diff --git a/docker/job-started-hook.sh b/docker/job-started-hook.sh index d124126..186259d 100644 --- a/docker/job-started-hook.sh +++ b/docker/job-started-hook.sh @@ -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 ;; diff --git a/skills-lock.json b/skills-lock.json index fafbe14..ca935c0 100644 --- a/skills-lock.json +++ b/skills-lock.json @@ -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" } } } diff --git a/src/account-runner-scheduler.ts b/src/account-runner-scheduler.ts index 4c16ba8..37f70f1 100644 --- a/src/account-runner-scheduler.ts +++ b/src/account-runner-scheduler.ts @@ -4,8 +4,11 @@ import { z } from "zod"; import { prepareRunnerApplication, reconcileRunnerApplicationCapacity } from "./cloudflare-containers"; import type { WorkerEnvironment } from "./environment"; import { githubRunnerTokenFor, type GitHubRepositoryTarget } from "./github-repository"; -import { githubTokenForRunner } from "./github-app"; +import { githubTokenForRunner, type GitHubAppDependencies } from "./github-app"; +import { GitHubJobAssignmentClient, type GitHubJobDetail } from "./github-job-assignment"; import { deleteGitHubRunner } from "./provision"; +import { assignResourceTraceRunner } from "./resource-traces"; +import { startRunnerProvisioningWorkflows } from "./runner-provisioning-workflow"; import { runnerProfileSchema, RUNNER_PROFILE_KEYS, @@ -26,6 +29,14 @@ import { } from "./scheduler-policy"; const SCHEDULER_ALARM_DELAY_MS = 60_000; +// A running Container's job-started hook polls runner-cache assignment about +// once per second. The GitHub self-heal below is an escape hatch for lost +// `in_progress` deliveries, so bound it instead of calling the API per poll. +const RUNNER_ASSIGNMENT_GITHUB_RECONCILE_INTERVAL_MS = 10_000; +// Each reconcile probes candidate jobs through the GitHub API with a 5 s +// per-request timeout; cap the whole sweep so it stays well inside the +// job-started hook's 30 s budget even when the candidate list is full. +const RUNNER_ASSIGNMENT_GITHUB_RECONCILE_BUDGET_MS = 8_000; const RUNNER_COMPLETION_GRACE_MS = 30_000; // GitHub can deliver workflow_job: completed before actions/cache finishes its // post-job upload through the runner's local results proxy. Keep the one-job @@ -98,6 +109,7 @@ interface JobRow { updated_at: number; runner_id: number | null; runner_attempt: number; + provisioning_owner_scoped: number; failure_reason: string | null; container_stopped_at: number | null; container_exit_code: number | null; @@ -156,6 +168,7 @@ interface JitRunnerRow { github_owner: string; github_repository: string; profile_key: string; + github_runner_id: number | null; assigned_job_id: string | null; assignment_observed: number; created_at: number; @@ -342,6 +355,12 @@ function jobMatchesRunnerClaim(job: JobRow, input: SchedulerRunnerClaimInput): b * mutations, while Workflows handle only runner start and rollout polling. */ export class AccountRunnerScheduler extends DurableObject { + private readonly runnerReconcileAttempts = new Map< + string, + { timestamp: number; cursor: number; runId?: string; page: number } + >(); + assignmentDependencies?: GitHubAppDependencies; + constructor(ctx: DurableObjectState, env: WorkerEnvironment) { super(ctx, env); ctx.blockConcurrencyWhile(async () => { @@ -497,6 +516,7 @@ export class AccountRunnerScheduler extends DurableObject { ["cache_write_allowed", "INTEGER NOT NULL DEFAULT 0"], ["github_assignment_observed", "INTEGER NOT NULL DEFAULT 0"], ["runner_attempt", "INTEGER NOT NULL DEFAULT 1"], + ["provisioning_owner_scoped", "INTEGER NOT NULL DEFAULT 0"], ["container_stopped_at", "INTEGER"], ["container_exit_code", "INTEGER"], ["container_stop_reason", "TEXT"], @@ -510,6 +530,14 @@ export class AccountRunnerScheduler extends DurableObject { this.ctx.storage.sql.exec(`ALTER TABLE scheduler_jobs ADD COLUMN ${name} ${definition}`); } } + + const jitRunnerColumns = this.rows<{ name: string }>("PRAGMA table_info(scheduler_jit_runners)"); + const jitRunnerColumnMigrations = [["github_runner_id", "INTEGER"]] as const; + for (const [name, definition] of jitRunnerColumnMigrations) { + if (!jitRunnerColumns.some((column) => column.name === name)) { + this.ctx.storage.sql.exec(`ALTER TABLE scheduler_jit_runners ADD COLUMN ${name} ${definition}`); + } + } // Jobs that were queued before the pool became multi-repository belong to // the original POC repository and must retain a usable cleanup credential. if (this.env.LEGACY_GITHUB_OWNER !== undefined && this.env.LEGACY_GITHUB_REPOSITORY !== undefined) { @@ -968,11 +996,14 @@ export class AccountRunnerScheduler extends DurableObject { return { accepted: true, queueReason: stored === undefined ? undefined : this.queueReason(stored), admissions }; } - async claimProvisioning(jobId: string): Promise { + async claimProvisioning(jobId: string, runnerName?: string): Promise { const job = this.job(jobId); if (job === undefined) { return { kind: "missing" }; } + if (runnerName !== undefined && job.runner_name !== runnerName) { + return { kind: "cancelled" }; + } if (job.status === "queued" || job.status === "admitted") { if (job.slot_id === null) { return { kind: "wait" }; @@ -993,7 +1024,8 @@ export class AccountRunnerScheduler extends DurableObject { return { kind: "wait" }; } this.ctx.storage.sql.exec( - "UPDATE scheduler_jobs SET status = 'provisioning', updated_at = ? WHERE job_id = ?", + "UPDATE scheduler_jobs SET status = 'provisioning', provisioning_owner_scoped = ?, updated_at = ? WHERE job_id = ?", + runnerName === undefined ? 0 : 1, now(), jobId, ); @@ -1116,9 +1148,9 @@ export class AccountRunnerScheduler extends DurableObject { return admissions; } - async canStart(jobId: string): Promise { + async canStart(jobId: string, runnerName?: string): Promise { const job = this.job(jobId); - return job?.status === "provisioning"; + return job?.status === "provisioning" && (runnerName === undefined || job.runner_name === runnerName); } async runnerProvisioned(jobId: string, runnerName: string, runnerId: number): Promise { @@ -1143,13 +1175,14 @@ export class AccountRunnerScheduler extends DurableObject { const timestamp = now(); this.ctx.storage.sql.exec( `INSERT INTO scheduler_jit_runners - (runner_name, source_job_id, github_owner, github_repository, profile_key, assigned_job_id, assignment_observed, created_at, updated_at) - VALUES (?, ?, ?, ?, ?, NULL, 0, ?, ?) + (runner_name, source_job_id, github_owner, github_repository, profile_key, github_runner_id, assigned_job_id, assignment_observed, created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, NULL, 0, ?, ?) ON CONFLICT(runner_name) DO UPDATE SET source_job_id = excluded.source_job_id, github_owner = excluded.github_owner, github_repository = excluded.github_repository, profile_key = excluded.profile_key, + github_runner_id = excluded.github_runner_id, assigned_job_id = NULL, assignment_observed = 0, updated_at = excluded.updated_at`, @@ -1158,6 +1191,7 @@ export class AccountRunnerScheduler extends DurableObject { job.github_owner, job.github_repository, job.profile_key, + runnerId, timestamp, timestamp, ); @@ -1191,6 +1225,7 @@ export class AccountRunnerScheduler extends DurableObject { async cacheAssignment( runnerName: string, repository: string, + runId?: string, ): Promise< | { jobId: string; @@ -1246,15 +1281,228 @@ export class AccountRunnerScheduler extends DurableObject { repo, now() - RUNNER_CACHE_POST_JOB_GRACE_MS, )[0]; - if (job === undefined || job.github_assignment_observed === 0 || job.cache_scope === "") { + if (job !== undefined && job.github_assignment_observed === 1 && job.cache_scope !== "") { + return { + jobId: job.job_id, + cacheScope: storedCacheScope(job.cache_scope, job.cache_fallback_scope, job.cache_write_allowed), + }; + } + if (runId !== undefined && !/^[1-9][0-9]*$/u.test(runId)) { + return undefined; + } + return this.reconcileRunnerAssignmentFromGitHub(runnerName, owner, repo, runId); + } + + /** + * `workflow_job: in_progress` is the authoritative runner-to-job assignment, + * but GitHub does not re-deliver webhooks the Worker drops or fails to + * acknowledge, so a lost delivery leaves a running job polling for its + * assignment until the job-started hook times out. New hooks provide their + * actual workflow run, whose jobs can be listed without scanning unrelated + * queued jobs. Older hooks retain a deadline-bounded candidate scan. Both + * paths require GitHub to confirm this runner and use the webhook's + * ownership-transfer and resource-attribution paths. + */ + private async reconcileRunnerAssignmentFromGitHub( + runnerName: string, + owner: string, + repo: string, + runId?: string, + ): Promise<{ jobId: string; cacheScope: SchedulerCacheScope } | undefined> { + const runner = this.rows( + `SELECT * FROM scheduler_jit_runners + WHERE runner_name = ? + AND lower(github_owner) = lower(?) + AND lower(github_repository) = lower(?)`, + runnerName, + owner, + repo, + )[0]; + if (runner === undefined) { + return undefined; + } + const attempt: { timestamp: number; cursor: number; runId?: string; page: number } = + this.runnerReconcileAttempts.get(runnerName) ?? { + timestamp: 0, + cursor: Number.MAX_SAFE_INTEGER, + page: 1, + }; + const timestamp = this.assignmentNow(); + if (timestamp - attempt.timestamp < RUNNER_ASSIGNMENT_GITHUB_RECONCILE_INTERVAL_MS) { + return undefined; + } + // Entries live only for the duration of a Durable Object instance, but a + // long-lived scheduler would otherwise accumulate one entry per claimed + // runner forever; drop state for runners that have not polled recently. + for (const [name, entry] of this.runnerReconcileAttempts) { + if (timestamp - entry.timestamp > 60 * 60 * 1000) { + this.runnerReconcileAttempts.delete(name); + } + } + attempt.timestamp = timestamp; + if (attempt.runId !== runId) { + attempt.runId = runId; + attempt.page = 1; + } + this.runnerReconcileAttempts.set(runnerName, attempt); + const deadline = timestamp + RUNNER_ASSIGNMENT_GITHUB_RECONCILE_BUDGET_MS; + + const target: GitHubRepositoryTarget = { owner, repository: repo }; + const client = new GitHubJobAssignmentClient(this.env, target, this.assignmentDependencies); + if (runId !== undefined) { + const sourceJob = this.job(runner.source_job_id); + if (sourceJob === undefined) { + return undefined; + } + while (this.assignmentNow() < deadline) { + const remaining = deadline - this.assignmentNow(); + if (remaining <= 0) { + return undefined; + } + // eslint-disable-next-line no-await-in-loop -- page requests stop as soon as this runner is found or the sweep deadline expires. + const page = await client.workflowJobs( + runId, + attempt.page, + sourceJob.github_installation_id, + AbortSignal.timeout(remaining), + ); + if (page === undefined) { + return undefined; + } + const detail = page.jobs.find((job) => job.status === "in_progress" && job.runner_name === runnerName); + if (detail !== undefined) { + return this.applyReconciledAssignment(runner, String(detail.id), detail, target); + } + if (page.jobs.length < 100 || attempt.page * 100 >= page.total_count) { + attempt.page = 1; + return undefined; + } + attempt.page += 1; + } + return undefined; + } + + // SQL pages bound memory, not the number of candidates reachable per + // sweep. Keep probing pages while the deadline permits, and resume after + // the last actually probed row when a slow request exhausts the budget. + while (this.assignmentNow() < deadline) { + const candidates = this.rows( + `SELECT * FROM scheduler_jobs + WHERE lower(github_owner) = lower(?) + AND lower(github_repository) = lower(?) + AND profile_key = ? + AND github_assignment_observed = 0 + AND status IN ('queued', 'admitted', 'provisioning', 'running') + AND CAST(job_id AS INTEGER) < ? + ORDER BY CAST(job_id AS INTEGER) DESC + LIMIT 9`, + owner, + repo, + runner.profile_key, + attempt.cursor, + ); + // One lookahead row tells us whether a full eight-row page is also + // the final page, even if its last probe consumes the whole budget. + for (const candidate of candidates.slice(0, 8)) { + const remaining = deadline - this.assignmentNow(); + if (remaining <= 0) { + return undefined; + } + const signal = AbortSignal.timeout(remaining); + // eslint-disable-next-line no-await-in-loop -- sequential probes share a credential and stop at the first match or deadline. + const detail = await (this.jobDetailOverride === undefined + ? client.jobDetail(candidate.job_id, candidate.github_installation_id, signal) + : this.jobDetailOverride(candidate.job_id, target, candidate.github_installation_id)); + attempt.cursor = Number(candidate.job_id); + if (detail?.status === "in_progress" && detail.runner_name === runnerName) { + // eslint-disable-next-line no-await-in-loop -- assignment validation follows the authoritative API response. + const result = await this.applyReconciledAssignment(runner, candidate.job_id, detail, target); + if (result !== undefined) { + return result; + } + } + } + if (candidates.length <= 8) { + // Reset only after fully consuming this page; an interrupted page + // must retain its cursor so its unprobed tail is revisited. + attempt.cursor = Number.MAX_SAFE_INTEGER; + return undefined; + } + } + return undefined; + } + + private assignmentNow(): number { + return this.assignmentDependencies?.now() ?? now(); + } + + private async applyReconciledAssignment( + runner: JitRunnerRow, + jobId: string, + detail: GitHubJobDetail, + target: GitHubRepositoryTarget, + ): Promise<{ jobId: string; cacheScope: SchedulerCacheScope } | undefined> { + const fresh = this.job(jobId); + if ( + fresh === undefined || + !["queued", "admitted", "provisioning", "running"].includes(fresh.status) || + fresh.github_owner.toLowerCase() !== target.owner.toLowerCase() || + fresh.github_repository.toLowerCase() !== target.repository.toLowerCase() || + fresh.profile_key !== runner.profile_key || + detail.runner_id === 0 + ) { + return undefined; + } + const result = await this.workflowJobStarted({ + jobId, + runnerName: runner.runner_name, + runnerId: detail.runner_id ?? runner.github_runner_id ?? undefined, + target, + profile: parseProfile(fresh.profile_json), + }); + if (!result.accepted) { + this.recordEvent("jit-runner-assignment-unreconciled", { + jobId, + detail: { runnerName: runner.runner_name, sourceJobId: runner.source_job_id }, + }); + return undefined; + } + this.runnerReconcileAttempts.delete(runner.runner_name); + this.recordEvent("github-assignment-reconciled", { + jobId, + detail: { + runnerName: runner.runner_name, + runnerId: detail.runner_id ?? runner.github_runner_id, + via: "github-api", + }, + }); + this.ctx.waitUntil(startRunnerProvisioningWorkflows(this.env, result.admissions)); + await assignResourceTraceRunner(this.env, { + runnerName: runner.runner_name, + jobId, + repository: `${target.owner}/${target.repository}`, + }); + const assigned = this.job(jobId); + if (assigned === undefined || assigned.github_assignment_observed !== 1 || assigned.cache_scope === "") { return undefined; } return { - jobId: job.job_id, - cacheScope: storedCacheScope(job.cache_scope, job.cache_fallback_scope, job.cache_write_allowed), + jobId, + cacheScope: storedCacheScope(assigned.cache_scope, assigned.cache_fallback_scope, assigned.cache_write_allowed), }; } + /** + * Authoritative state of a single workflow job. The `jobDetailOverride` + * seam exists so tests can answer without real API credentials or network + * access. + */ + jobDetailOverride?: ( + jobId: string, + target: GitHubRepositoryTarget, + installationId: number | null, + ) => Promise; + /** @deprecated Use cacheAssignment, which follows a JIT runner's actual job. */ async cacheScope(runnerName: string, repository: string, jobId: string): Promise { const [owner, repo, extra] = repository.split("/"); @@ -1287,13 +1535,51 @@ export class AccountRunnerScheduler extends DurableObject { * put the displaced job back through normal admission with a fresh JIT name. */ async workflowJobStarted(input: SchedulerRunnerClaimInput): Promise { - this.recordJitRunnerAssignment(input); + const assignmentRecorded = this.recordJitRunnerAssignment(input); const runnerOwner = this.rows( `SELECT * FROM scheduler_jobs WHERE runner_name = ? AND status IN ('provisioning', 'running', 'releasing')`, input.runnerName, )[0]; if (runnerOwner === undefined) { + // The runner is ours and GitHub says it executes a job we know, but no + // scheduler job still owns the runner — typically a mutual cross- + // assignment where both runners were displaced onto each other's jobs. + // Attach the runner to the job actually running so the displaced job is + // not requeued behind a runner that will never exist. + const actualJob = this.job(input.jobId); + // `provisioning` is excluded deliberately: an in-flight provisioning + // workflow for that job would observe `canStart() === false` after the + // rename and tear the adopted job down via provisioningFailed. + if ( + assignmentRecorded && + actualJob !== undefined && + jobMatchesRunnerClaim(actualJob, input) && + (actualJob.status === "queued" || actualJob.status === "admitted" || actualJob.status === "running") + ) { + this.ctx.storage.sql.exec( + `UPDATE scheduler_jobs + SET status = 'running', runner_name = ?, runner_id = ?, failure_reason = NULL, + container_stopped_at = NULL, container_exit_code = NULL, container_stop_reason = NULL, + recovery_due_at = NULL, runner_cleanup_state = 'none', runner_cleanup_due_at = NULL, + github_assignment_observed = 1, updated_at = ? + WHERE job_id = ?`, + input.runnerName, + input.runnerId ?? null, + now(), + actualJob.job_id, + ); + this.recordEvent("github-job-started-adopted-unowned-runner", { + jobId: input.jobId, + slotId: actualJob.slot_id ?? undefined, + detail: { runnerName: input.runnerName, runnerId: input.runnerId }, + }); + return { accepted: true, admissions: [] }; + } + this.recordEvent("github-job-started-without-owned-job", { + jobId: input.jobId, + detail: { runnerName: input.runnerName, runnerId: input.runnerId }, + }); return { accepted: false, admissions: [] }; } @@ -1391,7 +1677,7 @@ export class AccountRunnerScheduler extends DurableObject { return { accepted: true, admissions }; } - private recordJitRunnerAssignment(input: SchedulerRunnerClaimInput): void { + private recordJitRunnerAssignment(input: SchedulerRunnerClaimInput): boolean { const runner = this.rows( "SELECT * FROM scheduler_jit_runners WHERE runner_name = ?", input.runnerName, @@ -1405,7 +1691,15 @@ export class AccountRunnerScheduler extends DurableObject { runner.profile_key !== input.profile.key || !jobMatchesRunnerClaim(assignedJob, input) ) { - return; + this.recordEvent("jit-runner-assignment-unmatched", { + jobId: input.jobId, + detail: { + runnerName: input.runnerName, + runnerRegistered: runner !== undefined, + jobRegistered: assignedJob !== undefined, + }, + }); + return false; } this.ctx.storage.sql.exec( `UPDATE scheduler_jit_runners @@ -1415,6 +1709,7 @@ export class AccountRunnerScheduler extends DurableObject { now(), input.runnerName, ); + return true; } private releaseJob(job: JobRow, status: "completed" | "cancelled" | "failed", reason?: string): void { @@ -1465,9 +1760,29 @@ export class AccountRunnerScheduler extends DurableObject { }); } - async provisioningFailed(jobId: string, reason: string): Promise { + async provisioningFailed(jobId: string, reason: string, runnerName?: string): Promise { const job = this.job(jobId); - if (job === undefined || !activeJobStates.has(job.status)) { + // Only the workflow that owns the current provisioning attempt may fail + // the job. A stale or retried workflow — identifiable by the runner name + // from its claim, which embeds the `-rN` attempt suffix — must never tear + // down a job that already moved on (adopted onto a different JIT runner + // after a GitHub cross-assignment) or was re-queued and claimed again. + // During deployment, old workflows still send two arguments. An + // unscoped failure is accepted only for an attempt claimed through the + // legacy protocol. New workflows persist scoped ownership when claiming, + // so an old workflow cannot fail a replacement attempt. + const ownsAttempt = + job !== undefined && + (runnerName === undefined + ? job.provisioning_owner_scoped === 0 && job.github_assignment_observed === 0 + : job.runner_name === runnerName); + if (job === undefined || job.status !== "provisioning" || !ownsAttempt) { + if (job !== undefined && activeJobStates.has(job.status)) { + this.recordEvent("provisioning-failed-ignored", { + jobId, + detail: { status: job.status, runnerName, currentRunnerName: job.runner_name, reason }, + }); + } return { accepted: false, admissions: [] }; } this.releaseJob(job, "failed", reason); @@ -1904,7 +2219,7 @@ export class AccountRunnerScheduler extends DurableObject { // eslint-disable-next-line no-await-in-loop -- deterministic workflow IDs make each admission idempotent. await this.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 { diff --git a/src/github-job-assignment.ts b/src/github-job-assignment.ts new file mode 100644 index 0000000..7ef422d --- /dev/null +++ b/src/github-job-assignment.ts @@ -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; +export type GitHubWorkflowJobs = z.infer; + +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(); + + 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 { + 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 { + 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 { + 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; + } +} diff --git a/src/provision.ts b/src/provision.ts index 95c89e6..78fcd9b 100644 --- a/src/provision.ts +++ b/src/provision.ts @@ -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}`; } diff --git a/src/runner-cache.ts b/src/runner-cache.ts index de3ec69..44933d2 100644 --- a/src/runner-cache.ts +++ b/src/runner-cache.ts @@ -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; + cacheAssignment?( + runnerName: string, + repository: string, + runId?: string, + ): Promise; cacheScope?(runnerName: string, repository: string, jobId: string): Promise; /** @deprecated Kept temporarily for direct handler consumers upgrading to scoped access. */ canWriteCache?(runnerName: string, repository: string): Promise; @@ -455,11 +459,14 @@ interface ResolvedRunnerCacheClaim { async function resolveRunnerCacheClaim( writeAuthorizer: RunnerCacheWriteAuthorizer, claim: RunnerCacheClaim, + runId?: string, ): Promise { 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) && @@ -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 diff --git a/src/runner-provisioning-workflow.ts b/src/runner-provisioning-workflow.ts index 033fbb4..cb70da7 100644 --- a/src/runner-provisioning-workflow.ts +++ b/src/runner-provisioning-workflow.ts @@ -18,6 +18,7 @@ import { runnerContainerFor } from "./runner-container-router"; export interface RunnerProvisioningWorkflowParameters { jobId: string; + runnerName?: string; } const apiStepConfig = { @@ -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 { @@ -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, @@ -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 { @@ -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; @@ -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" }; } @@ -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; diff --git a/tests/durable/account-runner-scheduler.test.ts b/tests/durable/account-runner-scheduler.test.ts index eee9841..0a20b59 100644 --- a/tests/durable/account-runner-scheduler.test.ts +++ b/tests/durable/account-runner-scheduler.test.ts @@ -1,5 +1,5 @@ import { env, runInDurableObject } from "cloudflare:test"; -import { describe, expect, it } from "vite-plus/test"; +import { beforeAll, describe, expect, it, vi } from "vite-plus/test"; import type { SchedulerJobInput } from "../../src/account-runner-scheduler"; import { RUNNER_PROFILES } from "../../src/runner-profiles"; @@ -7,6 +7,27 @@ import { RUNNER_PROFILES } from "../../src/runner-profiles"; const target = { owner: "biw", repository: "runner-poc" }; const profile = RUNNER_PROFILES["standard-3"]; +beforeAll(async () => { + await env.RESOURCE_METRICS.prepare( + `CREATE TABLE IF NOT EXISTS resource_trace_assignments ( + runner_name TEXT PRIMARY KEY, job_id TEXT NOT NULL, repository TEXT NOT NULL, assigned_at INTEGER NOT NULL + )`, + ).run(); +}); + +async function assignmentClock( + scheduler: DurableObjectStub, + fetch: typeof globalThis.fetch = async () => { + throw new Error("Unexpected GitHub request"); + }, +) { + const clock = { time: 1_700_000_000_000 }; + await runInDurableObject(scheduler, async (instance) => { + instance.assignmentDependencies = { fetch, now: () => clock.time }; + }); + return clock; +} + function job(jobId: string, runnerName: string, cacheScope: string): SchedulerJobInput { return { jobId, @@ -123,6 +144,215 @@ describe("AccountRunnerScheduler JIT cache assignments", () => { await expect(scheduler.cacheAssignment(queuedJob.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); }); + it("recovers a lost in_progress delivery by asking GitHub which job the runner executes", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("github-reconcile"); + const queuedJob = job("500", "cf-standard-3-job-500", "refs/pull/500/merge"); + const reassignedJob = job("600", "cf-standard-3-job-600", "refs/pull/600/merge"); + + await scheduler.submit(queuedJob); + await scheduler.submit(reassignedJob); + await provisionRunner(scheduler, queuedJob.jobId, queuedJob.runnerName, 5_001); + + // GitHub assigned the provisioned runner to a different job and its + // in_progress webhook never reached the Worker. The claim must still + // resolve through the GitHub API self-heal instead of timing the job out. + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async (jobId) => + jobId === reassignedJob.jobId + ? { status: "in_progress", runner_id: 5_001, runner_name: queuedJob.runnerName } + : { status: "queued", runner_id: null, runner_name: null }; + }); + + await expect(scheduler.cacheAssignment(queuedJob.runnerName, "biw/runner-poc")).resolves.toEqual({ + jobId: reassignedJob.jobId, + cacheScope: { + scope: reassignedJob.cacheScope?.scope, + fallbackScope: "refs/heads/main", + writeAllowed: true, + }, + }); + + // The ownership transfer must run through the same path as the webhook: + // the displaced job is requeued under a fresh retry runner name and the + // actual job takes over the runner reservation. + await runInDurableObject(scheduler, async (_instance, state) => { + // SAFETY: the query selects exactly these two columns and every row carries them. + const displaced = state.storage.sql + .exec(`SELECT status, runner_name FROM scheduler_jobs WHERE job_id = ?`, queuedJob.jobId) + .toArray()[0] as { status: string; runner_name: string }; + expect(["queued", "admitted", "provisioning"]).toContain(displaced.status); + // The displaced job is requeued under a fresh `-rN` runner name (the + // provisioning workflow may fail in the test environment and retry). + expect(displaced.runner_name).toMatch(/cf-standard-3-job-500-r\d+/u); + // SAFETY: the query selects exactly these three columns and every row carries them. + const actual = state.storage.sql + .exec( + `SELECT status, runner_name, github_assignment_observed FROM scheduler_jobs WHERE job_id = ?`, + reassignedJob.jobId, + ) + .toArray()[0] as { status: string; runner_name: string; github_assignment_observed: number }; + expect(actual.status).toBe("running"); + expect(actual.runner_name).toBe(queuedJob.runnerName); + expect(actual.github_assignment_observed).toBe(1); + }); + }); + + it("adopts the running job when its runner was already displaced by a mutual cross-assignment", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("mutual-cross-assign"); + const jobA = job("900", "cf-standard-3-job-900", "refs/pull/900/merge"); + const jobB = job("901", "cf-standard-3-job-901", "refs/pull/901/merge"); + + await scheduler.submit(jobA); + await scheduler.submit(jobB); + await provisionRunner(scheduler, jobA.jobId, jobA.runnerName, 9_001); + await provisionRunner(scheduler, jobB.jobId, jobB.runnerName, 9_002); + + // GitHub cross-assigns both runners: RA executes B, RB executes A. + await scheduler.workflowJobStarted({ + jobId: jobB.jobId, + runnerName: jobA.runnerName, + runnerId: 9_001, + target, + profile, + }); + // RB's owner row (job B) was just moved onto RA, so nothing owns RB. The + // scheduler must still adopt job A onto RB instead of leaving A queued + // behind the never-provisioned `RA-r1` runner. + await scheduler.workflowJobStarted({ + jobId: jobA.jobId, + runnerName: jobB.runnerName, + runnerId: 9_002, + target, + profile, + }); + + await expect(scheduler.cacheAssignment(jobB.runnerName, "biw/runner-poc")).resolves.toEqual({ + jobId: jobA.jobId, + cacheScope: { + scope: jobA.cacheScope?.scope, + fallbackScope: "refs/heads/main", + writeAllowed: true, + }, + }); + await runInDurableObject(scheduler, async (_instance, state) => { + // SAFETY: the query selects exactly these two columns and every row carries them. + const adopted = state.storage.sql + .exec(`SELECT status, runner_name FROM scheduler_jobs WHERE job_id = ?`, jobA.jobId) + .toArray()[0] as { status: string; runner_name: string }; + expect(adopted.status).toBe("running"); + expect(adopted.runner_name).toBe(jobB.runnerName); + }); + + // A late provisioning workflow for job A's abandoned `RA-r1` runner must + // not tear the adopted job down: provisioningFailed only applies while + // the job is still provisioning that exact runner name. + await scheduler.provisioningFailed(jobA.jobId, "stale workflow", `${jobA.runnerName}-r1`); + + await runInDurableObject(scheduler, async (_instance, state) => { + // SAFETY: the query selects exactly these two columns and every row carries them. + const surviving = state.storage.sql + .exec(`SELECT status, runner_name FROM scheduler_jobs WHERE job_id = ?`, jobA.jobId) + .toArray()[0] as { status: string; runner_name: string }; + expect(surviving.status).toBe("running"); + expect(surviving.runner_name).toBe(jobB.runnerName); + }); + }); + + it("ignores a failure report from a superseded provisioning attempt", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("stale-provisioning-failure"); + const queuedJob = job("950", "cf-standard-3-job-950", "refs/pull/950/merge"); + const other = job("951", "cf-standard-3-job-951", "refs/pull/951/merge"); + + await scheduler.submit(queuedJob); + await scheduler.submit(other); + await provisionRunner(scheduler, queuedJob.jobId, queuedJob.runnerName, 9_501); + await provisionRunner(scheduler, other.jobId, other.runnerName, 9_502); + + // Cross-assignment: GitHub puts `other` on queuedJob's runner, requeueing + // queuedJob under `cf-standard-3-job-950-r1`; the new attempt then claims + // provisioning again. + await scheduler.workflowJobStarted({ + jobId: other.jobId, + runnerName: queuedJob.runnerName, + runnerId: 9_501, + target, + profile, + }); + // The fresh attempt claims provisioning again under the `-r1` runner + // name (claimProvisioning may return `wait` while slot capacity is being + // applied, so drive the state directly). + await runInDurableObject(scheduler, async (_instance, state) => { + state.storage.sql.exec( + `UPDATE scheduler_jobs SET status = 'provisioning', updated_at = ? WHERE job_id = ?`, + Date.now(), + queuedJob.jobId, + ); + }); + + // The superseded workflow reports its failure late; the job must not be + // failed because its current runner name no longer matches. + await scheduler.provisioningFailed(queuedJob.jobId, "stale workflow", queuedJob.runnerName); + + await runInDurableObject(scheduler, async (_instance, state) => { + // SAFETY: the query selects exactly these two columns and every row carries them. + const surviving = state.storage.sql + .exec(`SELECT status, runner_name FROM scheduler_jobs WHERE job_id = ?`, queuedJob.jobId) + .toArray()[0] as { status: string; runner_name: string }; + // The job may legitimately retry further (the test environment fails + // real provisioning attempts), but the stale report must not kill it. + expect(surviving.status).not.toBe("failed"); + expect(surviving.runner_name).toMatch(/cf-standard-3-job-950-r\d+/u); + }); + }); + + it("reconciles a running job whose in_progress webhook was lost", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("github-reconcile-running"); + const queuedJob = job("800", "cf-standard-3-job-800", "refs/pull/800/merge"); + + await scheduler.submit(queuedJob); + await provisionRunner(scheduler, queuedJob.jobId, queuedJob.runnerName, 8_001); + // runnerStarted moves the job to `running` before GitHub's in_progress + // webhook arrives; a lost delivery must still resolve. + await scheduler.runnerStarted(queuedJob.runnerName); + + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async () => ({ + status: "in_progress", + runner_id: 8_001, + runner_name: queuedJob.runnerName, + }); + }); + + await expect(scheduler.cacheAssignment(queuedJob.runnerName, "biw/runner-poc")).resolves.toEqual({ + jobId: queuedJob.jobId, + cacheScope: { + scope: queuedJob.cacheScope?.scope, + fallbackScope: "refs/heads/main", + writeAllowed: true, + }, + }); + }); + + it("does not resolve an assignment when GitHub reports the job on the runner is finished", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("github-reconcile-empty"); + const queuedJob = job("700", "cf-standard-3-job-700", "refs/pull/700/merge"); + + await scheduler.submit(queuedJob); + await provisionRunner(scheduler, queuedJob.jobId, queuedJob.runnerName, 7_001); + + await runInDurableObject(scheduler, async (instance) => { + // The runner name matches but the job already completed, so the status + // check itself must deny the claim. + instance.jobDetailOverride = async () => ({ + status: "completed", + runner_id: 7_001, + runner_name: queuedJob.runnerName, + }); + }); + + await expect(scheduler.cacheAssignment(queuedJob.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + }); + it("does not record an assignment across repository or machine-profile boundaries", async () => { const scheduler = env.RUNNER_SCHEDULER.getByName("isolated-assignment"); const queuedJob = job("400", "cf-standard-3-job-400", "refs/pull/400/merge"); @@ -139,4 +369,339 @@ describe("AccountRunnerScheduler JIT cache assignments", () => { await expect(scheduler.cacheAssignment(queuedJob.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); }); + + it("reaches a legacy assignment beyond 24 newer candidates in one bounded sweep", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("large-legacy-backlog"); + const source = job("1000", "cf-standard-3-job-1000", "refs/pull/1000/merge"); + await scheduler.submit(source); + for (let id = 1001; id <= 1040; id += 1) { + // eslint-disable-next-line no-await-in-loop -- prepare a deterministic queue with the executing job last. + await scheduler.submit(job(String(id), `cf-standard-3-job-${id}`, `refs/pull/${id}/merge`)); + } + await provisionRunner(scheduler, source.jobId, source.runnerName, 10_001); + const clock = await assignmentClock(scheduler); + const probes: string[] = []; + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async (jobId) => { + probes.push(jobId); + clock.time += 10; + return jobId === source.jobId + ? { status: "in_progress", runner_id: 10_001, runner_name: source.runnerName } + : { status: "queued", runner_id: null, runner_name: null }; + }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(probes).toHaveLength(41); + expect(clock.time - 1_700_000_000_000).toBeLessThan(8_000); + }); + + it("wraps a fully probed final page on the next eligible retry without an empty sweep", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("reconcile-cursor-wrap"); + const source = job("1100", "cf-standard-3-job-1100", "refs/pull/1100/merge"); + await scheduler.submit(source); + await provisionRunner(scheduler, source.jobId, source.runnerName, 11_001); + const clock = await assignmentClock(scheduler); + let probes = 0; + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async () => { + probes += 1; + return probes === 1 + ? { status: "queued", runner_id: null, runner_name: null } + : { status: "in_progress", runner_id: 11_001, runner_name: source.runnerName }; + }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + expect(probes).toBe(1); + clock.time += 10_000; + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(probes).toBe(2); + }); + + it("resumes the unprobed tail when the deadline interrupts a short final page", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("reconcile-cursor-deadline"); + const source = job("1200", "cf-standard-3-job-1200", "refs/pull/1200/merge"); + await scheduler.submit(source); + await scheduler.submit(job("1201", "cf-standard-3-job-1201", "refs/pull/1201/merge")); + await scheduler.submit(job("1202", "cf-standard-3-job-1202", "refs/pull/1202/merge")); + await provisionRunner(scheduler, source.jobId, source.runnerName, 12_001); + const clock = await assignmentClock(scheduler); + const probes: string[] = []; + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async (jobId) => { + probes.push(jobId); + if (probes.length === 1) { + clock.time += 8_000; + } + return jobId === source.jobId + ? { status: "in_progress", runner_id: 12_001, runner_name: source.runnerName } + : { status: "queued", runner_id: null, runner_name: null }; + }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + expect(probes).toEqual(["1202"]); + clock.time += 10_000; + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(probes).toEqual(["1202", "1201", "1200"]); + }); + + it("wraps a full final page when its last probe exhausts the sweep budget", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("reconcile-full-final-page"); + const source = job("1250", "cf-standard-3-job-1250", "refs/pull/1250/merge"); + await scheduler.submit(source); + for (let id = 1251; id <= 1257; id += 1) { + // eslint-disable-next-line no-await-in-loop -- prepare exactly one full candidate page. + await scheduler.submit(job(String(id), `cf-standard-3-job-${id}`, `refs/pull/${id}/merge`)); + } + await provisionRunner(scheduler, source.jobId, source.runnerName, 12_501); + const clock = await assignmentClock(scheduler); + let probes = 0; + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async (jobId) => { + probes += 1; + clock.time += 1_000; + return probes > 8 && jobId === source.jobId + ? { status: "in_progress", runner_id: 12_501, runner_name: source.runnerName } + : { status: "queued", runner_id: null, runner_name: null }; + }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + expect(probes).toBe(8); + clock.time += 10_000; + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(probes).toBe(16); + }); + + it("does not spend candidate probes on stopped or releasing jobs", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("reconcile-terminal-candidates"); + const source = job("1300", "cf-standard-3-job-1300", "refs/pull/1300/merge"); + await scheduler.submit(source); + await scheduler.submit(job("1301", "cf-standard-3-job-1301", "refs/pull/1301/merge")); + await scheduler.submit(job("1302", "cf-standard-3-job-1302", "refs/pull/1302/merge")); + await provisionRunner(scheduler, source.jobId, source.runnerName, 13_001); + await assignmentClock(scheduler); + const probes: string[] = []; + await runInDurableObject(scheduler, async (instance, state) => { + state.storage.sql.exec("UPDATE scheduler_jobs SET status = 'releasing' WHERE job_id = '1301'"); + state.storage.sql.exec("UPDATE scheduler_jobs SET status = 'stopped-awaiting-completion' WHERE job_id = '1302'"); + instance.jobDetailOverride = async (jobId) => { + probes.push(jobId); + return { status: "in_progress", runner_id: 13_001, runner_name: source.runnerName }; + }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(probes).toEqual([source.jobId]); + }); + + it("does not resurrect a job completed while GitHub was being queried", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("reconcile-completed-during-fetch"); + const source = job("1400", "cf-standard-3-job-1400", "refs/pull/1400/merge"); + await scheduler.submit(source); + await provisionRunner(scheduler, source.jobId, source.runnerName, 14_001); + await assignmentClock(scheduler); + await runInDurableObject(scheduler, async (instance) => { + instance.jobDetailOverride = async () => { + await instance.workflowJobCompleted(source.jobId); + return { status: "in_progress", runner_id: 14_001, runner_name: source.runnerName }; + }; + }); + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc")).resolves.toBeUndefined(); + expect((await scheduler.status()).jobs.find((entry) => entry.jobId === source.jobId)?.status).toBe("releasing"); + }); + + it("uses the actual workflow run despite a larger repository backlog and persists resource attribution", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("workflow-run-reconcile"); + const source = { ...job("1500", "cf-standard-3-job-1500", "refs/pull/1500/merge"), installationId: undefined }; + const actual = job("1499", "cf-standard-3-job-1499", "refs/pull/1499/merge"); + await scheduler.submit(source); + await scheduler.submit(actual); + for (let id = 1501; id <= 1540; id += 1) { + // eslint-disable-next-line no-await-in-loop -- populate unrelated jobs which must not consume API probes. + await scheduler.submit(job(String(id), `cf-standard-3-job-${id}`, `refs/pull/${id}/merge`)); + } + await provisionRunner(scheduler, source.jobId, source.runnerName, 15_001); + const fetch = vi.fn().mockImplementation(async () => + Response.json({ + total_count: 1, + jobs: [{ id: 1499, status: "in_progress", runner_id: 15_001, runner_name: source.runnerName }], + }), + ); + await assignmentClock(scheduler, fetch); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc", "12345")).resolves.toMatchObject({ + jobId: actual.jobId, + cacheScope: { scope: actual.cacheScope?.scope }, + }); + expect(fetch).toHaveBeenCalledExactlyOnceWith( + "https://api.github.com/repos/biw/runner-poc/actions/runs/12345/jobs?filter=latest&per_page=100&page=1", + expect.objectContaining({ headers: expect.objectContaining({ Authorization: "Bearer assignment-test-token" }) }), + ); + await expect( + env.RESOURCE_METRICS.prepare("SELECT job_id, repository FROM resource_trace_assignments WHERE runner_name = ?") + .bind(source.runnerName) + .first(), + ).resolves.toEqual({ job_id: actual.jobId, repository: "biw/runner-poc" }); + }); + + it("resumes workflow-run pagination after the sweep deadline", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("workflow-run-pagination"); + const source = { ...job("1600", "cf-standard-3-job-1600", "refs/pull/1600/merge"), installationId: undefined }; + await scheduler.submit(source); + await provisionRunner(scheduler, source.jobId, source.runnerName, 16_001); + const clock = { time: 1_700_000_000_000 }; + const fetch = vi + .fn() + .mockImplementationOnce(async () => { + clock.time += 8_000; + return Response.json({ + total_count: 101, + jobs: Array.from({ length: 100 }, (_value, index) => ({ + id: 10_000 + index, + status: "queued", + runner_id: 0, + runner_name: null, + })), + }); + }) + .mockImplementationOnce(async () => + Response.json({ + total_count: 101, + jobs: [{ id: 1600, status: "in_progress", runner_id: 16_001, runner_name: source.runnerName }], + }), + ); + await runInDurableObject(scheduler, async (instance) => { + instance.assignmentDependencies = { fetch, now: () => clock.time }; + }); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc", "12346")).resolves.toBeUndefined(); + expect(fetch).toHaveBeenCalledTimes(1); + clock.time += 10_000; + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc", "12346")).resolves.toMatchObject({ + jobId: source.jobId, + }); + expect(fetch.mock.calls[1]?.[0]).toContain("page=2"); + }); + + it.each(["repository", "profile", "terminal", "unknown", "runner-id"])( + "rejects a workflow-run assignment with an invalid %s boundary", + async (boundary) => { + const scheduler = env.RUNNER_SCHEDULER.getByName(`workflow-run-boundary-${boundary}`); + const source = { ...job("1650", "cf-standard-3-job-1650", "refs/pull/1650/merge"), installationId: undefined }; + const actual = { + ...job("1651", "cf-standard-3-job-1651", "refs/pull/1651/merge"), + target: boundary === "repository" ? { owner: "biw", repository: "other-repository" } : target, + profile: boundary === "profile" ? RUNNER_PROFILES["standard-2"] : profile, + }; + await scheduler.submit(source); + if (boundary !== "unknown") { + await scheduler.submit(actual); + } + if (boundary === "terminal") { + await scheduler.workflowJobCompleted(actual.jobId); + } + await provisionRunner(scheduler, source.jobId, source.runnerName, 16_501); + const fetch = vi.fn().mockImplementation(async () => + Response.json({ + total_count: 1, + jobs: [ + { + id: 1651, + status: "in_progress", + runner_id: boundary === "runner-id" ? 0 : 16_501, + runner_name: source.runnerName, + }, + ], + }), + ); + await assignmentClock(scheduler, fetch); + + await expect(scheduler.cacheAssignment(source.runnerName, "biw/runner-poc", "12347")).resolves.toBeUndefined(); + expect(fetch).toHaveBeenCalledTimes(1); + await expect( + env.RESOURCE_METRICS.prepare("SELECT job_id FROM resource_trace_assignments WHERE runner_name = ?") + .bind(source.runnerName) + .first(), + ).resolves.toBeNull(); + }, + ); + + it.each(["legacy", "scoped"] as const)("releases a failed %s provisioning attempt", async (protocol) => { + const scheduler = env.RUNNER_SCHEDULER.getByName(`compatible-failure-${protocol}`); + const source = job("1700", "cf-standard-3-job-1700", "refs/pull/1700/merge"); + await scheduler.submit(source); + await runInDurableObject(scheduler, async (_instance, state) => { + state.storage.sql.exec("UPDATE scheduler_slots SET applied_max_instances = desired_max_instances"); + }); + const claim = await (protocol === "scoped" + ? scheduler.claimProvisioning(source.jobId, source.runnerName) + : scheduler.claimProvisioning(source.jobId)); + expect(claim).toMatchObject({ kind: "provision" }); + if (protocol === "scoped") { + await scheduler.provisioningFailed(source.jobId, "container startup failed", source.runnerName); + } else { + await scheduler.provisioningFailed(source.jobId, "container startup failed"); + } + const status = await scheduler.status(); + expect(status.jobs.find((entry) => entry.jobId === source.jobId)?.status).toBe("failed"); + expect(status.slots.find((entry) => entry.slotId === "preset:standard-3")?.reservedCount).toBe(0); + }); + + it("accepts an old two-argument failure for an already provisioning legacy retry", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("legacy-inflight-failure"); + const source = job("1710", "cf-standard-3-job-1710", "refs/pull/1710/merge"); + await scheduler.submit(source); + await runInDurableObject(scheduler, async (_instance, state) => { + // Model a persisted attempt created before workflows sent runner names. + state.storage.sql.exec( + "UPDATE scheduler_jobs SET status = 'provisioning', runner_name = ?, runner_attempt = 2 WHERE job_id = ?", + `${source.runnerName}-r2`, + source.jobId, + ); + }); + + await scheduler.provisioningFailed(source.jobId, "legacy retry failed"); + const status = await scheduler.status(); + expect(status.jobs.find((entry) => entry.jobId === source.jobId)?.status).toBe("failed"); + expect(status.slots.find((entry) => entry.slotId === "preset:standard-3")?.reservedCount).toBe(0); + }); + + it("does not let an unscoped old workflow fail a replacement owned by a new workflow", async () => { + const scheduler = env.RUNNER_SCHEDULER.getByName("legacy-stale-failure"); + const source = job("1800", "cf-standard-3-job-1800", "refs/pull/1800/merge"); + const replacementName = `${source.runnerName}-r2`; + await scheduler.submit(source); + await runInDurableObject(scheduler, async (_instance, state) => { + state.storage.sql.exec("UPDATE scheduler_slots SET applied_max_instances = desired_max_instances"); + state.storage.sql.exec( + "UPDATE scheduler_jobs SET runner_name = ?, runner_attempt = 2 WHERE job_id = ?", + replacementName, + source.jobId, + ); + }); + await expect(scheduler.claimProvisioning(source.jobId, replacementName)).resolves.toMatchObject({ + kind: "provision", + }); + await scheduler.provisioningFailed(source.jobId, "old workflow failed"); + await expect(scheduler.canStart(source.jobId, source.runnerName)).resolves.toBe(false); + await expect(scheduler.canStart(source.jobId, replacementName)).resolves.toBe(true); + await expect(scheduler.claimProvisioning(source.jobId, source.runnerName)).resolves.toEqual({ kind: "cancelled" }); + expect((await scheduler.status()).jobs.find((entry) => entry.jobId === source.jobId)?.status).toBe("provisioning"); + await scheduler.provisioningFailed(source.jobId, "current workflow failed", replacementName); + expect((await scheduler.status()).jobs.find((entry) => entry.jobId === source.jobId)?.status).toBe("failed"); + }); }); diff --git a/tests/durable/runner-provisioning-workflow.test.ts b/tests/durable/runner-provisioning-workflow.test.ts index 58ced5a..f52d244 100644 --- a/tests/durable/runner-provisioning-workflow.test.ts +++ b/tests/durable/runner-provisioning-workflow.test.ts @@ -38,7 +38,7 @@ describe("delayed runner repository authorization", () => { { jobId: "next", runnerName: "cf-standard-3-job-next", workflowId: "runner-next" }, ]; const provisioningFailed = vi - .fn<(jobId: string, reason: string) => Promise<{ admissions: SchedulerAdmission[] }>>() + .fn<(jobId: string, reason: string, runnerName: string) => Promise<{ admissions: SchedulerAdmission[] }>>() .mockResolvedValue({ admissions }); const dependencies: EligibilityReleaseDependencies = { authorize: vi.fn().mockResolvedValue({ @@ -58,14 +58,18 @@ describe("delayed runner repository authorization", () => { target: plan.target, installationId: plan.installationId, }); - expect(provisioningFailed).toHaveBeenCalledExactlyOnceWith(plan.jobId, `Repository visibility is ${visibility}`); + expect(provisioningFailed).toHaveBeenCalledExactlyOnceWith( + plan.jobId, + `Repository visibility is ${visibility}`, + plan.runnerName, + ); expect(dependencies.startProvisioning).toHaveBeenCalledExactlyOnceWith(testEnvironment, admissions); }, ); it("keeps a private reservation without releasing or starting another workflow", async () => { const provisioningFailed = - vi.fn<(jobId: string, reason: string) => Promise<{ admissions: SchedulerAdmission[] }>>(); + vi.fn<(jobId: string, reason: string, runnerName: string) => Promise<{ admissions: SchedulerAdmission[] }>>(); const dependencies: EligibilityReleaseDependencies = { authorize: vi.fn().mockResolvedValue({ kind: "private" }), startProvisioning: vi.fn(), diff --git a/tests/github-job-assignment.test.ts b/tests/github-job-assignment.test.ts new file mode 100644 index 0000000..3c5c79f --- /dev/null +++ b/tests/github-job-assignment.test.ts @@ -0,0 +1,121 @@ +import { generateKeyPairSync } from "node:crypto"; + +import { afterEach, describe, expect, it, vi } from "vite-plus/test"; + +import { GitHubJobAssignmentClient } from "../src/github-job-assignment"; + +const target = { owner: "biw", repository: "runner-poc" }; +const legacyEnvironment = { + LEGACY_GITHUB_OWNER: target.owner, + LEGACY_GITHUB_REPOSITORY: target.repository, + GITHUB_RUNNER_TOKEN: "test-token", +}; +const privateKey = generateKeyPairSync("rsa", { modulusLength: 2048 }) + .privateKey.export({ + type: "pkcs8", + format: "pem", + }) + .toString(); +const appEnvironment = { GITHUB_APP_ID: "123", GITHUB_APP_PRIVATE_KEY: privateKey }; + +afterEach(() => vi.restoreAllMocks()); + +describe("GitHub assignment HTTP lookups", () => { + it("mints one installation token for multiple candidate requests and uses GitHub's real job endpoint", async () => { + const fetch = vi + .fn() + .mockResolvedValueOnce(Response.json({ token: "installation-token" })) + .mockResolvedValueOnce(Response.json({ status: "queued", runner_id: 0, runner_name: null })) + .mockResolvedValueOnce(Response.json({ status: "in_progress", runner_id: 42, runner_name: "cf-runner" })); + const client = new GitHubJobAssignmentClient(appEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + const signal = new AbortController().signal; + + await expect(client.jobDetail("100", 42, signal)).resolves.toEqual({ + status: "queued", + runner_id: 0, + runner_name: null, + }); + await expect(client.jobDetail("101", 42, signal)).resolves.toEqual({ + status: "in_progress", + runner_id: 42, + runner_name: "cf-runner", + }); + expect(fetch).toHaveBeenCalledTimes(3); + expect(fetch.mock.calls[0]?.[0]).toBe("https://api.github.com/app/installations/42/access_tokens"); + expect(fetch.mock.calls[1]?.[0]).toBe("https://api.github.com/repos/biw/runner-poc/actions/jobs/100"); + expect(fetch.mock.calls[2]?.[0]).toBe("https://api.github.com/repos/biw/runner-poc/actions/jobs/101"); + expect(fetch.mock.calls[1]?.[1]?.headers).toMatchObject({ Authorization: "Bearer installation-token" }); + }); + + it("lists the latest workflow attempt with explicit pagination and parses runner identities", async () => { + const page = { + total_count: 101, + jobs: [{ id: 100, status: "in_progress", runner_id: 42, runner_name: "cf-runner" }], + }; + const fetch = vi.fn().mockResolvedValue(Response.json(page)); + const client = new GitHubJobAssignmentClient(legacyEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + + await expect(client.workflowJobs("12345", 2, null, new AbortController().signal)).resolves.toEqual(page); + expect(fetch).toHaveBeenCalledExactlyOnceWith( + "https://api.github.com/repos/biw/runner-poc/actions/runs/12345/jobs?filter=latest&per_page=100&page=2", + expect.objectContaining({ headers: expect.objectContaining({ Authorization: "Bearer test-token" }) }), + ); + }); + + it.each([ + ["HTTP failure", new Response("rate limited", { status: 429 })], + ["invalid JSON", new Response("{unfinished")], + ["invalid runner identity", Response.json({ status: "in_progress", runner_id: "42", runner_name: "cf-runner" })], + ])("fails closed and records a %s", async (_scenario, response) => { + const error = vi.spyOn(console, "error").mockImplementation(() => undefined); + const fetch = vi.fn().mockResolvedValue(response); + const client = new GitHubJobAssignmentClient(legacyEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + + await expect(client.jobDetail("100", null, new AbortController().signal)).resolves.toBeUndefined(); + expect(error).toHaveBeenCalled(); + }); + + it("rejects a malformed workflow-jobs response", async () => { + const error = vi.spyOn(console, "error").mockImplementation(() => undefined); + const fetch = vi.fn().mockResolvedValue(Response.json({ total_count: 1, jobs: [{}] })); + const client = new GitHubJobAssignmentClient(legacyEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + + await expect(client.workflowJobs("12345", 1, null, new AbortController().signal)).resolves.toBeUndefined(); + expect(error).toHaveBeenCalledWith("GitHub assignment lookup returned invalid workflow jobs", { + runId: "12345", + page: 1, + }); + }); + + it("aborts token issuance at the sweep deadline without starting a job request", async () => { + vi.spyOn(console, "error").mockImplementation(() => undefined); + const controller = new AbortController(); + const fetch = vi.fn( + async (_input, init) => + new Promise((_resolve, reject) => { + if (init?.signal?.aborted) { + reject(init.signal.reason); + } else { + init?.signal?.addEventListener("abort", () => reject(init.signal?.reason), { once: true }); + } + }), + ); + const client = new GitHubJobAssignmentClient(appEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + const lookup = client.jobDetail("100", 42, controller.signal); + await vi.waitFor(() => expect(fetch).toHaveBeenCalledTimes(1)); + controller.abort(new DOMException("Sweep deadline exceeded", "TimeoutError")); + + await expect(lookup).resolves.toBeUndefined(); + expect(fetch).toHaveBeenCalledTimes(1); + expect(fetch.mock.calls[0]?.[1]?.signal).toBe(controller.signal); + }); + + it("does not perform network requests after the sweep deadline", async () => { + const fetch = vi.fn(); + const client = new GitHubJobAssignmentClient(legacyEnvironment, target, { fetch, now: () => 1_700_000_000_000 }); + const signal = AbortSignal.abort(); + await expect(client.jobDetail("100", null, signal)).resolves.toBeUndefined(); + await expect(client.workflowJobs("12345", 1, null, signal)).resolves.toBeUndefined(); + expect(fetch).not.toHaveBeenCalled(); + }); +}); diff --git a/tests/job-started-hook.test.ts b/tests/job-started-hook.test.ts index 78b6407..b40ccd8 100644 --- a/tests/job-started-hook.test.ts +++ b/tests/job-started-hook.test.ts @@ -43,6 +43,31 @@ async function runHook(configurationPath: string | undefined, environment: Recor } describe("runner cache assignment hook", () => { + it("includes the executing workflow run when polling for an assignment", async () => { + const requests: string[] = []; + const server = createServer((request, response) => { + requests.push(request.url ?? ""); + response.writeHead(200).end(); + }); + const port = await listen(server); + const directory = await mkdtemp(join(tmpdir(), "runner-job-hook-test-")); + const configurationPath = join(directory, "cache-assignment"); + await writeFile(configurationPath, `http://127.0.0.1:${port}/v1/runner-cache\nBearer runner-capability\n`, { + mode: 0o600, + }); + try { + await expect(runHook(configurationPath, { GITHUB_RUN_ID: "12345" })).resolves.toEqual({ + code: 0, + stdout: "", + stderr: "", + }); + expect(requests).toEqual(["/v1/runner-cache-v2/assignment?run_id=12345"]); + } finally { + await close(server); + await rm(directory, { recursive: true, force: true }); + } + }); + it("waits for the GitHub assignment webhook before a cache-enabled job begins", async () => { let attempts = 0; const server = createServer((_request, response) => { diff --git a/tests/runner-cache.test.ts b/tests/runner-cache.test.ts index fc5864f..cea4f3f 100644 --- a/tests/runner-cache.test.ts +++ b/tests/runner-cache.test.ts @@ -524,6 +524,41 @@ describe("runner R2 cache", () => { await expect(assigned.json()).resolves.toEqual({ ok: true }); }); + it("passes the workflow run hint only after authenticating the runner capability", async () => { + const environment = actionCacheEnvironment(actionCacheBucket()); + const cacheAssignment = vi + .fn<(runnerName: string, repository: string, runId?: string) => Promise>() + .mockResolvedValue({ jobId: "assigned-job", cacheScope: { scope: "refs/pull/2/merge", writeAllowed: true } }); + const endpoint = "https://runner.example/v1/runner-cache-v2/assignment?run_id=12345"; + const unauthorized = await handleRunnerCacheV2Request(new Request(endpoint), environment, { cacheAssignment }); + expect(unauthorized.status).toBe(401); + expect(cacheAssignment).not.toHaveBeenCalled(); + const assigned = await handleRunnerCacheV2Request( + new Request(endpoint, { + headers: { Authorization: await authorization() }, + }), + environment, + { cacheAssignment }, + ); + expect(assigned.status).toBe(200); + expect(cacheAssignment).toHaveBeenCalledExactlyOnceWith("cf-standard-3-job-42", "biw/example", "12345"); + }); + + it.each(["0", "-1", "1/path", "not-a-run"])("rejects the malformed workflow run hint %s", async (runId) => { + const environment = actionCacheEnvironment(actionCacheBucket()); + const cacheAssignment = + vi.fn<(runnerName: string, repository: string, runId?: string) => Promise>(); + const response = await handleRunnerCacheV2Request( + new Request(`https://runner.example/v1/runner-cache-v2/assignment?run_id=${encodeURIComponent(runId)}`, { + headers: { Authorization: await authorization() }, + }), + environment, + { cacheAssignment }, + ); + expect(response.status).toBe(400); + expect(cacheAssignment).not.toHaveBeenCalled(); + }); + it("implements CacheService v2 lookups and direct archive uploads in R2", async () => { const bucket = actionCacheBucket(); const environment = actionCacheEnvironment(bucket); diff --git a/vite.config.ts b/vite.config.ts index e5eaa38..e60762b 100644 --- a/vite.config.ts +++ b/vite.config.ts @@ -87,6 +87,9 @@ export default defineConfig({ miniflare: { bindings: { RUNNER_CACHE_MAX_BYTES: "20", + LEGACY_GITHUB_OWNER: "biw", + LEGACY_GITHUB_REPOSITORY: "runner-poc", + GITHUB_RUNNER_TOKEN: "assignment-test-token", }, }, wrangler: { configPath: "./wrangler.jsonc" },