diff --git a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts index 179756d80..c272e852d 100644 --- a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts +++ b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts @@ -551,10 +551,14 @@ describe("advisor", () => { it("coalesces multiple onTurnEnd calls while a prompt is in-flight", async () => { const promptInputs: string[] = []; const { promise: firstPromptPromise, resolve: finishFirstPrompt } = Promise.withResolvers(); + const { promise: secondPromptDone, resolve: finishSecondPrompt } = Promise.withResolvers(); + let promptCalls = 0; const agent: AdvisorAgent = { prompt: async input => { promptInputs.push(input); - await firstPromptPromise; + promptCalls++; + if (promptCalls === 1) await firstPromptPromise; + else finishSecondPrompt(); }, abort: () => {}, reset: () => {}, @@ -575,34 +579,24 @@ describe("advisor", () => { messages.push({ role: "user", content: "second", timestamp: 2 } as AgentMessage); runtime.onTurnEnd(); await Promise.resolve(); - expect(promptInputs).toHaveLength(1); + expect(promptInputs).toHaveLength(1); // second prompt not started yet finishFirstPrompt(); - await Promise.resolve(); - await Promise.resolve(); + await secondPromptDone; expect(promptInputs).toHaveLength(2); expect(promptInputs[1]).toContain("second"); }); - it("budgets only the batch sent after async context maintenance", async () => { + it("coalesces late-arriving deltas into the batch after context maintenance", async () => { const promptInputs: string[] = []; const { promise: firstMaintainStarted, resolve: startFirstMaintain } = Promise.withResolvers(); const { promise: finishFirstMaintain, resolve: releaseFirstMaintain } = Promise.withResolvers(); - const { promise: firstPromptStarted, resolve: startFirstPrompt } = Promise.withResolvers(); - const { promise: secondPromptStarted, resolve: startSecondPrompt } = Promise.withResolvers(); - const { promise: finishFirstPrompt, resolve: releaseFirstPrompt } = Promise.withResolvers(); + const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers(); let maintainCalls = 0; - let promptCalls = 0; const agent: AdvisorAgent = { prompt: async input => { promptInputs.push(input); - promptCalls++; - if (promptCalls === 1) { - startFirstPrompt(); - await finishFirstPrompt; - } else if (promptCalls === 2) { - startSecondPrompt(); - } + startPrompt(); }, abort: () => {}, reset: () => {}, @@ -625,19 +619,168 @@ describe("advisor", () => { runtime.onTurnEnd(); await firstMaintainStarted; + + // Second turn arrives while first maintainContext is still awaiting. messages.push({ role: "user", content: "second", timestamp: 2 } as AgentMessage); runtime.onTurnEnd(); releaseFirstMaintain(false); - await firstPromptStarted; + await promptStarted; + + // Both deltas land in a single prompt — late arrival coalesced before agent.prompt(). expect(promptInputs).toHaveLength(1); expect(promptInputs[0]).toContain("first"); - expect(promptInputs[0]).not.toContain("second"); + expect(promptInputs[0]).toContain("second"); + // The loop re-checked maintenance for the expanded batch. + expect(maintainCalls).toBe(2); + }); - releaseFirstPrompt(); - await secondPromptStarted; - expect(promptInputs).toHaveLength(2); - expect(promptInputs[1]).toContain("second"); + it("late-arriving delta that triggers reprime: full replay and correct turn accounting", async () => { + const promptInputs: string[] = []; + const { promise: firstMaintainStarted, resolve: startFirstMaintain } = Promise.withResolvers(); + const { promise: finishFirstMaintain, resolve: releaseFirstMaintain } = Promise.withResolvers(); + const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers(); + let resetCount = 0; + let maintainCalls = 0; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + startPrompt(); + }, + abort: () => {}, + reset: () => { + resetCount++; + }, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "turn1", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + maintainContext: async () => { + maintainCalls++; + if (maintainCalls === 1) { + startFirstMaintain(); + return await finishFirstMaintain; + } + // Second call (for the merged batch) → reprime. + return true; + }, + }; + const runtime = new AdvisorRuntime(agent, host); + + runtime.onTurnEnd(); + await firstMaintainStarted; + + messages.push({ role: "user", content: "turn2", timestamp: 2 } as AgentMessage); + runtime.onTurnEnd(); + + releaseFirstMaintain(false); + await promptStarted; + + // Full replay includes both turns. + expect(promptInputs).toHaveLength(1); + expect(promptInputs[0]).toContain("turn1"); + expect(promptInputs[0]).toContain("turn2"); + // Reprime resets the advisor agent. + expect(resetCount).toBeGreaterThan(0); + }); + + it("tags in-progress turns with [in progress] heading", async () => { + const promptInputs: string[] = []; + const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers(); + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + startPrompt(); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "hello", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host); + + runtime.onTurnEnd(messages, { willContinue: true }); + await promptStarted; + + expect(promptInputs).toHaveLength(1); + expect(promptInputs[0]).toContain("[in progress — more steps follow]"); + }); + + it("uses plain heading when willContinue is false or absent", async () => { + const promptInputs: string[] = []; + const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers(); + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + startPrompt(); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "done", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host); + + runtime.onTurnEnd(messages); + await promptStarted; + + expect(promptInputs).toHaveLength(1); + expect(promptInputs[0]).toContain("### Session update\n"); + expect(promptInputs[0]).not.toContain("[in progress"); + }); + + it("hasFreshBacklog is true only while pending queue is non-empty during a prompt", async () => { + const { promise: firstPromptStarted, resolve: startFirstPrompt } = Promise.withResolvers(); + const { promise: firstPromptDone, resolve: finishFirstPrompt } = Promise.withResolvers(); + const { promise: secondPromptDone, resolve: finishSecondPrompt } = Promise.withResolvers(); + let promptCalls = 0; + const agent: AdvisorAgent = { + prompt: async () => { + promptCalls++; + if (promptCalls === 1) { + startFirstPrompt(); + await firstPromptDone; + } else { + finishSecondPrompt(); + } + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "a", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host); + + runtime.onTurnEnd(); + await firstPromptStarted; + + // No late arrivals — false while first prompt runs with empty pending. + expect(runtime.hasFreshBacklog).toBe(false); + + // Push a second turn while the first prompt is still in-flight. + messages.push({ role: "user", content: "b", timestamp: 2 } as AgentMessage); + runtime.onTurnEnd(); + expect(runtime.hasFreshBacklog).toBe(true); + + finishFirstPrompt(); + await secondPromptDone; + + // After the second turn is fully drained, pending is empty again. + expect(runtime.hasFreshBacklog).toBe(false); }); it("sends the batch when context maintenance fails", async () => { @@ -661,9 +804,18 @@ describe("advisor", () => { expect(promptInputs[0]).toContain("first"); }); - it("excludes advisor custom messages from the rendered delta", () => { + it("excludes advisor custom messages from the rendered delta", async () => { const promptInputs: string[] = []; - const agent = makeAgent(promptInputs); + const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers(); + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + startPrompt(); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; const messages: AgentMessage[] = [ { role: "user", content: "hello", timestamp: 1 } as AgentMessage, { role: "custom", customType: "advisor", content: "note", display: true, timestamp: 2 } as AgentMessage, @@ -674,6 +826,7 @@ describe("advisor", () => { }; const runtime = new AdvisorRuntime(agent, host); runtime.onTurnEnd(); + await promptStarted; expect(promptInputs).toHaveLength(1); expect(promptInputs[0]).toContain("hello"); expect(promptInputs[0]).not.toContain("note"); @@ -882,7 +1035,7 @@ describe("advisor", () => { expect(promptInputs[1]).not.toContain("except the single plan file named below"); }); - it("renders the watched delta with a heading, watched-role labels, and no inner ## headings", () => { + it("renders the watched delta with a heading, watched-role labels, and no inner ## headings", async () => { const promptInputs: string[] = []; const agent = makeAgent(promptInputs); const messages: AgentMessage[] = [ @@ -920,6 +1073,7 @@ describe("advisor", () => { }; const runtime = new AdvisorRuntime(agent, host); runtime.onTurnEnd(); + await Promise.resolve(); expect(promptInputs).toHaveLength(1); const prompt = promptInputs[0]; expect(prompt).toContain("### Session update"); @@ -932,7 +1086,7 @@ describe("advisor", () => { expect(prompt.split("**agent**:").length - 1).toBe(1); }); - it("handles compaction shrink without prompting", () => { + it("handles compaction shrink without prompting", async () => { const promptInputs: string[] = []; const agent = makeAgent(promptInputs); let messages: AgentMessage[] = [ @@ -945,6 +1099,7 @@ describe("advisor", () => { }; const runtime = new AdvisorRuntime(agent, host); runtime.onTurnEnd(); + await Promise.resolve(); expect(promptInputs).toHaveLength(1); messages = [{ role: "user", content: "a", timestamp: 1 } as AgentMessage]; @@ -954,7 +1109,18 @@ describe("advisor", () => { it("reset re-primes the advisor with the full current transcript", async () => { const promptInputs: string[] = []; - const agent = makeAgent(promptInputs); + const { promise: secondPromptDone, resolve: finishSecond } = Promise.withResolvers(); + let promptCalls = 0; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + promptCalls++; + if (promptCalls === 2) finishSecond(); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage]; const host: AdvisorRuntimeHost = { snapshotMessages: () => messages, @@ -972,7 +1138,7 @@ describe("advisor", () => { runtime.reset(); runtime.onTurnEnd(); - await Promise.resolve(); + await secondPromptDone; // The next turn replays the full post-compaction transcript, not just new tail. expect(promptInputs).toHaveLength(2); expect(promptInputs[1]).toContain("summary-bbb"); @@ -980,10 +1146,16 @@ describe("advisor", () => { it("triggers a re-prime and full replay when maintainContext returns true", async () => { const promptInputs: string[] = []; + const { promise: firstPromptDone, resolve: finishFirst } = Promise.withResolvers(); + const { promise: secondPromptDone, resolve: finishSecond } = Promise.withResolvers(); + let promptCalls = 0; let resetCount = 0; const agent: AdvisorAgent = { prompt: async input => { promptInputs.push(input); + promptCalls++; + if (promptCalls === 1) finishFirst(); + else finishSecond(); }, abort: () => {}, reset: () => { @@ -1003,21 +1175,20 @@ describe("advisor", () => { }; const runtime = new AdvisorRuntime(agent, host); - // First turn: normal incremental prompt + // First turn: normal incremental prompt. runtime.onTurnEnd(messages); - await Promise.resolve(); + await firstPromptDone; expect(promptInputs).toHaveLength(1); expect(promptInputs[0]).toContain("aaa"); expect(resetCount).toBe(0); - // Second turn: maintainContext resolves true, triggering a re-prime + // Second turn: maintainContext returns true → re-prime. shouldRePrime = true; messages.push({ role: "user", content: "bbb", timestamp: 2 } as AgentMessage); runtime.onTurnEnd(messages); - await Promise.resolve(); - await Promise.resolve(); + await secondPromptDone; - // The reset cleared history and prompted a full replay (so the batch contains both aaa and bbb) + // Full replay includes both aaa and bbb. expect(promptInputs).toHaveLength(2); expect(promptInputs[1]).toContain("aaa"); expect(promptInputs[1]).toContain("bbb"); diff --git a/packages/coding-agent/src/advisor/runtime.ts b/packages/coding-agent/src/advisor/runtime.ts index 8c88f86ae..4c4587634 100644 --- a/packages/coding-agent/src/advisor/runtime.ts +++ b/packages/coding-agent/src/advisor/runtime.ts @@ -107,11 +107,21 @@ export class AdvisorRuntime { return this.#backlog; } - onTurnEnd(messages?: AgentMessage[]): void { + /** + * True while the advisor model is processing a batch AND newer primary turns + * have already arrived — `#pending` is non-empty during `agent.prompt()`. + * Used by the delivery path to annotate advice that was generated without + * seeing those newer turns. + */ + get hasFreshBacklog(): boolean { + return this.#pending.length > 0; + } + + onTurnEnd(messages?: AgentMessage[], opts?: { willContinue?: boolean }): void { if (this.disposed) return; const all = messages ?? this.host.snapshotMessages(); this.#latestMessages = all; - const render = this.#renderDelta(all); + const render = this.#renderDelta(all, opts?.willContinue ?? false); if (render) { this.#pending.push({ text: render, turns: 1 }); this.#backlog++; @@ -200,7 +210,7 @@ export class AdvisorRuntime { this.#wakeAllWaiters(); } - #renderDelta(messages?: AgentMessage[]): string | null { + #renderDelta(messages?: AgentMessage[], wip = false): string | null { const all = messages ?? this.#latestMessages ?? this.host.snapshotMessages(); if (all.length < this.#lastCount) { this.#lastCount = all.length; @@ -209,7 +219,7 @@ export class AdvisorRuntime { } const delta = all .slice(this.#lastCount) - .filter(m => !(m.role === "custom" && (m as { customType?: string }).customType === "advisor")) + .filter(m => !(m.role === "custom" && m.customType === "advisor")) .map(m => this.#dedupContextMessage(m)); this.#lastCount = all.length; if (delta.length === 0) return null; @@ -223,7 +233,10 @@ export class AdvisorRuntime { expandEditDiffs: true, }); if (!md.trim()) return null; - return `### Session update\n\n${md}`; + const heading = wip + ? "### Session update [in progress — more steps follow]" + : "### Session update"; + return `${heading}\n\n${md}`; } /** @@ -236,14 +249,13 @@ export class AdvisorRuntime { */ #dedupContextMessage(msg: AgentMessage): AgentMessage { if (msg.role !== "custom") return msg; - const type = (msg as { customType?: string }).customType; - if (!type || !PRIMARY_CONTEXT_CUSTOM_TYPES.has(type)) return msg; - const content = (msg as { content?: unknown }).content; - if (typeof content !== "string") return msg; - if (this.#seenContext.get(type) === content) { - return { ...(msg as object), content: "(unchanged — still in effect)" } as AgentMessage; + // Narrowed to CustomMessage: customType and content are properly typed. + if (!PRIMARY_CONTEXT_CUSTOM_TYPES.has(msg.customType)) return msg; + if (typeof msg.content !== "string") return msg; + if (this.#seenContext.get(msg.customType) === msg.content) { + return { ...msg, content: "(unchanged — still in effect)" }; } - this.#seenContext.set(type, content); + this.#seenContext.set(msg.customType, msg.content); return msg; } @@ -284,46 +296,71 @@ export class AdvisorRuntime { } } + /** + * Collect all currently pending deltas into one batch, running + * `maintainContext` for correct token budgeting. Loops until the pending + * queue is stable (nothing new arrived during a maintenance check) or a + * reprime is triggered. Every `await` inside the loop has an epoch guard so + * a reset/dispose mid-await cannot leak a stale batch into the post-reset + * conversation. + * + * Returns `null` when the epoch was invalidated — caller should `continue`. + * Returns `{ batch: null, finalTurns }` when there is nothing to render but + * backlog still needs to be decremented. + */ + async #collectAndMaintainBatch( + epoch: number, + ): Promise<{ batch: string | null; finalTurns: number } | null> { + const initial = this.#pending.splice(0); + let batchText = initial.map(b => b.text).join("\n\n"); + let turns = initial.reduce((sum, b) => sum + b.turns, 0); + + while (true) { + if (this.host.maintainContext) { + const incomingTokens = estimateTokens({ role: "user", content: batchText, timestamp: Date.now() }); + let shouldReprime = false; + try { + shouldReprime = await this.host.maintainContext(incomingTokens); + } catch (err) { + logger.debug("advisor context maintenance failed", { err: String(err) }); + } + // Epoch guard — a reset/dispose during the maintainContext await + // invalidates this batch. + if (this.#epoch !== epoch) return null; + + if (shouldReprime) { + // Tally deltas that arrived during this await before #resetAdvisorContext + // wipes #pending, so finalTurns stays accurate for backlog accounting. + turns += this.#pending.reduce((sum, b) => sum + b.turns, 0); + this.#resetAdvisorContext(false, false); + return { batch: this.#renderDelta(this.#latestMessages), finalTurns: turns }; + } + } + + // Coalesce any deltas that arrived while we were awaiting maintenance. + // If none arrived the batch is stable and we're done; otherwise merge + // and re-check the maintenance budget for the expanded batch. + const late = this.#pending.splice(0); + if (late.length === 0) break; + batchText = [batchText, ...late.map(b => b.text)].join("\n\n"); + turns += late.reduce((sum, b) => sum + b.turns, 0); + } + + return { batch: batchText || null, finalTurns: turns }; + } + async #drain(): Promise { if (this.#busy) return; this.#busy = true; try { while (!this.disposed && this.#pending.length) { - const popped = this.#pending.splice(0); const epoch = this.#epoch; - // Each delta already opens with a `### Session update` heading, so - // join with a blank line rather than a `---` rule. - const candidateBatch = popped.map(b => b.text).join("\n\n"); - const turnsCovered = popped.reduce((sum, b) => sum + b.turns, 0); - const incomingTokens = estimateTokens({ - role: "user", - content: candidateBatch, - timestamp: Date.now(), - }); + const result = await this.#collectAndMaintainBatch(epoch); - let shouldReprime = false; - if (this.host.maintainContext) { - try { - shouldReprime = await this.host.maintainContext(incomingTokens); - } catch (err) { - logger.debug("advisor context maintenance failed", { err: String(err) }); - } - } - // A reset/dispose during context maintenance invalidates this batch. - if (this.#epoch !== epoch) continue; + // Epoch was invalidated during batch collection; restart the loop. + if (result === null) continue; - let batch: string | null; - let finalTurns: number; - if (shouldReprime) { - // Promotion could not fit the advisor's context — re-prime. - const newTurns = this.#pending.reduce((sum, b) => sum + b.turns, 0); - this.#resetAdvisorContext(false, false); - batch = this.#renderDelta(this.#latestMessages); - finalTurns = turnsCovered + newTurns; - } else { - batch = candidateBatch; - finalTurns = turnsCovered; - } + const { batch, finalTurns } = result; if (this.disposed || batch === null) { this.#backlog = Math.max(0, this.#backlog - finalTurns); @@ -333,33 +370,27 @@ export class AdvisorRuntime { let success = false; // Capture the advisor's message count BEFORE the prompt so a failure can - // roll back the user batch + synthetic assistant-error turn `Agent.#runLoop` - // appends to internal state. Without this, a retry would replay the - // failed batch on top of the stale turns and the dropped-after-3 path - // would leak orphan failures into the next successful run's context. + // roll back the user batch + synthetic assistant-error turn Agent.#runLoop + // appends to internal state. Without this, a retry would replay the failed + // batch on top of stale turns and the dropped-after-3 path would leak + // orphan failures into the next successful run's context. const messageSnapshot = this.agent.state.messages.length; try { // Reset the host's per-update advisor state (one-advise-per-update - // gate) before each model cycle, so the new batch starts with a - // fresh budget. Dedupe history persists across cycles. + // gate) before each model cycle so the new batch starts fresh. this.host.beginAdvisorUpdate?.(); await this.agent.prompt(batch); - // `Agent.#runLoop` catches provider/stream failures internally and - // resolves `prompt()` cleanly with the assistant turn ending in - // `stopReason: "error"` and the message recorded on `state.error`. - // Treat that as a failed turn so OpenRouter ZDR-style endpoint - // rejections trip the retry/notify path instead of looking like a - // successful empty cycle. + // Agent.#runLoop catches provider/stream failures internally and + // resolves prompt() cleanly with stopReason: "error". Treat that + // as a failed turn so endpoint rejections trip the retry path. const promptError = this.agent.state.error; if (promptError) throw new Error(promptError); success = true; this.#consecutiveFailures = 0; this.#failureNotified = false; } catch (err) { - // reset()/dispose() aborts the in-flight prompt; the rejection is the - // reset itself, not a transient advisor failure. Drop the stale batch - // (reset already cleared #pending and rewound the cursor) instead of - // requeuing it into the post-reset conversation. + // reset()/dispose() aborts the in-flight prompt; treat it as a + // reset, not a transient failure — drop the stale batch. if (this.#epoch !== epoch) continue; this.#rollbackFailedTurn(messageSnapshot); logger.debug("advisor turn failed", { err: String(err) }); @@ -368,8 +399,7 @@ export class AdvisorRuntime { } catch (hookErr) { logger.debug("advisor onTurnError hook failed", { err: String(hookErr) }); } - // The hook awaits; a reset during it invalidates this batch like the - // prompt await above — drop it instead of requeueing stale content. + // Epoch guard after the async error hook. if (this.#epoch !== epoch) continue; this.#consecutiveFailures++; if (this.#consecutiveFailures >= 3) { @@ -383,9 +413,9 @@ export class AdvisorRuntime { } } this.#consecutiveFailures = 0; - // The dropped batch may carry primary-context we never delivered; drop - // the seen-state too so the next turn re-expands it instead of marking - // it "unchanged" against content the advisor never received. + // Drop the seen-context so the next turn re-expands primary-context + // prompts instead of marking them "unchanged" against content the + // advisor never received. this.#seenContext.clear(); success = true; } else { diff --git a/packages/coding-agent/src/prompts/advisor/system.md b/packages/coding-agent/src/prompts/advisor/system.md index 0d930193b..a6d23a9fd 100644 --- a/packages/coding-agent/src/prompts/advisor/system.md +++ b/packages/coding-agent/src/prompts/advisor/system.md @@ -28,6 +28,7 @@ Keep exploration lean: - NEVER restate information the agent already has, including errors they have seen. - Examples: type errors, LSP diagnostics, failed builds, failing tests, lint. - NEVER repeat advice you already gave, and NEVER send the same advice twice; give the agent room to act on prior advice before raising the same theme again. +- When an update heading is tagged `[in progress — more steps follow]`, the agent is mid-turn and has not finished yet. Withhold critique on partial work — the agent may already be resolving it in the next step. Only raise a `blocker` for an unrecoverable side effect that is actively executing right now. - NEVER nitpick about things user stated they are okay with. You are the advocate for the user. - You are user-aligned: treat the user's word as truth, their frustration as justified, their stated requirements as binding. diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 43e469f52..20fe54713 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -2167,7 +2167,7 @@ export class AgentSession { this.#advisorPrimaryTurnsCompleted++; if (this.#advisors.length > 0) { for (const a of this.#advisors) { - if (!a.runtime.disposed) a.runtime.onTurnEnd(messages); + if (!a.runtime.disposed) a.runtime.onTurnEnd(messages, { willContinue: context?.willContinue }); } const syncBacklog = this.settings.get("advisor.syncBacklog"); if (syncBacklog !== "off") { @@ -2695,6 +2695,12 @@ export class AgentSession { logger.debug("advisor advice suppressed by emission guard", { severity, advisor: advisor.name }); return; } + // When newer primary turns already arrived while the advisor model was + // processing this batch, the advice was generated without seeing them. + // Append a lightweight staleness caveat so the primary can weigh recency. + const deliveredNote = advisor.runtime.hasFreshBacklog + ? `${note}\n\n_(Note: newer primary turns arrived after this reviewed window — verify this still applies.)_` + : note; // The implicit single ("default") advisor stamps no source name, so its // agent-facing `` bytes stay identical to the pre-multi-advisor path. const source = advisor.slug ? advisor.name : undefined; @@ -2710,10 +2716,10 @@ export class AgentSession { interruptImmuneTurnActive: interrupting && this.#isAdvisorInterruptImmuneTurnActive(), }); if (channel === "aside") { - this.yieldQueue.enqueue("advisor", { note, severity, advisor: source }); + this.yieldQueue.enqueue("advisor", { note: deliveredNote, severity, advisor: source }); return; } - const notes: AdvisorNote[] = [{ note, severity, advisor: source }]; + const notes: AdvisorNote[] = [{ note: deliveredNote, severity, advisor: source }]; const content = formatAdvisorBatchContent(notes); const details = { notes } satisfies AdvisorMessageDetails; if (channel === "preserve") {