From 116ee167755f6460b6edfdc2684b1cd2e606fd91 Mon Sep 17 00:00:00 2001 From: AgentRun Codex Date: Mon, 20 Jul 2026 09:33:56 +0200 Subject: [PATCH] fix: expose runner source fetch progress --- deploy/runtime/boot/agentrun-boot-report.mjs | 29 +++ deploy/runtime/boot/agentrun-runner-boot.sh | 65 +++++ deploy/runtime/boot/agentrun-runner.sh | 2 +- src/mgr/runner-boot-observation.ts | 245 ++++++++++++++++++ src/mgr/server.ts | 10 + src/selftest/cases/24-runner-startup-retry.ts | 35 +++ 6 files changed, 385 insertions(+), 1 deletion(-) create mode 100644 deploy/runtime/boot/agentrun-boot-report.mjs create mode 100644 deploy/runtime/boot/agentrun-runner-boot.sh create mode 100644 src/mgr/runner-boot-observation.ts diff --git a/deploy/runtime/boot/agentrun-boot-report.mjs b/deploy/runtime/boot/agentrun-boot-report.mjs new file mode 100644 index 0000000..e3e1847 --- /dev/null +++ b/deploy/runtime/boot/agentrun-boot-report.mjs @@ -0,0 +1,29 @@ +const managerUrl = process.env.AGENTRUN_MGR_URL; +const runnerJobId = process.env.AGENTRUN_RUNNER_JOB_ID; +const phase = process.argv[2]; +const attempt = Number(process.argv[3]); +if (!managerUrl || !runnerJobId || !phase || !Number.isSafeInteger(attempt) || attempt < 1) process.exit(2); + +const body = { + phase, + attempt, + runId: process.env.AGENTRUN_RUN_ID, + commandId: process.env.AGENTRUN_COMMAND_ID, + attemptId: process.env.AGENTRUN_ATTEMPT_ID, + runnerId: process.env.AGENTRUN_RUNNER_ID, +}; +if (process.argv[4]) body.code = process.argv[4]; +if (process.argv[5]) body.summary = process.argv[5]; +const headers = { "content-type": "application/json" }; +if (process.env.AGENTRUN_API_KEY) headers.authorization = `Bearer ${process.env.AGENTRUN_API_KEY}`; +const response = await fetch(new URL(`/api/v1/runner-jobs/${encodeURIComponent(runnerJobId)}/boot-observations`, managerUrl), { + method: "POST", + headers, + body: JSON.stringify(body), +}); +const envelope = await response.json(); +if (!response.ok || envelope?.ok !== true) { + process.stderr.write(`${JSON.stringify({ ok: false, event: "agentrun-boot-report", status: response.status, failureKind: envelope?.failureKind ?? "infra-failed", message: envelope?.message ?? "manager boot observation failed" })}\n`); + process.exit(1); +} +process.stdout.write(`${JSON.stringify(envelope.data)}\n`); diff --git a/deploy/runtime/boot/agentrun-runner-boot.sh b/deploy/runtime/boot/agentrun-runner-boot.sh new file mode 100644 index 0000000..7f0dd8c --- /dev/null +++ b/deploy/runtime/boot/agentrun-runner-boot.sh @@ -0,0 +1,65 @@ +#!/bin/sh +set -eu + +entrypoint="${1:-src/runner/main.ts}" +repo_url="${AGENTRUN_BOOT_REPO_URL:-http://git-mirror-http.devops-infra.svc.cluster.local/pikasTech/agentrun.git}" +commit="${AGENTRUN_BOOT_COMMIT:-${AGENTRUN_SOURCE_COMMIT:-}}" +app_root="${AGENTRUN_APP_ROOT:-/workspace/agentrun}" +reporter=/opt/agentrun/deploy/runtime/boot/agentrun-boot-report.mjs + +if [ -z "$commit" ] || ! printf '%s' "$commit" | grep -Eq '^[0-9a-f]{40}$'; then + printf '{"ok":false,"event":"agentrun-boot","failureKind":"source-commit-invalid","message":"AGENTRUN_BOOT_COMMIT must be a full git SHA"}\n' >&2 + exit 64 +fi + +mkdir -p "$(dirname "$app_root")" +rm -rf "$app_root" +mkdir -p "$app_root" +cd "$app_root" +git init -q +git remote add origin "$repo_url" + +attempt=1 +fetch_log=/tmp/agentrun-boot-fetch.log +while :; do + start_response=$(node "$reporter" source-fetch-started "$attempt" 2>/tmp/agentrun-boot-report.log || true) + timeout_ms=$(printf '%s' "$start_response" | node -e 'let s="";process.stdin.on("data",d=>s+=d).on("end",()=>{try{const n=JSON.parse(s).attemptTimeoutMs;process.stdout.write(Number.isFinite(n)?String(n):"")}catch{}})') + timeout_seconds="" + if [ -n "$timeout_ms" ]; then timeout_seconds=$(( (timeout_ms + 999) / 1000 )); fi + + if { [ -n "$timeout_seconds" ] && timeout "${timeout_seconds}s" git fetch --depth=1 origin "$commit" >"$fetch_log" 2>&1; } || { [ -z "$timeout_seconds" ] && git fetch --depth=1 origin "$commit" >"$fetch_log" 2>&1; }; then + node "$reporter" source-fetch-succeeded "$attempt" >/tmp/agentrun-boot-report.log 2>&1 || true + break + fi + + message=$(tail -n 20 "$fetch_log" | tr '\n' ' ' | sed 's/[[:space:]]\+/ /g') + failure_kind=git-mirror-fetch-failed + if printf '%s' "$message" | grep -Eiq 'not our ref|unadvertised|Server does not allow request'; then + failure_kind=git-mirror-exact-commit-unavailable + elif printf '%s' "$message" | grep -Eiq 'Could not resolve host|Connection timed out|Failed to connect|timed out'; then + failure_kind=git-mirror-network-failed + elif printf '%s' "$message" | grep -Eiq 'Authentication failed|Permission denied|repository.*not found'; then + failure_kind=git-mirror-auth-failed + fi + failure_response=$(node "$reporter" source-fetch-failed "$attempt" "$failure_kind" 2>/tmp/agentrun-boot-report.log || true) + action=$(printf '%s' "$failure_response" | node -e 'let s="";process.stdin.on("data",d=>s+=d).on("end",()=>{try{process.stdout.write(JSON.parse(s).action||"")}catch{}})') + if [ "$action" != "retry" ]; then + printf '{"ok":false,"event":"agentrun-boot","failureKind":"%s","repoUrl":"%s","commit":"%s","message":%s}\n' "$failure_kind" "$repo_url" "$commit" "$(node -e 'console.log(JSON.stringify(process.argv[1] || ""))' "$message")" >&2 + exit 65 + fi + backoff_ms=$(printf '%s' "$failure_response" | node -e 'let s="";process.stdin.on("data",d=>s+=d).on("end",()=>{try{const n=JSON.parse(s).backoffMs;process.stdout.write(Number.isFinite(n)?String(n):"0")}catch{process.stdout.write("0")}})') + sleep_seconds=$(node -e 'process.stdout.write(String(Math.max(0, Number(process.argv[1]) || 0) / 1000))' "$backoff_ms") + sleep "$sleep_seconds" + attempt=$((attempt + 1)) +done + +git checkout -q --detach "$commit" +actual=$(git rev-parse HEAD) +if [ "$actual" != "$commit" ]; then + printf '{"ok":false,"event":"agentrun-boot","failureKind":"source-commit-mismatch","expected":"%s","actual":"%s"}\n' "$commit" "$actual" >&2 + exit 66 +fi + +ln -sfn /opt/agentrun/node_modules "$app_root/node_modules" +printf '{"ok":true,"event":"agentrun-boot","repoUrl":"%s","commit":"%s","entrypoint":"%s","nodeModules":"/opt/agentrun/node_modules"}\n' "$repo_url" "$commit" "$entrypoint" +exec bun "$entrypoint" diff --git a/deploy/runtime/boot/agentrun-runner.sh b/deploy/runtime/boot/agentrun-runner.sh index 564f783..3bec89f 100644 --- a/deploy/runtime/boot/agentrun-runner.sh +++ b/deploy/runtime/boot/agentrun-runner.sh @@ -1,3 +1,3 @@ #!/bin/sh set -eu -exec /opt/agentrun/deploy/runtime/boot/agentrun-boot.sh src/runner/main.ts +exec /opt/agentrun/deploy/runtime/boot/agentrun-runner-boot.sh src/runner/main.ts diff --git a/src/mgr/runner-boot-observation.ts b/src/mgr/runner-boot-observation.ts new file mode 100644 index 0000000..f147f0d --- /dev/null +++ b/src/mgr/runner-boot-observation.ts @@ -0,0 +1,245 @@ +import { AgentRunError } from "../common/errors.js"; +import type { JsonRecord, RunnerJobRecord } from "../common/types.js"; +import { agentRunBusinessTraceId } from "../common/otel-trace.js"; +import { nowIso, stableHash } from "../common/validation.js"; +import type { AgentRunStore } from "./store.js"; +import { isTerminalCommandState, isTerminalRunStatus } from "./store.js"; +import type { RunnerStartupRecoveryOptions } from "./runner-reconciler.js"; + +type BootPhase = "source-fetch-started" | "source-fetch-succeeded" | "source-fetch-failed"; + +interface BootFailure { + code: string; + summary: string; + retryable: boolean; +} + +export async function recordRunnerBootObservation(input: { + store: AgentRunStore; + runnerJobId: string; + body: JsonRecord; + policy?: RunnerStartupRecoveryOptions; +}): Promise { + const phase = bootPhase(input.body.phase); + const runId = requiredString(input.body, "runId"); + const commandId = requiredString(input.body, "commandId"); + const attemptId = requiredString(input.body, "attemptId"); + const runnerId = requiredString(input.body, "runnerId"); + const attempt = positiveInteger(input.body.attempt, "attempt"); + const job = await requireRunnerJob(input.store, { runnerJobId: input.runnerJobId, runId, commandId, attemptId, runnerId }); + const policy = input.policy ?? singleAttemptPolicy(); + const previous = recordValue(job.result.bootRecovery) ?? {}; + const firstObservedAt = stringValue(previous.firstObservedAt) ?? nowIso(); + const observedAt = nowIso(); + const maxAttempts = policy.maxAttempts; + const run = await input.store.getRun(runId); + const command = await input.store.getCommand(commandId); + const common: JsonRecord = { + failureDomain: "infrastructure", + component: "git-mirror", + attempt, + maxAttempts, + firstObservedAt, + observedAt, + runId, + commandId, + runnerJobId: job.id, + attemptId, + runnerId, + traceId: agentRunBusinessTraceId(run, command), + policyFallback: input.policy ? false : true, + valuesPrinted: false, + }; + + if (phase === "source-fetch-started") { + if (attempt > 1) { + const previousFailure = bootFailure(recordValue(previous.failure)?.code, recordValue(previous.failure)?.summary); + await appendTransition(input.store, job, "runner-startup-retry-started", { + ...common, + ...failurePayload(previousFailure), + retryPhase: "retryStarted", + backoffMs: 0, + nextRetryAt: null, + }); + } + await appendTransition(input.store, job, "runner-source-fetch-started", { + ...common, + code: "git-mirror-fetch-in-progress", + summary: "正在从 Git mirror 获取 Runner 源码", + retryable: true, + retryPhase: attempt > 1 ? "retryStarted" : "initialAttempt", + backoffMs: 0, + nextRetryAt: null, + }); + const next = await input.store.updateRunnerJobResult(job.id, { + bootRecovery: { state: "fetching", attempt, maxAttempts, firstObservedAt, observedAt, valuesPrinted: false }, + }); + return { + action: "fetch", + attempt, + maxAttempts, + attemptTimeoutMs: attemptTimeoutMs(policy, firstObservedAt, observedAt, attempt), + policyFallback: input.policy ? false : true, + runnerJobUpdatedAt: next.updatedAt, + valuesPrinted: false, + }; + } + + if (phase === "source-fetch-succeeded") { + await appendTransition(input.store, job, "runner-source-fetch-completed", { + ...common, + code: "git-mirror-fetch-completed", + summary: "Git mirror 源码获取完成", + retryable: false, + retryPhase: attempt > 1 ? "retryRecovered" : "initialAttemptCompleted", + backoffMs: 0, + nextRetryAt: null, + }); + if (attempt > 1) { + const previousFailure = bootFailure(recordValue(previous.failure)?.code, recordValue(previous.failure)?.summary); + await appendTransition(input.store, job, "runner-startup-retry-recovered", { + ...common, + ...failurePayload(previousFailure), + retryPhase: "retryRecovered", + backoffMs: 0, + nextRetryAt: null, + }); + } + const next = await input.store.updateRunnerJobResult(job.id, { + bootRecovery: { state: "completed", attempt, maxAttempts, firstObservedAt, observedAt, valuesPrinted: false }, + }); + return { action: "continue", attempt, maxAttempts, runnerJobUpdatedAt: next.updatedAt, valuesPrinted: false }; + } + + const failure = bootFailure(input.body.code, input.body.summary ?? input.body.message); + const deadlineAt = new Date(Date.parse(firstObservedAt) + policy.totalDeadlineMs).toISOString(); + const deadlineExceeded = Date.parse(observedAt) >= Date.parse(deadlineAt); + const retry = failure.retryable && !deadlineExceeded && attempt < maxAttempts; + await appendTransition(input.store, job, "runner-startup-failure-observed", { + ...common, + ...failurePayload(failure), + retryPhase: "failureObserved", + backoffMs: 0, + nextRetryAt: null, + deadlineAt, + willRetry: retry, + }); + if (retry) { + const backoffMs = startupBackoffMs(policy, attempt, job.id); + const nextRetryAt = new Date(Date.parse(observedAt) + backoffMs).toISOString(); + await appendTransition(input.store, job, "runner-startup-retry-scheduled", { + ...common, + ...failurePayload(failure), + retryPhase: "retryScheduled", + backoffMs, + nextRetryAt, + deadlineAt, + willRetry: true, + }); + const next = await input.store.updateRunnerJobResult(job.id, { + bootRecovery: { state: "scheduled", attempt, maxAttempts, firstObservedAt, observedAt, nextRetryAt, deadlineAt, failure, valuesPrinted: false }, + }); + return { action: "retry", attempt, nextAttempt: attempt + 1, maxAttempts, backoffMs, nextRetryAt, runnerJobUpdatedAt: next.updatedAt, valuesPrinted: false }; + } + + await appendTransition(input.store, job, "runner-startup-retry-exhausted", { + ...common, + ...failurePayload(failure), + retryPhase: "retryExhausted", + backoffMs: 0, + nextRetryAt: null, + deadlineAt, + willRetry: false, + }); + const failureMessage = failure.summary; + if (!isTerminalCommandState(command.state)) { + await input.store.finishCommand(commandId, { terminalStatus: "failed", failureKind: "infra-failed", failureMessage }); + } + const currentRun = await input.store.getRun(runId); + if (!isTerminalRunStatus(currentRun.status)) { + await input.store.finishRun(runId, { terminalStatus: "failed", failureKind: "infra-failed", failureMessage }); + } + const next = await input.store.updateRunnerJobResult(job.id, { + bootRecovery: { state: "exhausted", attempt, maxAttempts, firstObservedAt, observedAt, deadlineAt, failure, valuesPrinted: false }, + }); + return { action: "terminal", attempt, maxAttempts, terminalStatus: "failed", failureKind: "infra-failed", code: failure.code, summary: failure.summary, runnerJobUpdatedAt: next.updatedAt, valuesPrinted: false }; +} + +async function requireRunnerJob(store: AgentRunStore, identity: { runnerJobId: string; runId: string; commandId: string; attemptId: string; runnerId: string }): Promise { + const job = (await store.listRunnerJobs(identity.runId, identity.commandId)).find((item) => item.id === identity.runnerJobId); + if (!job) throw new AgentRunError("schema-invalid", `runner job ${identity.runnerJobId} was not found for boot observation`, { httpStatus: 404 }); + for (const key of ["runId", "commandId", "attemptId", "runnerId"] as const) { + if (job[key] !== identity[key]) throw new AgentRunError("runner-lease-conflict", `runner boot ${key} does not match runner job`, { httpStatus: 409, details: { runnerJobId: job.id, field: key, valuesPrinted: false } }); + } + return job; +} + +async function appendTransition(store: AgentRunStore, job: RunnerJobRecord, phase: string, payload: JsonRecord): Promise { + const transitionKey = stableHash({ runnerJobId: job.id, phase, attempt: payload.attempt, code: payload.code, nextRetryAt: payload.nextRetryAt ?? null }); + const events = await store.listEventsForCommand(job.runId, job.commandId, 2_000); + if (events.some((event) => event.payload.transitionKey === transitionKey)) return; + await store.appendEvent(job.runId, "backend_status", { phase, ...payload, transitionKey }); +} + +function bootFailure(codeValue: unknown, summaryValue: unknown): BootFailure { + const code = typeof codeValue === "string" && codeValue.trim().length > 0 ? codeValue.trim() : "git-mirror-fetch-failed"; + const summary = typeof summaryValue === "string" && summaryValue.trim().length > 0 ? summaryValue.trim().slice(0, 500) : summaryForCode(code); + const retryable = code === "git-mirror-network-failed" || code === "git-mirror-fetch-failed"; + return { code, summary, retryable }; +} + +function summaryForCode(code: string): string { + if (code === "git-mirror-network-failed") return "Git mirror 网络暂时不可达"; + if (code === "git-mirror-auth-failed") return "Git mirror 鉴权失败"; + if (code === "git-mirror-exact-commit-unavailable") return "Git mirror 中不存在指定源码提交"; + if (code === "source-commit-mismatch") return "Runner 获取的源码提交不匹配"; + return "Git mirror 源码获取失败"; +} + +function failurePayload(failure: BootFailure): JsonRecord { + return { code: failure.code, summary: failure.summary, message: failure.summary, failureKind: "infra-failed", retryable: failure.retryable, willRetry: failure.retryable }; +} + +function bootPhase(value: unknown): BootPhase { + if (value === "source-fetch-started" || value === "source-fetch-succeeded" || value === "source-fetch-failed") return value; + throw new AgentRunError("schema-invalid", "runner boot observation phase is not supported", { httpStatus: 400 }); +} + +function requiredString(record: JsonRecord, field: string): string { + const value = record[field]; + if (typeof value !== "string" || value.trim().length === 0) throw new AgentRunError("schema-invalid", `${field} is required`, { httpStatus: 400 }); + return value.trim(); +} + +function positiveInteger(value: unknown, field: string): number { + if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 1) throw new AgentRunError("schema-invalid", `${field} must be a positive integer`, { httpStatus: 400 }); + return value; +} + +function recordValue(value: unknown): JsonRecord | null { + return value && typeof value === "object" && !Array.isArray(value) ? value as JsonRecord : null; +} + +function stringValue(value: unknown): string | null { + return typeof value === "string" && value.trim().length > 0 ? value : null; +} + +function singleAttemptPolicy(): RunnerStartupRecoveryOptions { + return { maxAttempts: 1, initialBackoffMs: 1_000, maxBackoffMs: 1_000, multiplier: 2, totalDeadlineMs: 60_000, jitterRatio: 0 }; +} + +function attemptTimeoutMs(policy: RunnerStartupRecoveryOptions, firstObservedAt: string, observedAt: string, attempt: number): number { + const elapsed = Math.max(0, Date.parse(observedAt) - Date.parse(firstObservedAt)); + const remaining = Math.max(1_000, policy.totalDeadlineMs - elapsed); + const attemptsRemaining = Math.max(1, policy.maxAttempts - attempt + 1); + return Math.max(1_000, Math.floor(remaining / attemptsRemaining)); +} + +function startupBackoffMs(policy: RunnerStartupRecoveryOptions, attempt: number, seed: string): number { + const exponential = Math.min(policy.maxBackoffMs, Math.round(policy.initialBackoffMs * (policy.multiplier ** Math.max(0, attempt - 1)))); + if (policy.jitterRatio <= 0) return exponential; + const hash = Number.parseInt(stableHash({ seed, attempt }).slice(0, 8), 16); + const unit = Number.isFinite(hash) ? hash / 0xffffffff : 0.5; + const factor = 1 - policy.jitterRatio + (2 * policy.jitterRatio * unit); + return Math.max(1, Math.min(policy.maxBackoffMs, Math.round(exponential * factor))); +} diff --git a/src/mgr/server.ts b/src/mgr/server.ts index 909a49b..35065b6 100644 --- a/src/mgr/server.ts +++ b/src/mgr/server.ts @@ -26,6 +26,7 @@ import { listToolCredentials, setGithubSshToolCredential, showToolCredential } f import { aipodSpecFromInput, applyAipodSpec, deleteAipodSpec, listAipodSpecs, renderAipodSpecByName, showAipodSpec } from "../common/aipod-specs.js"; import { staticWorkReadyCapabilitySummary } from "../common/work-ready.js"; import { agentRunBusinessTraceId, agentRunDiagnosticOtelTraceContext, emitAgentRunOtelSpan, type AgentRunDiagnosticOtelTraceContext } from "../common/otel-trace.js"; +import { recordRunnerBootObservation } from "./runner-boot-observation.js"; function pvcOptions(defaults: { kubectlCommand?: string } | undefined): SessionPvcOptions { return defaults?.kubectlCommand ? { kubectlCommand: defaults.kubectlCommand } : {}; @@ -1084,6 +1085,15 @@ async function route({ method, url, body, signal, store, sourceCommit, authSumma if (method === "POST" && path === "/api/v1/reconciler/kafka-event-outbox") { return await relayKafkaEventOutboxOnce({ store, options: kafkaOutboxRelayOptions }) as JsonValue; } + const runnerBootObservationMatch = path.match(/^\/api\/v1\/runner-jobs\/([^/]+)\/boot-observations$/u); + if (method === "POST" && runnerBootObservationMatch) { + return await recordRunnerBootObservation({ + store, + runnerJobId: runnerBootObservationMatch[1] ?? "", + body: asRecord(body ?? {}, "runnerBootObservation"), + ...(runnerReconcilerOptions.startupRecovery ? { policy: runnerReconcilerOptions.startupRecovery } : {}), + }) as JsonValue; + } const commandShowMatch = path.match(/^\/api\/v1\/runs\/([^/]+)\/commands\/([^/]+)$/u); if (method === "GET" && commandShowMatch) return await store.getCommand(commandShowMatch[2] ?? "") as unknown as JsonValue; const commandResultMatch = path.match(/^\/api\/v1\/runs\/([^/]+)\/commands\/([^/]+)\/result$/u); diff --git a/src/selftest/cases/24-runner-startup-retry.ts b/src/selftest/cases/24-runner-startup-retry.ts index 1767b0c..256b03f 100644 --- a/src/selftest/cases/24-runner-startup-retry.ts +++ b/src/selftest/cases/24-runner-startup-retry.ts @@ -4,6 +4,7 @@ import path from "node:path"; import type { JsonRecord } from "../../common/types.js"; import { reconcileRunnerJobsOnce, type RunnerStartupRecoveryOptions } from "../../mgr/runner-reconciler.js"; import { MemoryAgentRunStore } from "../../mgr/store.js"; +import { recordRunnerBootObservation } from "../../mgr/runner-boot-observation.js"; import type { SelfTestCase, SelfTestContext } from "../harness.js"; const policy: RunnerStartupRecoveryOptions = { @@ -27,6 +28,7 @@ const selfTest: SelfTestCase = async (context) => { await assertDeterministicFailureDoesNotRetry(context, fixturePath, kubectlCommand); await assertFailedJobWithoutTerminalOutboxFails(context, fixturePath, kubectlCommand); await assertSucceededJobWithoutTerminalOutboxFails(context, fixturePath, kubectlCommand); + await assertRunnerBootObservationLifecycle(); } finally { if (previousFixture === undefined) delete process.env.AGENTRUN_SELFTEST_STARTUP_FIXTURE_PATH; else process.env.AGENTRUN_SELFTEST_STARTUP_FIXTURE_PATH = previousFixture; @@ -40,10 +42,43 @@ const selfTest: SelfTestCase = async (context) => { "failed-job-without-terminal-outbox-fails", "succeeded-job-without-terminal-outbox-fails", "semantic-failure-events-carry-workbench-fields", + "runner-git-mirror-boot-is-visible-before-claim", ], }; }; +async function assertRunnerBootObservationLifecycle(): Promise { + const fact = createRunnerFact("boot-observation"); + const identity = { runId: fact.runId, commandId: fact.commandId, attemptId: "attempt-boot-observation", runnerId: "runner-boot-observation" }; + const started = await recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, phase: "source-fetch-started", attempt: 1 }, policy: { ...policy, maxAttempts: 2, initialBackoffMs: 10 } }); + assert.equal(started.action, "fetch"); + assert.equal(typeof started.attemptTimeoutMs, "number"); + await recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, phase: "source-fetch-started", attempt: 1 }, policy: { ...policy, maxAttempts: 2, initialBackoffMs: 10 } }); + assert.equal(fact.store.listEvents(fact.runId, 0, 100).filter((event) => event.payload.phase === "runner-source-fetch-started").length, 1); + await assert.rejects( + recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, runnerId: "runner-other", phase: "source-fetch-started", attempt: 1 }, policy: { ...policy, maxAttempts: 2 } }), + (error: unknown) => error instanceof Error && error.message.includes("runnerId does not match"), + ); + const failed = await recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, phase: "source-fetch-failed", attempt: 1, code: "git-mirror-network-failed", summary: "Git mirror network unavailable" }, policy: { ...policy, maxAttempts: 2, initialBackoffMs: 10 } }); + assert.equal(failed.action, "retry"); + await recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, phase: "source-fetch-started", attempt: 2 }, policy: { ...policy, maxAttempts: 2, initialBackoffMs: 10 } }); + const succeeded = await recordRunnerBootObservation({ store: fact.store, runnerJobId: fact.runnerJobId, body: { ...identity, phase: "source-fetch-succeeded", attempt: 2 }, policy: { ...policy, maxAttempts: 2, initialBackoffMs: 10 } }); + assert.equal(succeeded.action, "continue"); + const events = fact.store.listEvents(fact.runId, 0, 100); + for (const expected of ["runner-source-fetch-started", "runner-startup-failure-observed", "runner-startup-retry-scheduled", "runner-startup-retry-started", "runner-source-fetch-completed", "runner-startup-retry-recovered"]) phase(events, expected); + assert.equal(fact.store.getCommand(fact.commandId).state, "pending"); + assert.equal(fact.store.getRun(fact.runId).status, "pending"); + + const terminalFact = createRunnerFact("boot-terminal"); + const terminalIdentity = { runId: terminalFact.runId, commandId: terminalFact.commandId, attemptId: "attempt-boot-terminal", runnerId: "runner-boot-terminal" }; + await recordRunnerBootObservation({ store: terminalFact.store, runnerJobId: terminalFact.runnerJobId, body: { ...terminalIdentity, phase: "source-fetch-started", attempt: 1 }, policy: { ...policy, maxAttempts: 3 } }); + const terminal = await recordRunnerBootObservation({ store: terminalFact.store, runnerJobId: terminalFact.runnerJobId, body: { ...terminalIdentity, phase: "source-fetch-failed", attempt: 1, code: "git-mirror-auth-failed" }, policy: { ...policy, maxAttempts: 3 } }); + assert.equal(terminal.action, "terminal"); + assert.equal(terminalFact.store.getCommand(terminalFact.commandId).state, "failed"); + assert.equal(terminalFact.store.getRun(terminalFact.runId).status, "failed"); + assert.equal(terminalFact.store.listEvents(terminalFact.runId, 0, 100).some((event) => event.type === "terminal_status"), true); +} + async function assertTransientRegistryFailureExhausts(context: SelfTestContext, fixturePath: string, kubectlCommand: string): Promise { const fact = createRunnerFact("transient"); await writeFixture(fixturePath, "ImagePullBackOff", "failed to pull image: dial tcp 10.43.249.83:5000: connect: connection refused");