diff --git a/package.json b/package.json index fe3010a..bf6737d 100644 --- a/package.json +++ b/package.json @@ -7,6 +7,7 @@ "check": "tsc --noEmit", "self-test": "bun run src/selftest/run.ts", "self-test:postgres-kafka-order": "bun run src/selftest/integration/postgres-kafka-order.ts", + "self-test:postgres-runner-retention-fence": "bun run src/selftest/integration/postgres-runner-retention-fence.ts", "test": "bun run src/selftest/run.ts", "cli": "bun scripts/agentrun-cli.ts" }, diff --git a/src/mgr/kubernetes-runner-job.ts b/src/mgr/kubernetes-runner-job.ts index facb6ea..ebdf63e 100644 --- a/src/mgr/kubernetes-runner-job.ts +++ b/src/mgr/kubernetes-runner-job.ts @@ -79,6 +79,12 @@ export interface CreateRunnerJobInput extends JsonRecord { } export async function createKubernetesRunnerJob(options: { store: AgentRunStore; runId: string; input: CreateRunnerJobInput; defaults: RunnerJobDefaults; signal?: AbortSignal }): Promise { + if (!options.defaults.retention) return await createKubernetesRunnerJobUnderFence(options); + const namespace = optionalString(options.input.namespace) ?? options.defaults.namespace; + return await options.store.withRunnerCreateFence(namespace, async () => await createKubernetesRunnerJobUnderFence(options)); +} + +async function createKubernetesRunnerJobUnderFence(options: { store: AgentRunStore; runId: string; input: CreateRunnerJobInput; defaults: RunnerJobDefaults; signal?: AbortSignal }): Promise { throwIfAborted(options.signal, "runner dispatch attempt"); const commandId = stringField(options.input, "commandId"); const run = await options.store.getRun(options.runId); diff --git a/src/mgr/postgres-store.ts b/src/mgr/postgres-store.ts index a222087..1c61132 100644 --- a/src/mgr/postgres-store.ts +++ b/src/mgr/postgres-store.ts @@ -5,7 +5,7 @@ import { AgentRunError } from "../common/errors.js"; import { redactJson } from "../common/redaction.js"; import type { BackendProfile, BackendTurnResult, CancelRequestRecord, CancelStage, CancelTargetKind, CommandRecord, CommandState, CreateCommandInput, CreateQueueTaskInput, CreateRunInput, EventType, FailureKind, JsonRecord, JsonValue, KafkaEventOutboxRecord, ListGcExpiredSessionsInput, QueueCommanderSnapshot, QueueReadCursorRecord, QueueStats, QueueTaskListResult, QueueTaskRecord, QueueTaskState, RunEvent, RunnerDispatchCompletion, RunnerDispatchIntentRecord, RunnerJobRecord, RunnerRecord, RunRecord, RunStatus, SessionEventPage, SessionListResult, SessionReadCursorRecord, SessionRecord, SessionRef, SessionStoragePatch, SessionSummary, TerminalStatus, UpsertSessionInput } from "../common/types.js"; import { newId, nowIso, stableHash } from "../common/validation.js"; -import type { AgentRunStore, DurableQueueClaim, EventOutboxConfig, ListQueueTasksInput, ListSessionsInput, SaveRunnerJobInput, SessionEventPageInput, StoreHealth, UpdateQueueTaskAttemptInput } from "./store.js"; +import type { AgentRunStore, DurableQueueClaim, EventOutboxConfig, ListQueueTasksInput, ListSessionsInput, RunnerRetentionFenceFacts, RunnerRetentionFenceInput, SaveRunnerJobInput, SessionEventPageInput, StoreHealth, UpdateQueueTaskAttemptInput } from "./store.js"; import { assertSessionBoundary, buildQueueStats, buildQueueTaskSummary, buildSessionSummary, cancelStagePayload, clampQueueLimit, clampSessionLimit, commandStateFromTerminal, fenceLateEventForCancelledRun, isLeaseExpired, isSessionOutputEvent, isTerminalCommandState, isTerminalQueueTaskState, isTerminalRunStatus, lateWriteRejectedPayload, parseQueueCursor, parseSessionCursor, queueTaskMatchesCommander, queueTaskPayloadHash, queueTaskSort, sessionListFilters, sessionMatchesListState, sessionRefFromRecord, sessionSort, sessionTitleFromCommand, statusFromTerminal, summarizeResourceBundleRef, summarizeSessionRef, titleFromMetadata } from "./store.js"; import { backendCapabilitiesSqlValues, mergeBackendCapability } from "../common/backend-profiles.js"; import { normalizeRunEventPayload, requireEventType, requireExternallyAppendableEventType, userMessagePayloadForCommand } from "../common/events.js"; @@ -458,12 +458,14 @@ export async function createPostgresAgentRunStore(options: PostgresStoreOptions) export class PostgresAgentRunStore implements AgentRunStore { private readonly pool: Pool; + private readonly runnerCreateFencePool: Pool; private readonly eventOutboxConfig: EventOutboxConfig; private migrationReady = false; private appliedMigrationId: string | null = null; constructor(options: PostgresStoreOptions) { this.pool = new Pool({ connectionString: options.connectionString, application_name: "agentrun-mgr-v01", ...(options.poolMax === undefined ? {} : { max: options.poolMax }) }); + this.runnerCreateFencePool = new Pool({ connectionString: options.connectionString, application_name: "agentrun-mgr-runner-capacity-fence", max: 1 }); this.eventOutboxConfig = options.eventOutbox ?? { enabled: false, topic: null, source: null }; } @@ -881,6 +883,50 @@ CREATE TABLE IF NOT EXISTS agentrun_schema_migrations ( }); } + async withRunnerCreateFence(namespace: string, operation: () => Promise): Promise { + const client = await this.runnerCreateFencePool.connect(); + try { + await client.query("BEGIN"); + await client.query("SELECT pg_advisory_xact_lock(hashtext($1), hashtext($2))", ["agentrun-runner-create", namespace]); + const result = await operation(); + await client.query("COMMIT"); + return result; + } catch (error) { + await client.query("ROLLBACK"); + throw error; + } finally { + client.release(); + } + } + + async withRunnerRetentionFence(input: RunnerRetentionFenceInput, operation: (facts: RunnerRetentionFenceFacts) => Promise): Promise { + return this.withTransaction(async (client) => { + const run = await this.requireRunForUpdate(client, input.runId); + const commandRows = await client.query("SELECT * FROM agentrun_commands WHERE run_id = $1 ORDER BY seq ASC FOR UPDATE", [run.id]); + const commands = commandRows.rows.slice(0, 100).map(commandFromRow); + const commandRow = commandRows.rows.find((item) => String(item.id) === input.commandId); + if (!commandRow) throw new AgentRunError("schema-invalid", `command ${input.commandId} does not belong to run ${input.runId}`, { httpStatus: 409 }); + const command = commandFromRow(commandRow); + const eventRows = await client.query("SELECT * FROM agentrun_events WHERE run_id = $1 ORDER BY seq ASC LIMIT 500", [run.id]); + const result = await operation({ run, command, commands, events: eventRows.rows.map(eventFromRow) }); + if (!input.terminalizeStalePending || isTerminalRunStatus(run.status)) return result; + + const at = nowIso(); + const updatedCommands = await client.query("UPDATE agentrun_commands SET state = 'failed', updated_at = $2 WHERE run_id = $1 AND state NOT IN ('completed', 'failed', 'cancelled') RETURNING *", [run.id, at]); + await fenceActiveRunnerDispatchIntents(client, { runId: run.id }, input.reason, at); + const updatedRun = await client.query("UPDATE agentrun_runs SET status = 'failed', terminal_status = 'failed', failure_kind = 'infra-failed', failure_message = $2, claimed_by = NULL, lease_expires_at = NULL, updated_at = $3 WHERE id = $1 RETURNING *", [run.id, input.reason, at]); + await client.query("DELETE FROM agentrun_leases WHERE run_id = $1", [run.id]); + for (const row of updatedCommands.rows) { + const item = commandFromRow(row); + await this.appendEventWithLockedRun(client, run.id, "backend_status", { phase: "command-terminal", commandId: item.id, state: "failed", terminalStatus: "failed", failureKind: "infra-failed", message: input.reason, retentionFence: true }); + } + await this.appendEventWithLockedRun(client, run.id, "terminal_status", { terminalStatus: "failed", failureKind: "infra-failed", message: input.reason, retentionFence: true }); + const next = runFromRow(updatedRun.rows[0]); + await this.touchSessionForRun(client, next, { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: run.id, terminalStatus: "failed", failureKind: "infra-failed", lastActivityAt: at }, { bumpVersion: true, at }); + return result; + }); + } + async claimRun(runId: string, runnerId: string, leaseMs: number): Promise { return this.withTransaction(async (client) => { const run = await this.requireRunForUpdate(client, runId); @@ -1396,7 +1442,7 @@ CREATE TABLE IF NOT EXISTS agentrun_schema_migrations ( } async close(): Promise { - await this.pool.end(); + await Promise.all([this.pool.end(), this.runnerCreateFencePool.end()]); } private async createCancelRequest(client: PoolClient, input: { targetKind: CancelTargetKind; targetId: string; run: RunRecord; command: CommandRecord | null; reason: string; at: string; stage: CancelStage }): Promise { diff --git a/src/mgr/runner-retention.ts b/src/mgr/runner-retention.ts index 9aa280d..411d67f 100644 --- a/src/mgr/runner-retention.ts +++ b/src/mgr/runner-retention.ts @@ -12,6 +12,7 @@ export interface RunnerRetentionOptions { maxRunners: number; cleanupOrder: "oldest-inactive-last-active-first"; activeHeartbeatMaxAgeMs: number; + stalePendingMaxAgeMs: number; matchLabels: Record; jobNamePrefixes: string[]; ageBasedCleanup: { @@ -26,6 +27,7 @@ export interface RunnerRetentionSummary extends JsonRecord { enabled: boolean; namespace: string; maxRunners: number; + stalePendingMaxAgeMs: number; incomingRunnerCount: number; liveRunnerCountBefore: number; liveRunnerCountAfter: number; @@ -44,35 +46,73 @@ export interface RunnerRetentionSummary extends JsonRecord { valuesPrinted: false; } -type CandidateKind = "idle" | "terminal" | "inactive"; +type CandidateKind = "idle" | "terminal" | "inactive" | "stale-pending"; + +interface RunnerPodObservation extends JsonRecord { + name: string; + uid: string | null; + resourceVersion: string | null; + phase: string | null; + createdAt: string | null; + ageMs: number | null; + stale: boolean; + podStartTime: string | null; + ready: boolean; + runnerContainerState: "waiting" | "running" | "terminated" | "unknown"; + runnerContainerLastState: "running" | "terminated" | "none"; + runnerContainerStarted: boolean; + runnerContainerReady: boolean; + runnerContainerRestartCount: number; + waiting: JsonRecord[]; + waitingEvidenceSources: string[]; + eventFactsAvailable: boolean; + factsAvailable: boolean; + valuesPrinted: false; +} interface K8sObject { kind?: string; metadata?: { name?: string; namespace?: string; + uid?: string; + resourceVersion?: string; creationTimestamp?: string; labels?: Record; annotations?: Record; ownerReferences?: Array<{ kind?: string; name?: string }>; }; status?: JsonRecord; + involvedObject?: { kind?: string; name?: string }; + reason?: string; + message?: string; + type?: string; + count?: number; + firstTimestamp?: string; + lastTimestamp?: string; + eventTime?: string; + series?: JsonRecord; } interface RunnerResourceEntry { resourceKind: "job" | "pod"; name: string; jobName: string; + resourceUid: string | null; + resourceVersion: string | null; runId: string | null; commandId: string | null; runnerId: string | null; createdAtMs: number; terminal: boolean; podPhase: string | null; + podObservations: RunnerPodObservation[]; protectedActive: boolean; protectedReason: string | null; candidateKind: CandidateKind; + classificationReason: string; lastActiveAtMs: number; + dbFactsFingerprint: string | null; } interface Snapshot { @@ -95,6 +135,8 @@ export async function enforceRunnerRetentionBeforeCreate(input: { store: AgentRu const selected = candidates.slice(0, requiredDeleteCount); const protectedActiveRunners = [...before.jobs, ...before.pods].filter((item) => item.protectedActive); const protectedActiveRunnerCount = protectedActiveRunners.length; + const inactiveCandidates = boundedResourceSummaries(candidates); + const protectedSummaries = boundedResourceSummaries(protectedActiveRunners); if (requiredDeleteCount > 0 && selected.length < requiredDeleteCount) { throw new AgentRunError("infra-failed", "runner retention cannot safely make room for a new runner", { httpStatus: 409, @@ -102,12 +144,16 @@ export async function enforceRunnerRetentionBeforeCreate(input: { store: AgentRu reason: "runner-retention-no-safe-candidate", namespace: input.options.namespace, maxRunners: input.options.maxRunners, + stalePendingMaxAgeMs: input.options.stalePendingMaxAgeMs, incomingRunnerCount, liveRunnerCountBefore: before.liveRunnerCount, requiredDeleteCount, inactiveCandidateCount: candidates.length, + inactiveCandidates: inactiveCandidates.items, + inactiveCandidateOmittedCount: inactiveCandidates.omittedCount, protectedActiveRunnerCount, - protectedActiveRunners: protectedActiveRunners.map(protectedRunnerSummary), + protectedActiveRunners: protectedSummaries.items, + protectedActiveRunnerOmittedCount: protectedSummaries.omittedCount, activeRunRisk: protectedActiveRunnerCount > 0, cleanupPolicy: cleanupPolicySummary(), valuesPrinted: false, @@ -119,20 +165,56 @@ export async function enforceRunnerRetentionBeforeCreate(input: { store: AgentRu let deletedRunnerPodCount = 0; let deletedAssociatedResourceCount = 0; for (const item of selected) { + await revalidateSelectedRunner(input.store, input.options, item); + if (!item.runId || !item.commandId || !item.dbFactsFingerprint) { + throw new AgentRunError("infra-failed", "runner retention selected resource has no fenceable database identity", { + httpStatus: 409, + details: { reason: "runner-retention-db-fence-unavailable", namespace: input.options.namespace, resourceKind: item.resourceKind, name: item.name, valuesPrinted: false }, + }); + } + await input.store.withRunnerRetentionFence({ + runId: item.runId, + commandId: item.commandId, + terminalizeStalePending: item.candidateKind === "stale-pending", + reason: `runner retention retired stale Kubernetes ${item.resourceKind}/${item.name}`, + }, async (facts) => { + const currentDbFingerprint = databaseFactsFingerprint(facts.run, facts.command, facts.commands, facts.events); + if (currentDbFingerprint !== item.dbFactsFingerprint) throw databaseFenceRejected(input.options, item, currentDbFingerprint, facts.run); + await deleteRunnerResourceCas(input.options, item); + }); if (item.resourceKind === "job") { - await kubectlRun(input.options.kubectlCommand ?? "kubectl", ["delete", "job", item.name, "-n", input.options.namespace, "--ignore-not-found=true"], input.options.signal); deletedRunnerJobCount++; deletedAssociatedResourceCount += await deleteAssociatedResources(input.options, item.jobName); } else { - await kubectlRun(input.options.kubectlCommand ?? "kubectl", ["delete", "pod", item.name, "-n", input.options.namespace, "--ignore-not-found=true"], input.options.signal); deletedRunnerPodCount++; } } const after = selected.length > 0 ? await runnerSnapshot(input.store, input.options) : before; + if (after.liveRunnerCount > targetBeforeCreate) { + throw new AgentRunError("infra-failed", "runner retention did not make enough room for a new runner", { + httpStatus: 409, + details: { + reason: "runner-retention-postcondition-failed", + namespace: input.options.namespace, + maxRunners: input.options.maxRunners, + incomingRunnerCount, + targetBeforeCreate, + liveRunnerCountBefore: before.liveRunnerCount, + liveRunnerCountAfter: after.liveRunnerCount, + selectedRunnerCount: selected.length, + deletedRunnerJobCount, + deletedRunnerPodCount, + deletedAssociatedResourceCount, + selected: selected.map(resourceSummary), + valuesPrinted: false, + }, + }); + } return { enabled: true, namespace: input.options.namespace, maxRunners: input.options.maxRunners, + stalePendingMaxAgeMs: input.options.stalePendingMaxAgeMs, incomingRunnerCount, liveRunnerCountBefore: before.liveRunnerCount, liveRunnerCountAfter: after.liveRunnerCount, @@ -140,8 +222,11 @@ export async function enforceRunnerRetentionBeforeCreate(input: { store: AgentRu nonTerminalRunnerPodCountBefore: before.nonTerminalRunnerPodCount, orphanNonTerminalRunnerPodCountBefore: before.orphanNonTerminalRunnerPodCount, inactiveCandidateCount: candidates.length, + inactiveCandidates: inactiveCandidates.items, + inactiveCandidateOmittedCount: inactiveCandidates.omittedCount, protectedActiveRunnerCount, - protectedActiveRunners: protectedActiveRunners.map(protectedRunnerSummary), + protectedActiveRunners: protectedSummaries.items, + protectedActiveRunnerOmittedCount: protectedSummaries.omittedCount, selectedRunnerCount: selected.length, deletedRunnerJobCount, deletedRunnerPodCount, @@ -149,7 +234,7 @@ export async function enforceRunnerRetentionBeforeCreate(input: { store: AgentRu activeRunRisk: protectedActiveRunnerCount > 0, reason: requiredDeleteCount === 0 ? "within-limit" : "pre-create-retention", cleanupPolicy: cleanupPolicySummary(), - selected: selected.map((item) => ({ resourceKind: item.resourceKind, name: item.name, jobName: item.jobName, candidateKind: item.candidateKind, protectedActive: false, lastActiveAt: new Date(item.lastActiveAtMs).toISOString(), valuesPrinted: false })), + selected: selected.map(resourceSummary), valuesPrinted: false, } satisfies RunnerRetentionSummary; } @@ -162,17 +247,49 @@ async function runnerSnapshot(store: AgentRunStore, options: RunnerRetentionOpti const failure = observed.find((item): item is PromiseRejectedResult => item.status === "rejected"); if (failure) throw failure.reason; const [jobObjects, podObjects] = observed.map((item) => (item as PromiseFulfilledResult).value); + const matchingPods = podObjects.filter((item) => { + const name = objectName(item); + if (!name) return false; + return matchesJobPrefix(podJobName(item) ?? name, options.jobNamePrefixes); + }); + const podsByJob = new Map(); + const eventsByPod = new Map(); + let eventFactsAvailable = true; + let eventObjects: K8sObject[] = []; + if (matchingPods.length > 0) { + try { + eventObjects = await kubectlGetPodEvents(options); + } catch (error) { + if (options.signal?.aborted) throw error; + eventFactsAvailable = false; + } + } + for (const item of eventObjects) { + const podName = podEventName(item); + if (!podName) continue; + const current = eventsByPod.get(podName) ?? []; + current.push(item); + eventsByPod.set(podName, current); + } + for (const item of matchingPods) { + const name = objectName(item); + if (!name) continue; + const jobName = podJobName(item) ?? name; + const current = podsByJob.get(jobName) ?? []; + current.push(item); + podsByJob.set(jobName, current); + } const jobs: RunnerResourceEntry[] = []; for (const item of jobObjects) { const name = objectName(item); if (!name || !matchesJobPrefix(name, options.jobNamePrefixes)) continue; - jobs.push(await jobEntry(store, item, options)); + jobs.push(await jobEntry(store, item, options, podsByJob.get(name) ?? [], eventsByPod, eventFactsAvailable)); } const jobNames = new Set(jobs.map((item) => item.jobName)); const pods: RunnerResourceEntry[] = []; let nonTerminalRunnerPodCount = 0; let orphanNonTerminalRunnerPodCount = 0; - for (const item of podObjects) { + for (const item of matchingPods) { const name = objectName(item); if (!name) continue; const jobName = podJobName(item) ?? name; @@ -181,7 +298,7 @@ async function runnerSnapshot(store: AgentRunStore, options: RunnerRetentionOpti if (!terminal) nonTerminalRunnerPodCount++; if (jobNames.has(jobName)) continue; if (!terminal) orphanNonTerminalRunnerPodCount++; - pods.push(await podEntry(store, item, options, jobName)); + pods.push(await podEntry(store, item, options, jobName, eventsByPod.get(name) ?? [], eventFactsAvailable)); } return { jobs, @@ -193,9 +310,11 @@ async function runnerSnapshot(store: AgentRunStore, options: RunnerRetentionOpti }; } -async function jobEntry(store: AgentRunStore, item: K8sObject, options: RunnerRetentionOptions): Promise { +async function jobEntry(store: AgentRunStore, item: K8sObject, options: RunnerRetentionOptions, pods: K8sObject[], eventsByPod: Map, eventFactsAvailable: boolean): Promise { const name = objectName(item) ?? "unknown"; const annotations = item.metadata?.annotations ?? {}; + const podObservations = pods.map((pod) => observePod(pod, options.stalePendingMaxAgeMs, eventsByPod.get(objectName(pod) ?? "") ?? [], eventFactsAvailable)); + const podIdentityMatch = pods.every((pod) => resourceIdentityMatches(annotations, pod.metadata?.annotations ?? {})); const activity = await activityFor(store, { runId: annotations["agentrun.pikastech.local/run-id"] ?? null, commandId: annotations["agentrun.pikastech.local/command-id"] ?? null, @@ -203,21 +322,28 @@ async function jobEntry(store: AgentRunStore, item: K8sObject, options: RunnerRe terminal: jobTerminal(item), createdAtMs: objectCreatedAtMs(item), activeHeartbeatMaxAgeMs: options.activeHeartbeatMaxAgeMs, + stalePendingMaxAgeMs: options.stalePendingMaxAgeMs, + podObservations, + podIdentityMatch, }); return { resourceKind: "job", name, jobName: name, + resourceUid: item.metadata?.uid ?? null, + resourceVersion: item.metadata?.resourceVersion ?? null, createdAtMs: objectCreatedAtMs(item), terminal: jobTerminal(item), podPhase: null, + podObservations, ...activity, }; } -async function podEntry(store: AgentRunStore, item: K8sObject, options: RunnerRetentionOptions, jobName: string): Promise { +async function podEntry(store: AgentRunStore, item: K8sObject, options: RunnerRetentionOptions, jobName: string, events: K8sObject[], eventFactsAvailable: boolean): Promise { const name = objectName(item) ?? "unknown"; const annotations = item.metadata?.annotations ?? {}; + const podObservations = [observePod(item, options.stalePendingMaxAgeMs, events, eventFactsAvailable)]; const activity = await activityFor(store, { runId: annotations["agentrun.pikastech.local/run-id"] ?? null, commandId: annotations["agentrun.pikastech.local/command-id"] ?? null, @@ -225,42 +351,174 @@ async function podEntry(store: AgentRunStore, item: K8sObject, options: RunnerRe terminal: podTerminal(item), createdAtMs: objectCreatedAtMs(item), activeHeartbeatMaxAgeMs: options.activeHeartbeatMaxAgeMs, + stalePendingMaxAgeMs: options.stalePendingMaxAgeMs, + podObservations, + podIdentityMatch: true, }); return { resourceKind: "pod", name, jobName, + resourceUid: item.metadata?.uid ?? null, + resourceVersion: item.metadata?.resourceVersion ?? null, createdAtMs: objectCreatedAtMs(item), terminal: podTerminal(item), podPhase: stringPath(item, ["status", "phase"]), + podObservations, ...activity, }; } -async function activityFor(store: AgentRunStore, input: { runId: string | null; commandId: string | null; runnerId: string | null; terminal: boolean; createdAtMs: number; activeHeartbeatMaxAgeMs: number }): Promise> { - let run: RunRecord | null = null; - let command: CommandRecord | null = null; - let openCommand: CommandRecord | null = null; - if (input.runId) { - try { run = await store.getRun(input.runId); } catch { run = null; } +type ActivityFields = Pick; + +async function activityFor(store: AgentRunStore, input: { + runId: string | null; + commandId: string | null; + runnerId: string | null; + terminal: boolean; + createdAtMs: number; + activeHeartbeatMaxAgeMs: number; + stalePendingMaxAgeMs: number; + podObservations: RunnerPodObservation[]; + podIdentityMatch: boolean; +}): Promise { + if (!input.runId || !input.commandId || !input.runnerId) return protectedActivity(input, "identity-facts-unavailable", input.createdAtMs, null); + if (!input.podIdentityMatch) return protectedActivity(input, "pod-identity-mismatch", input.createdAtMs, null); + + let run: RunRecord; + let command: CommandRecord; + try { + [run, command] = await Promise.all([store.getRun(input.runId), store.getCommand(input.commandId)]); + } catch { + return protectedActivity(input, "db-facts-unavailable", input.createdAtMs, null); } - if (input.commandId) { - try { command = await store.getCommand(input.commandId); } catch { command = null; } + + const lastActiveAtMs = Math.max(input.createdAtMs, isoMs(run.updatedAt), isoMs(command.updatedAt), isoMs(run.leaseExpiresAt)); + if (command.runId !== run.id) return protectedActivity(input, "db-identity-mismatch", lastActiveAtMs, null); + let commands: CommandRecord[]; + let events: import("../common/types.js").RunEvent[]; + try { + [commands, events] = await Promise.all([store.listCommands(run.id, 0, 100), store.listEvents(run.id, 0, 500)]); + } catch { + return protectedActivity(input, "db-progress-facts-unavailable", lastActiveAtMs, null); } - if (run && !isTerminalRunStatus(run.status) && !input.terminal) { - try { - openCommand = (await store.listCommands(run.id, 0, 100)).find((item) => !isTerminalCommandState(item.state)) ?? null; - } catch { - openCommand = null; + const dbFactsFingerprint = databaseFactsFingerprint(run, command, commands, events); + const freshClaim = Boolean(run.claimedBy && heartbeatFresh(run.leaseExpiresAt, input.activeHeartbeatMaxAgeMs)); + if (freshClaim) return protectedActivity(input, "fresh-active-run", lastActiveAtMs, dbFactsFingerprint); + const activePodReason = activePodProtectionReason(input.podObservations); + if (activePodReason) return protectedActivity(input, activePodReason, lastActiveAtMs, dbFactsFingerprint); + if (isTerminalRunStatus(run.status)) return candidateActivity(input, "terminal", "terminal-run", lastActiveAtMs, dbFactsFingerprint); + if (isTerminalCommandState(command.state)) return candidateActivity(input, input.terminal ? "terminal" : "idle", input.terminal ? "terminal-runner-resource" : "terminal-command-idle-resource", lastActiveAtMs, dbFactsFingerprint); + if (input.terminal) return candidateActivity(input, "terminal", "terminal-runner-resource", lastActiveAtMs, dbFactsFingerprint); + + const podBlockReason = stalePendingPodBlockReason(input.podObservations); + if (podBlockReason) return protectedActivity(input, podBlockReason, lastActiveAtMs, dbFactsFingerprint); + if (!hasNoBusinessProgress(run, command, commands, events, input.stalePendingMaxAgeMs)) { + return protectedActivity(input, "business-progress-observed-or-unverifiable", lastActiveAtMs, dbFactsFingerprint); + } + return candidateActivity(input, "stale-pending", "stale-pending-no-business-progress", lastActiveAtMs, dbFactsFingerprint); +} + +async function revalidateSelectedRunner(store: AgentRunStore, options: RunnerRetentionOptions, selected: RunnerResourceEntry): Promise { + const snapshot = await runnerSnapshot(store, options); + const resources = selected.resourceKind === "job" ? snapshot.jobs : snapshot.pods; + const current = resources.find((item) => item.name === selected.name) ?? null; + const selectedFingerprint = runnerSafetyFingerprint(selected); + const currentFingerprint = current ? runnerSafetyFingerprint(current) : null; + if (current && !current.protectedActive && currentFingerprint === selectedFingerprint) return; + + throw new AgentRunError("infra-failed", "runner retention selected resource changed before deletion", { + httpStatus: 409, + details: { + reason: "runner-retention-selected-runner-changed", + changeReason: current === null ? "resource-missing" : current.protectedActive ? "protected-active" : "safety-fingerprint-changed", + namespace: options.namespace, + resourceKind: selected.resourceKind, + name: selected.name, + selectedFingerprint, + currentFingerprint, + selected: resourceSummary(selected), + current: current ? resourceSummary(current) : null, + valuesPrinted: false, + }, + }); +} + +function databaseFenceRejected(options: RunnerRetentionOptions, item: RunnerResourceEntry, currentDbFingerprint: string, run: RunRecord): AgentRunError { + return new AgentRunError("infra-failed", "runner retention database facts changed before atomic deletion", { + httpStatus: 409, + details: { + reason: "runner-retention-db-fence-rejected", + changeReason: run.claimedBy || run.status === "claimed" || run.status === "running" ? "run-became-active" : "database-fingerprint-changed", + namespace: options.namespace, + resourceKind: item.resourceKind, + name: item.name, + selectedDbFingerprint: item.dbFactsFingerprint, + currentDbFingerprint, + runStatus: run.status, + claimedBy: run.claimedBy, + leaseExpiresAt: run.leaseExpiresAt, + valuesPrinted: false, + }, + }); +} + +async function deleteRunnerResourceCas(options: RunnerRetentionOptions, item: RunnerResourceEntry): Promise { + if (item.resourceKind === "job") { + for (const pod of [...item.podObservations].sort((left, right) => left.name.localeCompare(right.name))) { + await deleteKubernetesObjectCas(options, item, "pod", pod.name, pod.uid, pod.resourceVersion); } + const currentJob = await kubectlGetObject(options, "job", item.name); + const expectedIdentity = { + "agentrun.pikastech.local/run-id": item.runId ?? "", + "agentrun.pikastech.local/command-id": item.commandId ?? "", + "agentrun.pikastech.local/runner-id": item.runnerId ?? "", + }; + if (currentJob.metadata?.uid !== item.resourceUid || !resourceIdentityMatches(expectedIdentity, currentJob.metadata?.annotations ?? {})) { + throw new AgentRunError("infra-failed", "runner retention Job identity changed after Pod CAS deletion", { + httpStatus: 409, + details: { reason: "runner-retention-k8s-cas-rejected", changeReason: "job-identity-changed", namespace: options.namespace, resourceKind: "job", name: item.name, expectedUid: item.resourceUid, currentUid: currentJob.metadata?.uid ?? null, valuesPrinted: false }, + }); + } + await deleteKubernetesObjectCas(options, item, "job", item.name, currentJob.metadata?.uid ?? null, currentJob.metadata?.resourceVersion ?? null); + return; } - const lastActiveAtMs = Math.max(input.createdAtMs, isoMs(run?.updatedAt), isoMs(command?.updatedAt), isoMs(openCommand?.updatedAt), isoMs(run?.leaseExpiresAt)); - if (run && !isTerminalRunStatus(run.status) && !input.terminal && (openCommand || (command && !isTerminalCommandState(command.state)))) { - const freshClaim = input.runnerId && run.claimedBy === input.runnerId && heartbeatFresh(run.leaseExpiresAt, input.activeHeartbeatMaxAgeMs); - return { runId: input.runId, commandId: input.commandId, runnerId: input.runnerId, protectedActive: true, protectedReason: freshClaim ? "fresh-active-run" : openCommand ? "nonterminal-run-open-command" : "nonterminal-runner-resource", candidateKind: "inactive", lastActiveAtMs }; + await deleteKubernetesObjectCas(options, item, item.resourceKind, item.name, item.resourceUid, item.resourceVersion); +} + +async function deleteKubernetesObjectCas(options: RunnerRetentionOptions, selected: RunnerResourceEntry, resourceKind: "job" | "pod", name: string, uid: string | null, resourceVersion: string | null): Promise { + if (!uid || !resourceVersion) { + throw new AgentRunError("infra-failed", "runner retention cannot bind Kubernetes delete preconditions", { + httpStatus: 409, + details: { reason: "runner-retention-k8s-cas-unavailable", namespace: options.namespace, selectedResourceKind: selected.resourceKind, selectedName: selected.name, resourceKind, name, uidPresent: Boolean(uid), resourceVersionPresent: Boolean(resourceVersion), valuesPrinted: false }, + }); } - const candidateKind: CandidateKind = command && isTerminalCommandState(command.state) && !input.terminal ? "idle" : input.terminal || (run !== null && isTerminalRunStatus(run.status)) ? "terminal" : "inactive"; - return { runId: input.runId, commandId: input.commandId, runnerId: input.runnerId, protectedActive: false, protectedReason: null, candidateKind, lastActiveAtMs }; + const resourcePath = resourceKind === "job" + ? `/apis/batch/v1/namespaces/${encodeURIComponent(options.namespace)}/jobs/${encodeURIComponent(name)}` + : `/api/v1/namespaces/${encodeURIComponent(options.namespace)}/pods/${encodeURIComponent(name)}`; + const deleteOptions = JSON.stringify({ apiVersion: "v1", kind: "DeleteOptions", propagationPolicy: "Background", preconditions: { uid, resourceVersion } }); + const result = await kubectlRun(options.kubectlCommand ?? "kubectl", ["delete", "--raw", resourcePath, "-f", "-"], options.signal, deleteOptions); + if (result.code === 0) return; + const conflict = /\b(?:Conflict|PreconditionFailed|resourceVersion|UID)\b/iu.test(`${result.stderr}\n${result.stdout}`); + throw new AgentRunError("infra-failed", `Kubernetes CAS delete ${resourceKind} failed with code ${result.code}`, { + httpStatus: conflict ? 409 : 502, + details: redactJson({ + reason: "runner-retention-k8s-cas-rejected", + changeReason: conflict ? "object-identity-or-version-changed" : "delete-request-failed", + namespace: options.namespace, + selectedResourceKind: selected.resourceKind, + selectedName: selected.name, + resourceKind, + name, + expectedUid: uid, + expectedResourceVersion: resourceVersion, + exitCode: result.code, + signal: result.signal, + stderr: redactText(result.stderr.slice(-2000)), + stdout: redactText(result.stdout.slice(-1000)), + valuesPrinted: false, + }), + }); } async function kubectlGetList(options: RunnerRetentionOptions, resource: string, selector: string): Promise { @@ -274,6 +532,21 @@ async function kubectlGetList(options: RunnerRetentionOptions, resource: string, return Array.isArray(items) ? items.filter((entry) => typeof entry === "object" && entry !== null && !Array.isArray(entry)).map((entry) => entry as unknown as K8sObject) : []; } +async function kubectlGetObject(options: RunnerRetentionOptions, resource: "job" | "pod", name: string): Promise { + const result = await kubectlRun(options.kubectlCommand ?? "kubectl", ["get", resource, name, "-n", options.namespace, "-o", "json"], options.signal); + if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl get ${resource}/${name} failed with code ${result.code}`, { httpStatus: 409, details: redactJson({ reason: "runner-retention-k8s-cas-rejected", changeReason: "resource-missing-before-cas", namespace: options.namespace, resourceKind: resource, name, stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); + return parseJsonObject(result.stdout, `kubectl get ${resource}/${name}`) as unknown as K8sObject; +} + +async function kubectlGetPodEvents(options: RunnerRetentionOptions): Promise { + const args = ["get", "events", "-n", options.namespace, "--field-selector", "involvedObject.kind=Pod,type=Warning", "-o", "json"]; + const result = await kubectlRun(options.kubectlCommand ?? "kubectl", args, options.signal); + if (result.code !== 0) throw new AgentRunError("infra-failed", `kubectl get events failed with code ${result.code}`, { httpStatus: 502, details: redactJson({ stderr: redactText(result.stderr.slice(-2000)), stdout: redactText(result.stdout.slice(-1000)), valuesPrinted: false }) }); + const parsed = parseJsonObject(result.stdout, "kubectl get events"); + const items = parsed.items; + return Array.isArray(items) ? items.filter((entry) => typeof entry === "object" && entry !== null && !Array.isArray(entry)).map((entry) => entry as unknown as K8sObject) : []; +} + async function deleteAssociatedResources(options: RunnerRetentionOptions, jobName: string): Promise { let deleted = 0; const selector = `agentrun.pikastech.local/runner-job=${jobName}`; @@ -285,21 +558,25 @@ async function deleteAssociatedResources(options: RunnerRetentionOptions, jobNam return deleted; } -async function kubectlRun(kubectlCommand: string, args: string[], abortSignal?: AbortSignal): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { +async function kubectlRun(kubectlCommand: string, args: string[], abortSignal?: AbortSignal, stdin?: string): Promise<{ code: number | null; signal: NodeJS.Signals | null; stdout: string; stderr: string }> { throwIfAborted(abortSignal, "kubectl"); - const child = spawn(kubectlCommand, args, { stdio: ["ignore", "pipe", "pipe"] }); + const child = spawn(kubectlCommand, args, { stdio: [stdin === undefined ? "ignore" : "pipe", "pipe", "pipe"] }); let stdout = ""; let stderr = ""; - child.stdout.setEncoding("utf8"); - child.stderr.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { stdout += String(chunk); }); - child.stderr.on("data", (chunk) => { stderr += String(chunk); }); + const stdoutStream = child.stdout; + const stderrStream = child.stderr; + if (!stdoutStream || !stderrStream) throw new AgentRunError("infra-failed", "kubectl output pipes were not created", { httpStatus: 503 }); + stdoutStream.setEncoding("utf8"); + stderrStream.setEncoding("utf8"); + stdoutStream.on("data", (chunk) => { stdout += String(chunk); }); + stderrStream.on("data", (chunk) => { stderr += String(chunk); }); + if (stdin !== undefined) child.stdin?.end(stdin); const result = await waitForChildProcess(child, { ...(abortSignal ? { signal: abortSignal } : {}), label: "kubectl" }); return { ...result, stdout, stderr }; } function compareCandidates(left: RunnerResourceEntry, right: RunnerResourceEntry): number { - const priority = (kind: CandidateKind): number => kind === "idle" ? 0 : kind === "terminal" ? 1 : 2; + const priority = (kind: CandidateKind): number => kind === "idle" ? 0 : kind === "terminal" ? 1 : kind === "stale-pending" ? 2 : 3; const byKind = priority(left.candidateKind) - priority(right.candidateKind); if (byKind !== 0) return byKind; const byActive = left.lastActiveAtMs - right.lastActiveAtMs; @@ -307,28 +584,34 @@ function compareCandidates(left: RunnerResourceEntry, right: RunnerResourceEntry return left.name.localeCompare(right.name); } -function protectedRunnerSummary(item: RunnerResourceEntry): JsonRecord { - return { +function resourceSummary(item: RunnerResourceEntry): JsonRecord { + return redactJson({ resourceKind: item.resourceKind, name: item.name, jobName: item.jobName, + resourceUid: item.resourceUid, + resourceVersion: item.resourceVersion, runId: item.runId, commandId: item.commandId, runnerId: item.runnerId, protectedReason: item.protectedReason, + candidateKind: item.candidateKind, + classificationReason: item.classificationReason, terminal: item.terminal, podPhase: item.podPhase, + podObservations: item.podObservations.slice(0, 5), + podObservationOmittedCount: Math.max(0, item.podObservations.length - 5), lastActiveAt: new Date(item.lastActiveAtMs).toISOString(), valuesPrinted: false, - }; + }) as JsonRecord; } function cleanupPolicySummary(): JsonRecord { return { mode: "db-ledger-and-k8s-observation", forceActive: false, - activeProtection: "protect non-terminal DB run/command while Kubernetes runner resource is non-terminal; fresh heartbeat is not required during manager rolling recovery", - deletionOrder: "idle, terminal, inactive by lastActiveAt", + activeProtection: "protect non-terminal DB run/command unless unclaimed no-progress facts and stale Pending Pod observations are complete", + deletionOrder: "idle, terminal, stale-pending, inactive by lastActiveAt", valuesPrinted: false, }; } @@ -347,6 +630,13 @@ function podJobName(item: K8sObject): string | null { return item.metadata?.ownerReferences?.find((owner) => owner.kind === "Job")?.name ?? null; } +function resourceIdentityMatches(left: Record, right: Record): boolean { + for (const key of ["agentrun.pikastech.local/run-id", "agentrun.pikastech.local/command-id", "agentrun.pikastech.local/runner-id"] as const) { + if (!left[key] || !right[key] || left[key] !== right[key]) return false; + } + return true; +} + function jobTerminal(item: K8sObject): boolean { const conditions = item.status?.conditions; if (!Array.isArray(conditions)) return false; @@ -361,6 +651,259 @@ function podTerminal(item: K8sObject): boolean { return phase === "Succeeded" || phase === "Failed"; } +function observePod(item: K8sObject, stalePendingMaxAgeMs: number, events: K8sObject[], eventFactsAvailable: boolean): RunnerPodObservation { + const createdAtMs = objectCreatedAtMs(item); + const ageMs = createdAtMs > 0 ? Math.max(0, Date.now() - createdAtMs) : null; + const phase = stringPath(item, ["status", "phase"]); + const conditions = jsonRecordArray(item.status?.conditions); + const statuses = jsonRecordArray(item.status?.containerStatuses); + const runner = statuses.find((status) => status.name === "runner") ?? null; + const runnerState = jsonRecord(runner?.state); + const runnerContainerState: RunnerPodObservation["runnerContainerState"] = runnerState?.waiting ? "waiting" : runnerState?.running ? "running" : runnerState?.terminated ? "terminated" : "unknown"; + const ready = conditions.some((condition) => condition.type === "Ready" && condition.status === "True"); + const waiting = [ + ...statuses.flatMap((status) => waitingSummary(status, "container", typeof status.name === "string" ? status.name : null)), + ...jsonRecordArray(item.status?.initContainerStatuses).flatMap((status) => waitingSummary(status, "init-container", typeof status.name === "string" ? status.name : null)), + ...conditions.flatMap((condition) => condition.status === "False" && typeof condition.reason === "string" && condition.reason.length > 0 + ? [{ source: "condition", name: typeof condition.type === "string" ? condition.type : null, reason: condition.reason, message: boundedRedactedText(condition.message), deterministic: condition.reason === "Unschedulable", valuesPrinted: false }] + : []), + ...events.flatMap(eventWaitingSummary), + ]; + const runnerContainerRestartCount = typeof runner?.restartCount === "number" ? runner.restartCount : 0; + const runnerLastState = jsonRecord(runner?.lastState); + const runnerContainerLastState: RunnerPodObservation["runnerContainerLastState"] = runnerLastState?.running ? "running" : runnerLastState?.terminated ? "terminated" : "none"; + const runnerContainerStarted = runner?.started === true || runnerContainerRestartCount > 0 || runnerContainerState === "running" || runnerContainerState === "terminated" || Boolean(runnerLastState?.running) || Boolean(runnerLastState?.terminated); + const runnerContainerReady = runner?.ready === true; + return { + name: objectName(item) ?? "unknown", + uid: item.metadata?.uid ?? null, + resourceVersion: item.metadata?.resourceVersion ?? null, + phase, + createdAt: createdAtMs > 0 ? new Date(createdAtMs).toISOString() : null, + ageMs, + stale: ageMs !== null && ageMs >= stalePendingMaxAgeMs, + podStartTime: boundedRedactedText(item.status?.startTime), + ready, + runnerContainerState, + runnerContainerLastState, + runnerContainerStarted, + runnerContainerReady, + runnerContainerRestartCount, + waiting, + waitingEvidenceSources: [...new Set(waiting.map((entry) => typeof entry.source === "string" ? entry.source : "unknown"))], + eventFactsAvailable, + factsAvailable: phase !== null && createdAtMs > 0, + valuesPrinted: false, + }; +} + +function waitingSummary(status: JsonRecord, source: string, name: string | null): JsonRecord[] { + const state = jsonRecord(status.state); + const waiting = jsonRecord(state?.waiting); + if (!waiting || typeof waiting.reason !== "string" || waiting.reason.length === 0) return []; + return [{ source, name, reason: waiting.reason, message: boundedRedactedText(waiting.message), deterministic: waiting.reason !== "ContainerCreating" && waiting.reason !== "PodInitializing", valuesPrinted: false }]; +} + +function eventWaitingSummary(event: K8sObject): JsonRecord[] { + if (event.type !== "Warning" || typeof event.reason !== "string" || event.reason.length === 0) return []; + return [{ + source: "event", + name: typeof event.metadata?.name === "string" ? event.metadata.name : null, + reason: event.reason, + message: boundedRedactedText(event.message), + deterministic: true, + count: typeof event.count === "number" ? event.count : null, + lastObservedAt: boundedRedactedText(event.series?.lastObservedTime ?? event.eventTime ?? event.lastTimestamp ?? event.firstTimestamp), + valuesPrinted: false, + }]; +} + +function podEventName(item: K8sObject): string | null { + if (item.involvedObject?.kind !== "Pod") return null; + return typeof item.involvedObject.name === "string" && item.involvedObject.name.length > 0 ? item.involvedObject.name : null; +} + +function stalePendingPodBlockReason(observations: RunnerPodObservation[]): string | null { + if (observations.length === 0) return "pod-facts-unavailable"; + if (observations.some((item) => !item.factsAvailable)) return "pod-facts-unavailable"; + const activeReason = activePodProtectionReason(observations); + if (activeReason) return activeReason; + if (observations.some((item) => item.phase !== "Pending")) return "runner-pod-not-pending"; + if (observations.some((item) => !item.stale)) return "stale-pending-window-active"; + if (observations.some((item) => !item.waiting.some((entry) => entry.deterministic === true) && !item.eventFactsAvailable)) return "pod-waiting-event-facts-unavailable"; + if (observations.some((item) => item.waiting.length > 0 && !item.waiting.some((entry) => entry.deterministic === true))) return "pod-deterministic-waiting-reason-unavailable"; + if (observations.some((item) => item.waiting.length === 0)) return "pod-waiting-reason-unavailable"; + return null; +} + +function activePodProtectionReason(observations: RunnerPodObservation[]): string | null { + if (observations.some((item) => item.phase === "Running" || item.ready || item.runnerContainerReady || item.runnerContainerStarted || item.runnerContainerState === "running" || item.runnerContainerState === "terminated")) return "runner-pod-started-running-or-ready"; + return null; +} + +function hasNoBusinessProgress(run: RunRecord, command: CommandRecord, commands: CommandRecord[], events: import("../common/types.js").RunEvent[], stalePendingMaxAgeMs: number): boolean { + const cutoff = Date.now() - stalePendingMaxAgeMs; + if (run.status !== "pending" || run.claimedBy !== null || run.leaseExpiresAt !== null || run.cancelRequestId !== null || run.cancelRequestedAt !== null) return false; + if (command.state !== "pending" || command.acknowledgedAt !== null || command.cancelRequestId !== null || command.cancelRequestedAt !== null) return false; + if (!sameValidTimestamp(run.createdAt, run.updatedAt) || !sameValidTimestamp(command.createdAt, command.updatedAt)) return false; + if (isoMs(run.createdAt) > cutoff || isoMs(command.createdAt) > cutoff) return false; + if (commands.length !== 1 || commands[0]?.id !== command.id) return false; + if (events.length === 0 || events.length >= 500) return false; + const allowedPhases = new Set(["run-created", "command-created", "runner-job-created", "runner-dispatch-completed"]); + const userMessages = events.filter((event) => event.type === "user_message"); + if (userMessages.length > 1) return false; + const phases = events.flatMap((event) => { + if (event.type === "user_message") { + const payload = event.payload; + const valid = payload.commandId === command.id + && typeof payload.userMessageId === "string" && payload.userMessageId.length > 0 + && typeof payload.text === "string" && payload.text.length > 0 + && typeof payload.source === "string" && payload.source.length > 0; + return valid ? [] : [null]; + } + const phase = typeof event.payload.phase === "string" ? event.payload.phase : null; + return event.type === "backend_status" && phase !== null && allowedPhases.has(phase) ? [phase] : [null]; + }); + if (phases.includes(null)) return false; + if (!phases.includes("run-created") || !phases.includes("command-created") || !phases.includes("runner-job-created")) return false; + return events.every((event) => event.payload.commandId === undefined || event.payload.commandId === command.id); +} + +function protectedActivity(input: { runId: string | null; commandId: string | null; runnerId: string | null }, reason: string, lastActiveAtMs: number, dbFactsFingerprint: string | null): ActivityFields { + return { runId: input.runId, commandId: input.commandId, runnerId: input.runnerId, protectedActive: true, protectedReason: reason, candidateKind: "inactive", classificationReason: reason, lastActiveAtMs, dbFactsFingerprint }; +} + +function candidateActivity(input: { runId: string | null; commandId: string | null; runnerId: string | null }, kind: CandidateKind, reason: string, lastActiveAtMs: number, dbFactsFingerprint: string): ActivityFields { + return { runId: input.runId, commandId: input.commandId, runnerId: input.runnerId, protectedActive: false, protectedReason: null, candidateKind: kind, classificationReason: reason, lastActiveAtMs, dbFactsFingerprint }; +} + +function databaseFactsFingerprint(run: RunRecord, command: CommandRecord, commands: CommandRecord[], events: import("../common/types.js").RunEvent[]): string { + const commandFacts = (item: CommandRecord): JsonRecord => ({ + id: item.id, + runId: item.runId, + seq: item.seq, + type: item.type, + state: item.state, + payloadHash: item.payloadHash, + cancelEpoch: item.cancelEpoch, + cancelRequestId: item.cancelRequestId, + cancelRequestedAt: item.cancelRequestedAt, + cancelReason: item.cancelReason, + createdAt: item.createdAt, + updatedAt: item.updatedAt, + acknowledgedAt: item.acknowledgedAt, + dispatchIntentHash: stableHash(item.dispatchIntent ?? null), + valuesPrinted: false, + }); + return stableHash({ + run: { + id: run.id, + status: run.status, + terminalStatus: run.terminalStatus, + failureKind: run.failureKind, + failureMessageHash: run.failureMessage ? stableHash(run.failureMessage) : null, + cancelEpoch: run.cancelEpoch, + cancelRequestId: run.cancelRequestId, + cancelRequestedAt: run.cancelRequestedAt, + cancelReason: run.cancelReason, + createdAt: run.createdAt, + updatedAt: run.updatedAt, + claimedBy: run.claimedBy, + leaseExpiresAt: run.leaseExpiresAt, + valuesPrinted: false, + }, + command: commandFacts(command), + commands: [...commands].sort((left, right) => left.seq - right.seq || left.id.localeCompare(right.id)).map(commandFacts), + events: [...events].sort((left, right) => left.seq - right.seq || left.id.localeCompare(right.id)).map((event) => ({ + id: event.id, + runId: event.runId, + seq: event.seq, + type: event.type, + createdAt: event.createdAt, + payloadHash: stableHash(event.payload), + valuesPrinted: false, + })), + valuesPrinted: false, + }); +} + +function runnerSafetyFingerprint(item: RunnerResourceEntry): string { + return stableHash({ + schema: "agentrun-runner-retention-safety-v1", + resourceKind: item.resourceKind, + name: item.name, + jobName: item.jobName, + resourceUid: item.resourceUid, + resourceVersion: item.resourceVersion, + runId: item.runId, + commandId: item.commandId, + runnerId: item.runnerId, + createdAtMs: item.createdAtMs, + terminal: item.terminal, + podPhase: item.podPhase, + protectedActive: item.protectedActive, + protectedReason: item.protectedReason, + candidateKind: item.candidateKind, + classificationReason: item.classificationReason, + lastActiveAtMs: item.lastActiveAtMs, + dbFactsFingerprint: item.dbFactsFingerprint, + podObservations: [...item.podObservations].sort((left, right) => left.name.localeCompare(right.name)).map(podSafetyFacts), + valuesPrinted: false, + }); +} + +function podSafetyFacts(item: RunnerPodObservation): JsonRecord { + const waiting = item.waiting.map((entry) => ({ + source: typeof entry.source === "string" ? entry.source : null, + name: typeof entry.name === "string" ? entry.name : null, + reason: typeof entry.reason === "string" ? entry.reason : null, + deterministic: entry.deterministic === true, + valuesPrinted: false, + })).sort((left, right) => JSON.stringify(left).localeCompare(JSON.stringify(right))); + return { + name: item.name, + uid: item.uid, + resourceVersion: item.resourceVersion, + phase: item.phase, + createdAt: item.createdAt, + stale: item.stale, + podStartTime: item.podStartTime, + ready: item.ready, + runnerContainerState: item.runnerContainerState, + runnerContainerLastState: item.runnerContainerLastState, + runnerContainerStarted: item.runnerContainerStarted, + runnerContainerReady: item.runnerContainerReady, + runnerContainerRestartCount: item.runnerContainerRestartCount, + waiting, + waitingEvidenceSources: [...item.waitingEvidenceSources].sort(), + eventFactsAvailable: item.eventFactsAvailable, + factsAvailable: item.factsAvailable, + valuesPrinted: false, + }; +} + +function boundedResourceSummaries(items: RunnerResourceEntry[]): { items: JsonRecord[]; omittedCount: number } { + const limit = 20; + return { items: items.slice(0, limit).map(resourceSummary), omittedCount: Math.max(0, items.length - limit) }; +} + +function jsonRecord(value: unknown): JsonRecord | null { + return typeof value === "object" && value !== null && !Array.isArray(value) ? value as JsonRecord : null; +} + +function jsonRecordArray(value: unknown): JsonRecord[] { + return Array.isArray(value) ? value.map(jsonRecord).filter((item): item is JsonRecord => item !== null) : []; +} + +function boundedRedactedText(value: unknown): string | null { + if (typeof value !== "string" || value.length === 0) return null; + return redactText(value).slice(0, 512); +} + +function sameValidTimestamp(left: string, right: string): boolean { + return isoMs(left) > 0 && isoMs(left) === isoMs(right); +} + function heartbeatFresh(leaseExpiresAt: string | null, activeHeartbeatMaxAgeMs: number): boolean { const leaseMs = isoMs(leaseExpiresAt); if (leaseMs <= 0) return false; diff --git a/src/mgr/server.ts b/src/mgr/server.ts index 6f38124..2fc39ef 100644 --- a/src/mgr/server.ts +++ b/src/mgr/server.ts @@ -95,6 +95,7 @@ function runnerRetentionOptionsForRequest(defaults: RunnerRetentionOptions | und const maxRunners = optionalPositiveInteger("AGENTRUN_RUNNER_RETENTION_MAX_RUNNERS", process.env.AGENTRUN_RUNNER_RETENTION_MAX_RUNNERS); if (maxRunners === undefined) return undefined; const activeHeartbeatMaxAgeMs = requiredPositiveInteger("AGENTRUN_RUNNER_RETENTION_ACTIVE_HEARTBEAT_MAX_AGE_MS", process.env.AGENTRUN_RUNNER_RETENTION_ACTIVE_HEARTBEAT_MAX_AGE_MS); + const stalePendingMaxAgeMs = requiredPositiveInteger("AGENTRUN_RUNNER_RETENTION_STALE_PENDING_MAX_AGE_MS", process.env.AGENTRUN_RUNNER_RETENTION_STALE_PENDING_MAX_AGE_MS); const cleanupOrder = process.env.AGENTRUN_RUNNER_RETENTION_CLEANUP_ORDER ?? "oldest-inactive-last-active-first"; if (cleanupOrder !== "oldest-inactive-last-active-first") throw new AgentRunError("schema-invalid", "AGENTRUN_RUNNER_RETENTION_CLEANUP_ORDER is unsupported", { httpStatus: 500 }); const jobNamePrefixes = stringListEnv("AGENTRUN_RUNNER_RETENTION_JOB_NAME_PREFIXES", process.env.AGENTRUN_RUNNER_RETENTION_JOB_NAME_PREFIXES); @@ -103,6 +104,7 @@ function runnerRetentionOptionsForRequest(defaults: RunnerRetentionOptions | und maxRunners, cleanupOrder, activeHeartbeatMaxAgeMs, + stalePendingMaxAgeMs, matchLabels: jsonRecordEnv("AGENTRUN_RUNNER_RETENTION_MATCH_LABELS_JSON", process.env.AGENTRUN_RUNNER_RETENTION_MATCH_LABELS_JSON), jobNamePrefixes: jobNamePrefixes.length > 0 ? jobNamePrefixes : [jobNamePrefix], ageBasedCleanup: { diff --git a/src/mgr/store.ts b/src/mgr/store.ts index 64ffbc9..99c9199 100644 --- a/src/mgr/store.ts +++ b/src/mgr/store.ts @@ -32,6 +32,20 @@ export interface DurableQueueClaim { attemptCount: number; } +export interface RunnerRetentionFenceFacts { + run: RunRecord; + command: CommandRecord; + commands: CommandRecord[]; + events: RunEvent[]; +} + +export interface RunnerRetentionFenceInput { + runId: string; + commandId: string; + terminalizeStalePending: boolean; + reason: string; +} + export interface AgentRunStore { health(): MaybePromise; createRun(input: CreateRunInput): MaybePromise; @@ -57,6 +71,8 @@ export interface AgentRunStore { getRunnerJobByIdempotencyKey(runId: string, idempotencyKey: string, payloadHash: string): MaybePromise; saveRunnerJob(input: SaveRunnerJobInput): MaybePromise; updateRunnerJobResult(runnerJobId: string, patch: JsonRecord): MaybePromise; + withRunnerCreateFence(namespace: string, operation: () => Promise): MaybePromise; + withRunnerRetentionFence(input: RunnerRetentionFenceInput, operation: (facts: RunnerRetentionFenceFacts) => Promise): MaybePromise; claimRun(runId: string, runnerId: string, leaseMs: number): MaybePromise; heartbeat(runId: string, runnerId: string, leaseMs: number): MaybePromise; ackCommand(commandId: string): MaybePromise; @@ -179,6 +195,8 @@ export class MemoryAgentRunStore implements AgentRunStore { private readonly kafkaEventOutbox = new Map(); private readonly queueTasks = new Map(); private readonly queueReadCursors = new Map(); + private readonly runnerCreateFenceTails = new Map>(); + private readonly runnerRetentionFences = new Set(); private queueVersion = 0; private sessionVersion = 0; private kafkaOutboxSeq = 0; @@ -429,7 +447,40 @@ export class MemoryAgentRunStore implements AgentRunStore { return next; } + async withRunnerCreateFence(namespace: string, operation: () => Promise): Promise { + const predecessor = this.runnerCreateFenceTails.get(namespace) ?? Promise.resolve(); + let release = (): void => {}; + const current = new Promise((resolve) => { release = resolve; }); + const tail = predecessor.then(async () => current); + this.runnerCreateFenceTails.set(namespace, tail); + await predecessor; + try { + return await operation(); + } finally { + release(); + if (this.runnerCreateFenceTails.get(namespace) === tail) this.runnerCreateFenceTails.delete(namespace); + } + } + + async withRunnerRetentionFence(input: RunnerRetentionFenceInput, operation: (facts: RunnerRetentionFenceFacts) => Promise): Promise { + if (this.runnerRetentionFences.has(input.runId)) throw runnerRetentionFenceConflict(input.runId); + const run = this.getRun(input.runId); + const command = this.getCommand(input.commandId); + if (command.runId !== run.id) throw new AgentRunError("schema-invalid", `command ${input.commandId} does not belong to run ${input.runId}`, { httpStatus: 409 }); + this.runnerRetentionFences.add(input.runId); + try { + const commands = this.listCommands(run.id, 0, 100); + const events = this.listEvents(run.id, 0, 500); + const result = await operation({ run, command, commands, events }); + if (input.terminalizeStalePending && !isTerminalRunStatus(run.status)) this.terminalizeRunnerRetention(run, commands, input.reason); + return result; + } finally { + this.runnerRetentionFences.delete(input.runId); + } + } + claimRun(runId: string, runnerId: string, leaseMs: number): RunRecord { + if (this.runnerRetentionFences.has(runId)) throw runnerRetentionFenceConflict(runId); const run = this.getRun(runId); if (isTerminalRunStatus(run.status)) throw new AgentRunError(run.failureKind ?? (run.status === "cancelled" ? "cancelled" : "schema-invalid"), `run ${runId} is already terminal: ${run.status}`, { httpStatus: 409 }); if (run.claimedBy && run.claimedBy !== runnerId && !isLeaseExpired(run.leaseExpiresAt)) throw new AgentRunError("runner-lease-conflict", `run ${runId} is already claimed`, { httpStatus: 409 }); @@ -440,6 +491,7 @@ export class MemoryAgentRunStore implements AgentRunStore { } heartbeat(runId: string, runnerId: string, leaseMs: number): RunRecord { + if (this.runnerRetentionFences.has(runId)) throw runnerRetentionFenceConflict(runId); const run = this.getRun(runId); if (isTerminalRunStatus(run.status)) return run; if (run.claimedBy !== runnerId) throw new AgentRunError("runner-lease-conflict", `run ${runId} is not claimed by ${runnerId}`, { httpStatus: 409 }); @@ -812,6 +864,20 @@ export class MemoryAgentRunStore implements AgentRunStore { return next; } + private terminalizeRunnerRetention(run: RunRecord, commands: CommandRecord[], reason: string): void { + const at = nowIso(); + for (const command of commands.filter((item) => !isTerminalCommandState(item.state))) { + this.commands.set(command.id, { ...command, state: "failed", updatedAt: at }); + this.cancelRunnerDispatchIntentForCommand(command.id, reason, at); + const prepared = this.prepareEvent(run.id, "backend_status", { phase: "command-terminal", commandId: command.id, state: "failed", terminalStatus: "failed", failureKind: "infra-failed", message: reason, retentionFence: true }); + this.commitPreparedEvent(prepared); + } + const next = this.updateRun(run.id, { status: "failed", terminalStatus: "failed", failureKind: "infra-failed", failureMessage: reason, claimedBy: null, leaseExpiresAt: null }); + this.cancelRunnerDispatchIntentsForRun(run.id, reason, at); + this.commitPreparedEvent(this.prepareEvent(run.id, "terminal_status", { terminalStatus: "failed", failureKind: "infra-failed", message: reason, retentionFence: true })); + this.touchSessionForRun(next, { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: run.id, terminalStatus: "failed", failureKind: "infra-failed", lastActivityAt: at }, { bumpVersion: true, at }); + } + private requireRunnerDispatchIntent(intentId: string): RunnerDispatchIntentRecord { const intent = this.runnerDispatchIntents.get(intentId); if (!intent) throw new AgentRunError("schema-invalid", `runner dispatch intent ${intentId} was not found`, { httpStatus: 404 }); @@ -1251,6 +1317,13 @@ function staleClaimError(component: string, claim: DurableQueueClaim): AgentRunE return new AgentRunError("runner-lease-conflict", `${component} claim is stale`, { httpStatus: 409, details: { id: claim.id, attemptCount: claim.attemptCount, valuesPrinted: false } }); } +function runnerRetentionFenceConflict(runId: string): AgentRunError { + return new AgentRunError("runner-lease-conflict", `run ${runId} is fenced by runner retention`, { + httpStatus: 409, + details: { reason: "runner-retention-db-fenced", runId, valuesPrinted: false }, + }); +} + function isActiveRunnerDispatchIntent(intent: RunnerDispatchIntentRecord): boolean { return intent.state === "pending" || intent.state === "dispatching" || intent.state === "retry"; } diff --git a/src/selftest/cases/20-runner-k8s-job.ts b/src/selftest/cases/20-runner-k8s-job.ts index ecb0d1b..79a3a58 100644 --- a/src/selftest/cases/20-runner-k8s-job.ts +++ b/src/selftest/cases/20-runner-k8s-job.ts @@ -278,6 +278,7 @@ process.exit(1); const retentionKubectl = path.join(context.tmp, "fake-kubectl-retention.js"); const retentionCreatedManifest = path.join(context.tmp, "created-retention-runner-job.json"); const retentionDeleted = path.join(context.tmp, "retention-deleted.json"); + const previousRetentionDbFacts = process.env.AGENTRUN_SELFTEST_RETENTION_DB_FACTS; await writeFile(retentionDeleted, "[]\n"); await writeFile(retentionKubectl, `#!/usr/bin/env bun const args = Bun.argv.slice(2); @@ -291,38 +292,59 @@ async function readStdin() { for await (const chunk of Bun.stdin.stream()) chunks.push(chunk); return Buffer.concat(chunks.map((chunk) => Buffer.from(chunk))).toString("utf8"); } -function runnerJob(name, createdAt, commandId) { +function runnerJob(name, createdAt, identity) { return { apiVersion: "batch/v1", kind: "Job", metadata: { name, namespace: "agentrun-v02", + uid: "uid-" + name, + resourceVersion: "1", creationTimestamp: createdAt, labels: { "app.kubernetes.io/part-of": "agentrun", "app.kubernetes.io/name": "agentrun-runner", "app.kubernetes.io/component": "runner", "agentrun.pikastech.local/lane": "v0.2", "job-name": name }, - annotations: { "agentrun.pikastech.local/run-id": "run-old", "agentrun.pikastech.local/command-id": commandId, "agentrun.pikastech.local/runner-id": "runner-old", token: null, secret: null }, + annotations: { "agentrun.pikastech.local/run-id": identity.runId, "agentrun.pikastech.local/command-id": identity.commandId, "agentrun.pikastech.local/runner-id": "runner-old", token: null, secret: null }, }, status: {}, }; } +function retentionJobs() { + const facts = JSON.parse(process.env.AGENTRUN_SELFTEST_RETENTION_DB_FACTS || "[]"); + return [ + runnerJob("agentrun-v02-runner-old-a", "2026-01-01T00:00:00Z", facts[0]), + runnerJob("agentrun-v02-runner-old-b", "2026-01-01T00:01:00Z", facts[1]), + ]; +} if (args[0] === "get" && args[1] === "jobs") { const deleted = await readDeleted(); const deletedJobs = new Set(deleted.filter((item) => item.resource === "job").map((item) => item.name)); - const items = [ - runnerJob("agentrun-v02-runner-old-a", "2026-01-01T00:00:00Z", "cmd-old-a"), - runnerJob("agentrun-v02-runner-old-b", "2026-01-01T00:01:00Z", "cmd-old-b"), - ].filter((item) => !deletedJobs.has(item.metadata.name)); + const items = retentionJobs().filter((item) => !deletedJobs.has(item.metadata.name)); console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items })); process.exit(0); } +if (args[0] === "get" && args[1] === "job" && args[2]) { + const deleted = await readDeleted(); + const item = retentionJobs().find((candidate) => candidate.metadata.name === args[2]); + if (!item || deleted.some((entry) => entry.resource === "job" && entry.name === args[2])) process.exit(1); + console.log(JSON.stringify(item)); + process.exit(0); +} if (args[0] === "get" && args[1] === "pods") { console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items: [] })); process.exit(0); } +if (args[0] === "get" && args[1] === "events") { console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items: [] })); process.exit(0); } if (args[0] === "delete") { const deleted = await readDeleted(); - const resource = args[1]; - const name = args[2] && !args[2].startsWith("-") ? args[2] : args.join(" "); + const rawIndex = args.indexOf("--raw"); + const rawPath = rawIndex >= 0 ? args[rawIndex + 1] : null; + const resource = rawPath?.includes("/jobs/") ? "job" : args[1]; + const name = rawPath ? decodeURIComponent(rawPath.split("/").filter(Boolean).at(-1)) : args[2] && !args[2].startsWith("-") ? args[2] : args.join(" "); + if (rawPath) { + const body = JSON.parse(await readStdin()); + const current = retentionJobs().find((candidate) => candidate.metadata.name === name); + if (!current || body?.preconditions?.uid !== current.metadata.uid || body?.preconditions?.resourceVersion !== current.metadata.resourceVersion) process.exit(1); + } deleted.push({ resource, name, args }); await writeDeleted(deleted); console.log(String(resource) + "/" + String(name).replace(/[^a-z0-9-]+/g, "-") + " deleted"); @@ -339,11 +361,12 @@ console.error("unsupported fake retention kubectl args: " + args.join(" ")); process.exit(1); `); await chmod(retentionKubectl, 0o755); + const retentionStore = new MemoryAgentRunStore(); const serverWithRetention = await startManagerServer({ port: 0, host: "127.0.0.1", sourceCommit: "self-test", - store: new MemoryAgentRunStore(), + store: retentionStore, runnerJobDefaults: { namespace: "agentrun-v02", managerUrl: "http://agentrun-mgr.agentrun-v02.svc.cluster.local:8080", @@ -358,6 +381,7 @@ process.exit(1); maxRunners: 1, cleanupOrder: "oldest-inactive-last-active-first", activeHeartbeatMaxAgeMs: 900_000, + stalePendingMaxAgeMs: 900_000, matchLabels: { "app.kubernetes.io/part-of": "agentrun", "app.kubernetes.io/name": "agentrun-runner", "app.kubernetes.io/component": "runner" }, jobNamePrefixes: ["agentrun-v02-runner"], ageBasedCleanup: { enabled: false, maxAgeHours: null }, @@ -366,6 +390,15 @@ process.exit(1); }); try { const retentionClient = new ManagerClient(serverWithRetention.baseUrl); + const oldItems = [ + await createRunWithCommand(retentionClient, context, "terminal old runner a", "selftest-runner-retention-old-a", 15_000), + await createRunWithCommand(retentionClient, context, "terminal old runner b", "selftest-runner-retention-old-b", 15_000), + ]; + for (const item of oldItems) { + retentionStore.finishCommand(item.commandId, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + retentionStore.finishRun(item.runId, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + } + process.env.AGENTRUN_SELFTEST_RETENTION_DB_FACTS = JSON.stringify(oldItems); const retentionItem = await createRunWithCommand(retentionClient, context, "retention job create", "selftest-runner-retention", 15_000); const retentionCreated = await retentionClient.post(`/api/v1/runs/${retentionItem.runId}/runner-jobs`, { commandId: retentionItem.commandId, attemptId: "attempt_selftest_retention" }) as JsonRecord; const retention = (retentionCreated.retention as JsonRecord).preCreateCleanup as JsonRecord; @@ -386,6 +419,8 @@ process.exit(1); assert.ok(deletedItems.some((item) => String((item.args as string[]).join(" ")).includes("agentrun.pikastech.local/runner-job=agentrun-v02-runner-old-a")), "old runner associated resources should be selected by runner-job label"); assertNoSecretLeak(retentionCreated); } finally { + if (previousRetentionDbFacts === undefined) delete process.env.AGENTRUN_SELFTEST_RETENTION_DB_FACTS; + else process.env.AGENTRUN_SELFTEST_RETENTION_DB_FACTS = previousRetentionDbFacts; await new Promise((resolve) => serverWithRetention.server.close(() => resolve())); } const protectedStore = new MemoryAgentRunStore(); @@ -418,6 +453,7 @@ if (args[0] === "get" && args[1] === "jobs") { process.exit(0); } if (args[0] === "get" && args[1] === "pods") { console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items: [] })); process.exit(0); } +if (args[0] === "get" && args[1] === "events") { console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items: [] })); process.exit(0); } if (args[0] === "delete") { const deleted = await readDeleted(); deleted.push({ resource: args[1], name: args[2], args }); await writeDeleted(deleted); console.log(String(args[1]) + "/" + String(args[2]) + " deleted"); process.exit(0); } console.error("unsupported fake protected kubectl args: " + args.join(" ")); process.exit(1); @@ -432,6 +468,7 @@ process.exit(1); maxRunners: 1, cleanupOrder: "oldest-inactive-last-active-first", activeHeartbeatMaxAgeMs: 0, + stalePendingMaxAgeMs: 900_000, matchLabels: { "app.kubernetes.io/part-of": "agentrun", "app.kubernetes.io/name": "agentrun-runner", "app.kubernetes.io/component": "runner" }, jobNamePrefixes: ["agentrun-v02-runner"], ageBasedCleanup: { enabled: false, maxAgeHours: null }, diff --git a/src/selftest/cases/22-runner-retention-stale-pending.ts b/src/selftest/cases/22-runner-retention-stale-pending.ts new file mode 100644 index 0000000..547361e --- /dev/null +++ b/src/selftest/cases/22-runner-retention-stale-pending.ts @@ -0,0 +1,665 @@ +import assert from "node:assert/strict"; +import { chmod, readFile, writeFile } from "node:fs/promises"; +import path from "node:path"; +import { AgentRunError } from "../../common/errors.js"; +import type { CommandRecord, JsonRecord, RunEvent, RunRecord } from "../../common/types.js"; +import { enforceRunnerRetentionBeforeCreate, type RunnerRetentionOptions, type RunnerRetentionSummary } from "../../mgr/runner-retention.js"; +import { MemoryAgentRunStore, type RunnerRetentionFenceFacts, type RunnerRetentionFenceInput } from "../../mgr/store.js"; +import type { SelfTestCase } from "../harness.js"; + +const stalePendingMaxAgeMs = 15 * 60_000; +const labels = { + "app.kubernetes.io/part-of": "agentrun", + "app.kubernetes.io/name": "agentrun-runner", + "app.kubernetes.io/component": "runner", +}; + +interface RetentionFact { + run: RunRecord; + command: CommandRecord; + events: RunEvent[]; +} + +interface RetentionFixture { + jobs: JsonRecord[]; + pods: JsonRecord[]; + events?: JsonRecord[]; + eventsAfterFirstRead?: JsonRecord[]; + racePods?: JsonRecord[]; + raceJobs?: JsonRecord[]; + eventsError?: boolean; + trackEventReads?: boolean; + churnEventMetadata?: boolean; + raceAfterEvents?: boolean; + raceOnCasDelete?: boolean; + bumpJobResourceVersionAfterPodDelete?: boolean; + bumpedJobResourceVersion?: string; + deleteErrorResources?: string[]; + retainDeletedJobs?: boolean; +} + +class RetentionFactStore extends MemoryAgentRunStore { + private readonly facts: RetentionFact[]; + private beforeFence: (() => void) | null = null; + + constructor(facts: RetentionFact[]) { + super(); + this.facts = facts; + } + + override getRun(runId: string): RunRecord { + const run = this.facts.find((item) => item.run.id === runId)?.run; + if (!run) throw new AgentRunError("schema-invalid", `run ${runId} was not found`, { httpStatus: 404 }); + return run; + } + + override getCommand(commandId: string): CommandRecord { + const command = this.facts.find((item) => item.command.id === commandId)?.command; + if (!command) throw new AgentRunError("schema-invalid", `command ${commandId} was not found`, { httpStatus: 404 }); + return command; + } + + override listCommands(runId: string, afterSeq: number, limit: number): CommandRecord[] { + this.getRun(runId); + return this.facts.map((item) => item.command).filter((command) => command.runId === runId && command.seq > afterSeq).slice(0, limit); + } + + override listEvents(runId: string, afterSeq: number, limit: number): RunEvent[] { + const fact = this.facts.find((item) => item.run.id === runId); + if (!fact) throw new AgentRunError("schema-invalid", `run ${runId} was not found`, { httpStatus: 404 }); + return fact.events.filter((event) => event.seq > afterSeq).slice(0, limit); + } + + setBeforeFence(callback: (() => void) | null): void { + this.beforeFence = callback; + } + + claimFact(runId: string, runnerId: string): void { + const fact = this.facts.find((item) => item.run.id === runId); + if (!fact) throw new Error(`missing retention fact ${runId}`); + const at = new Date(Date.now() + 60_000).toISOString(); + fact.run = { ...fact.run, status: "claimed", claimedBy: runnerId, leaseExpiresAt: at, updatedAt: at }; + fact.events.push(runEvent(runId, fact.events.length + 1, at, "backend_status", { phase: "run-claimed", runnerId, commandId: fact.command.id })); + } + + override async withRunnerRetentionFence(input: RunnerRetentionFenceInput, operation: (facts: RunnerRetentionFenceFacts) => Promise): Promise { + const beforeFence = this.beforeFence; + this.beforeFence = null; + beforeFence?.(); + const fact = this.facts.find((item) => item.run.id === input.runId && item.command.id === input.commandId); + if (!fact) throw new Error(`missing retention fence facts ${input.runId}/${input.commandId}`); + const result = await operation({ run: fact.run, command: fact.command, commands: this.listCommands(fact.run.id, 0, 100), events: this.listEvents(fact.run.id, 0, 500) }); + if (input.terminalizeStalePending) { + const at = new Date().toISOString(); + fact.command = { ...fact.command, state: "failed", updatedAt: at }; + fact.run = { ...fact.run, status: "failed", terminalStatus: "failed", failureKind: "infra-failed", failureMessage: input.reason, claimedBy: null, leaseExpiresAt: null, updatedAt: at }; + } + return result; + } +} + +const selfTest: SelfTestCase = async (context) => { + const kubectlCommand = path.join(context.root, "src/selftest/fixtures/fake-kubectl-retention-stale-pending.mjs"); + const fixturePath = path.join(context.tmp, "runner-retention-stale-pending-fixture.json"); + const statePath = path.join(context.tmp, "runner-retention-stale-pending-state.json"); + const previousFixturePath = process.env.AGENTRUN_SELFTEST_RETENTION_FIXTURE_PATH; + const previousStatePath = process.env.AGENTRUN_SELFTEST_RETENTION_STATE_PATH; + process.env.AGENTRUN_SELFTEST_RETENTION_FIXTURE_PATH = fixturePath; + process.env.AGENTRUN_SELFTEST_RETENTION_STATE_PATH = statePath; + await chmod(kubectlCommand, 0o755); + + try { + await assertTwentyStalePendingSelectsOnlyOldest({ fixturePath, statePath, kubectlCommand }); + await assertLegacyNoUserMessageCandidate({ fixturePath, statePath, kubectlCommand }); + await assertEventOnlyFailedMountCandidate({ fixturePath, statePath, kubectlCommand }); + await assertProtectedMatrix({ fixturePath, statePath, kubectlCommand }); + await assertPendingToRunningRaceFailsClosed({ fixturePath, statePath, kubectlCommand }); + await assertDbClaimAfterRecheckFailsClosed({ fixturePath, statePath, kubectlCommand }); + await assertPodCasRaceFailsClosed({ fixturePath, statePath, kubectlCommand }); + await assertDeleteFailureFailsClosed({ fixturePath, statePath, kubectlCommand }); + await assertDeletionPostconditionFailsClosed({ fixturePath, statePath, kubectlCommand }); + await assertMemoryDbFenceRejectsConcurrentClaim(); + await assertRunnerCreateFenceSerializesNamespace(); + } finally { + restoreEnv("AGENTRUN_SELFTEST_RETENTION_FIXTURE_PATH", previousFixturePath); + restoreEnv("AGENTRUN_SELFTEST_RETENTION_STATE_PATH", previousStatePath); + } + + return { + name: "runner-retention-stale-pending", + tests: [ + "stale-pending-typed-candidate", + "formal-user-message-is-not-runner-progress", + "legacy-no-user-message-candidate", + "user-message-command-binding-protected", + "incoming-selects-only-oldest", + "no-incoming-does-not-delete", + "fresh-pending-protected", + "running-ready-protected", + "fresh-claim-protected", + "tool-execution-protected", + "facts-unavailable-protected", + "waiting-reason-required", + "diagnostics-redacted", + "failed-mount-event-only-evidence", + "container-started-pending-protected", + "pod-identity-mismatch-protected", + "event-facts-unavailable-protected", + "terminal-ledger-running-pod-protected", + "selected-runner-race-fails-closed", + "db-claim-after-recheck-fails-closed", + "pod-resource-version-cas-fails-closed", + "runner-delete-cas-exit-code-checked", + "runner-delete-postcondition-checked", + "memory-db-fence-rejects-concurrent-claim", + "namespace-runner-create-fence-prevents-capacity-overshoot", + ], + }; +}; + +export default selfTest; + +async function assertTwentyStalePendingSelectsOnlyOldest(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const now = Date.now(); + const facts = Array.from({ length: 20 }, (_, index) => pendingFact(index, new Date(now - (60 - index) * 60_000).toISOString())); + const fixture = { + jobs: facts.map((fact, index) => runnerJob(fact, index)), + pods: facts.map((fact, index) => runnerPod(fact, index, "pending")), + events: facts.map((fact, index) => failedMountEvent(fact, index)), + trackEventReads: true, + churnEventMetadata: true, + bumpJobResourceVersionAfterPodDelete: true, + bumpedJobResourceVersion: "2", + }; + await writeFixture(input, fixture); + const store = new RetentionFactStore(facts); + const options = retentionOptions(input.kubectlCommand, 20); + + const withinLimit = await enforceRunnerRetentionBeforeCreate({ store, options, incomingRunnerCount: 0 }); + assert.equal(withinLimit.liveRunnerCountBefore, 20); + assert.equal(withinLimit.inactiveCandidateCount, 20); + assert.equal(withinLimit.selectedRunnerCount, 0); + assert.equal(deletionOperations(await readDeletionState(input.statePath)).length, 0); + + const cleanup = await enforceRunnerRetentionBeforeCreate({ store, options, incomingRunnerCount: 1 }); + assert.equal(cleanup.liveRunnerCountBefore, 20); + assert.equal(cleanup.liveRunnerCountAfter, 19); + assert.equal(cleanup.inactiveCandidateCount, 20); + assert.equal(cleanup.selectedRunnerCount, 1); + assert.equal(cleanup.deletedRunnerJobCount, 1); + assert.equal(cleanup.deletedRunnerPodCount, 0); + assert.equal(cleanup.deletedAssociatedResourceCount, 2); + const selected = firstSummary(cleanup, "selected"); + assert.equal(selected.name, "agentrun-v02-runner-stale-00"); + assert.equal(selected.candidateKind, "stale-pending"); + assert.equal(selected.classificationReason, "stale-pending-no-business-progress"); + assert.equal(selected.runId, facts[0]?.run.id); + assert.equal(selected.commandId, facts[0]?.command.id); + const observation = ((selected.podObservations as JsonRecord[])[0] ?? {}) as JsonRecord; + assert.equal(observation.phase, "Pending"); + assert.equal(observation.stale, true); + assert.equal(observation.runnerContainerStarted, false); + assert.equal(observation.runnerContainerReady, false); + assert.equal(observation.podStartTime, facts[0]?.run.createdAt); + assert.ok(Number(observation.ageMs) >= stalePendingMaxAgeMs); + assert.equal((((observation.waiting as JsonRecord[])[0] ?? {}) as JsonRecord).reason, "ContainerCreating"); + assert.ok((observation.waiting as JsonRecord[]).some((item) => item.source === "event" && item.reason === "FailedMount")); + const serialized = JSON.stringify(cleanup); + assert.equal(serialized.includes("abc123-secret-value"), false); + assert.equal(serialized.includes("REDACTED"), true); + const deleted = await readDeletionState(input.statePath); + assert.deepEqual(deleted.filter((item) => item.resource === "job").map((item) => item.name), ["agentrun-v02-runner-stale-00"]); + const jobCas = deleted.find((item) => item.resource === "cas-request" && item.targetKind === "job") ?? {}; + assert.equal(jobCas.expectedResourceVersion, "2"); + assert.equal(deleted.filter((item) => item.resource === "read" && item.name === "events").length, 4); +} + +async function assertProtectedMatrix(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const oldAt = new Date(Date.now() - 60 * 60_000).toISOString(); + const freshAt = new Date(Date.now() - 60_000).toISOString(); + const base = pendingFact(100, oldAt); + await assertProtected(input, [pendingFact(101, freshAt)], [runnerPod(pendingFact(101, freshAt), 101, "pending")], "stale-pending-window-active"); + await assertProtected(input, [base], [runnerPod(base, 100, "running")], "runner-pod-started-running-or-ready"); + await assertProtected(input, [base], [runnerPod(base, 100, "ready")], "runner-pod-started-running-or-ready"); + await assertProtected(input, [base], [runnerPod(base, 100, "started-pending")], "runner-pod-started-running-or-ready"); + const claimed = pendingFact(102, oldAt, "claimed"); + await assertProtected(input, [claimed], [runnerPod(claimed, 102, "pending")], "fresh-active-run"); + const toolExecuting = pendingFact(103, oldAt, "tool-executing"); + await assertProtected(input, [toolExecuting], [runnerPod(toolExecuting, 103, "failed-config")], "business-progress-observed-or-unverifiable"); + const unknown = pendingFact(104, oldAt); + await assertProtected(input, [], [runnerPod(unknown, 104, "pending")], "db-facts-unavailable", unknown); + await assertProtected(input, [base], [runnerPod(base, 100, "no-waiting")], "pod-waiting-reason-unavailable"); + await assertProtected(input, [base], [], "pod-facts-unavailable"); + await assertProtected(input, [base], [runnerPod(base, 100, "pending")], "pod-waiting-event-facts-unavailable", base, { eventsError: true }); + await assertProtected(input, [base], [runnerPod(base, 100, "pending"), runnerPod(base, 101, "running", 100)], "runner-pod-started-running-or-ready"); + const terminalRunning = pendingFact(105, oldAt); + terminalRunning.run.status = "completed"; + terminalRunning.run.terminalStatus = "completed"; + terminalRunning.command.state = "completed"; + await assertProtected(input, [terminalRunning], [runnerPod(terminalRunning, 105, "running")], "runner-pod-started-running-or-ready"); + const identityMismatchPod = runnerPod(base, 100, "pending"); + ((identityMismatchPod.metadata as JsonRecord).annotations as JsonRecord)["agentrun.pikastech.local/command-id"] = "cmd-mismatched"; + await assertProtected(input, [base], [identityMismatchPod], "pod-identity-mismatch"); + const mismatchedUserMessage = pendingFact(106, oldAt); + const userMessage = mismatchedUserMessage.events.find((event) => event.type === "user_message"); + assert.ok(userMessage); + userMessage.payload.commandId = "cmd-mismatched"; + await assertProtected(input, [mismatchedUserMessage], [runnerPod(mismatchedUserMessage, 106, "failed-config")], "business-progress-observed-or-unverifiable"); +} + +async function assertLegacyNoUserMessageCandidate(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(89, new Date(Date.now() - 60 * 60_000).toISOString()); + fact.events = fact.events + .filter((event) => event.type !== "user_message") + .map((event, index) => ({ ...event, id: `evt-${fact.run.id}-${index + 1}`, seq: index + 1 })); + await writeFixture(input, { + jobs: [runnerJob(fact, 89)], + pods: [runnerPod(fact, 89, "failed-config")], + }); + const summary = await enforceRunnerRetentionBeforeCreate({ store: new RetentionFactStore([fact]), options: retentionOptions(input.kubectlCommand, 1), incomingRunnerCount: 1 }); + assert.equal(summary.selectedRunnerCount, 1); + assert.equal(firstSummary(summary, "selected").candidateKind, "stale-pending"); +} + +async function assertEventOnlyFailedMountCandidate(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(90, new Date(Date.now() - 60 * 60_000).toISOString()); + const fixture = { + jobs: [runnerJob(fact, 90)], + pods: [runnerPod(fact, 90, "event-only")], + events: [failedMountEvent(fact, 90)], + }; + await writeFixture(input, fixture); + const summary = await enforceRunnerRetentionBeforeCreate({ store: new RetentionFactStore([fact]), options: retentionOptions(input.kubectlCommand, 1), incomingRunnerCount: 0 }); + assert.equal(summary.inactiveCandidateCount, 1); + assert.equal(summary.selectedRunnerCount, 0); + const candidate = firstSummary(summary, "inactiveCandidates"); + assert.equal(candidate.candidateKind, "stale-pending"); + const observation = ((candidate.podObservations as JsonRecord[])[0] ?? {}) as JsonRecord; + assert.equal(observation.runnerContainerState, "unknown"); + assert.deepEqual((observation.waiting as JsonRecord[]).map((item) => [item.source, item.reason]), [["event", "FailedMount"]]); + assert.equal((await readDeletionState(input.statePath)).length, 0); +} + +async function assertPendingToRunningRaceFailsClosed(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(91, new Date(Date.now() - 60 * 60_000).toISOString()); + const runningPod = runnerPod(fact, 91, "running"); + (runningPod.metadata as JsonRecord).resourceVersion = "2"; + await writeFixture(input, { + jobs: [runnerJob(fact, 91)], + pods: [runnerPod(fact, 91, "pending")], + events: [failedMountEvent(fact, 91)], + raceAfterEvents: true, + racePods: [runningPod], + }); + + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ + store: new RetentionFactStore([fact]), + options: retentionOptions(input.kubectlCommand, 1), + incomingRunnerCount: 1, + })); + assert.equal(error.details?.reason, "runner-retention-selected-runner-changed"); + assert.equal(error.details?.changeReason, "protected-active"); + const current = (error.details?.current ?? {}) as JsonRecord; + assert.equal(current.protectedReason, "runner-pod-started-running-or-ready"); + assert.equal(deletionOperations(await readDeletionState(input.statePath)).length, 0); +} + +async function assertDbClaimAfterRecheckFailsClosed(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(94, new Date(Date.now() - 60 * 60_000).toISOString()); + await writeFixture(input, { jobs: [runnerJob(fact, 94)], pods: [runnerPod(fact, 94, "pending")], events: [failedMountEvent(fact, 94)] }); + const store = new RetentionFactStore([fact]); + store.setBeforeFence(() => store.claimFact(fact.run.id, "runner-raced-after-recheck")); + + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ store, options: retentionOptions(input.kubectlCommand, 1), incomingRunnerCount: 1 })); + assert.equal(error.details?.reason, "runner-retention-db-fence-rejected"); + assert.equal(error.details?.changeReason, "run-became-active"); + assert.equal(error.details?.claimedBy, "runner-raced-after-recheck"); + const state = await readDeletionState(input.statePath); + assert.equal(deletionOperations(state).length, 0); + assert.equal(state.filter((item) => item.resource === "cas-request").length, 0); +} + +async function assertPodCasRaceFailsClosed(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(95, new Date(Date.now() - 60 * 60_000).toISOString()); + const runningPod = runnerPod(fact, 95, "running"); + (runningPod.metadata as JsonRecord).resourceVersion = "2"; + await writeFixture(input, { + jobs: [runnerJob(fact, 95)], + pods: [runnerPod(fact, 95, "pending")], + events: [failedMountEvent(fact, 95)], + raceOnCasDelete: true, + racePods: [runningPod], + }); + + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ store: new RetentionFactStore([fact]), options: retentionOptions(input.kubectlCommand, 1), incomingRunnerCount: 1 })); + assert.equal(error.details?.reason, "runner-retention-k8s-cas-rejected"); + assert.equal(error.details?.changeReason, "object-identity-or-version-changed"); + assert.equal(error.details?.resourceKind, "pod"); + assert.equal(error.details?.expectedResourceVersion, "1"); + const state = await readDeletionState(input.statePath); + assert.equal(deletionOperations(state).length, 0); + assert.equal(state.filter((item) => item.resource === "cas-request" && item.targetKind === "pod").length, 1); +} + +async function assertDeleteFailureFailsClosed(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(92, new Date(Date.now() - 60 * 60_000).toISOString()); + await writeFixture(input, { + jobs: [runnerJob(fact, 92)], + pods: [runnerPod(fact, 92, "pending")], + events: [failedMountEvent(fact, 92)], + deleteErrorResources: ["job"], + }); + + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ + store: new RetentionFactStore([fact]), + options: retentionOptions(input.kubectlCommand, 1), + incomingRunnerCount: 1, + })); + assert.equal(error.details?.reason, "runner-retention-k8s-cas-rejected"); + assert.equal(error.details?.changeReason, "delete-request-failed"); + assert.equal(error.details?.resourceKind, "job"); + assert.equal(error.details?.exitCode, 42); + const state = await readDeletionState(input.statePath); + assert.deepEqual(deletionOperations(state).filter((item) => item.resource === "job"), []); + assert.deepEqual(deletionOperations(state).filter((item) => item.resource === "pod").map((item) => item.name), ["agentrun-v02-runner-stale-92-pod-92"]); +} + +async function assertDeletionPostconditionFailsClosed(input: { fixturePath: string; statePath: string; kubectlCommand: string }): Promise { + const fact = pendingFact(93, new Date(Date.now() - 60 * 60_000).toISOString()); + await writeFixture(input, { + jobs: [runnerJob(fact, 93)], + pods: [runnerPod(fact, 93, "pending")], + events: [failedMountEvent(fact, 93)], + retainDeletedJobs: true, + }); + + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ + store: new RetentionFactStore([fact]), + options: retentionOptions(input.kubectlCommand, 1), + incomingRunnerCount: 1, + })); + assert.equal(error.details?.reason, "runner-retention-postcondition-failed"); + assert.equal(error.details?.targetBeforeCreate, 0); + assert.equal(error.details?.liveRunnerCountBefore, 1); + assert.equal(error.details?.liveRunnerCountAfter, 1); + assert.equal(error.details?.deletedRunnerJobCount, 1); + assert.deepEqual(deletionOperations(await readDeletionState(input.statePath)).filter((item) => item.resource === "job").map((item) => item.name), ["agentrun-v02-runner-stale-93"]); +} + +async function assertMemoryDbFenceRejectsConcurrentClaim(): Promise { + const store = new MemoryAgentRunStore(); + const run = store.createRun({ tenantId: "unidesk", projectId: "pikasTech/agentrun", workspaceRef: { kind: "host-path", path: "/tmp/agentrun-retention-memory" }, providerId: "NC01", backendProfile: "codex", executionPolicy: { sandbox: "workspace-write", approval: "never", timeoutMs: 60_000, network: "default", secretScope: { allowCredentialEcho: false, providerCredentials: [] } }, traceSink: null }); + const command = store.createCommand(run.id, { type: "turn", payload: { prompt: "memory retention fence" } }); + let enterFence = (): void => {}; + let releaseFence = (): void => {}; + const entered = new Promise((resolve) => { enterFence = resolve; }); + const release = new Promise((resolve) => { releaseFence = resolve; }); + const fenced = Promise.resolve(store.withRunnerRetentionFence({ runId: run.id, commandId: command.id, terminalizeStalePending: true, reason: "retention memory self-test" }, async () => { + enterFence(); + await release; + })); + await entered; + assert.throws( + () => store.claimRun(run.id, "runner-during-retention", 60_000), + (error) => error instanceof AgentRunError && error.failureKind === "runner-lease-conflict" && error.details?.reason === "runner-retention-db-fenced", + ); + releaseFence(); + await fenced; + assert.equal(store.getRun(run.id).status, "failed"); + assert.equal(store.getCommand(command.id).state, "failed"); + assert.throws(() => store.claimRun(run.id, "runner-after-retention", 60_000), (error) => error instanceof AgentRunError && error.httpStatus === 409); +} + +async function assertRunnerCreateFenceSerializesNamespace(): Promise { + const store = new MemoryAgentRunStore(); + const maxRunners = 20; + let liveRunnerCount = 19; + let enterFirst = (): void => {}; + let releaseFirst = (): void => {}; + const firstEntered = new Promise((resolve) => { enterFirst = resolve; }); + const firstRelease = new Promise((resolve) => { releaseFirst = resolve; }); + let secondEntered = false; + const first = Promise.resolve(store.withRunnerCreateFence("agentrun-nc01-v02", async () => { + assert.ok(liveRunnerCount + 1 <= maxRunners); + enterFirst(); + await firstRelease; + liveRunnerCount++; + return "first"; + })); + await firstEntered; + const second = Promise.resolve(store.withRunnerCreateFence("agentrun-nc01-v02", async () => { + secondEntered = true; + if (liveRunnerCount + 1 > maxRunners) throw new AgentRunError("infra-failed", "capacity fence rejected second incoming runner", { httpStatus: 409 }); + liveRunnerCount++; + return "second"; + })); + await Promise.resolve(); + assert.equal(secondEntered, false); + releaseFirst(); + const results = await Promise.allSettled([first, second]); + assert.equal(results[0]?.status, "fulfilled"); + assert.equal(results[1]?.status, "rejected"); + assert.equal(secondEntered, true); + assert.equal(liveRunnerCount, maxRunners); +} + +async function assertProtected( + input: { fixturePath: string; statePath: string; kubectlCommand: string }, + facts: RetentionFact[], + pods: JsonRecord[], + expectedReason: string, + identityFact: RetentionFact = facts[0] as RetentionFact, + fixtureOptions: { events?: JsonRecord[]; eventsError?: boolean } = {}, +): Promise { + const firstPodMetadata = (pods[0]?.metadata ?? {}) as JsonRecord; + const firstPodLabels = (firstPodMetadata.labels ?? {}) as JsonRecord; + const jobName = typeof firstPodLabels["job-name"] === "string" ? firstPodLabels["job-name"] : "agentrun-v02-runner-stale-999"; + const job = runnerJob(identityFact, 999); + const jobMetadata = job.metadata as JsonRecord; + jobMetadata.name = jobName; + jobMetadata.labels = { ...((jobMetadata.labels ?? {}) as JsonRecord), "job-name": jobName }; + const fixture = { jobs: [job], pods, events: fixtureOptions.events ?? [], ...(fixtureOptions.eventsError === true ? { eventsError: true } : {}) }; + await writeFixture(input, fixture); + const error = await retentionError(() => enforceRunnerRetentionBeforeCreate({ store: new RetentionFactStore(facts), options: retentionOptions(input.kubectlCommand, 1), incomingRunnerCount: 1 })); + assert.equal(error.details?.reason, "runner-retention-no-safe-candidate"); + assert.equal(error.details?.inactiveCandidateCount, 0); + const protectedRunner = (((error.details?.protectedActiveRunners as JsonRecord[] | undefined) ?? [])[0] ?? {}) as JsonRecord; + assert.equal(protectedRunner.protectedReason, expectedReason, JSON.stringify(error.details)); + assert.equal((await readDeletionState(input.statePath)).length, 0); +} + +function pendingFact(index: number, createdAt: string, variant: "pending" | "claimed" | "tool-executing" = "pending"): RetentionFact { + const suffix = String(index).padStart(2, "0"); + const runId = `run-retention-${suffix}`; + const commandId = `cmd-retention-${suffix}`; + const progressedAt = new Date(Date.parse(createdAt) + 60_000).toISOString(); + const claimed = variant === "claimed"; + const toolExecuting = variant === "tool-executing"; + const run: RunRecord = { + id: runId, + tenantId: "unidesk", + projectId: "pikasTech/agentrun", + workspaceRef: { kind: "host-path", path: "/tmp/agentrun-retention-selftest" }, + sessionRef: null, + resourceBundleRef: null, + providerId: "NC01", + backendProfile: "codex", + executionPolicy: { sandbox: "workspace-write", approval: "never", timeoutMs: 60_000, network: "default", secretScope: { allowCredentialEcho: false, providerCredentials: [] } }, + traceSink: null, + status: claimed ? "claimed" : "pending", + terminalStatus: null, + failureKind: null, + failureMessage: null, + cancelEpoch: 0, + cancelRequestId: null, + cancelRequestedAt: null, + cancelReason: null, + createdAt, + updatedAt: claimed ? progressedAt : createdAt, + claimedBy: claimed ? `runner-retention-${suffix}` : null, + leaseExpiresAt: claimed ? new Date(Date.now() + 60_000).toISOString() : null, + }; + const command: CommandRecord = { + id: commandId, + runId, + seq: 1, + type: "turn", + payload: { prompt: "retention self-test" }, + idempotencyKey: `retention-${suffix}`, + payloadHash: `hash-${suffix}`, + state: toolExecuting ? "acknowledged" : "pending", + cancelEpoch: 0, + cancelRequestId: null, + cancelRequestedAt: null, + cancelReason: null, + createdAt, + updatedAt: toolExecuting ? progressedAt : createdAt, + acknowledgedAt: toolExecuting ? progressedAt : null, + }; + const events: RunEvent[] = [ + runEvent(runId, 1, createdAt, "backend_status", { phase: "run-created" }), + runEvent(runId, 2, createdAt, "user_message", { commandId, traceId: `trc-retention-${suffix}`, userMessageId: `msg-retention-${suffix}`, text: "retention self-test", source: "agentrun-turn-command-payload", valuesPrinted: false }), + runEvent(runId, 3, createdAt, "backend_status", { phase: "command-created", commandId }), + runEvent(runId, 4, createdAt, "backend_status", { phase: "runner-job-created", commandId }), + runEvent(runId, 5, createdAt, "backend_status", { phase: "runner-dispatch-completed", commandId }), + ...(claimed ? [runEvent(runId, 6, progressedAt, "backend_status", { phase: "run-claimed", commandId })] : []), + ...(toolExecuting ? [runEvent(runId, 6, progressedAt, "tool_call", { commandId, tool: "shell", phase: "started" })] : []), + ]; + return { run, command, events }; +} + +function runEvent(runId: string, seq: number, createdAt: string, type: RunEvent["type"], payload: JsonRecord): RunEvent { + return { id: `evt-${runId}-${seq}`, runId, seq, type, payload, createdAt }; +} + +function runnerJob(fact: RetentionFact, index: number): JsonRecord { + const suffix = String(index).padStart(2, "0"); + const name = `agentrun-v02-runner-stale-${suffix}`; + return { + apiVersion: "batch/v1", + kind: "Job", + metadata: { + name, + namespace: "agentrun-nc01-v02", + uid: `job-uid-${suffix}`, + resourceVersion: "1", + creationTimestamp: fact.run.createdAt, + labels: { ...labels, "job-name": name }, + annotations: identityAnnotations(fact), + }, + status: {}, + }; +} + +function runnerPod(fact: RetentionFact, index: number, state: "pending" | "running" | "ready" | "started-pending" | "failed-config" | "event-only" | "no-waiting", jobIndex = index): JsonRecord { + const suffix = String(index).padStart(2, "0"); + const jobName = `agentrun-v02-runner-stale-${String(jobIndex).padStart(2, "0")}`; + const name = `${jobName}-pod-${suffix}`; + const waitingStatus = { + name: "runner", + ready: false, + started: state === "started-pending", + restartCount: 0, + state: { waiting: { reason: state === "failed-config" ? "CreateContainerConfigError" : "ContainerCreating", message: "FailedMount: secret agentrun-provider-test not found; Authorization: Bearer abc123-secret-value" } }, + }; + const runningStatus = { + name: "runner", + ready: true, + started: true, + restartCount: 0, + state: { running: { startedAt: fact.run.createdAt } }, + }; + return { + apiVersion: "v1", + kind: "Pod", + metadata: { + name, + namespace: "agentrun-nc01-v02", + uid: `pod-uid-${suffix}`, + resourceVersion: "1", + creationTimestamp: fact.run.createdAt, + labels: { ...labels, "job-name": jobName }, + annotations: identityAnnotations(fact), + ownerReferences: [{ kind: "Job", name: jobName }], + }, + status: { + phase: state === "running" ? "Running" : "Pending", + startTime: fact.run.createdAt, + conditions: state === "no-waiting" || state === "event-only" ? [] : [{ type: "Ready", status: state === "running" || state === "ready" ? "True" : "False", reason: state === "running" || state === "ready" ? null : "ContainersNotReady", message: "runner container is waiting" }], + containerStatuses: state === "no-waiting" || state === "event-only" ? [] : [state === "running" ? runningStatus : { ...waitingStatus, ready: state === "ready" }], + }, + }; +} + +function failedMountEvent(fact: RetentionFact, index: number): JsonRecord { + const suffix = String(index).padStart(2, "0"); + const jobName = `agentrun-v02-runner-stale-${suffix}`; + return { + apiVersion: "v1", + kind: "Event", + metadata: { name: `${jobName}.failed-mount`, namespace: "agentrun-nc01-v02" }, + involvedObject: { kind: "Pod", name: `${jobName}-pod-${suffix}` }, + type: "Warning", + reason: "FailedMount", + message: "MountVolume.SetUp failed for volume provider: secret agentrun-provider-test not found; token=abc123-secret-value", + count: 12, + firstTimestamp: fact.run.createdAt, + lastTimestamp: fact.run.createdAt, + }; +} + +function identityAnnotations(fact: RetentionFact): JsonRecord { + return { + "agentrun.pikastech.local/run-id": fact.run.id, + "agentrun.pikastech.local/command-id": fact.command.id, + "agentrun.pikastech.local/runner-id": `runner-${fact.run.id}`, + }; +} + +function retentionOptions(kubectlCommand: string, maxRunners: number): RunnerRetentionOptions { + return { + namespace: "agentrun-nc01-v02", + maxRunners, + cleanupOrder: "oldest-inactive-last-active-first", + activeHeartbeatMaxAgeMs: 900_000, + stalePendingMaxAgeMs, + matchLabels: labels, + jobNamePrefixes: ["agentrun-v02-runner"], + ageBasedCleanup: { enabled: false, maxAgeHours: null }, + kubectlCommand, + }; +} + +async function writeFixture(input: { fixturePath: string; statePath: string }, fixture: RetentionFixture): Promise { + await writeFile(input.fixturePath, `${JSON.stringify(fixture, null, 2)}\n`, "utf8"); + await writeFile(input.statePath, "[]\n", "utf8"); +} + +async function readDeletionState(statePath: string): Promise { + return JSON.parse(await readFile(statePath, "utf8")) as JsonRecord[]; +} + +function deletionOperations(state: JsonRecord[]): JsonRecord[] { + return state.filter((item) => item.resource === "job" || item.resource === "pod" || item.resource === "secret" || item.resource === "pvc"); +} + +async function retentionError(run: () => Promise): Promise { + try { + await run(); + } catch (error) { + if (error instanceof AgentRunError) return error; + throw error; + } + throw new Error("expected runner retention to fail closed"); +} + +function firstSummary(summary: RunnerRetentionSummary, key: string): JsonRecord { + const items = summary[key]; + assert.ok(Array.isArray(items) && items.length > 0, `${key} summary is required`); + return items[0] as JsonRecord; +} + +function restoreEnv(key: string, value: string | undefined): void { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; +} diff --git a/src/selftest/fixtures/fake-kubectl-retention-stale-pending.mjs b/src/selftest/fixtures/fake-kubectl-retention-stale-pending.mjs new file mode 100755 index 0000000..c589587 --- /dev/null +++ b/src/selftest/fixtures/fake-kubectl-retention-stale-pending.mjs @@ -0,0 +1,142 @@ +#!/usr/bin/env bun + +const args = Bun.argv.slice(2); +const fixturePath = process.env.AGENTRUN_SELFTEST_RETENTION_FIXTURE_PATH; +const statePath = process.env.AGENTRUN_SELFTEST_RETENTION_STATE_PATH; + +if (!fixturePath || !statePath) { + console.error("AGENTRUN_SELFTEST_RETENTION_FIXTURE_PATH and AGENTRUN_SELFTEST_RETENTION_STATE_PATH are required"); + process.exit(2); +} + +async function readJson(path, fallback) { + try { + return JSON.parse(await Bun.file(path).text()); + } catch { + return fallback; + } +} + +function podJobName(item) { + return item?.metadata?.labels?.["job-name"] + ?? item?.metadata?.ownerReferences?.find((owner) => owner?.kind === "Job")?.name + ?? item?.metadata?.name; +} + +const fixture = await readJson(fixturePath, { jobs: [], pods: [] }); +const state = await readJson(statePath, []); +const eventReadCount = state.filter((item) => item.resource === "read" && item.name === "events").length; +const deletedJobs = fixture.retainDeletedJobs === true + ? new Set() + : new Set(state.filter((item) => item.resource === "job").map((item) => item.name)); +const deletedPods = new Set(state.filter((item) => item.resource === "pod").map((item) => item.name)); + +function currentJob(name) { + const base = fixture.jobs.find((item) => item?.metadata?.name === name); + if (!base || fixture.bumpJobResourceVersionAfterPodDelete !== true) return base; + const podDeleted = fixture.pods.some((pod) => podJobName(pod) === name && deletedPods.has(pod?.metadata?.name)); + if (!podDeleted) return base; + return { ...base, metadata: { ...(base.metadata ?? {}), resourceVersion: fixture.bumpedJobResourceVersion ?? "2" } }; +} + +if (args[0] === "get" && args[1] === "job" && args[2]) { + const item = currentJob(String(args[2])); + if (!item || deletedJobs.has(item?.metadata?.name)) { + console.error(`Error from server (NotFound): job/${String(args[2])} was not found`); + process.exit(1); + } + console.log(JSON.stringify(item)); + process.exit(0); +} + +if (args[0] === "get" && args[1] === "jobs") { + const items = fixture.jobs.filter((item) => !deletedJobs.has(item?.metadata?.name)); + console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items })); + process.exit(0); +} + +if (args[0] === "get" && args[1] === "pods") { + const sourcePods = fixture.raceAfterEvents === true && eventReadCount > 0 && Array.isArray(fixture.racePods) + ? fixture.racePods + : fixture.pods; + const items = sourcePods.filter((item) => !deletedJobs.has(podJobName(item)) && !deletedPods.has(item?.metadata?.name)); + console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items })); + process.exit(0); +} + +if (args[0] === "get" && args[1] === "events") { + if (fixture.eventsError === true) { + console.error("events access denied by self-test fixture"); + process.exit(1); + } + const sourceItems = eventReadCount > 0 && Array.isArray(fixture.eventsAfterFirstRead) + ? fixture.eventsAfterFirstRead + : fixture.events ?? []; + const items = fixture.churnEventMetadata === true + ? sourceItems.map((item) => ({ + ...item, + count: Number(item?.count ?? 0) + eventReadCount, + lastTimestamp: new Date(Date.parse(item?.lastTimestamp ?? "2026-01-01T00:00:00.000Z") + eventReadCount * 1000).toISOString(), + message: `${String(item?.message ?? "warning")} occurrence=${eventReadCount + 1}`, + })) + : sourceItems; + if (fixture.trackEventReads === true || fixture.raceAfterEvents === true || Array.isArray(fixture.eventsAfterFirstRead)) { + state.push({ resource: "read", name: "events", sequence: eventReadCount + 1, valuesPrinted: false }); + await Bun.write(statePath, `${JSON.stringify(state, null, 2)}\n`); + } + console.log(JSON.stringify({ apiVersion: "v1", kind: "List", items })); + process.exit(0); +} + +if (args[0] === "delete") { + const rawIndex = args.indexOf("--raw"); + if (rawIndex >= 0) { + const resourcePath = String(args[rawIndex + 1] ?? ""); + const resource = resourcePath.includes("/jobs/") ? "job" : resourcePath.includes("/pods/") ? "pod" : "unknown"; + const name = decodeURIComponent(resourcePath.split("/").filter(Boolean).at(-1) ?? "unknown"); + const deleteOptions = await readStdinJson(); + const expected = deleteOptions?.preconditions ?? {}; + const source = fixture.raceOnCasDelete === true && Array.isArray(fixture.racePods) ? fixture.racePods : fixture.pods; + const current = resource === "job" + ? (fixture.raceOnCasDelete === true && Array.isArray(fixture.raceJobs) ? fixture.raceJobs.find((item) => item?.metadata?.name === name) : currentJob(name)) + : source.find((item) => item?.metadata?.name === name); + state.push({ resource: "cas-request", name, targetKind: resource, expectedUid: expected.uid ?? null, expectedResourceVersion: expected.resourceVersion ?? null, valuesPrinted: false }); + await Bun.write(statePath, `${JSON.stringify(state, null, 2)}\n`); + if (!expected.uid || !expected.resourceVersion || !current || current?.metadata?.uid !== expected.uid || current?.metadata?.resourceVersion !== expected.resourceVersion) { + console.error(`Error from server (Conflict): ${resource}/${name} resourceVersion or UID precondition failed`); + process.exit(1); + } + if (Array.isArray(fixture.deleteErrorResources) && fixture.deleteErrorResources.includes(resource)) { + console.error(`self-test forced ${resource} deletion failure`); + process.exit(42); + } + state.push({ resource, name, selector: null, cas: true, valuesPrinted: false }); + await Bun.write(statePath, `${JSON.stringify(state, null, 2)}\n`); + console.log(JSON.stringify({ kind: "Status", apiVersion: "v1", status: "Success", details: { name, kind: resource } })); + process.exit(0); + } + const resource = String(args[1] ?? "unknown"); + const directName = args[2] && !String(args[2]).startsWith("-") ? String(args[2]) : null; + const selectorIndex = args.indexOf("-l"); + const selector = selectorIndex >= 0 ? String(args[selectorIndex + 1] ?? "") : null; + const name = directName ?? selector ?? "selected"; + if (Array.isArray(fixture.deleteErrorResources) && fixture.deleteErrorResources.includes(resource)) { + console.error(`self-test forced ${resource} deletion failure`); + process.exit(42); + } + state.push({ resource, name, selector, valuesPrinted: false }); + await Bun.write(statePath, `${JSON.stringify(state, null, 2)}\n`); + console.log(`${resource}/${name.replace(/[^a-z0-9.-]+/giu, "-")} deleted`); + process.exit(0); +} + +async function readStdinJson() { + try { + return JSON.parse(await Bun.stdin.text()); + } catch { + return null; + } +} + +console.error(`unsupported fake retention kubectl args: ${args.join(" ")}`); +process.exit(1); diff --git a/src/selftest/integration/postgres-runner-retention-fence.ts b/src/selftest/integration/postgres-runner-retention-fence.ts new file mode 100644 index 0000000..52e2872 --- /dev/null +++ b/src/selftest/integration/postgres-runner-retention-fence.ts @@ -0,0 +1,80 @@ +import assert from "node:assert/strict"; +import { AgentRunError } from "../../common/errors.js"; +import { createPostgresAgentRunStore } from "../../mgr/postgres-store.js"; + +const connectionString = process.env.AGENTRUN_SELFTEST_DATABASE_URL?.trim(); +if (!connectionString) throw new Error("AGENTRUN_SELFTEST_DATABASE_URL is required; this integration test never falls back to memory"); + +const primary = await createPostgresAgentRunStore({ connectionString }); +const concurrent = await createPostgresAgentRunStore({ connectionString }); + +try { + await assertNamespaceCreateFence(); + await assertRetentionFenceBlocksClaim(); + console.log(JSON.stringify({ ok: true, tests: ["postgres-namespace-create-fence", "postgres-run-retention-fence-blocks-claim"] })); +} finally { + await Promise.all([primary.close(), concurrent.close()]); +} + +async function assertNamespaceCreateFence(): Promise { + let firstEntered = (): void => {}; + let releaseFirst = (): void => {}; + const entered = new Promise((resolve) => { firstEntered = resolve; }); + const release = new Promise((resolve) => { releaseFirst = resolve; }); + let secondEntered = false; + const first = primary.withRunnerCreateFence("agentrun-postgres-retention-selftest", async () => { + firstEntered(); + await release; + return "first"; + }); + await entered; + const second = concurrent.withRunnerCreateFence("agentrun-postgres-retention-selftest", async () => { + secondEntered = true; + return "second"; + }); + await delay(50); + assert.equal(secondEntered, false); + releaseFirst(); + assert.deepEqual(await Promise.all([first, second]), ["first", "second"]); +} + +async function assertRetentionFenceBlocksClaim(): Promise { + const run = await primary.createRun({ + tenantId: "selftest", + projectId: "pikasTech/agentrun", + workspaceRef: { kind: "host-path", path: "/tmp/agentrun-postgres-retention-selftest" }, + providerId: "NC01", + backendProfile: "codex", + executionPolicy: { sandbox: "workspace-write", approval: "never", timeoutMs: 60_000, network: "default", secretScope: { allowCredentialEcho: false, providerCredentials: [] } }, + traceSink: null, + }); + const command = await primary.createCommand(run.id, { type: "turn", payload: { prompt: "postgres retention fence self-test" } }); + let fenceEntered = (): void => {}; + let releaseFence = (): void => {}; + const entered = new Promise((resolve) => { fenceEntered = resolve; }); + const release = new Promise((resolve) => { releaseFence = resolve; }); + const fenced = primary.withRunnerRetentionFence({ runId: run.id, commandId: command.id, terminalizeStalePending: true, reason: "postgres retention fence self-test" }, async (facts) => { + assert.equal(facts.run.id, run.id); + assert.equal(facts.command.id, command.id); + fenceEntered(); + await release; + }); + await entered; + let claimSettled = false; + const claim = concurrent.claimRun(run.id, "runner-postgres-retention-race", 60_000) + .then(() => null, (error: unknown) => error) + .finally(() => { claimSettled = true; }); + await delay(50); + assert.equal(claimSettled, false); + releaseFence(); + await fenced; + const claimError = await claim; + assert.ok(claimError instanceof AgentRunError); + assert.equal(claimError.httpStatus, 409); + assert.equal((await primary.getRun(run.id)).status, "failed"); + assert.equal((await primary.getCommand(command.id)).state, "failed"); +} + +async function delay(ms: number): Promise { + await new Promise((resolve) => setTimeout(resolve, ms)); +}