diff --git a/packages/coding-agent/.changes/post-compaction-idle.md b/packages/coding-agent/.changes/post-compaction-idle.md new file mode 100644 index 0000000000..c059f46e0b --- /dev/null +++ b/packages/coding-agent/.changes/post-compaction-idle.md @@ -0,0 +1 @@ +- Removed the delay before continuing sessions after compaction. diff --git a/packages/coding-agent/.changes/rlm-activity-change-waiter.md b/packages/coding-agent/.changes/rlm-activity-change-waiter.md new file mode 100644 index 0000000000..754d5df11a --- /dev/null +++ b/packages/coding-agent/.changes/rlm-activity-change-waiter.md @@ -0,0 +1 @@ +- Wait for RLM session activity changes without zero-delay polling. diff --git a/packages/coding-agent/src/core/agent-session.ts b/packages/coding-agent/src/core/agent-session.ts index 9b0f660219..33a9b4023c 100644 --- a/packages/coding-agent/src/core/agent-session.ts +++ b/packages/coding-agent/src/core/agent-session.ts @@ -830,11 +830,12 @@ function createAgentMessageDeferred(): AgentMessageDeferred { /** One-shot settlement for a scheduled post-compaction continuation; a settled failure is never re-exposed to later waiters. */ interface PostCompactionContinuationSettlement extends AgentMessageDeferred { + continueAfterSessionInput: boolean; settled: boolean; } function createPostCompactionContinuationSettlement(): PostCompactionContinuationSettlement { - return { ...createAgentMessageDeferred(), settled: false }; + return { ...createAgentMessageDeferred(), continueAfterSessionInput: false, settled: false }; } export interface ModelCycleResult { @@ -1065,7 +1066,7 @@ export class AgentSession { private _pendingSessionActionFenceWaiters = 0; private readonly _sessionActionCommitContext = new AsyncLocalStorage(); private readonly _sessionActionCommitDisposeAbortController = new AbortController(); - // Checkpoint and handoff waiters share lifecycle-edge notifications to avoid polling. + // Checkpoint, handoff, and activity waiters share lifecycle-edge notifications to avoid polling. private readonly _sessionInputCheckpointWaiters = new Set<() => void>(); private _pendingNextTurnMessages: CustomMessage[] = []; @@ -1198,7 +1199,6 @@ export class AgentSession { private _compactAutoRefinePending = false; private _turnIntervalAutoRefinePending = false; private _postCompactionContinuationScheduled = false; - private _postCompactionContinuationTimer: ReturnType | undefined; private _postCompactionContinuationSettlement: PostCompactionContinuationSettlement | undefined; private _postCompactionContinuationMessages: AgentMessage[] = []; private _scheduledPostCompactionContinuationMessages: AgentMessage[] = []; @@ -2506,6 +2506,7 @@ export class AgentSession { if (this._refineInFlight === applySettled) { this._refineInFlight = undefined; } + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); } } @@ -2714,6 +2715,7 @@ export class AgentSession { if (this._refineInFlight === applySettled) { this._refineInFlight = undefined; } + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); } } @@ -3704,6 +3706,7 @@ export class AgentSession { this._retryResolve(); this._retryResolve = undefined; this._retryPromise = undefined; + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); } } @@ -6583,6 +6586,19 @@ export class AgentSession { for (const resolve of waiters) resolve(); } + private _waitForSessionActivityChange(signal: AbortSignal): Promise { + return new Promise((resolve) => { + const finish = () => { + this._sessionInputCheckpointWaiters.delete(finish); + signal.removeEventListener("abort", finish); + resolve(); + }; + this._sessionInputCheckpointWaiters.add(finish); + signal.addEventListener("abort", finish, { once: true }); + if (signal.aborted) finish(); + }); + } + private _observeSessionActionDeferral(action: QueuedSessionAction): { deferred: Promise; stop(): void; @@ -6794,10 +6810,29 @@ export class AgentSession { } async waitForIdle(): Promise { - while (true) { + await this._waitForIdleOrSettlement(); + } + + /** + * {@link waitForIdle} loop; with a settlement, returns once that settlement is + * superseded so a cancelled post-compaction runner cannot keep a checkpoint + * waiter registered (a leaked waiter holds hasPendingAdmissionWaiters true and + * blocks daemon passivation). + */ + private async _waitForIdleOrSettlement(settlement?: PostCompactionContinuationSettlement): Promise { + while (settlement === undefined || this._postCompactionContinuationSettlement === settlement) { if (this._actionStore.queuedActions().length > 0) { if (this._sessionInputPumpSuspended || this._queuedWorkPauses.size > 0) { - await new Promise((resolve) => this._sessionInputCheckpointWaiters.add(resolve)); + let wake = () => {}; + const changed = new Promise((resolve) => { + wake = resolve; + this._sessionInputCheckpointWaiters.add(resolve); + }); + try { + await (settlement ? Promise.race([changed, settlement.promise]) : changed); + } finally { + this._sessionInputCheckpointWaiters.delete(wake); + } continue; } this._scheduleSessionInputPump(); @@ -7311,6 +7346,7 @@ export class AgentSession { throw new Error("Cannot compact without aborting while the agent is running."); } const hadPostCompactionContinue = this._postCompactionContinuationScheduled; + const continueAfterSessionInput = this._postCompactionContinuationSettlement?.continueAfterSessionInput ?? false; this._disconnectFromAgent(); if (!options.skipAbort) await this.abort(); let didCompact = false; @@ -7375,11 +7411,12 @@ export class AgentSession { this._compactionOperation = undefined; } resolveCompactionOperation(); + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); if (didCompact) { this._discardPendingAutoRefine({ cancelPostCompactionContinue: true }); if (hadPostCompactionContinue) { - this._schedulePostCompactionContinue(); + this._schedulePostCompactionContinue(continueAfterSessionInput); } // Queued agent or session-owned inputs resume the loop; defer refine // behind them instead of interleaving it before their turns. @@ -7498,20 +7535,17 @@ export class AgentSession { } private _settlePostCompactionContinue(error?: Error): void { - if (!error && (this._postCompactionContinuationScheduled || this._postCompactionContinuationTimer)) return; + if (!error && this._postCompactionContinuationScheduled) return; const settlement = this._postCompactionContinuationSettlement; if (!settlement || settlement.settled) return; settlement.settled = true; this._postCompactionContinuationSettlement = undefined; if (error) settlement.reject(error); else settlement.resolve(); + this._notifySessionInputCheckpointChange(); } private _cancelPostCompactionContinue(): void { - if (this._postCompactionContinuationTimer) { - clearTimeout(this._postCompactionContinuationTimer); - this._postCompactionContinuationTimer = undefined; - } this._postCompactionContinuationScheduled = false; this._scheduledPostCompactionContinuationMessages = []; this._settlePostCompactionContinue(); @@ -7607,75 +7641,126 @@ export class AgentSession { this._scheduleAutoRefine("compact"); } - private _schedulePostCompactionContinue(): void { - if (this._postCompactionContinuationScheduled) { - return; - } + private _schedulePostCompactionContinue(continueAfterSessionInput = false): void { if (!this._postCompactionContinuationSettlement || this._postCompactionContinuationSettlement.settled) { this._postCompactionContinuationSettlement = createPostCompactionContinuationSettlement(); } + const settlement = this._postCompactionContinuationSettlement; + settlement.continueAfterSessionInput ||= continueAfterSessionInput; + if (this._postCompactionContinuationScheduled) { + return; + } this._postCompactionContinuationScheduled = true; this._scheduledPostCompactionContinuationMessages = [...this._postCompactionContinuationMessages]; - this._postCompactionContinuationTimer = setTimeout(() => { - this._postCompactionContinuationTimer = undefined; - void this._runScheduledPostCompactionContinue() - .catch(() => undefined) - .finally(() => this._settlePostCompactionContinue()); - }, 100); + void this._runScheduledPostCompactionContinue(settlement) + .catch(() => undefined) + .finally(() => { + if (this._postCompactionContinuationSettlement === settlement) { + this._settlePostCompactionContinue(); + } + }); } private _sessionOwnsScheduledContinuations(continuationMessages: AgentMessage[]): boolean { return continuationMessages.some((message) => this._postCompactionContinuationMessages.includes(message)); } - private async _runScheduledPostCompactionContinue(): Promise { - await this._waitForRefineIdle(); - if (!this._postCompactionContinuationScheduled) { - return; - } - if (this.isStreaming || this.isCompacting || this.isRetrying || this._queuedWorkPauses.size > 0) { - this._postCompactionContinuationScheduled = false; - this._schedulePostCompactionContinue(); - return; + private async _waitForQueuedWorkResume(settlement: PostCompactionContinuationSettlement): Promise { + while (this._queuedWorkPauses.size > 0 && this._postCompactionContinuationSettlement === settlement) { + let resume = () => {}; + const resumed = new Promise((resolve) => { + resume = resolve; + this._sessionInputCheckpointWaiters.add(resolve); + }); + try { + await Promise.race([resumed, settlement.promise]); + } finally { + this._sessionInputCheckpointWaiters.delete(resume); + } } + } - const continuationMessages = [...this._scheduledPostCompactionContinuationMessages]; - if (continuationMessages.length > 0 && !this._sessionOwnsScheduledContinuations(continuationMessages)) { - this._cancelPostCompactionContinue(); - this._scheduleAutoRefineAfterAgentEnd(); - return; - } - // An empty queue is not idle while the scheduler still owns active work. - if (this.unfinishedActionCount > 0 || this._sessionInputPumpRequested) { - this._scheduleSessionInputPump(); - await this._sessionInputPump; - if (this._postCompactionContinuationScheduled) { - this._postCompactionContinuationScheduled = false; - const shouldReschedule = - continuationMessages.length === 0 - ? this.unfinishedActionCount > 0 - : this._sessionOwnsScheduledContinuations(continuationMessages); - if (shouldReschedule) { - this._schedulePostCompactionContinue(); - } else { - this._scheduledPostCompactionContinuationMessages = []; + private async _runScheduledPostCompactionContinue(settlement: PostCompactionContinuationSettlement): Promise { + while (this._postCompactionContinuationScheduled && this._postCompactionContinuationSettlement === settlement) { + await this.agent.waitForIdle(); + await this.waitForRetry(); + await this._waitForRefineIdle(); + await this._waitForQueuedWorkResume(settlement); + const compactionOperation = this._compactionOperation; + if (compactionOperation) { + await Promise.race([compactionOperation, settlement.promise]); + continue; + } + + const commitFence = await this._acquireSessionActionCommitFence(); + let continuation: Promise | undefined; + let continuationMessages: AgentMessage[] = []; + let waitForSessionInput = false; + try { + await this.agent.waitForIdle(); + if ( + !this._postCompactionContinuationScheduled || + this._postCompactionContinuationSettlement !== settlement + ) { + return; + } + + if (this._queuedWorkPauses.size > 0 || this._compactionOperation || this._refineInFlight) { + continue; + } + + continuationMessages = [...this._scheduledPostCompactionContinuationMessages]; + if (continuationMessages.length > 0 && !this._sessionOwnsScheduledContinuations(continuationMessages)) { + this._cancelPostCompactionContinue(); this._scheduleAutoRefineAfterAgentEnd(); + return; + } + if (this.unfinishedActionCount > 0 || this._sessionInputPumpRequested) { + this._scheduleSessionInputPump(); + waitForSessionInput = true; + } else { + this._postCompactionContinuationScheduled = false; + continuation = this.agent.continue(); } + } finally { + commitFence.release(); } - return; - } - this._postCompactionContinuationScheduled = false; - try { - await this.agent.continue(); - this._forgetConsumedPostCompactionContinuations(continuationMessages); - } catch (error) { - const code = error instanceof AgentContinueError ? error.code : undefined; - if (code === "busy") { - this._schedulePostCompactionContinue(); - } else if (code !== "nothing-to-continue") { - // "nothing-to-continue" means the turn already completed; anything else must reject headless idle waiters. - this._settlePostCompactionContinue(this._asError(error)); + if (waitForSessionInput) { + await this._waitForIdleOrSettlement(settlement); + if (this._postCompactionContinuationSettlement !== settlement) return; + const shouldContinue = + (settlement.continueAfterSessionInput && continuationMessages.length === 0) || + this._sessionOwnsScheduledContinuations(continuationMessages); + if (shouldContinue) { + this._scheduledPostCompactionContinuationMessages = [...this._postCompactionContinuationMessages]; + continue; + } + this._postCompactionContinuationScheduled = false; + this._scheduledPostCompactionContinuationMessages = []; + this._scheduleAutoRefineAfterAgentEnd(); + return; + } + + try { + await continuation; + if (this._postCompactionContinuationSettlement === settlement) { + this._forgetConsumedPostCompactionContinuations(continuationMessages); + } + return; + } catch (error) { + const code = error instanceof AgentContinueError ? error.code : undefined; + if (code === "busy") { + if (this._postCompactionContinuationSettlement === settlement) { + this._postCompactionContinuationScheduled = true; + this._scheduledPostCompactionContinuationMessages = [...this._postCompactionContinuationMessages]; + } + continue; + } + if (code !== "nothing-to-continue" && this._postCompactionContinuationSettlement === settlement) { + this._settlePostCompactionContinue(this._asError(error)); + } + return; } } } @@ -8019,6 +8104,7 @@ export class AgentSession { if (this._refineInFlight === applySettled) { this._refineInFlight = undefined; } + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); } } @@ -8489,12 +8575,17 @@ export class AgentSession { (reason === "requested" || reason === "threshold") && (shouldContinueAfterCompaction || this.agent.hasQueuedMessages() || this.hasPendingSessionWork) ) { - this._schedulePostCompactionContinue(); + this._schedulePostCompactionContinue(shouldContinueAfterCompaction); } }; this._emit({ type: "compaction_start", reason, customInstructions }); this._autoCompactionAbortController = new AbortController(); + let resolveCompactionOperation: () => void = () => {}; + const compactionOperation = new Promise((resolve) => { + resolveCompactionOperation = resolve; + }); + this._compactionOperation = compactionOperation; try { const authResult = this.model ? await this._modelRegistry.getApiKeyAndHeaders(this.model) : undefined; @@ -8541,13 +8632,13 @@ export class AgentSession { this.agent.state.messages = messages.slice(0, -1); } - this._schedulePostCompactionContinue(); + this._schedulePostCompactionContinue(true); this._scheduleAutoRefineAfterCompaction(willContinueAfterCompaction); return true; } else if (shouldContinueAfterCompaction || hasQueuedMessages) { // Compaction can intentionally stop a tool loop between turns. // Queued follow-up/steering/custom messages can also be waiting. - this._schedulePostCompactionContinue(); + this._schedulePostCompactionContinue(shouldContinueAfterCompaction); this._scheduleAutoRefineAfterCompaction(willContinueAfterCompaction); } else { this._scheduleAutoRefineAfterCompaction(willContinueAfterCompaction); @@ -8597,6 +8688,11 @@ export class AgentSession { return false; } finally { this._autoCompactionAbortController = undefined; + if (this._compactionOperation === compactionOperation) { + this._compactionOperation = undefined; + } + resolveCompactionOperation(); + this._notifySessionInputCheckpointChange(); this._scheduleSessionInputPump(); } } @@ -10036,12 +10132,9 @@ export class AgentSession { try { while (true) { await wait(this.waitForHeadlessIdle()); - // Strong RLM quiescence also owns session-level work (bash, refine, - // branch mutation, and manual compaction) that interactive waitForIdle - // intentionally ignores. Yield a macrotask while such work is active so - // recursive parent/child barriers cannot form a microtask busy-loop. + // Strong RLM quiescence also owns work that interactive waitForIdle ignores. if (this.isSessionActive || this._hasDeferredRlmTerminalNotices()) { - await wait(new Promise((resolve) => setTimeout(resolve, 0))); + await wait(this._waitForSessionActivityChange(cancellation.signal)); continue; } const unsettledRuns = [...this._unsettledRlmChildRuns].filter((run) => !run.settled); @@ -10941,6 +11034,7 @@ export class AgentSession { return result; } finally { this._bashAbortController = undefined; + this._notifySessionInputCheckpointChange(); } } @@ -10985,6 +11079,7 @@ export class AgentSession { ); } finally { this._userBashRunning = false; + this._notifySessionInputCheckpointChange(); } // Emitted after the slot is released so clients never observe a bash_end // while the session still rejects new commands as already running. @@ -11444,6 +11539,7 @@ export class AgentSession { this._branchSummaryOperation = undefined; } resolveBranchSummaryOperation(); + this._notifySessionInputCheckpointChange(); } } diff --git a/packages/coding-agent/test/agent-session-recursion.test.ts b/packages/coding-agent/test/agent-session-recursion.test.ts index 2c373352ba..f0461a2dcc 100644 --- a/packages/coding-agent/test/agent-session-recursion.test.ts +++ b/packages/coding-agent/test/agent-session-recursion.test.ts @@ -1426,7 +1426,7 @@ describe("AgentSession rlm recursion", () => { }); }); - it("strong quiescence yields to a real gated child bash while interactive idle remains resolved", async () => { + it("strong quiescence waits for a gated child bash activity change", async () => { const child = createSession({ rlmSessionDir: join(tempDir, "bash-active-child") }); const bashStarted = deferred(); const bashCompletion = deferred(); @@ -1446,7 +1446,7 @@ describe("AgentSession rlm recursion", () => { const originalHeadlessIdle = child.waitForHeadlessIdle.bind(child); let headlessIdleCalls = 0; vi.spyOn(child, "waitForHeadlessIdle").mockImplementation(async () => { - if (++headlessIdleCalls > 100) throw new Error("RLM quiescence spun without yielding a macrotask"); + headlessIdleCalls++; await originalHeadlessIdle(); }); @@ -1459,11 +1459,12 @@ describe("AgentSession rlm recursion", () => { sleep(20).then(() => "timer" as const), ]); expect(firstBoundary).toBe("timer"); + expect(headlessIdleCalls).toBe(1); bashCompletion.resolve(); await bash; await expect(quiescence).resolves.toBeUndefined(); - expect(headlessIdleCalls).toBeLessThan(100); + expect(headlessIdleCalls).toBe(2); }); it("rechecks parent self-activity after a child quiescence boundary", async () => { diff --git a/packages/coding-agent/test/suite/agent-session-compaction.test.ts b/packages/coding-agent/test/suite/agent-session-compaction.test.ts index 234179b9ec..f19ce5b304 100644 --- a/packages/coding-agent/test/suite/agent-session-compaction.test.ts +++ b/packages/coding-agent/test/suite/agent-session-compaction.test.ts @@ -4,6 +4,7 @@ import { type AssistantMessage, fauxAssistantMessage, type Model, type ToolResul import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { SessionManager } from "../../src/core/session-manager.js"; import { createHarness, getMessageText, type Harness } from "./harness.js"; +import { createDeferred } from "./scheduling.js"; type SessionWithCompactionInternals = { _checkCompaction: ( @@ -267,6 +268,96 @@ describe("AgentSession compaction characterization", () => { } }); + it("waits for active manual compaction before continuing", async () => { + const compactionStarted = createDeferred(); + const compactionRelease = createDeferred(); + const harness = await createHarness({ + settings: { compaction: { keepRecentTokens: 1 } }, + extensionFactories: [ + (pi) => { + pi.on("session_before_compact", async (event) => { + compactionStarted.resolve(); + await compactionRelease.promise; + return { + compaction: { + summary: "summary from extension", + firstKeptEntryId: event.preparation.firstKeptEntryId, + tokensBefore: event.preparation.tokensBefore, + details: { source: "extension" }, + }, + }; + }); + }, + ], + }); + harnesses.push(harness); + const internals = harness.session as unknown as { + _schedulePostCompactionContinue(): void; + }; + harness.setResponses([fauxAssistantMessage("one"), fauxAssistantMessage("two")]); + await harness.session.prompt("first"); + await harness.session.prompt("second"); + const pause = harness.session.acquireQueuedWorkPause(); + const continueAgent = vi.spyOn(harness.session.agent, "continue").mockResolvedValue(); + internals._schedulePostCompactionContinue(); + + const compaction = harness.session.compact(undefined, { skipAbort: true }); + await compactionStarted.promise; + pause.release(); + await new Promise(setImmediate); + expect(continueAgent).not.toHaveBeenCalled(); + + compactionRelease.resolve(); + await compaction; + await harness.session.waitForHeadlessIdle(); + expect(continueAgent).toHaveBeenCalledTimes(1); + }); + + it("waits for active auto-compaction before continuing", async () => { + const compactionStarted = createDeferred(); + const compactionRelease = createDeferred(); + const harness = await createHarness({ + settings: { compaction: { keepRecentTokens: 1 } }, + extensionFactories: [ + (pi) => { + pi.on("session_before_compact", async (event) => { + compactionStarted.resolve(); + await compactionRelease.promise; + return { + compaction: { + summary: "summary from extension", + firstKeptEntryId: event.preparation.firstKeptEntryId, + tokensBefore: event.preparation.tokensBefore, + details: { source: "extension" }, + }, + }; + }); + }, + ], + }); + harnesses.push(harness); + const internals = harness.session as unknown as SessionWithCompactionInternals & { + _schedulePostCompactionContinue(): void; + }; + harness.setResponses([fauxAssistantMessage("one"), fauxAssistantMessage("two")]); + await harness.session.prompt("first"); + await harness.session.prompt("second"); + const pause = harness.session.acquireQueuedWorkPause(); + const continueAgent = vi.spyOn(harness.session.agent, "continue").mockResolvedValue(); + internals._schedulePostCompactionContinue(); + + const compaction = internals._runAutoCompaction("threshold", false); + await compactionStarted.promise; + pause.release(); + await new Promise(setImmediate); + expect(continueAgent).not.toHaveBeenCalled(); + + compactionRelease.resolve(); + await compaction; + await harness.session.waitForHeadlessIdle(); + expect(continueAgent).toHaveBeenCalledTimes(1); + }); + it("treats session-owned queued inputs as queued work after compaction", async () => { const harness = await createHarness({ settings: { compaction: { keepRecentTokens: 1 } }, @@ -310,6 +401,49 @@ describe("AgentSession compaction characterization", () => { } }); + it("releases the runner's suspended idle wait when the continuation is cancelled", async () => { + const harness = await createHarness(); + harnesses.push(harness); + const session = harness.session; + const internals = session as unknown as { + _schedulePostCompactionContinue(): void; + _cancelPostCompactionContinue(): void; + _sessionInputCheckpointWaiters: Set<() => void>; + _sessionInputPumpSuspended: boolean; + }; + // A queued follow-up held back by a pause, then a pump suspension (the + // requestAbort teardown state): the queue stays populated but undispatchable. + const pause = session.acquireQueuedWorkPause(); + await session.followUp("queued across abort"); + expect(session.queuedActionCount).toBe(1); + session.requestAbort(); + pause.release(); + expect(internals._sessionInputPumpSuspended).toBe(true); + expect(session.queuedActionCount).toBe(1); + + // The runner passes its pre-dispatch guards (no pauses, agent idle) and + // parks in the session idle wait on a checkpoint waiter. + internals._schedulePostCompactionContinue(); + await vi.waitFor(() => { + expect(internals._sessionInputCheckpointWaiters.size).toBeGreaterThan(0); + }); + expect(session.hasPendingAdmissionWaiters).toBe(true); + + // Cancelling the continuation must release that waiter: a stuck waiter + // keeps hasPendingAdmissionWaiters true and blocks daemon passivation. + // The settle's own notify empties the set for a moment; a leaked runner + // re-parks within a microtask, so settle real time before asserting. + internals._cancelPostCompactionContinue(); + await new Promise((resolve) => setTimeout(resolve, 100)); + try { + expect(internals._sessionInputCheckpointWaiters.size).toBe(0); + expect(session.hasPendingAdmissionWaiters).toBe(false); + } finally { + session.clearQueue(); + session.resumeQueuedWork(); + } + }); + it("defers post-compaction refine behind a preparing session action", async () => { const preparationReached = vi.fn(); let releasePreparation = () => {}; @@ -875,12 +1009,19 @@ describe("AgentSession compaction characterization", () => { response: "continuation handled", tracked: true, }, - ])("does not continue again after the session pump handles $name", async ({ text, response, tracked }) => { + { + name: "an empty resume request", + text: "concurrent input", + response: "concurrent input handled", + tracked: false, + continueAfterSessionInput: true, + }, + ])("settles $name after the session pump runs", async ({ text, response, tracked, continueAfterSessionInput }) => { vi.useFakeTimers(); const harness = await createHarness(); harnesses.push(harness); const sessionInternals = harness.session as unknown as { - _schedulePostCompactionContinue(): void; + _schedulePostCompactionContinue(continueAfterSessionInput?: boolean): void; _postCompactionContinuationMessages: AgentMessage[]; _postCompactionContinuationScheduled: boolean; _createPreparedTurnAction( @@ -906,10 +1047,10 @@ describe("AgentSession compaction characterization", () => { ); const continueSpy = vi.spyOn(harness.session.agent, "continue"); - sessionInternals._schedulePostCompactionContinue(); + sessionInternals._schedulePostCompactionContinue(continueAfterSessionInput); await vi.advanceTimersByTimeAsync(200); - expect(continueSpy).not.toHaveBeenCalled(); + expect(continueSpy).toHaveBeenCalledTimes(continueAfterSessionInput ? 1 : 0); expect(sessionInternals._postCompactionContinuationScheduled).toBe(false); expect(sessionInternals._postCompactionContinuationMessages).toEqual([]); expect(harness.session.messages.at(-1)).toMatchObject({ @@ -919,7 +1060,6 @@ describe("AgentSession compaction characterization", () => { }); it("keeps autonomous threshold continuations when post-compaction continue must retry", async () => { - vi.useFakeTimers(); const harness = await createHarness({ autonomous: { enabled: true, @@ -933,6 +1073,7 @@ describe("AgentSession compaction characterization", () => { harnesses.push(harness); const sessionInternals = harness.session as unknown as { _schedulePostCompactionContinue(): void; + _cancelPostCompactionContinue(): void; _postCompactionContinuationMessages: AgentMessage[]; _postCompactionContinuationScheduled: boolean; }; @@ -945,16 +1086,82 @@ describe("AgentSession compaction characterization", () => { harness.session.agent.state.messages = [ { role: "user", content: [{ type: "text", text: "hello" }], timestamp: Date.now() - 1000 }, ]; + const activeRunSettled = createDeferred(); const continueSpy = vi .spyOn(harness.session.agent, "continue") .mockRejectedValueOnce(new AgentContinueError("busy", "already processing")); + vi.spyOn(harness.session.agent, "waitForIdle").mockImplementation(() => + continueSpy.mock.calls.length === 0 ? Promise.resolve() : activeRunSettled.promise, + ); sessionInternals._schedulePostCompactionContinue(); - await vi.advanceTimersByTimeAsync(100); + await vi.waitFor(() => expect(continueSpy).toHaveBeenCalledTimes(1)); - expect(continueSpy).toHaveBeenCalledTimes(1); expect(sessionInternals._postCompactionContinuationMessages).toEqual([queuedMessage]); expect(sessionInternals._postCompactionContinuationScheduled).toBe(true); + sessionInternals._cancelPostCompactionContinue(); + activeRunSettled.resolve(); + }); + + it("keeps replacement continuation messages when a cancelled continue settles late", async () => { + const harness = await createHarness(); + harnesses.push(harness); + const sessionInternals = harness.session as unknown as { + _schedulePostCompactionContinue(): void; + _cancelPostCompactionContinue(): void; + _postCompactionContinuationMessages: AgentMessage[]; + }; + const queuedMessage = { + role: "user", + content: [{ type: "text", text: "autonomous follow-up" }], + timestamp: Date.now(), + } satisfies AgentMessage; + sessionInternals._postCompactionContinuationMessages = [queuedMessage]; + const staleRun = createDeferred(); + const replacementRun = createDeferred(); + const continueSpy = vi + .spyOn(harness.session.agent, "continue") + .mockReturnValueOnce(staleRun.promise) + .mockReturnValueOnce(replacementRun.promise); + + sessionInternals._schedulePostCompactionContinue(); + await vi.waitFor(() => expect(continueSpy).toHaveBeenCalledTimes(1)); + sessionInternals._cancelPostCompactionContinue(); + sessionInternals._schedulePostCompactionContinue(); + await vi.waitFor(() => expect(continueSpy).toHaveBeenCalledTimes(2)); + + staleRun.resolve(); + await new Promise(setImmediate); + expect(sessionInternals._postCompactionContinuationMessages).toEqual([queuedMessage]); + + replacementRun.resolve(); + await harness.session.waitForHeadlessIdle(); + expect(sessionInternals._postCompactionContinuationMessages).toEqual([]); + }); + + it("waits for an in-flight refine application before continuing", async () => { + const harness = await createHarness(); + harnesses.push(harness); + const sessionInternals = harness.session as unknown as { + _schedulePostCompactionContinue(): void; + _refineInFlight: Promise | undefined; + }; + const continueSpy = vi.spyOn(harness.session.agent, "continue").mockResolvedValue(); + const pause = harness.session.acquireQueuedWorkPause(); + sessionInternals._schedulePostCompactionContinue(); + await new Promise(setImmediate); + + // Refine enters its apply phase while the runner waits out the pause. + const refineApply = createDeferred(); + sessionInternals._refineInFlight = refineApply.promise; + pause.release(); + await new Promise(setImmediate); + expect(continueSpy).not.toHaveBeenCalled(); + + sessionInternals._refineInFlight = undefined; + refineApply.resolve(); + await harness.session.waitForHeadlessIdle(); + expect(continueSpy).toHaveBeenCalledTimes(1); }); it("clears queued autonomous threshold continuations when autonomous mode is disabled", async () => { diff --git a/packages/coding-agent/test/suite/agent-session-queue.test.ts b/packages/coding-agent/test/suite/agent-session-queue.test.ts index 3c19d7cfe3..14e31404a1 100644 --- a/packages/coding-agent/test/suite/agent-session-queue.test.ts +++ b/packages/coding-agent/test/suite/agent-session-queue.test.ts @@ -39,7 +39,7 @@ type AutoRefineInternals = { _scheduleAutoRefine(reason: AutoRefineReason): void; _scheduleAutoRefineAfterCompaction(willContinueAfterCompaction: boolean): void; _scheduleAutoRefineAfterAgentEnd(): void; - _schedulePostCompactionContinue(): void; + _schedulePostCompactionContinue(continueAfterSessionInput?: boolean): void; _invalidatePendingAutoRefineForBranchChange(): Promise; _cancelPostCompactionContinue(): void; _assistantTurnsSinceAutoRefine: number; @@ -327,34 +327,71 @@ describe("AgentSession queue characterization", () => { } }); - it("retries a scheduled post-compaction continuation when another run starts first", async () => { - vi.useFakeTimers(); + it("waits for the active run to settle before retrying a post-compaction continuation", async () => { const harness = await createAutoRefineHarness({ settings: { autoRefine: { enabled: true, turnInterval: 25, cooldownMs: 0 } }, }); harnesses.push(harness); const internals = harness.session as unknown as AutoRefineInternals; + const activeRunSettled = createDeferred(); const continueAgent = vi .spyOn(harness.session.agent, "continue") .mockRejectedValueOnce( new AgentContinueError("busy", "Agent is already processing. Wait for completion before continuing."), ) .mockResolvedValueOnce(); + vi.spyOn(harness.session.agent, "waitForIdle").mockImplementation(() => + continueAgent.mock.calls.length === 0 ? Promise.resolve() : activeRunSettled.promise, + ); - try { - internals._schedulePostCompactionContinue(); - await vi.advanceTimersByTimeAsync(100); + internals._schedulePostCompactionContinue(); + await vi.waitFor(() => expect(continueAgent).toHaveBeenCalledTimes(1)); + expect(internals._postCompactionContinuationScheduled).toBe(true); - expect(continueAgent).toHaveBeenCalledTimes(1); - expect(internals._postCompactionContinuationScheduled).toBe(true); + activeRunSettled.resolve(); + await vi.waitFor(() => expect(continueAgent).toHaveBeenCalledTimes(2)); + expect(internals._postCompactionContinuationScheduled).toBe(false); + }); - await vi.advanceTimersByTimeAsync(100); + it("does not let a failed cancelled continuation reject its replacement", async () => { + const harness = await createAutoRefineHarness(); + harnesses.push(harness); + const internals = harness.session as unknown as AutoRefineInternals; + const cancelledRun = createDeferred(); + const replacementRun = createDeferred(); + const continueAgent = vi + .spyOn(harness.session.agent, "continue") + .mockReturnValueOnce(cancelledRun.promise) + .mockReturnValueOnce(replacementRun.promise); - expect(continueAgent).toHaveBeenCalledTimes(2); - expect(internals._postCompactionContinuationScheduled).toBe(false); - } finally { - vi.useRealTimers(); - } + internals._schedulePostCompactionContinue(); + await vi.waitFor(() => expect(continueAgent).toHaveBeenCalledTimes(1)); + internals._cancelPostCompactionContinue(); + internals._schedulePostCompactionContinue(); + await vi.waitFor(() => expect(continueAgent).toHaveBeenCalledTimes(2)); + const idle = harness.session.waitForHeadlessIdle(); + + cancelledRun.reject(new Error("cancelled continuation failed")); + await new Promise(setImmediate); + replacementRun.resolve(); + + await expect(idle).resolves.toBeUndefined(); + }); + + it("waits for a queued-work pause to release before post-compaction continuation", async () => { + const harness = await createAutoRefineHarness(); + harnesses.push(harness); + const internals = harness.session as unknown as AutoRefineInternals; + const pause = harness.session.acquireQueuedWorkPause(); + const continueAgent = vi.spyOn(harness.session.agent, "continue").mockResolvedValue(); + + internals._schedulePostCompactionContinue(); + await new Promise(setImmediate); + expect(continueAgent).not.toHaveBeenCalled(); + + pause.release(); + await harness.session.waitForHeadlessIdle(); + expect(continueAgent).toHaveBeenCalledTimes(1); }); it("cancels scheduled post-compaction continuation on branch changes", async () => { @@ -407,24 +444,22 @@ describe("AgentSession queue characterization", () => { }); it("keeps scheduled post-compaction continuation when session-input pump compaction skips without aborting", async () => { - vi.useFakeTimers(); const harness = await createAutoRefineHarness({ settings: { autoRefine: { enabled: true, turnInterval: 25, cooldownMs: 0 } }, }); harnesses.push(harness); const internals = harness.session as unknown as AutoRefineInternals; - try { - internals._schedulePostCompactionContinue(); + const idle = createDeferred(); + vi.spyOn(harness.session.agent, "waitForIdle").mockReturnValue(idle.promise); + internals._schedulePostCompactionContinue(); - await expect(harness.session.compact(undefined, { skipAbort: true })).rejects.toThrow( - "Session is too short to compact", - ); + await expect(harness.session.compact(undefined, { skipAbort: true })).rejects.toThrow( + "Session is too short to compact", + ); - expect(internals._postCompactionContinuationScheduled).toBe(true); - } finally { - internals._cancelPostCompactionContinue(); - vi.useRealTimers(); - } + expect(internals._postCompactionContinuationScheduled).toBe(true); + internals._cancelPostCompactionContinue(); + idle.resolve(); }); it("auto-refine pending review uses the in-progress guard and catches refine failures", async () => {