Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
64 changes: 64 additions & 0 deletions src/node/services/agentSession.goalAutoPause.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown> };
const confirm = internal.confirmTurnUseLease.bind(session);
const entered = Promise.withResolvers<void>();
const release = Promise.withResolvers<void>();
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.
Expand Down
33 changes: 30 additions & 3 deletions src/node/services/agentSession.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> | null = null;

Expand Down Expand Up @@ -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<T>(
attempt: PreparationAttempt,
Expand Down Expand Up @@ -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({
Expand All @@ -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);
}
Expand Down Expand Up @@ -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 ||
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 = () =>
Expand Down Expand Up @@ -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;
Expand Down
122 changes: 121 additions & 1 deletion src/node/services/goalAdvancement.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<void> | 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 () => {
Expand All @@ -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,
Expand All @@ -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", {
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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<number> {
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<boolean>;
};
const load = internal.loadAutoRetryEnabledPreference.bind(session);
const readEntered = Promise.withResolvers<void>();
const readRelease = Promise.withResolvers<void>();
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<unknown> };
const confirm = internal.confirmTurnUseLease.bind(session);
const entered = Promise.withResolvers<void>();
const releaseSend = Promise.withResolvers<void>();
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();
Expand Down
Loading