fix: resolved streaming stability and reference isolation issues
- Prevent object reference sharing between agent snapshots and stream events by deep-cloning tool-call arguments. - Stabilize GFM tables and Mermaid diagrams during streaming by delaying transcript block commits until content finalization. - Implement session resume safety to prevent crashes when working directories are missing. - Add comprehensive test suites to verify streaming commit stability and immutable snapshot isolation.
This commit is contained in:
@@ -1,6 +1,10 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
### Fixed
|
||||
|
||||
- Ensure deep-cloning of tool-call arguments respects own enumerable properties
|
||||
- Prevent direct object references between agent message snapshots and streaming events
|
||||
|
||||
## [16.1.0] - 2026-06-19
|
||||
|
||||
@@ -886,4 +890,4 @@ Initial release under @oh-my-pi scope. See previous releases at [badlogic/pi-mon
|
||||
### Changed
|
||||
|
||||
- `Agent` constructor now has all options optional (empty options use defaults).
|
||||
- `queueMessage()` is now synchronous (no longer returns a Promise).
|
||||
- `queueMessage()` is now synchronous (no longer returns a Promise).
|
||||
@@ -2096,3 +2096,138 @@ describe("agentLoopContinue with AgentMessage", () => {
|
||||
expect(finalMessage.errorMessage).toBe("Deadline exceeded");
|
||||
});
|
||||
});
|
||||
|
||||
describe("agentLoop streaming snapshots", () => {
|
||||
it("deep-clones tool-call arguments into message_update snapshots, copying only own enumerable properties", async () => {
|
||||
const context: AgentContext = {
|
||||
systemPrompt: ["You are helpful."],
|
||||
messages: [],
|
||||
tools: [],
|
||||
};
|
||||
const config: AgentLoopConfig = {
|
||||
model: createMockModel().model,
|
||||
convertToLlm: identityConverter,
|
||||
};
|
||||
|
||||
// Arguments carry a nested object, a nested array, primitives, and an
|
||||
// INHERITED enumerable property. cloneUnknown must deep-clone the own
|
||||
// nested structures (fresh references), pass primitives through by value,
|
||||
// and copy only OWN enumerable keys — the inherited key must not leak in
|
||||
// (this pins the `Object.hasOwn` guard that replaced `Object.entries`).
|
||||
const inheritedProto = { inheritedKey: "from-prototype" };
|
||||
const innerArray = [2, 3];
|
||||
const base: Record<string, unknown> = Object.create(inheritedProto);
|
||||
const sourceArgs: Record<string, unknown> = Object.assign(base, {
|
||||
nestedObj: { a: 1, b: "two" },
|
||||
nestedArr: [1, innerArray, { c: 4 }],
|
||||
num: 42,
|
||||
str: "hi",
|
||||
flag: true,
|
||||
nul: null,
|
||||
});
|
||||
const toolCall = { type: "toolCall" as const, id: "tc-clone", name: "noop", arguments: sourceArgs };
|
||||
|
||||
// Turn 0 streams the tool call; the unknown tool produces an error result
|
||||
// and the loop calls the model again — turn 1 returns plain text so the
|
||||
// loop terminates instead of spinning forever.
|
||||
let turn = 0;
|
||||
const streamFn = () => {
|
||||
const stream = new AssistantMessageEventStream();
|
||||
if (turn++ === 0) {
|
||||
const partial = createAssistantMessage([toolCall], "toolUse");
|
||||
stream.push({ type: "start", partial });
|
||||
stream.push({ type: "toolcall_start", contentIndex: 0, partial });
|
||||
stream.push({ type: "toolcall_delta", contentIndex: 0, delta: "{}", partial });
|
||||
stream.push({ type: "toolcall_end", contentIndex: 0, toolCall, partial });
|
||||
stream.push({ type: "done", reason: "toolUse", message: partial });
|
||||
} else {
|
||||
const partial = createAssistantMessage([{ type: "text", text: "done" }], "stop");
|
||||
stream.push({ type: "start", partial });
|
||||
stream.push({ type: "text_delta", contentIndex: 0, delta: "done", partial });
|
||||
stream.push({ type: "done", reason: "stop", message: partial });
|
||||
}
|
||||
return stream;
|
||||
};
|
||||
|
||||
const events: AgentEvent[] = [];
|
||||
const stream = agentLoop([createUserMessage("call noop")], context, config, undefined, streamFn);
|
||||
for await (const event of stream) {
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
const toolUpdate = events.find(
|
||||
(e): e is Extract<AgentEvent, { type: "message_update" }> =>
|
||||
e.type === "message_update" &&
|
||||
e.message.role === "assistant" &&
|
||||
e.message.content.some(c => c.type === "toolCall"),
|
||||
);
|
||||
expect(toolUpdate).toBeDefined();
|
||||
if (toolUpdate?.message.role !== "assistant") throw new Error("missing tool-call update");
|
||||
const block = toolUpdate.message.content.find(c => c.type === "toolCall");
|
||||
if (block?.type !== "toolCall") throw new Error("missing tool-call block");
|
||||
const cloned: Record<string, unknown> = block.arguments;
|
||||
|
||||
// Fresh top-level object, not the source reference.
|
||||
expect(cloned).not.toBe(sourceArgs);
|
||||
// Own enumerable keys only — the inherited property must not appear.
|
||||
expect("inheritedKey" in cloned).toBe(false);
|
||||
expect(Object.hasOwn(cloned, "inheritedKey")).toBe(false);
|
||||
// Nested object: deep-cloned (equal value, distinct reference).
|
||||
expect(cloned.nestedObj).toEqual({ a: 1, b: "two" });
|
||||
expect(cloned.nestedObj).not.toBe(sourceArgs.nestedObj);
|
||||
// Nested array: deep-cloned through the Array.isArray fast path, recursively.
|
||||
expect(Array.isArray(cloned.nestedArr)).toBe(true);
|
||||
expect(cloned.nestedArr).toEqual([1, [2, 3], { c: 4 }]);
|
||||
expect(cloned.nestedArr).not.toBe(sourceArgs.nestedArr);
|
||||
expect((cloned.nestedArr as unknown[])[1]).not.toBe(innerArray);
|
||||
// Primitives pass through by value.
|
||||
expect(cloned.num).toBe(42);
|
||||
expect(cloned.str).toBe("hi");
|
||||
expect(cloned.flag).toBe(true);
|
||||
expect(cloned.nul).toBeNull();
|
||||
});
|
||||
|
||||
it("shares one immutable snapshot between message and assistantMessageEvent.partial on message_update", async () => {
|
||||
const context: AgentContext = {
|
||||
systemPrompt: ["You are helpful."],
|
||||
messages: [],
|
||||
tools: [],
|
||||
};
|
||||
const config: AgentLoopConfig = {
|
||||
model: createMockModel().model,
|
||||
convertToLlm: identityConverter,
|
||||
};
|
||||
|
||||
const livePartial = createAssistantMessage([{ type: "text", text: "Hi" }], "stop");
|
||||
const streamFn = () => {
|
||||
const stream = new AssistantMessageEventStream();
|
||||
stream.push({ type: "start", partial: livePartial });
|
||||
stream.push({ type: "text_start", contentIndex: 0, partial: livePartial });
|
||||
stream.push({ type: "text_delta", contentIndex: 0, delta: "Hi", partial: livePartial });
|
||||
stream.push({ type: "text_end", contentIndex: 0, content: "Hi", partial: livePartial });
|
||||
stream.push({ type: "done", reason: "stop", message: livePartial });
|
||||
return stream;
|
||||
};
|
||||
|
||||
const events: AgentEvent[] = [];
|
||||
const stream = agentLoop([createUserMessage("say hi")], context, config, undefined, streamFn);
|
||||
for await (const event of stream) {
|
||||
events.push(event);
|
||||
}
|
||||
|
||||
const update = events.find(
|
||||
(e): e is Extract<AgentEvent, { type: "message_update" }> => e.type === "message_update",
|
||||
);
|
||||
expect(update).toBeDefined();
|
||||
if (!update) throw new Error("missing message_update");
|
||||
const ame = update.assistantMessageEvent;
|
||||
expect("partial" in ame).toBe(true);
|
||||
if (!("partial" in ame)) throw new Error("expected a partial-bearing assistant event");
|
||||
// Alias contract: `message` and `assistantMessageEvent.partial` are ONE
|
||||
// shared snapshot...
|
||||
expect(update.message).toBe(ame.partial);
|
||||
// ...and that snapshot is an independent deep clone of the live streaming
|
||||
// partial, never the mutable partial object itself.
|
||||
expect(update.message).not.toBe(livePartial);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
# Changelog
|
||||
|
||||
## [Unreleased]
|
||||
|
||||
### Added
|
||||
|
||||
- Added a welcome-screen tip for the `/advisor` runtime. Tips ending in a `[NEW]` marker now render a bold rainbow `NEW!` tag (it shimmers across the welcome intro's animation frames, then settles into a still rainbow) and are weighted to surface more often in the random tip rotation.
|
||||
@@ -14,7 +13,9 @@
|
||||
|
||||
### Fixed
|
||||
|
||||
- Prevented stale fragments in scrollback when streaming GFM tables by delaying commit until final
|
||||
- Resuming a session whose project directory no longer exists (deleted or renamed worktree) no longer crashes with an unhandled `ENOENT … chdir` rejection. The resume now keeps the current working directory instead of trying to `chdir` into the missing path, across the in-session selector, the `--resume` startup picker, and `SessionManager.open`/`continueRecent`.
|
||||
- Fixed a streaming mermaid diagram stranding stale diagram fragments in native scrollback once the reply scrolled past the viewport (cleared only by a full repaint / Ctrl+L). While streaming, the assistant block defaulted to commit-stable, so the transcript advertised the diagram's scrolled-off rows as durable snapshot content and the renderer committed an intermediate layout to immutable terminal history; the diagram's later re-layout then froze that superseded fragment in scrollback. A still-streaming reply carrying a ` ```mermaid ` fence is now commit-unstable, so the diagram stays wholly in the repaintable live region and commits once, at its final layout, when the turn finalizes.
|
||||
|
||||
## [16.1.1] - 2026-06-19
|
||||
|
||||
|
||||
@@ -16,6 +16,24 @@ import { type CacheInvalidation, CacheInvalidationMarkerComponent } from "./cach
|
||||
*/
|
||||
const MAX_TRANSCRIPT_ERROR_LINES = 8;
|
||||
|
||||
/**
|
||||
* A fenced ` ```mermaid ` block anchored at a line start (≤3 leading spaces,
|
||||
* 3+ backticks, optional info-string whitespace) so prose that merely mentions
|
||||
* a mermaid fence inline does not match. See
|
||||
* {@link AssistantMessageComponent.isTranscriptBlockCommitStable}.
|
||||
*/
|
||||
const LIVE_MERMAID_FENCE = /^ {0,3}`{3,}[ \t]*mermaid\b/m;
|
||||
|
||||
/**
|
||||
* A GFM table delimiter row (`| --- | :--: |`, with or without bounding pipes)
|
||||
* anchored at a line start. The header row alone does not render a table — this
|
||||
* delimiter is what makes Markdown lay one out, and a streaming table re-aligns
|
||||
* its columns as rows arrive. Requires at least one column pipe so a bare
|
||||
* thematic break (`---`) does not match. See
|
||||
* {@link AssistantMessageComponent.isTranscriptBlockCommitStable}.
|
||||
*/
|
||||
const MARKDOWN_TABLE_DELIMITER = /^ {0,3}\|?(?:[ \t]*:?-+:?[ \t]*\|)+[ \t]*:?-*:?[ \t]*$/m;
|
||||
|
||||
/**
|
||||
* Frames for the streaming "thinking" pulse rendered in place of a hidden
|
||||
* thinking block while the model is still producing it. A single fixed-width
|
||||
@@ -36,6 +54,15 @@ export class AssistantMessageComponent extends Container {
|
||||
#convertedKittyImages = new Map<string, ImageContent>();
|
||||
#kittyConversionsInFlight = new Set<string>();
|
||||
#transcriptBlockFinalized: boolean;
|
||||
/**
|
||||
* True while a non-finalized text item carries reflowing Markdown — a
|
||||
* ` ```mermaid ` fence or a GFM table — whose layout re-flows every frame as
|
||||
* source arrives (a diagram reshaping, a table re-aligning its columns), so
|
||||
* no prefix is byte-stable until the message finalizes. See
|
||||
* {@link isTranscriptBlockCommitStable}. Recomputed in {@link updateContent}
|
||||
* ahead of the fast-path return, so it tracks every stream tick.
|
||||
*/
|
||||
#hasLiveReflowingMarkdown = false;
|
||||
/**
|
||||
* When true, the turn-ending `Error: …` line for `stopReason === "error"` is
|
||||
* suppressed because the same error is currently shown in the pinned banner
|
||||
@@ -192,6 +219,21 @@ export class AssistantMessageComponent extends Container {
|
||||
return this.#transcriptBlockFinalized;
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether this still-live block's scrolled-off rows may be committed to
|
||||
* immutable native scrollback (the {@link TranscriptContainer} durable-
|
||||
* snapshot path). Reflowing Markdown — a streaming mermaid diagram or a GFM
|
||||
* table — re-lays-out its body as source arrives (the diagram reshapes, the
|
||||
* table re-aligns its columns), so committing an intermediate layout strands
|
||||
* a stale fragment in native scrollback that only a full repaint (Ctrl+L) can
|
||||
* clear. While such content is still streaming the block therefore stays
|
||||
* wholly in the repaintable live region and commits once, at its final
|
||||
* layout, when the turn finalizes.
|
||||
*/
|
||||
isTranscriptBlockCommitStable(): boolean {
|
||||
return this.#transcriptBlockFinalized || !this.#hasLiveReflowingMarkdown;
|
||||
}
|
||||
|
||||
getTranscriptBlockVersion(): number {
|
||||
return this.#blockVersion;
|
||||
}
|
||||
@@ -418,6 +460,17 @@ export class AssistantMessageComponent extends Container {
|
||||
this.#lastMessage = message;
|
||||
this.#lastUpdateTransient = opts?.transient === true;
|
||||
|
||||
// Streaming reflowing Markdown (a mermaid diagram reshaping, a GFM table
|
||||
// re-aligning columns) re-lays-out its body each frame; see
|
||||
// isTranscriptBlockCommitStable. Detect it from raw text — a Markdown
|
||||
// parser only resolves these once the closing fence / delimiter row
|
||||
// arrives, but the stale native-scrollback commits happen mid-stream.
|
||||
this.#hasLiveReflowingMarkdown = message.content.some(
|
||||
content =>
|
||||
content.type === "text" &&
|
||||
(LIVE_MERMAID_FENCE.test(content.text) || MARKDOWN_TABLE_DELIMITER.test(content.text)),
|
||||
);
|
||||
|
||||
// Fast path: reuse Markdown children when shape is stable during streaming
|
||||
if (this.#tryFastPathUpdate(message)) return;
|
||||
|
||||
|
||||
@@ -98,6 +98,74 @@ describe("AssistantMessageComponent mermaid markdown", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("AssistantMessageComponent reflowing-markdown commit stability", () => {
|
||||
// A streaming reply is built empty then fed via updateContent (the live path);
|
||||
// passing a message to the constructor would mark it finalized.
|
||||
it("is commit-unstable while a streaming reply still carries a mermaid fence", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("Here is the flow:\n\n```mermaid\nflowchart TD\n A-->B"));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(false);
|
||||
});
|
||||
|
||||
it("becomes commit-stable once the mermaid reply finalizes", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("```mermaid\nflowchart TD\n A-->B\n```"));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(false);
|
||||
component.markTranscriptBlockFinalized();
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
|
||||
it("stays commit-stable for a streaming reply without a mermaid fence", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("A long normal reply.\n\n- one\n- two\n\nMore prose follows."));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
|
||||
it("does not trip on prose that mentions a mermaid fence inline", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("Wrap the diagram in a ```mermaid block to render it."));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
|
||||
it("is commit-unstable for intro prose followed by a streaming mermaid tail", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(
|
||||
createAssistantMessage("Intro prose above the diagram.\n\n```mermaid\nflowchart TD\n A-->B"),
|
||||
);
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(false);
|
||||
});
|
||||
|
||||
it("is commit-unstable while a streaming reply still renders a GFM table", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("Results:\n\n| Name | Score |\n| --- | --- |\n| a | 1 |"));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(false);
|
||||
});
|
||||
|
||||
it("becomes commit-stable once a table reply finalizes", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("| Name | Score |\n| --- | --- |\n| a | 1 |\n| b | 2 |"));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(false);
|
||||
component.markTranscriptBlockFinalized();
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
|
||||
it("stays commit-stable for pipe-heavy prose with no table delimiter row", () => {
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(
|
||||
createAssistantMessage("Weigh cost | benefit | risk before deciding, and note `a || b` short-circuits."),
|
||||
);
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
|
||||
it("stays commit-stable for a streaming table header before its delimiter row arrives", () => {
|
||||
// Header alone is just a paragraph with pipes — Markdown lays out no table,
|
||||
// and nothing re-flows, until the delimiter row streams in.
|
||||
const component = new AssistantMessageComponent();
|
||||
component.updateContent(createAssistantMessage("| Name | Score |"));
|
||||
expect(component.isTranscriptBlockCommitStable()).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("AssistantMessageComponent thinking renderers", () => {
|
||||
it("renders all extension outputs below visible thinking blocks in registration order", () => {
|
||||
const contexts: Array<{ contentIndex: number; thinkingIndex: number; text: string }> = [];
|
||||
|
||||
@@ -329,6 +329,44 @@ describe("TranscriptContainer", () => {
|
||||
// blockStart 2 + (7 rows - TAIL_VOLATILITY_ROWS holdback of 4) = 5.
|
||||
expect(container.getNativeScrollbackCommitSafeEnd()).toBe(5);
|
||||
});
|
||||
|
||||
it("withholds the durable snapshot commit for a streaming mermaid reply, then promotes it on finalize", () => {
|
||||
// A fenced mermaid diagram re-lays-out its whole body every frame as the
|
||||
// reply streams. If its scrolled-off rows were committed to immutable
|
||||
// native scrollback (the durable-snapshot path) the later re-layout would
|
||||
// strand a stale diagram fragment in history that only Ctrl+L clears, so
|
||||
// the live mermaid reply must advertise no durable snapshot end.
|
||||
const container = new TranscriptContainer();
|
||||
container.addChild(new StreamingBlock(["earlier turn"], true));
|
||||
const assistant = new AssistantMessageComponent();
|
||||
assistant.updateContent(
|
||||
makeAssistantMessage({ content: [{ type: "text", text: "```mermaid\nflowchart TD\n A-->B\n```" }] }),
|
||||
);
|
||||
container.addChild(assistant);
|
||||
|
||||
for (let frame = 0; frame < 4; frame++) container.render(60);
|
||||
expect(container.getNativeScrollbackSnapshotSafeEnd()).toBeUndefined();
|
||||
|
||||
// Finalized: the final layout is permanent and commits like any block.
|
||||
assistant.markTranscriptBlockFinalized();
|
||||
container.render(60);
|
||||
expect(container.getNativeScrollbackSnapshotSafeEnd()).toBeGreaterThan(0);
|
||||
});
|
||||
|
||||
it("still commits the durable snapshot of a streaming reply without a mermaid fence", () => {
|
||||
// Guard the narrow scope: an ordinary streaming reply stays commit-stable,
|
||||
// so its settled rows still reach native scrollback while it streams.
|
||||
const container = new TranscriptContainer();
|
||||
container.addChild(new StreamingBlock(["earlier turn"], true));
|
||||
const assistant = new AssistantMessageComponent();
|
||||
assistant.updateContent(
|
||||
makeAssistantMessage({ content: [{ type: "text", text: "A normal streamed answer with **bold** text." }] }),
|
||||
);
|
||||
container.addChild(assistant);
|
||||
|
||||
for (let frame = 0; frame < 4; frame++) container.render(60);
|
||||
expect(container.getNativeScrollbackSnapshotSafeEnd()).toBeGreaterThan(0);
|
||||
});
|
||||
it("does not re-render finalized rows already committed to native scrollback", () => {
|
||||
const container = new TranscriptContainer();
|
||||
const committed = new CountingFinalizedBlock(["committed"]);
|
||||
|
||||
Reference in New Issue
Block a user