Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions docker/job-started-hook.sh
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,13 @@ if [ -z "$cache_endpoint" ] || [ -z "$cache_authorization" ]; then
fi

assignment_endpoint="${cache_endpoint%/v1/runner-cache}/v1/runner-cache-v2/assignment"
# The actual workflow run narrows recovery to its jobs when a webhook is lost.
# Older images omit this hint and retain the scheduler's candidate fallback.
workflow_run_id=${GITHUB_RUN_ID:-}
case "$workflow_run_id" in
''|*[!0-9]*|0*) ;;
*) assignment_endpoint="$assignment_endpoint?run_id=$workflow_run_id" ;;
esac
assignment_max_attempts=${CF_RUNNER_CACHE_ASSIGNMENT_MAX_ATTEMPTS:-30}
case "$assignment_max_attempts" in
''|*[!0-9]*) assignment_max_attempts=30 ;;
Expand Down
4 changes: 2 additions & 2 deletions skills-lock.json
Original file line number Diff line number Diff line change
Expand Up @@ -11,13 +11,13 @@
"source": "dmmulroy/anti-slop",
"sourceType": "github",
"skillPath": "skills/install-anti-slop/SKILL.md",
"computedHash": "b903a7ac2c823d473a4b0046834fccbab0bc75ee4fda8be4f81efa895d475b01"
"computedHash": "4031728fbe75bdcad6ee3208fd52b5d66e167b056fefee1fa9758e9a6cb9c0c8"
},
"review-fix-address-bots": {
"source": "biw/skills",
"sourceType": "github",
"skillPath": "skills/review-fix-address-bots/SKILL.md",
"computedHash": "44f5a087d34e444abf84a783f0a6e84ffc6be6b471c10149a7750e6416495a19"
"computedHash": "b6f670a22af7cad23628756a1918a11151f77eaea248eec0948b0a6516ad6642"
}
}
}
347 changes: 331 additions & 16 deletions src/account-runner-scheduler.ts

Large diffs are not rendered by default.

133 changes: 133 additions & 0 deletions src/github-job-assignment.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
import { z } from "zod";

import { githubTokenForRunner, type GitHubAppDependencies, type GitHubAppEnvironment } from "./github-app";
import {
githubRunnerTokenFor,
type GitHubDynamicSecretEnvironment,
type GitHubRepositoryTarget,
} from "./github-repository";
import { githubHeaders } from "./provision";

const githubJobDetailSchema = z.object({
status: z.string(),
runner_id: z.number().int().nonnegative().nullable(),
runner_name: z.string().nullable(),
});
const workflowJobsSchema = z.object({
total_count: z.number().int().nonnegative(),
jobs: z.array(githubJobDetailSchema.extend({ id: z.number().int().positive() })),
});

export type GitHubJobDetail = z.infer<typeof githubJobDetailSchema>;
export type GitHubWorkflowJobs = z.infer<typeof workflowJobsSchema>;

const dependencies: GitHubAppDependencies = {
fetch: (input, init) => fetch(input, init),
now: () => Date.now(),
};

