fix app stream head replacement during init hydration (#663)

This commit is contained in:
Mohamed Boudra
2026-05-02 14:13:39 +08:00
committed by GitHub
parent bb889dae99
commit d46d92528a
2 changed files with 98 additions and 3 deletions

View File

@@ -910,6 +910,67 @@ describe("processAgentStreamEvents", () => {
lastActivityAt: new Date(3000),
});
});
it("keeps a live Claude assistant paragraph contiguous when init tail hydration lands mid-stream", () => {
const seq186Text = "Now let me write the updated tool-call-detail-parser.ts with all the sub";
const seq187Text = "agent additions.";
const liveSeq186 = processAgentStreamEvent({
...baseStreamInput,
event: makeTimelineEvent(seq186Text),
seq: 186,
epoch: "epoch-1",
currentTail: [],
currentHead: [],
currentCursor: undefined,
timestamp: new Date("2026-05-02T10:00:00.186Z"),
});
expect(getAssistantTexts(liveSeq186.head)).toEqual([seq186Text]);
expect(getAssistantTexts(liveSeq186.tail)).toEqual([]);
const initTailHydration = processTimelineResponse({
...baseTimelineInput,
currentTail: liveSeq186.tail,
currentHead: liveSeq186.head,
currentCursor: liveSeq186.cursor ?? undefined,
isInitializing: true,
hasActiveInitDeferred: true,
initRequestDirection: "tail",
payload: {
agentId: "agent-1",
direction: "tail",
reset: false,
epoch: "epoch-1",
startCursor: { seq: 186 },
endCursor: { seq: 186 },
entries: [makeTimelineEntry(186, seq186Text)],
error: null,
},
});
expect(getAssistantTexts(initTailHydration.tail)).toEqual([]);
expect(getAssistantTexts(initTailHydration.head)).toEqual([seq186Text]);
const liveSeq187 = processAgentStreamEvent({
...baseStreamInput,
event: makeTimelineEvent(seq187Text),
seq: 187,
epoch: "epoch-1",
currentTail: initTailHydration.tail,
currentHead: initTailHydration.head,
currentCursor: initTailHydration.cursor ?? undefined,
timestamp: new Date("2026-05-02T10:00:00.187Z"),
});
const finalAssistantItems = [...liveSeq187.tail, ...liveSeq187.head].filter(
(item): item is Extract<StreamItem, { kind: "assistant_message" }> =>
item.kind === "assistant_message",
);
expect(finalAssistantItems).toHaveLength(1);
expect(finalAssistantItems[0]?.text).toBe(`${seq186Text}${seq187Text}`);
});
});
describe("createAgentStreamReducerQueue", () => {

View File

@@ -99,16 +99,49 @@ interface TimelinePathResult {
sideEffects: TimelineReducerSideEffect[];
}
function preserveReplacePathAssistantHead(params: {
tail: StreamItem[];
currentHead: StreamItem[];
}): {
tail: StreamItem[];
head: StreamItem[];
} {
const { tail, currentHead } = params;
const liveAssistant = currentHead.findLast(
(item): item is Extract<StreamItem, { kind: "assistant_message" }> =>
item.kind === "assistant_message",
);
if (!liveAssistant) {
return { tail, head: [] };
}
const tailAssistant = tail.at(-1);
if (!tailAssistant || tailAssistant.kind !== "assistant_message") {
return { tail, head: currentHead };
}
if (!liveAssistant.text.startsWith(tailAssistant.text)) {
return { tail, head: [] };
}
return {
tail: tail.slice(0, -1),
head: [{ ...liveAssistant, text: tailAssistant.text }],
};
}
function applyTimelineReplacePath(args: {
timelineUnits: TimelineUnit[];
payload: ProcessTimelineResponseInput["payload"];
bootstrapPolicy: ReturnType<typeof deriveBootstrapTailTimelinePolicy>;
currentHead: StreamItem[];
toHydratedEvents: (
units: TimelineUnit[],
) => Array<{ event: AgentStreamEventPayload; timestamp: Date }>;
}): TimelinePathResult {
const { timelineUnits, payload, bootstrapPolicy, toHydratedEvents } = args;
const tail = hydrateStreamState(toHydratedEvents(timelineUnits), { source: "canonical" });
const { timelineUnits, payload, bootstrapPolicy, currentHead, toHydratedEvents } = args;
const hydratedTail = hydrateStreamState(toHydratedEvents(timelineUnits), { source: "canonical" });
const { tail, head } = preserveReplacePathAssistantHead({
tail: hydratedTail,
currentHead,
});
const cursor: TimelineCursor | null =
payload.startCursor && payload.endCursor
? {
@@ -121,7 +154,7 @@ function applyTimelineReplacePath(args: {
if (bootstrapPolicy.catchUpCursor) {
sideEffects.push({ type: "catch_up", cursor: bootstrapPolicy.catchUpCursor });
}
return { tail, head: [], cursor, cursorChanged: true, sideEffects };
return { tail, head, cursor, cursorChanged: true, sideEffects };
}
interface IncrementalAcceptResult {
@@ -307,6 +340,7 @@ export function processTimelineResponse(
timelineUnits,
payload,
bootstrapPolicy,
currentHead,
toHydratedEvents,
})
: applyTimelineIncrementalPath({