diff --git a/src/node/services/agentSession.goalAutoPause.test.ts b/src/node/services/agentSession.goalAutoPause.test.ts index f5315698e61..6bf556c7a25 100644 --- a/src/node/services/agentSession.goalAutoPause.test.ts +++ b/src/node/services/agentSession.goalAutoPause.test.ts @@ -535,6 +535,70 @@ describe("AgentSession goal safety hooks", () => { await session.dispose(); }); + test("a direct session send in its preflight defers redispatched follow-ups (#5506)", async () => { + // A user send made on the session itself (`xum run`) arms no WorkspaceService ticket, so the + // idle probes must count the session's own preflight too. Before the fix the follow-up was + // admitted, moved the turn, and the user's send was refused as a context mutation. + const workspaceId = "compaction-followup-direct-send-preflight"; + const { session, goalService, historyService, cleanup } = + await createSessionHarness(workspaceId); + cleanups.push(cleanup); + const created = await setGoalOk(goalService, { workspaceId, objective: "Direct send race" }); + const summary = createMuxMessage( + `summary-${crypto.randomUUID()}`, + "assistant", + "Compacted conversation.", + { + muxMetadata: { + type: "compaction-summary", + pendingFollowUp: { + text: "Continue working on the goal.", + agentId: "exec", + model: "openai:gpt-4o", + goalKind: GOAL_CONTINUATION_KIND, + goalId: created.goalId, + }, + }, + } + ); + expect((await historyService.appendToHistory(workspaceId, summary)).success).toBe(true); + + // Hold the user's send in prepareMessage's preflight (its turn-lease confirmation). + const internal = session as unknown as { confirmTurnUseLease(): Promise }; + const confirm = internal.confirmTurnUseLease.bind(session); + const entered = Promise.withResolvers(); + const release = Promise.withResolvers(); + spyOn(internal, "confirmTurnUseLease").mockImplementationOnce(async () => { + entered.resolve(); + await release.promise; + return confirm(); + }); + const sending = session.sendMessage("user sends directly", SEND_OPTIONS); + let dispatched: boolean | undefined; + try { + await entered.promise; + dispatched = await session.dispatchPendingCompactionFollowUpIfNeeded(); + } finally { + release.resolve(); + } + // Target assertion: the user's send is not refused by the follow-up... + expect((await sending).success).toBe(true); + // ...because the follow-up yielded to it. + expect(dispatched).toBe(false); + const history = await historyService.getLastMessages(workspaceId, 10); + expect(history.success).toBe(true); + if (history.success) { + expect( + history.data.some((message) => + message.parts.some( + (part) => part.type === "text" && part.text === "Continue working on the goal." + ) + ) + ).toBe(false); + } + await session.dispose(); + }); + test("a recovered budget wrap-up follow-up installs its missing reservation", async () => { // Codex P2 (PRRT_kwDOPxxmWM6cRJEE): a crash between wrap-up send // acceptance and tryMarkBudgetLimitInjected leaves the goal unmarked. diff --git a/src/node/services/agentSession.ts b/src/node/services/agentSession.ts index b9973702f41..9c0a65fb03b 100644 --- a/src/node/services/agentSession.ts +++ b/src/node/services/agentSession.ts @@ -1338,6 +1338,13 @@ export class AgentSession { } | null = null; /** setAutoRetryEnabled(false) calls still applying; a goal resume after an error waits them out. */ private autoRetryOptOutsInFlight = 0; + /** + * Bumped by each terminal-error record that gets past its early returns + * (recordGoalAdvancementAfterStreamError). A record that a later such record overtook during its + * preference read is dropped, so one failure hands over one + * resume and the predecessor's send options cannot replace the successor's (#5546). + */ + private goalAdvancementRecordGeneration = 0; private autoRetryStateVersion = 0; private autoRetryStateLoad: Promise | null = null; @@ -4103,6 +4110,15 @@ export class AgentSession { return this.manualSendsInPreflight > 0 || this.hasExternalManualSendPreflight?.() === true; } + /** + * A send that a redispatched follow-up's idle rule yields to: a user send in this session's own + * preflight (a direct sendMessage, e.g. `xum run`, arms no WorkspaceService ticket, #5506), or + * any send WorkspaceService holds in its preflight. + */ + private hasFollowUpBlockingSendPreflight(): boolean { + return this.manualSendsInPreflight > 0 || this.hasExternalSendPreflight?.() === true; + } + /** Correlated callbacks settle before publishing idle; teardown joins this whole physical lease. */ private async completePreparation( attempt: PreparationAttempt, @@ -7969,6 +7985,7 @@ export class AgentSession { if (failureType === "runtime_not_ready" || failureType === "runtime_start_failed") { const failedUserMessageId = this.activeStreamUserMessageId; + const failedContext = this.activeStreamContext; this.activeCompactionRequest = undefined; this.resetActiveStreamState(); await this.handleStreamFailureForAutoRetry({ @@ -7978,6 +7995,12 @@ export class AgentSession { if (!this.coordinator.isCurrentTurn(turn) || !this.coordinator.isCurrentOperation(operation)) return { success: false, error, failureHandled: true }; await this.updateStartupAutoRetryAbandonFromFailure(failureType, failedUserMessageId); + if (!this.coordinator.isCurrentTurn(turn) || !this.coordinator.isCurrentOperation(operation)) + return { success: false, error, failureHandled: true }; + // This failure has no stream error event, so handleStreamError's G4 settlement never runs: + // settle here, or a non-retryable runtime_not_ready leaves an active goal stranded and its + // failed kickoff installed (#5546). The hand-over waits for this preparation's idle. + await this.recordGoalAdvancementAfterStreamError(failureType, failedContext); } else { await this.handleStreamError(buildStreamErrorEventData(error, { acpPromptId }), operation); } @@ -9376,10 +9399,14 @@ export class AgentSession { const failedOptions = failed?.options; if (failedOptions?.agentId === "plan" || failedOptions?.agentId === "compact") return; if (this.config.findWorkspace(this.workspaceId)?.parentWorkspaceId != null) return; + // Bumped only by a record that will record: a failure that settles nothing (aborted, retried, + // plan or compact) must not make an earlier record stale. + const generation = ++this.goalAdvancementRecordGeneration; try { const fence = goalService.captureGoalAdvancementFence(this.workspaceId); const autoRetryEnabled = await this.loadAutoRetryEnabledPreference(); if ( + generation !== this.goalAdvancementRecordGeneration || !autoRetryEnabled || this.autoRetryOptOutsInFlight > 0 || this.coordinator.closing || @@ -11745,7 +11772,7 @@ export class AgentSession { // Codex P1 (PRRT_kwDOPxxmWM6cRJD-): a manual service-level send can sit // in its preflight (awaiting pricing/settings) without queueing or // holding the turn phase — it must win over the synthetic follow-up too. - const hasExternalPreflightSend = this.hasExternalSendPreflight?.() === true; + const hasExternalPreflightSend = this.hasFollowUpBlockingSendPreflight(); if ( enforceIdleRule && (hasQueuedMessages || hasActiveNonCompletingTurn || hasExternalPreflightSend) @@ -11809,7 +11836,7 @@ export class AgentSession { const idleRuleStale = enforceIdleRule ? () => this.hasPendingManualFollowUp() || - this.hasExternalSendPreflight?.() === true || + this.hasFollowUpBlockingSendPreflight() || (this.isBusy() && this.coordinator.phase !== "completing") : undefined; const followUpAdmissionStale = () => @@ -11965,7 +11992,7 @@ export class AgentSession { }); await this.skipIdleRuleFollowUp( lastMessage, - this.hasPendingManualFollowUp() || this.hasExternalSendPreflight?.() === true, + this.hasPendingManualFollowUp() || this.hasFollowUpBlockingSendPreflight(), this.isBusy() && this.coordinator.phase !== "completing" ); return false; diff --git a/src/node/services/goalAdvancement.test.ts b/src/node/services/goalAdvancement.test.ts index 81a9f87aa08..47c03b59707 100644 --- a/src/node/services/goalAdvancement.test.ts +++ b/src/node/services/goalAdvancement.test.ts @@ -5,7 +5,8 @@ import * as path from "path"; import { afterEach, beforeEach, describe, expect, mock, spyOn, test } from "bun:test"; import type { Config } from "@/node/config"; import type { GoalRecordV1 } from "@/common/types/goal"; -import { Ok } from "@/common/types/result"; +import { Err, Ok } from "@/common/types/result"; +import type { SendMessageError } from "@/common/types/errors"; import type { StreamErrorType } from "@/common/types/errors"; import type { StreamAbortEvent, StreamEndEvent } from "@/common/types/stream"; import { GOAL_STREAM_ERROR_RESUME_MAX_ATTEMPTS } from "@/constants/goals"; @@ -64,6 +65,8 @@ describe("goal advancement after automatic work ends or is abandoned (G4)", () = let failureType: StreamErrorType; /** When set, a failing stream reports its failure only once this settles. */ let failureGate: Promise | null; + /** When set, the next send fails before any stream with this error (and clears it). */ + let preStreamError: SendMessageError | null; let streamCalls: number; beforeEach(async () => { @@ -79,6 +82,7 @@ describe("goal advancement after automatic work ends or is abandoned (G4)", () = providerUp = false; failureType = "authentication"; failureGate = null; + preStreamError = null; streamCalls = 0; const harness = await createAgentSessionHarness({ workspaceId, @@ -88,6 +92,11 @@ describe("goal advancement after automatic work ends or is abandoned (G4)", () = aiServiceOverrides: { streamMessage: mock(() => { streamCalls += 1; + if (preStreamError != null) { + const error = preStreamError; + preStreamError = null; + return Promise.resolve(Err(error)); + } if (providerUp) { // A real stream reports its start before its handle settles. aiEmitter.emit("stream-start", { @@ -225,6 +234,28 @@ describe("goal advancement after automatic work ends or is abandoned (G4)", () = }); }); + test("G4 (#5546 item 3): a goal turn that fails before any stream with a non-retryable error resumes the goal", async () => { + await setGoalOk(service, { workspaceId, objective: "Ship G4" }); + // The kickoff's send fails before a stream exists: the container is gone. This failure + // never reaches handleStreamError (no stream error event), and RetryManager does not retry it. + preStreamError = { type: "runtime_not_ready", message: "Container not found" }; + const requestsBefore = requestDispatch.mock.calls.length; + expect(await dispatchAt(Date.now())).toBe(true); + await session.waitForIdle(); + expect(streamCalls).toBe(1); + expect(session.hasPendingAutoRetry()).toBe(false); + expect(await waitForRequests(requestsBefore)).toBe(1); + // Target assertion: the failure armed a bounded resume behind the error backoff. (The code + // left the failed kickoff installed, eligible at once: it re-fired with no backoff or bound.) + expect(await service.checkGoalContinuationEligibility(workspaceId, Date.now())).toMatchObject( + { eligible: false, reason: "error_backoff" } + ); + expect(await eligibilityAfterBackoff()).toMatchObject({ + eligible: true, + candidate: { source: "stream_error", sendOptions: { model: TEST_MODEL } }, + }); + }); + test("G4: a failed heartbeat turn resumes the goal with the goal's options, not the heartbeat's", async () => { await setGoalOk(service, { workspaceId, objective: "Ship G4" }); // The heartbeat turn runs while the kickoff candidate is gone (already consumed). @@ -966,6 +997,95 @@ describe("goal advancement after automatic work ends or is abandoned (G4)", () = ); }); + /** + * A failed turn's error record waits in its auto-retry preference read while a successor + * turn of `successorAgentId` starts and fails. Returns the dispatch requests made before it. + */ + async function recordOvertakenBy(successorAgentId: string): Promise { + const { release, requestsBefore } = await gatedFailingTurn(); + // Hold the first error's record in its auto-retry preference read (the only argument-less + // read on this path; RetryManager's own read passes its generation probe). + const internal = session as unknown as { + loadAutoRetryEnabledPreference(...args: unknown[]): Promise; + }; + const load = internal.loadAutoRetryEnabledPreference.bind(session); + const readEntered = Promise.withResolvers(); + const readRelease = Promise.withResolvers(); + let held = false; + spyOn(internal, "loadAutoRetryEnabledPreference").mockImplementation(async (...args) => { + if (!held && args.length === 0) { + held = true; + readEntered.resolve(); + await readRelease.promise; + } + return load(...args); + }); + failureGate = null; + release(); + await readEntered.promise; + const sent = await session.sendMessage( + "Peer message", + { model: TEST_MODEL, agentId: successorAgentId }, + { acceptanceOrigin: "automatic", synthetic: true, agentInitiated: true } + ); + expect(sent.success).toBe(true); + await session.waitForIdle(); + readRelease.resolve(); + await session.waitForIdle(); + return requestsBefore; + } + + test("G4 (#5546 item 1): an error record overtaken during its preference read hands over nothing", async () => { + // The successor's own error record hands over the resume. + const requestsBefore = await recordOvertakenBy("exec"); + expect(await waitForRequests(requestsBefore)).toBe(1); + // Target assertion: the predecessor's record is stale and hands over no second resume + // (the code consumed a second resume attempt for the same failure). + expect(await waitForRequests(requestsBefore + 1, 150)).toBe(0); + }); + + test("G4 (#5546 item 1): a successor whose failure records nothing leaves the earlier record", async () => { + // A failed plan turn records no advancement, so it must not make the earlier record stale. + const requestsBefore = await recordOvertakenBy("plan"); + // Target assertion: the predecessor's record still hands over its resume, once. + expect(await waitForRequests(requestsBefore)).toBe(1); + expect(await waitForRequests(requestsBefore + 1, 150)).toBe(0); + }); + + test("G4 (#5546 item 4): a dequeued manual send still in its preflight blocks the advancement", async () => { + const { release, requestsBefore } = await gatedFailingTurn(); + // The user's message waits behind the failing turn; the terminal error leaves it queued. + session.queueMessage("Do this next", { model: TEST_MODEL, agentId: "exec" }); + release(); + await session.waitForIdle(); + expect(await waitForRequests(requestsBefore, 100)).toBe(0); + // Hold the dequeued send in its preflight (turn-lease confirmation): the queue is empty, so + // only the preparation blocks the advancement. Its PREPARING phase does today, and + // userInputBlocksGoalAdvancement's preparing-manual-send check backs it up (that state is + // cleared before the preparation publishes idle). + const internal = session as unknown as { confirmTurnUseLease(): Promise }; + const confirm = internal.confirmTurnUseLease.bind(session); + const entered = Promise.withResolvers(); + const releaseSend = Promise.withResolvers(); + spyOn(internal, "confirmTurnUseLease").mockImplementationOnce(async () => { + entered.resolve(); + await releaseSend.promise; + return confirm(); + }); + providerUp = true; + session.drainQueuedMessagesIfIdle(); + try { + await entered.promise; + // Target assertion: the dequeue re-evaluated the advancement, which waits for the send. + expect(await waitForRequests(requestsBefore, 150)).toBe(0); + } finally { + releaseSend.resolve(); + } + await settle(() => Promise.resolve(streamCalls === 2), 1_000); + // The user's turn streamed: it owns the goal from here, not the error's resume. + expect(streamCalls).toBe(2); + }); + test("G4 control: a user Stop while blocked discards the pending advancement", async () => { const { withdraw, requestsBefore } = await errorBlockedByHeldBackWork(); await session.interruptStream();