/** One lookup sweep reuses each installation credential across all its requests. */
export class GitHubJobAssignmentClient {
private readonly tokens = new Map<number | null, string | undefined>();

constructor(
private readonly env: GitHubAppEnvironment & GitHubDynamicSecretEnvironment,
private readonly target: GitHubRepositoryTarget,
private readonly deps: GitHubAppDependencies = dependencies,
) {}

private async request(
path: string,
installationId: number | null,
signal: AbortSignal,
): Promise<Response | undefined> {
try {
if (signal.aborted) {
return undefined;
}
if (!this.tokens.has(installationId)) {
const token = await githubTokenForRunner(
this.env,
this.target,
installationId,
(target) => githubRunnerTokenFor(this.env, target),
{
fetch: (input, init) => this.deps.fetch(input, { ...init, signal }),
now: this.deps.now,
},
);
this.tokens.set(installationId, token);
}
const token = this.tokens.get(installationId);
if (token === undefined || signal.aborted) {
return undefined;
}
const response = await this.deps.fetch(
`https://api.github.com/repos/${encodeURIComponent(this.target.owner)}/${encodeURIComponent(this.target.repository)}/actions/${path}`,
{
headers: githubHeaders(token),
signal: AbortSignal.any([signal, AbortSignal.timeout(5_000)]),
},
);
if (!response.ok) {
console.error("GitHub assignment lookup failed", { path, status: response.status });
await response.body?.cancel();
return undefined;
}
return response;
} catch (error) {
console.error("GitHub assignment lookup failed", {
path,
reason: signal.aborted ? "deadline" : error instanceof Error ? error.name : "request-error",
});
return undefined;
}
}

async jobDetail(
jobId: string,
installationId: number | null,
signal: AbortSignal,
): Promise<GitHubJobDetail | undefined> {
const response = await this.request(`jobs/${encodeURIComponent(jobId)}`, installationId, signal);
if (response === undefined) {
return undefined;
}
try {
const parsed = githubJobDetailSchema.safeParse(await response.json());
if (parsed.success) {
return parsed.data;
}
} catch {
// Invalid JSON and unexpected response shapes both leave the claim unresolved.
}
console.error("GitHub assignment lookup returned an invalid job", { jobId });
return undefined;
}

async workflowJobs(
runId: string,
page: number,
installationId: number | null,
signal: AbortSignal,
): Promise<GitHubWorkflowJobs | undefined> {
const response = await this.request(
`runs/${encodeURIComponent(runId)}/jobs?filter=latest&per_page=100&page=${page}`,
installationId,
signal,
);
if (response === undefined) {
return undefined;
}
try {
const parsed = workflowJobsSchema.safeParse(await response.json());
if (parsed.success) {
return parsed.data;
}
} catch {
// Never accept an assignment from a partial or malformed response.
}
console.error("GitHub assignment lookup returned invalid workflow jobs", { runId, page });
return undefined;
}
}
4 changes: 2 additions & 2 deletions src/provision.ts
Original file line number Diff line number Diff line change
Expand Up @@ -66,14 +66,14 @@ const githubJitConfigSchema = z.object({
});
const githubErrorResponseSchema = z.object({ message: z.string().optional() });

