fix: 修复 Session 续跑终态竞态
This commit is contained in:
@@ -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,25 @@ CREATE TABLE IF NOT EXISTS agentrun_schema_migrations (
|
||||
return this.withTransaction(async (client) => await this.createCommandWithClient(client, runId, input));
|
||||
}
|
||||
|
||||
async admitSessionSteer(input: SessionSteerAdmissionInput): Promise<SessionSteerAdmission | null> {
|
||||
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 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?.activeRunId || !session.activeCommandId) return null;
|
||||
const runResult = await client.query("SELECT * FROM agentrun_runs WHERE id = $1 FOR UPDATE", [session.activeRunId]);
|
||||
const commandResult = await client.query("SELECT * FROM agentrun_commands WHERE id = $1 FOR UPDATE", [session.activeCommandId]);
|
||||
if (!runResult.rows[0] || !commandResult.rows[0]) return null;
|
||||
const run = runFromRow(runResult.rows[0]);
|
||||
const targetCommand = commandFromRow(commandResult.rows[0]);
|
||||
const receivable = receivableActiveSessionTurn(session, run, targetCommand);
|
||||
if (!receivable) return null;
|
||||
const command = await this.createCommandWithClient(client, run.id, input.command);
|
||||
return { run: await this.requireRunForUpdate(client, run.id), command, targetCommand, ...receivable };
|
||||
});
|
||||
}
|
||||
|
||||
async admitSessionTurn(input: SessionTurnAdmissionInput): Promise<SessionTurnAdmission> {
|
||||
assertSessionTurnAdmissionContract(input);
|
||||
return this.withTransaction(async (client) => {
|
||||
@@ -1102,7 +1121,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;
|
||||
});
|
||||
}
|
||||
|
||||
+14
-60
@@ -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<ReturnType<typeof activeReceivableCommand>>, request: JsonRecord, runBody: JsonRecord | null, reusedIdleRun: JsonRecord | null = null): JsonRecord {
|
||||
function sessionSendPlan(sessionId: string, decision: "steer" | "turn", active: Awaited<ReturnType<typeof activeReceivableCommand>>, 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<ReturnType<typeof activeReceivableCommand>>; 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<ReturnType<typeof activeReceivableCommand>>; 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<Awaited<ReturnType<typeof activ
|
||||
return { runId: active.run.id, commandId: active.command.id, commandState: active.command.state, runStatus: active.run.status, leaseExpiresAt: active.run.leaseExpiresAt, leaseExpired: active.leaseExpired, reason: active.reason, valuesPrinted: false };
|
||||
}
|
||||
|
||||
function idleRunSummary(idleRun: NonNullable<Awaited<ReturnType<typeof idleReceivableRun>>>): 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 };
|
||||
|
||||
+48
-1
@@ -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<StoreHealth>;
|
||||
createRun(input: CreateRunInput, identity?: CreateRunIdentity): MaybePromise<RunRecord>;
|
||||
@@ -68,6 +81,7 @@ export interface AgentRunStore {
|
||||
listEvents(runId: string, afterSeq: number, limit: number): MaybePromise<RunEvent[]>;
|
||||
listEventsForCommand(runId: string, commandId: string, limit: number): MaybePromise<RunEvent[]>;
|
||||
createCommand(runId: string, input: CreateCommandInput): MaybePromise<CommandRecord>;
|
||||
admitSessionSteer(input: SessionSteerAdmissionInput): MaybePromise<SessionSteerAdmission | null>;
|
||||
admitSessionTurn(input: SessionTurnAdmissionInput): MaybePromise<SessionTurnAdmission>;
|
||||
getCommand(commandId: string): MaybePromise<CommandRecord>;
|
||||
listCommands(runId: string, afterSeq: number, limit: number): MaybePromise<CommandRecord[]>;
|
||||
@@ -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", {
|
||||
|
||||
@@ -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<void> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
await new Promise<void>((resolve) => server.close(() => resolve()));
|
||||
}
|
||||
|
||||
@@ -31,6 +31,7 @@ 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-legacy-half-commit-recovery-after-restart",
|
||||
"postgres-session-turn-transaction-rollback",
|
||||
"postgres-session-resource-run-skips-json-null",
|
||||
@@ -86,13 +87,20 @@ async function assertConcurrentAdmissionAndRestartRecovery(): Promise<void> {
|
||||
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<void> {
|
||||
|
||||
Reference in New Issue
Block a user