diff --git a/src/mgr/postgres-store.ts b/src/mgr/postgres-store.ts index f27382e..4ddaff6 100644 --- a/src/mgr/postgres-store.ts +++ b/src/mgr/postgres-store.ts @@ -5,8 +5,8 @@ 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, QueueAttemptListResult, QueueAttemptRecord, QueueAttemptRef, QueueCommanderSnapshot, QueueReadCursorRecord, QueueRetryActivation, QueueRetryReservation, 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 { ActivateQueueRetryAttemptInput, AgentRunStore, CreateRunIdentity, DurableQueueClaim, EventOutboxConfig, ListQueueAttemptsInput, ListQueueTasksInput, ListSessionsInput, RecordQueueRetryAttemptFailureInput, ReserveQueueRetryAttemptInput, RunnerRetentionFenceFacts, RunnerRetentionFenceInput, SaveRunnerJobInput, SessionEventPageInput, SessionTurnAdmission, SessionTurnAdmissionInput, StoreHealth, UpdateQueueTaskAttemptInput } from "./store.js"; -import { activeSessionAdmissionConflict, assertQueueAttemptActivationReplay, assertQueueRetryable, assertQueueTaskPayloadHash, assertRunCreateReplay, assertSessionBoundary, assertSessionCommandReplay, assertSessionTurnAdmissionContract, buildQueueStats, buildQueueTaskSummary, buildSessionSummary, cancelStagePayload, clampQueueLimit, clampSessionLimit, commandStateFromTerminal, fenceLateEventForCancelledRun, invalidSessionAdmissionProjection, isLeaseExpired, isSessionOutputEvent, isTerminalCommandState, isTerminalQueueAttemptState, isTerminalQueueTaskState, isTerminalRunStatus, latestQueueAttempt, lateWriteRejectedPayload, nextQueueRetryIndex, parseQueueCursor, parseSessionCursor, queueAttemptRef, queueRetryReservation, queueRetryStateError, queueTaskMatchesCommander, queueTaskPayloadHash, queueTaskSort, runCreatedEventPayload, runnerDispatchFailureKind, runnerDispatchFailureReason, sessionListFilters, sessionMatchesListState, sessionRefFromRecord, sessionSort, sessionTitleFromCommand, statusFromTerminal, summarizeSessionRef, titleFromMetadata } from "./store.js"; +import type { ActivateQueueRetryAttemptInput, AgentRunStore, CreateRunIdentity, DurableQueueClaim, EventOutboxConfig, ListQueueAttemptsInput, ListQueueTasksInput, ListSessionsInput, RecordQueueRetryAttemptFailureInput, ReserveQueueRetryAttemptInput, RunnerRetentionFenceFacts, RunnerRetentionFenceInput, SaveRunnerJobInput, SessionEventPageInput, SessionSteerAdmission, SessionSteerAdmissionInput, SessionTurnAdmission, SessionTurnAdmissionInput, StoreHealth, UpdateQueueTaskAttemptInput } from "./store.js"; +import { activeSessionAdmissionConflict, assertQueueAttemptActivationReplay, assertQueueRetryable, assertQueueTaskPayloadHash, assertRunCreateReplay, assertSessionBoundary, assertSessionCommandReplay, assertSessionSteerAdmissionContract, assertSessionTurnAdmissionContract, buildQueueStats, buildQueueTaskSummary, buildSessionSummary, cancelStagePayload, clampQueueLimit, clampSessionLimit, commandStateFromTerminal, fenceLateEventForCancelledRun, invalidSessionAdmissionProjection, isLeaseExpired, isSessionOutputEvent, isTerminalCommandState, isTerminalQueueAttemptState, isTerminalQueueTaskState, isTerminalRunStatus, latestQueueAttempt, lateWriteRejectedPayload, nextQueueRetryIndex, parseQueueCursor, parseSessionCursor, queueAttemptRef, queueRetryReservation, queueRetryStateError, queueTaskMatchesCommander, queueTaskPayloadHash, queueTaskSort, receivableActiveSessionTurn, runCreatedEventPayload, runnerDispatchFailureKind, runnerDispatchFailureReason, sessionListFilters, sessionMatchesListState, sessionRefFromRecord, sessionSort, sessionTitleFromCommand, statusFromTerminal, summarizeSessionRef, titleFromMetadata } from "./store.js"; import { backendCapabilitiesSqlValues, mergeBackendCapability } from "../common/backend-profiles.js"; import { normalizeRunEventPayload, requireEventType, requireExternallyAppendableEventType, userMessagePayloadForCommand } from "../common/events.js"; import { buildKafkaEventOutboxRecord, canonicalAgentRunEventPartitionKey } from "./event-outbox.js"; @@ -611,6 +611,28 @@ CREATE TABLE IF NOT EXISTS agentrun_schema_migrations ( return this.withTransaction(async (client) => await this.createCommandWithClient(client, runId, input)); } + async admitSessionSteer(input: SessionSteerAdmissionInput): Promise { + assertSessionSteerAdmissionContract(input); + return this.withTransaction(async (client) => { + await client.query("SELECT pg_advisory_xact_lock(hashtext($1), hashtext($2))", ["agentrun-session-turn-admission", input.sessionId]); + const observedSessionResult = await client.query("SELECT * FROM agentrun_sessions WHERE session_id = $1", [input.sessionId]); + const observedSession = observedSessionResult.rows[0] ? sessionFromRow(observedSessionResult.rows[0]) : null; + if (!observedSession?.activeRunId || !observedSession.activeCommandId) return null; + const runResult = await client.query("SELECT * FROM agentrun_runs WHERE id = $1 FOR UPDATE", [observedSession.activeRunId]); + const commandResult = await client.query("SELECT * FROM agentrun_commands WHERE id = $1 FOR UPDATE", [observedSession.activeCommandId]); + if (!runResult.rows[0] || !commandResult.rows[0]) return null; + const run = runFromRow(runResult.rows[0]); + const targetCommand = commandFromRow(commandResult.rows[0]); + const sessionResult = await client.query("SELECT * FROM agentrun_sessions WHERE session_id = $1 FOR UPDATE", [input.sessionId]); + const session = sessionResult.rows[0] ? sessionFromRow(sessionResult.rows[0]) : null; + if (!session) return null; + const receivable = receivableActiveSessionTurn(session, run, targetCommand); + if (!receivable) return null; + const command = await this.createCommandWithClient(client, run.id, input.command); + return { run, command, targetCommand, ...receivable }; + }); + } + async admitSessionTurn(input: SessionTurnAdmissionInput): Promise { assertSessionTurnAdmissionContract(input); return this.withTransaction(async (client) => { @@ -1102,7 +1124,11 @@ CREATE TABLE IF NOT EXISTS agentrun_schema_migrations ( await fenceActiveRunnerDispatchIntents(client, { runId }, `run terminalized as ${status}`, at); if (result.threadId && run.sessionRef?.sessionId) await this.upsertSessionThread(client, run, result.threadId, result.turnId ?? null); await this.appendEventWithLockedRun(client, runId, "terminal_status", { terminalStatus: result.terminalStatus, failureKind: result.failureKind, message: result.failureMessage }); - await this.touchSessionForRun(client, run, { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: runId, terminalStatus: result.terminalStatus, failureKind: result.failureKind, lastActivityAt: run.updatedAt }, { bumpVersion: true, at: run.updatedAt }); + const sessionId = run.sessionRef?.sessionId; + const sessionResult = sessionId ? await client.query("SELECT * FROM agentrun_sessions WHERE session_id = $1 FOR UPDATE", [sessionId]) : null; + const session = sessionResult?.rows[0] ? sessionFromRow(sessionResult.rows[0]) : null; + const stillCurrent = !session || ((session.activeRunId === null || session.activeRunId === runId) && (session.lastRunId === null || session.lastRunId === runId)); + if (stillCurrent) await this.touchSessionForRun(client, run, { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: runId, terminalStatus: result.terminalStatus, failureKind: result.failureKind, lastActivityAt: run.updatedAt }, { bumpVersion: true, at: run.updatedAt }); return run; }); } diff --git a/src/mgr/server.ts b/src/mgr/server.ts index 68dc569..6f546c6 100644 --- a/src/mgr/server.ts +++ b/src/mgr/server.ts @@ -2,7 +2,7 @@ import type { Server } from "node:http"; import { createServer } from "node:http"; import type { AddressInfo } from "node:net"; import type { AgentRunStore, ListQueueTasksInput, ListSessionsInput, SessionEventPageInput } from "./store.js"; -import { assertSessionBoundary, openAgentRunStoreFromEnv } from "./store.js"; +import { assertSessionBoundary, openAgentRunStoreFromEnv, receivableActiveSessionTurn } from "./store.js"; import { agentRunKafkaConfig, type AgentRunKafkaConfig } from "../common/kafka-events.js"; import { AgentRunError, errorToJson } from "../common/errors.js"; import { asRecord, stableHash, validateBackendProfile, validateCreateCommand, validateCreateQueueTask, validateCreateRun, validateQueueTaskState, validateSessionListState, validateSessionRunnerJobInput } from "../common/validation.js"; @@ -1146,27 +1146,17 @@ async function sendToSession(input: { store: AgentRunStore; sessionId: string; b if (commandIdempotencyKey) commandBody.idempotencyKey = commandIdempotencyKey; const request = { method: "POST", path: `/api/v1/runs/${active.run.id}/commands`, commandType: "steer", payloadBytes: jsonByteLength(payload), valuesPrinted: false }; if (dryRun) return sessionSendPlan(input.sessionId, "steer", active, request, null); - const command = await input.store.createCommand(active.run.id, validateCreateCommand(commandBody)); - void emitAgentRunOtelSpan("command_created", active.run, process.env, { command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "steer", commandType: command.type, commandState: command.state } }); - void emitAgentRunOtelSpan("session_send", active.run, process.env, { command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "steer", reusedActiveRun: true } }); - return sessionSendResponse({ sessionId: input.sessionId, decision: "steer", run: active.run, command, runnerJob: null, runnerAdmission: null, activeBefore: active, dryRun: false }); + const admitted = await input.store.admitSessionSteer({ sessionId: input.sessionId, command: validateCreateCommand(commandBody) }); + if (admitted) { + const activeAtAdmission = { run: admitted.run, command: admitted.targetCommand, reason: admitted.reason, leaseExpired: admitted.leaseExpired }; + void emitAgentRunOtelSpan("command_created", admitted.run, process.env, { command: admitted.command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "steer", commandType: admitted.command.type, commandState: admitted.command.state } }); + void emitAgentRunOtelSpan("session_send", admitted.run, process.env, { command: admitted.command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "steer", reusedActiveRun: true } }); + return sessionSendResponse({ sessionId: input.sessionId, decision: "steer", run: admitted.run, command: admitted.command, runnerJob: null, runnerAdmission: null, activeBefore: activeAtAdmission, dryRun: false }); + } } const pendingAdmission = existing ? await matchingPendingSessionTurn(input.store, existing, payload, commandIdempotencyKey) : null; - const idleRun = existing ? await idleReceivableRun(input.store, existing) : null; - if (idleRun) { - const commandBody: JsonRecord = { type: "turn", payload }; - if (commandIdempotencyKey) commandBody.idempotencyKey = commandIdempotencyKey; - const request = { method: "POST", path: `/api/v1/runs/${idleRun.run.id}/commands`, commandType: "turn", payloadBytes: jsonByteLength(payload), reusedIdleRun: true, valuesPrinted: false }; - if (dryRun) return sessionSendPlan(input.sessionId, "turn", active, request, null, idleRunSummary(idleRun)); - const command = await input.store.createCommand(idleRun.run.id, validateCreateCommand(commandBody)); - const run = await input.store.getRun(idleRun.run.id); - void emitAgentRunOtelSpan("command_created", run, process.env, { command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "turn", reusedIdleRun: true, commandType: command.type, commandState: command.state } }); - void emitAgentRunOtelSpan("session_send", run, process.env, { command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "turn", createRunnerJob: false, reusedIdleRun: true } }); - return sessionSendResponse({ sessionId: input.sessionId, decision: "turn", run, command, runnerJob: null, runnerAdmission: null, activeBefore: active, reusedIdleRun: idleRunSummary(idleRun), dryRun: false }); - } - const runRecord = asRecord(record.run ?? record.runBase ?? null, "sessionSend.run"); const resourceRun = existing ? await input.store.getLatestSessionResourceRun(input.sessionId) : null; const runBody = inheritSessionRunResources(sessionSendRunBody(input.sessionId, runRecord), resourceRun); @@ -1184,7 +1174,7 @@ async function sendToSession(input: { store: AgentRunStore; sessionId: string; b runnerJobBytes: createRunnerJob ? jsonByteLength(runnerJobBody) : 0, valuesPrinted: false, }; - if (dryRun) return sessionSendPlan(input.sessionId, "turn", active, request, runBody, null); + if (dryRun) return sessionSendPlan(input.sessionId, "turn", active, request, runBody); if (createRunnerJob && !input.runnerDispatcherOptions.enabled) { throw new AgentRunError("infra-failed", "durable runner dispatch is not enabled for session turn admission", { httpStatus: 503, @@ -1202,7 +1192,7 @@ async function sendToSession(input: { store: AgentRunStore; sessionId: string; b const runnerAdmission = createRunnerJob ? runnerAdmissionSummary(input.sessionId, run, command, admitted.disposition, admitted.recoveredPriorPartialWrite) : runnerAdmissionNotRequested(input.sessionId, run, command, admitted.disposition); const runnerJob = createRunnerJob ? plannedRunnerJob(runnerAdmission) : null; void emitAgentRunOtelSpan("session_send", run, process.env, { command, startTimeMs: startedAt, kind: 2, attributes: { "http.method": "POST", "http.route": "/api/v1/sessions/:sessionId/send", "http.status_code": 200, decision: "turn", createRunnerJob } }); - return sessionSendResponse({ sessionId: input.sessionId, decision: "turn", run, command, runnerJob, runnerAdmission, activeBefore: active, reusedIdleRun: null, dryRun: false }); + return sessionSendResponse({ sessionId: input.sessionId, decision: "turn", run, command, runnerJob, runnerAdmission, activeBefore: null, dryRun: false }); } async function matchingPendingSessionTurn(store: AgentRunStore, session: SessionRecord, payload: JsonRecord, idempotencyKey: string | null): Promise<{ run: RunRecord; command: CommandRecord; intent: RunnerDispatchIntentRecord | null } | null> { @@ -1231,29 +1221,8 @@ function frozenRunContract(run: RunRecord): CreateRunInput { async function activeReceivableCommand(store: AgentRunStore, session: SessionRecord): Promise<{ run: RunRecord; command: CommandRecord; reason: string; leaseExpired: boolean } | null> { if (!session.activeRunId || !session.activeCommandId) return null; const [run, command] = await Promise.all([store.getRun(session.activeRunId), store.getCommand(session.activeCommandId)]); - if (run.sessionRef?.sessionId !== session.sessionId || command.runId !== run.id) return null; - if (runIsTerminal(run) || commandIsTerminal(command)) return null; - if (command.type !== "turn") return null; - const leaseExpired = run.leaseExpiresAt ? Date.parse(run.leaseExpiresAt) <= Date.now() : false; - if (leaseExpired) return null; - if (command.state !== "acknowledged") return null; - if (run.status !== "claimed" && run.status !== "running") return null; - return { run, command, reason: "active-turn-running", leaseExpired }; -} - -async function idleReceivableRun(store: AgentRunStore, session: SessionRecord): Promise<{ run: RunRecord; lastCommand: CommandRecord; reason: string; leaseExpired: boolean } | null> { - if (!session.lastRunId || !session.lastCommandId) return null; - try { - const [run, lastCommand] = await Promise.all([store.getRun(session.lastRunId), store.getCommand(session.lastCommandId)]); - if (run.sessionRef?.sessionId !== session.sessionId || lastCommand.runId !== run.id) return null; - if (runIsTerminal(run) || !commandIsTerminal(lastCommand)) return null; - const leaseExpired = run.leaseExpiresAt ? Date.parse(run.leaseExpiresAt) <= Date.now() : false; - if (leaseExpired) return null; - if (run.status !== "claimed" && run.status !== "running") return null; - return { run, lastCommand, reason: "idle-run-reusable", leaseExpired }; - } catch { - return null; - } + const receivable = receivableActiveSessionTurn(session, run, command); + return receivable ? { run, command, ...receivable } : null; } function sessionSendPayload(record: JsonRecord): JsonRecord { @@ -1295,7 +1264,7 @@ function inheritSessionRunResources(runBody: JsonRecord, resourceRun: RunRecord }; } -function sessionSendPlan(sessionId: string, decision: "steer" | "turn", active: Awaited>, request: JsonRecord, runBody: JsonRecord | null, reusedIdleRun: JsonRecord | null = null): JsonRecord { +function sessionSendPlan(sessionId: string, decision: "steer" | "turn", active: Awaited>, request: JsonRecord, runBody: JsonRecord | null): JsonRecord { return { action: "session-send-plan", dryRun: true, @@ -1304,7 +1273,6 @@ function sessionSendPlan(sessionId: string, decision: "steer" | "turn", active: decision, internalCommandType: decision, activeBefore: active ? activeBeforeSummary(active) : null, - reusedIdleRun, request, ...(runBody ? { run: { bodyBytes: jsonByteLength(runBody), sessionRef: summarizeSendSessionRef(runBody), valuesPrinted: false } } : {}), next: { confirm: managerActionDescriptor({ action: "send-session", operation: "send", resourceKind: "session", resourceName: sessionId, sessionId, inputKind: "prompt" }), note: "Remove --dry-run to perform the mutation. Manager will decide internal steer vs turn from durable session state." }, @@ -1356,7 +1324,7 @@ function plannedRunnerJob(admission: JsonRecord): JsonRecord { return { action: "runner-dispatch-admitted", state: admission.state, runId: admission.runId, commandId: admission.commandId, runnerJobId: admission.plannedRunnerJobId, dispatchIntentId: admission.dispatchIntentId, mutation: false, durable: true, valuesPrinted: false }; } -function sessionSendResponse(input: { sessionId: string; decision: "steer" | "turn"; run: RunRecord; command: CommandRecord; runnerJob: JsonValue; runnerAdmission: JsonRecord | null; activeBefore: Awaited>; reusedIdleRun?: JsonRecord | null; dryRun: false }): JsonRecord { +function sessionSendResponse(input: { sessionId: string; decision: "steer" | "turn"; run: RunRecord; command: CommandRecord; runnerJob: JsonValue; runnerAdmission: JsonRecord | null; activeBefore: Awaited>; dryRun: false }): JsonRecord { return { action: "session-send", dryRun: input.dryRun, @@ -1370,7 +1338,6 @@ function sessionSendResponse(input: { sessionId: string; decision: "steer" | "tu runnerJob: input.runnerJob, runnerAdmission: input.runnerAdmission, activeBefore: input.activeBefore ? activeBeforeSummary(input.activeBefore) : null, - reusedIdleRun: input.reusedIdleRun ?? null, pollActions: [ managerActionDescriptor({ action: "inspect-session", operation: "describe", resourceKind: "session", resourceName: input.sessionId, sessionId: input.sessionId, readerId: "cli" }), managerActionDescriptor({ action: "poll-trace", operation: "events", resourceKind: "run", resourceName: input.run.id, runId: input.run.id, commandId: input.command.id, sessionId: input.sessionId, afterSeq: 0, limit: 100 }), @@ -1404,19 +1371,6 @@ function activeBeforeSummary(active: NonNullable>>): JsonRecord { - return { - runId: idleRun.run.id, - commandId: idleRun.lastCommand.id, - runStatus: idleRun.run.status, - commandState: idleRun.lastCommand.state, - leaseExpiresAt: idleRun.run.leaseExpiresAt, - leaseExpired: idleRun.leaseExpired, - reason: idleRun.reason, - valuesPrinted: false, - }; -} - function summarizeSendSessionRef(runBody: JsonRecord): JsonRecord { const ref = asJsonRecord(runBody.sessionRef) ?? {}; return { sessionId: optionalString(ref.sessionId), conversationId: optionalString(ref.conversationId), threadId: optionalString(ref.threadId), metadataKeys: Object.keys(asJsonRecord(ref.metadata) ?? {}).sort(), valuesPrinted: false }; diff --git a/src/mgr/store.ts b/src/mgr/store.ts index e6217a2..7db8634 100644 --- a/src/mgr/store.ts +++ b/src/mgr/store.ts @@ -61,6 +61,19 @@ export interface SessionTurnAdmission extends JsonRecord { partialWrite: false; } +export interface SessionSteerAdmissionInput { + sessionId: string; + command: CreateCommandInput; +} + +export interface SessionSteerAdmission extends JsonRecord { + run: RunRecord; + command: CommandRecord; + targetCommand: CommandRecord; + reason: "active-turn-running"; + leaseExpired: false; +} + export interface AgentRunStore { health(): MaybePromise; createRun(input: CreateRunInput, identity?: CreateRunIdentity): MaybePromise; @@ -68,6 +81,7 @@ export interface AgentRunStore { listEvents(runId: string, afterSeq: number, limit: number): MaybePromise; listEventsForCommand(runId: string, commandId: string, limit: number): MaybePromise; createCommand(runId: string, input: CreateCommandInput): MaybePromise; + admitSessionSteer(input: SessionSteerAdmissionInput): MaybePromise; admitSessionTurn(input: SessionTurnAdmissionInput): MaybePromise; getCommand(commandId: string): MaybePromise; listCommands(runId: string, afterSeq: number, limit: number): MaybePromise; @@ -336,6 +350,19 @@ export class MemoryAgentRunStore implements AgentRunStore { return attachRunnerDispatchIntent(command, intent); } + admitSessionSteer(input: SessionSteerAdmissionInput): SessionSteerAdmission | null { + assertSessionSteerAdmissionContract(input); + const session = this.sessions.get(input.sessionId) ?? null; + if (!session?.activeRunId || !session.activeCommandId) return null; + const run = this.runs.get(session.activeRunId) ?? null; + const targetCommand = this.commands.get(session.activeCommandId) ?? null; + if (!run || !targetCommand) return null; + const receivable = receivableActiveSessionTurn(session, run, targetCommand); + if (!receivable) return null; + const command = this.createCommand(run.id, input.command); + return { run: this.getRun(run.id), command, targetCommand, ...receivable }; + } + admitSessionTurn(input: SessionTurnAdmissionInput): SessionTurnAdmission { assertSessionTurnAdmissionContract(input); const session = this.sessions.get(input.sessionId) ?? null; @@ -712,7 +739,10 @@ export class MemoryAgentRunStore implements AgentRunStore { this.cancelRunnerDispatchIntentsForRun(runId, `run terminalized as ${status}`, next.updatedAt); if (result.threadId && next.sessionRef?.sessionId) this.upsertSessionThread(next, result.threadId, result.turnId ?? null); this.appendEvent(runId, "terminal_status", { terminalStatus: result.terminalStatus, failureKind: result.failureKind, message: result.failureMessage }); - this.touchSessionForRun(this.getRun(runId), { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: runId, terminalStatus: result.terminalStatus, failureKind: result.failureKind, lastActivityAt: next.updatedAt }, { bumpVersion: true, at: next.updatedAt }); + const sessionId = next.sessionRef?.sessionId; + const session = sessionId ? this.sessions.get(sessionId) ?? null : null; + const stillCurrent = !session || ((session.activeRunId === null || session.activeRunId === runId) && (session.lastRunId === null || session.lastRunId === runId)); + if (stillCurrent) this.touchSessionForRun(this.getRun(runId), { executionState: "terminal", activeRunId: null, activeCommandId: null, lastRunId: runId, terminalStatus: result.terminalStatus, failureKind: result.failureKind, lastActivityAt: next.updatedAt }, { bumpVersion: true, at: next.updatedAt }); return next; } @@ -1637,6 +1667,23 @@ export function assertSessionTurnAdmissionContract(input: SessionTurnAdmissionIn if (input.command.type !== "turn") throw new AgentRunError("schema-invalid", "session turn admission only accepts turn commands", { httpStatus: 400 }); } +export function assertSessionSteerAdmissionContract(input: SessionSteerAdmissionInput): void { + if (!input.sessionId.trim()) throw new AgentRunError("schema-invalid", "session steer admission requires a non-empty sessionId", { httpStatus: 400 }); + if (input.command.type !== "steer") throw new AgentRunError("schema-invalid", "session steer admission only accepts steer commands", { httpStatus: 400 }); + if (input.command.dispatch) throw new AgentRunError("schema-invalid", "session steer admission does not accept runner dispatch", { httpStatus: 400 }); +} + +export function receivableActiveSessionTurn(session: SessionRecord, run: RunRecord, command: CommandRecord, now = Date.now()): { reason: "active-turn-running"; leaseExpired: false } | null { + if (session.activeRunId !== run.id || session.activeCommandId !== command.id) return null; + if (run.sessionRef?.sessionId !== session.sessionId || command.runId !== run.id) return null; + if (isTerminalRunStatus(run.status) || isTerminalCommandState(command.state)) return null; + if (command.type !== "turn" || command.state !== "acknowledged") return null; + if (run.status !== "claimed" && run.status !== "running") return null; + const leaseExpired = run.leaseExpiresAt ? Date.parse(run.leaseExpiresAt) <= now : false; + if (leaseExpired) return null; + return { reason: "active-turn-running", leaseExpired: false }; +} + export function assertSessionCommandReplay(existing: CommandRecord, input: CreateCommandInput): void { if (existing.type !== input.type || existing.payloadHash !== stableHash(input.payload) || (existing.idempotencyKey ?? null) !== (input.idempotencyKey ?? null)) { throw new AgentRunError("schema-invalid", "session turn admission replay does not match the existing command", { diff --git a/src/selftest/cases/66-session-turn-admission.ts b/src/selftest/cases/66-session-turn-admission.ts index 3fd2334..90daa50 100644 --- a/src/selftest/cases/66-session-turn-admission.ts +++ b/src/selftest/cases/66-session-turn-admission.ts @@ -8,7 +8,7 @@ import { ManagerClient } from "../../mgr/client.js"; import type { RunnerJobDefaults } from "../../mgr/kubernetes-runner-job.js"; import { dispatchRunnerIntentsOnce, type RunnerDispatcherOptions } from "../../mgr/runner-dispatcher.js"; import { startManagerServer } from "../../mgr/server.js"; -import { MemoryAgentRunStore, type SessionTurnAdmissionInput } from "../../mgr/store.js"; +import { MemoryAgentRunStore, type SessionSteerAdmission, type SessionSteerAdmissionInput, type SessionTurnAdmissionInput } from "../../mgr/store.js"; import type { SelfTestCase } from "../harness.js"; const image = "127.0.0.1:5000/agentrun/agentrun-mgr@sha256:1111111111111111111111111111111111111111111111111111111111111111"; @@ -28,6 +28,8 @@ const selfTest: SelfTestCase = async (context) => { await assertKubernetesCreateFailureTerminalizes(fakeKubectl); await assertSuccessfulSameSessionContinuity(fakeKubectl); await assertSessionFollowUpInheritsAipodResources(fakeKubectl); + await assertActiveSessionSendSteersAtomically(fakeKubectl); + await assertSessionFollowUpTerminalRaceCreatesNewRunner(fakeKubectl); } finally { restoreEnv("AGENTRUN_SELFTEST_KUBECTL_MODE", previousMode); } @@ -44,6 +46,8 @@ const selfTest: SelfTestCase = async (context) => { "kubectl-create-failure-terminalizes-and-releases-session", "successful-same-session-continuity", "session-follow-up-inherits-aipod-workspace-bundles-and-tool-credentials", + "active-session-send-steers-under-session-admission-lock", + "session-follow-up-terminal-race-creates-new-runner-without-empty-old-command", ], }; }; @@ -320,6 +324,82 @@ async function assertSessionFollowUpInheritsAipodResources(kubectlCommand: strin } } +async function assertSessionFollowUpTerminalRaceCreatesNewRunner(kubectlCommand: string): Promise { + const store = new TerminalizingSteerStore(); + const sessionId = "sess_artificer_terminal_race"; + const initial = store.createRun(artificerRunInput(sessionId)); + const initialCommand = store.createCommand(initial.id, validateCreateCommand({ type: "turn", payload: { prompt: "initial artificer work" } })); + store.claimRun(initial.id, "runner_artificer_terminal_race", 60_000); + store.ackCommand(initialCommand.id); + + const server = await startManagerServer({ + store, + sourceCommit: "selftest", + runnerJobDefaults: defaults(kubectlCommand), + runnerDispatcherOptions: { ...dispatcherOptions(3), intervalMs: 60_000 }, + kafkaOutboxRelayOptions: { enabled: false }, + }); + try { + const client = new ManagerClient(server.baseUrl); + const response = await client.post(`/api/v1/sessions/${sessionId}/send`, sessionSendBody(sessionId, "follow up during terminal convergence", "session-terminal-race-follow-up")) as JsonRecord; + const nextRun = response.run as RunRecord; + const nextCommand = response.command as CommandRecord; + const admission = response.runnerAdmission as JsonRecord; + assert.equal(response.decision, "turn"); + assert.equal(response.activeBefore, null); + assert.equal(Object.hasOwn(response, "reusedIdleRun"), false); + assert.notEqual(nextRun.id, initial.id); + assert.notEqual(nextCommand.id, initialCommand.id); + assert.equal(admission.state, "pending"); + assert.equal(admission.retryAuthority, "dispatcher"); + assert.equal(store.getRunnerDispatchIntent(nextCommand.id)?.state, "pending"); + assert.deepEqual(store.listCommands(initial.id, 0, 100).map((command) => command.id), [initialCommand.id]); + assert.equal(store.getCommand(initialCommand.id).state, "completed"); + + store.finishRun(initial.id, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + const session = store.getSession(sessionId); + assert.equal(session?.activeRunId, nextRun.id); + assert.equal(session?.activeCommandId, nextCommand.id); + assert.equal(session?.lastRunId, nextRun.id); + assert.equal(session?.lastCommandId, nextCommand.id); + } finally { + await closeServer(server.server); + } +} + +async function assertActiveSessionSendSteersAtomically(kubectlCommand: string): Promise { + const store = new MemoryAgentRunStore(); + const sessionId = "sess_artificer_active_steer"; + const run = store.createRun(artificerRunInput(sessionId)); + const turn = store.createCommand(run.id, validateCreateCommand({ type: "turn", payload: { prompt: "active artificer work" } })); + store.claimRun(run.id, "runner_artificer_active_steer", 60_000); + store.ackCommand(turn.id); + + const server = await startManagerServer({ + store, + sourceCommit: "selftest", + runnerJobDefaults: defaults(kubectlCommand), + runnerDispatcherOptions: { enabled: false }, + kafkaOutboxRelayOptions: { enabled: false }, + }); + try { + const client = new ManagerClient(server.baseUrl); + const response = await client.post(`/api/v1/sessions/${sessionId}/send`, sessionSendBody(sessionId, "steer active artificer turn", "session-active-steer")) as JsonRecord; + const command = response.command as CommandRecord; + const activeBefore = response.activeBefore as JsonRecord; + assert.equal(response.decision, "steer"); + assert.equal((response.run as RunRecord).id, run.id); + assert.equal(command.type, "steer"); + assert.equal(activeBefore.runId, run.id); + assert.equal(activeBefore.commandId, turn.id); + assert.equal(activeBefore.commandState, "acknowledged"); + assert.equal(response.runnerAdmission, null); + assert.equal(store.listCommands(run.id, 0, 100).length, 2); + } finally { + await closeServer(server.server); + } +} + function assertTerminalizedAdmission(store: MemoryAgentRunStore, runId: string, commandId: string, sessionId: string, reason: string): void { const intent = store.getRunnerDispatchIntent(commandId); assert.equal(intent?.state, "failed"); @@ -471,6 +551,16 @@ class FailingSessionCommandStore extends MemoryAgentRunStore { } } +class TerminalizingSteerStore extends MemoryAgentRunStore { + override admitSessionSteer(input: SessionSteerAdmissionInput): SessionSteerAdmission | null { + const session = this.getSession(input.sessionId); + if (session?.activeCommandId) { + this.finishCommand(session.activeCommandId, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + } + return super.admitSessionSteer(input); + } +} + async function closeServer(server: import("node:http").Server): Promise { await new Promise((resolve) => server.close(() => resolve())); } diff --git a/src/selftest/integration/postgres-session-turn-admission.ts b/src/selftest/integration/postgres-session-turn-admission.ts index 2fc79e9..c72ee2e 100644 --- a/src/selftest/integration/postgres-session-turn-admission.ts +++ b/src/selftest/integration/postgres-session-turn-admission.ts @@ -13,6 +13,7 @@ const concurrentSessionId = `ses_pg_admission_${suffix}`; const legacySessionId = `ses_pg_admission_legacy_${suffix}`; const rollbackSessionId = `ses_pg_admission_rollback_${suffix}`; const resourceSessionId = `ses_pg_admission_resource_${suffix}`; +const terminalRaceSessionId = `ses_pg_admission_terminal_race_${suffix}`; const eventOutbox = { enabled: true, topic: "agentrun.event.v1", source: "session-admission-selftest" } as const; const primary = await createPostgresAgentRunStore({ connectionString, poolMax: 4, eventOutbox }); const concurrent = await createPostgresAgentRunStore({ connectionString, poolMax: 4, eventOutbox }); @@ -24,6 +25,7 @@ try { await assertLegacyHalfCommitRecoveryAfterRestart(); await assertAdmissionTransactionRollback(); await assertLatestSessionResourceRunSkipsJsonNull(); + await assertTerminalAndFollowUpLockOrder(); console.log(JSON.stringify({ ok: true, tests: [ @@ -31,6 +33,8 @@ try { "postgres-session-turn-restart-durable-dispatch-intent", "postgres-session-turn-active-conflict", "postgres-session-turn-same-session-continuity", + "postgres-session-follow-up-terminal-race-preserves-new-active-run", + "postgres-session-terminal-follow-up-lock-order", "postgres-legacy-half-commit-recovery-after-restart", "postgres-session-turn-transaction-rollback", "postgres-session-resource-run-skips-json-null", @@ -43,6 +47,7 @@ try { await cleanupSession(legacySessionId).catch(() => undefined); await cleanupSession(rollbackSessionId).catch(() => undefined); await cleanupSession(resourceSessionId).catch(() => undefined); + await cleanupSession(terminalRaceSessionId).catch(() => undefined); await Promise.all(reopenedStores.map(async (store) => await store.close())); await Promise.all([primary.close(), concurrent.close(), control.end()]); } @@ -86,13 +91,20 @@ async function assertConcurrentAdmissionAndRestartRecovery(): Promise { assert.equal(replayed.disposition, "replayed"); await restarted.finishCommand(first.command.id, { terminalStatus: "completed", failureKind: null, failureMessage: null }); - await restarted.finishRun(first.run.id, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + const staleSteer = await restarted.admitSessionSteer({ + sessionId: concurrentSessionId, + command: validateCreateCommand({ type: "steer", payload: { prompt: "must not enter terminal run" }, idempotencyKey: `pg-admission-stale-steer-${suffix}` }), + }); + assert.equal(staleSteer, null); const next = await concurrent.admitSessionTurn(admissionInput(concurrentSessionId, "postgres same session next turn", `pg-admission-next-${suffix}`)); assert.equal(next.disposition, "created"); assert.notEqual(next.run.id, first.run.id); + await restarted.finishRun(first.run.id, { terminalStatus: "completed", failureKind: null, failureMessage: null }); const session = await primary.getSession(concurrentSessionId); assert.equal(session?.activeRunId, next.run.id); assert.equal(session?.activeCommandId, next.command.id); + assert.equal(session?.lastRunId, next.run.id); + assert.equal(session?.lastCommandId, next.command.id); } async function assertLegacyHalfCommitRecoveryAfterRestart(): Promise { @@ -149,6 +161,74 @@ async function assertLatestSessionResourceRunSkipsJsonNull(): Promise { assert.equal((await primary.getLatestSessionResourceRun(resourceSessionId))?.id, initial.id); } +async function assertTerminalAndFollowUpLockOrder(): Promise { + const run = await primary.createRun(validateCreateRun(runInput(terminalRaceSessionId))); + const command = await primary.createCommand(run.id, validateCreateCommand({ type: "turn", payload: { prompt: "terminal lock order" }, idempotencyKey: `pg-terminal-lock-order-${suffix}` })); + const runner = await primary.registerRunner({ runId: run.id, attemptId: `attempt_terminal_lock_${suffix}`, backendProfile: "codex", placement: "selftest", sourceCommit: "selftest" }); + await primary.claimRun(run.id, runner.id, 60_000); + await primary.ackCommand(command.id); + + const terminalClient = await control.connect(); + let steerPromise: Promise>> | null = null; + try { + await terminalClient.query("BEGIN"); + await terminalClient.query("SET LOCAL lock_timeout = '2s'"); + await terminalClient.query("SELECT id FROM agentrun_runs WHERE id = $1 FOR UPDATE", [run.id]); + await terminalClient.query("SELECT id FROM agentrun_commands WHERE id = $1 FOR UPDATE", [command.id]); + + steerPromise = concurrent.admitSessionSteer({ + sessionId: terminalRaceSessionId, + command: validateCreateCommand({ type: "steer", payload: { prompt: "follow up at terminal boundary" }, idempotencyKey: `pg-terminal-lock-steer-${suffix}` }), + }); + await waitForBlockedRunLock(); + + await terminalClient.query("SELECT session_id FROM agentrun_sessions WHERE session_id = $1 FOR UPDATE", [terminalRaceSessionId]); + await terminalClient.query("UPDATE agentrun_commands SET state = 'completed', updated_at = now() WHERE id = $1", [command.id]); + await terminalClient.query( + `UPDATE agentrun_sessions + SET execution_state = 'terminal', active_run_id = NULL, active_command_id = NULL, + last_run_id = $2, last_command_id = $3, terminal_status = 'completed', failure_kind = NULL, updated_at = now() + WHERE session_id = $1`, + [terminalRaceSessionId, run.id, command.id], + ); + await terminalClient.query("COMMIT"); + + assert.equal(await steerPromise, null); + assert.deepEqual((await primary.listCommands(run.id, 0, 100)).map((item) => item.id), [command.id]); + const next = await primary.admitSessionTurn(admissionInput(terminalRaceSessionId, "new run after terminal boundary", `pg-terminal-lock-next-${suffix}`)); + await concurrent.finishRun(run.id, { terminalStatus: "completed", failureKind: null, failureMessage: null }); + const session = await primary.getSession(terminalRaceSessionId); + assert.equal(session?.activeRunId, next.run.id); + assert.equal(session?.activeCommandId, next.command.id); + assert.equal(session?.lastRunId, next.run.id); + assert.equal(session?.lastCommandId, next.command.id); + } finally { + await terminalClient.query("ROLLBACK").catch(() => undefined); + terminalClient.release(); + if (steerPromise) await steerPromise.catch(() => undefined); + } +} + +async function waitForBlockedRunLock(): Promise { + const deadline = Date.now() + 2_000; + while (Date.now() < deadline) { + const result = await control.query<{ blocked: number }>( + `SELECT count(*)::int AS blocked + FROM pg_stat_activity + WHERE application_name = 'agentrun-mgr-v01' + AND wait_event_type = 'Lock' + AND query LIKE 'SELECT * FROM agentrun_runs WHERE id = $1 FOR UPDATE%'`, + ); + if ((result.rows[0]?.blocked ?? 0) > 0) return; + await delay(10); + } + throw new Error("session steer did not block on the run lock before the terminal path acquired the session lock"); +} + +async function delay(ms: number): Promise { + await new Promise((resolve) => setTimeout(resolve, ms)); +} + function admissionInput(sessionId: string, prompt: string, idempotencyKey: string): SessionTurnAdmissionInput { return { sessionId,