const githubHeaders = (token: string): HeadersInit => ({
export const githubHeaders = (token: string): HeadersInit => ({
Accept: "application/vnd.github+json",
Authorization: `Bearer ${token}`,
"User-Agent": "cloudflare-github-actions-runner",
"X-GitHub-Api-Version": "2022-11-28",
});

function githubRunnerUrl(target: GitHubRepositoryTarget, suffix = ""): string {
export function githubRunnerUrl(target: GitHubRepositoryTarget, suffix = ""): string {
return `https://api.github.com/repos/${encodeURIComponent(target.owner)}/${encodeURIComponent(target.repository)}/actions/runners${suffix}`;
}

Expand Down
20 changes: 17 additions & 3 deletions src/runner-cache.ts
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,11 @@ export interface RunnerCacheWriteAuthorizer {
* JIT runners are eligible for any compatible queued job, so the job that
* caused a runner to be provisioned is not necessarily the one it executes.
*/
cacheAssignment?(runnerName: string, repository: string): Promise<AssignedRunnerCacheScope | undefined>;
cacheAssignment?(
runnerName: string,
repository: string,
runId?: string,
): Promise<AssignedRunnerCacheScope | undefined>;
cacheScope?(runnerName: string, repository: string, jobId: string): Promise<RunnerCacheScope | undefined>;
/** @deprecated Kept temporarily for direct handler consumers upgrading to scoped access. */
canWriteCache?(runnerName: string, repository: string): Promise<boolean>;
Expand Down Expand Up @@ -455,11 +459,14 @@ interface ResolvedRunnerCacheClaim {
async function resolveRunnerCacheClaim(
writeAuthorizer: RunnerCacheWriteAuthorizer,
claim: RunnerCacheClaim,
runId?: string,
): Promise<ResolvedRunnerCacheClaim | undefined> {
const assignment =
writeAuthorizer.cacheAssignment === undefined
? undefined
: await writeAuthorizer.cacheAssignment(claim.runnerName, claim.repository);
: runId === undefined
? await writeAuthorizer.cacheAssignment(claim.runnerName, claim.repository)
: await writeAuthorizer.cacheAssignment(claim.runnerName, claim.repository, runId);
if (
assignment !== undefined &&
validRunnerCacheScope(assignment.cacheScope.scope) &&
Expand Down Expand Up @@ -996,7 +1003,14 @@ export async function handleRunnerCacheV2Request(
if (claim === undefined) {
return json({ error: "Unauthorized" }, 401);
}
const resolved = await resolveRunnerCacheClaim(writeAuthorizer, claim);
const runId =
request.method === "GET" && url.pathname === "/v1/runner-cache-v2/assignment"
? (url.searchParams.get("run_id") ?? undefined)
: undefined;
if (runId !== undefined && !/^[1-9][0-9]*$/u.test(runId)) {
return json({ error: "Invalid workflow run ID" }, 400);
}
const resolved = await resolveRunnerCacheClaim(writeAuthorizer, claim, runId);
// GitHub can start a JIT runner a moment before its authoritative
// `workflow_job: in_progress` webhook reaches this Worker. The runner's
// pre-job hook polls this endpoint, so actions/cache never observes that
Expand Down
22 changes: 16 additions & 6 deletions src/runner-provisioning-workflow.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import { runnerContainerFor } from "./runner-container-router";

export interface RunnerProvisioningWorkflowParameters {
jobId: string;
runnerName?: string;
}

const apiStepConfig = {
Expand All @@ -43,7 +44,7 @@ function schedulerFor(env: WorkerEnvironment) {
}

interface EligibilityReleaseScheduler {
provisioningFailed(jobId: string, reason: string): Promise<{ admissions: SchedulerAdmission[] }>;
provisioningFailed(jobId: string, reason: string, runnerName: string): Promise<{ admissions: SchedulerAdmission[] }>;
}

export interface EligibilityReleaseDependencies {
Expand Down Expand Up @@ -72,7 +73,11 @@ export async function releaseIfRepositoryIsIneligible(
return false;
}

const released = await scheduler.provisioningFailed(plan.jobId, `Repository visibility is ${eligibility.visibility}`);
const released = await scheduler.provisioningFailed(
plan.jobId,
`Repository visibility is ${eligibility.visibility}`,
plan.runnerName,
);
await dependencies.startProvisioning(env, released.admissions);
console.log("Cloudflare runner provisioning rejected by repository eligibility", {
jobId: plan.jobId,
Expand Down Expand Up @@ -137,7 +142,7 @@ export async function startRunnerProvisioningWorkflows(
// eslint-disable-next-line no-await-in-loop -- duplicate IDs are handled before the next job is scheduled.
await env.RUNNER_PROVISIONING_WORKFLOW.create({
id: admission.workflowId,
params: { jobId: admission.jobId },
params: { jobId: admission.jobId, runnerName: admission.runnerName },
retention: { successRetention: "1 day", errorRetention: "7 days" },
});
} catch {
Expand All @@ -164,7 +169,7 @@ export class RunnerProvisioningWorkflow extends WorkflowEntrypoint<
let plan: RunnerProvisioningPlan | undefined;
for (let attempt = 1; attempt <= 120; attempt += 1) {
// eslint-disable-next-line no-await-in-loop -- a custom slot may be configuring for an earlier job.
const claim = await scheduler.claimProvisioning(event.payload.jobId);
const claim = await scheduler.claimProvisioning(event.payload.jobId, event.payload.runnerName);
if (claim.kind === "provision") {
plan = claim;
break;
Expand Down Expand Up @@ -199,9 +204,13 @@ export class RunnerProvisioningWorkflow extends WorkflowEntrypoint<
await startRunnerProvisioningWorkflows(this.env, [...configuredAdmissions, ...capacityAdmissions]);
}

const mayStart = await scheduler.canStart(plan.jobId);
const mayStart = await scheduler.canStart(plan.jobId, plan.runnerName);
if (!mayStart) {
const released = await scheduler.provisioningFailed(plan.jobId, "GitHub completed before runner provisioning");
const released = await scheduler.provisioningFailed(
plan.jobId,
"GitHub completed before runner provisioning",
plan.runnerName,
);
await startRunnerProvisioningWorkflows(this.env, released.admissions);
return { kind: "cancelled" };
}
Expand Down Expand Up @@ -259,6 +268,7 @@ export class RunnerProvisioningWorkflow extends WorkflowEntrypoint<
const released = await scheduler.provisioningFailed(
plan.jobId,
error instanceof Error ? error.message : "Runner provisioning failed",
plan.runnerName,
);
await startRunnerProvisioningWorkflows(this.env, released.admissions);
throw error;
Expand Down
Loading
Loading