fix(advisor): reduce stale advisories via delta coalescing, WIP markers, and delivery annotation
Three related changes that address the pattern of the advisor flagging things
the primary already fixed:
Fix 1 — coalesce late-arriving deltas before agent.prompt (runtime.ts)
Refactored #drain into a reusable #collectAndMaintainBatch helper that loops
until the pending queue is stable (no new deltas arrive during a maintenance
check) before calling agent.prompt. Previously, any turn queued during the
maintainContext await was deferred a full extra model-call cycle; now it is
merged into the current batch after re-checking the token budget for the
expanded payload. Every await in the loop has an epoch guard so a
reset/dispose mid-await cannot leak a stale batch. finalTurns always counts
all merged turns so #backlog decrements correctly.
Fix 2 — hasFreshBacklog + delivery-time staleness annotation (runtime.ts, agent-session.ts)
Added AdvisorRuntime.hasFreshBacklog getter (true when #pending.length > 0
while agent.prompt is running — i.e., newer primary turns arrived after the
reviewed window). #routeAdvice checks it at delivery time and appends a
lightweight caveat to the note so the primary agent knows to verify before
acting. Uses #pending.length not #backlog, which is always > 0 mid-call.
Fix 3 — willContinue WIP marker in rendered delta + system prompt (runtime.ts, agent-session.ts, system.md)
onTurnEnd now accepts { willContinue } and passes it through to #renderDelta,
which tags the heading '[in progress — more steps follow]' for intermediate
turns. The agent-session.ts call site passes context.willContinue. The advisor
system prompt instructs the model to withhold critique on WIP updates.
Also fixed pre-existing inline casts in #renderDelta and #dedupContextMessage
that suppressed the type checker instead of using the narrowing already provided
by the role discriminant.
All 75 advisor tests pass; pre-existing type errors in cursor.ts are unrelated.
This commit is contained in:
@@ -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<void>();
|
||||
const { promise: secondPromptDone, resolve: finishSecondPrompt } = Promise.withResolvers<void>();
|
||||
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<void>();
|
||||
const { promise: finishFirstMaintain, resolve: releaseFirstMaintain } = Promise.withResolvers<boolean>();
|
||||
const { promise: firstPromptStarted, resolve: startFirstPrompt } = Promise.withResolvers<void>();
|
||||
const { promise: secondPromptStarted, resolve: startSecondPrompt } = Promise.withResolvers<void>();
|
||||
const { promise: finishFirstPrompt, resolve: releaseFirstPrompt } = Promise.withResolvers<void>();
|
||||
const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers<void>();
|
||||
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<void>();
|
||||
const { promise: finishFirstMaintain, resolve: releaseFirstMaintain } = Promise.withResolvers<boolean>();
|
||||
const { promise: promptStarted, resolve: startPrompt } = Promise.withResolvers<void>();
|
||||
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<void>();
|
||||
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<void>();
|
||||
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<void>();
|
||||
const { promise: firstPromptDone, resolve: finishFirstPrompt } = Promise.withResolvers<void>();
|
||||
const { promise: secondPromptDone, resolve: finishSecondPrompt } = Promise.withResolvers<void>();
|
||||
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<void>();
|
||||
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<void>();
|
||||
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<void>();
|
||||
const { promise: secondPromptDone, resolve: finishSecond } = Promise.withResolvers<void>();
|
||||
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");
|
||||
|
||||
@@ -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<void> {
|
||||
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 {
|
||||
|
||||
@@ -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.
|
||||
</communication>
|
||||
|
||||
@@ -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 `<advisory>` 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") {
|
||||
|
||||
Reference in New Issue
Block a user