Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
103 changes: 89 additions & 14 deletions apps/api/src/modules/sessions/domain/session-event-stream-fold.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,14 @@ import type { RuntimeEventId, SessionRunId } from "@mosoo/id";
// events. The merge rules mirror the pre-persistence compactor
// (runtime-event-compaction.ts): deltas append, snapshots prefer the longer
// prefix-matching text.
//
// Rows carry no message identity, so a stream whose driver-side identity
// fractured (a dropped message_start splits one reply across several
// started/delta groups, YEF-884) still folds into several fragment rows, and
// the trailing message.added snapshot lands after its stream already closed.
// To heal both shapes, closed fragment rows are kept as supersede candidates:
// a snapshot whose text prefix-matches the concatenated fragments replaces
// them with a single row instead of rendering a duplicate.

export interface StreamFoldableSessionEventRow {
content_text: string;
Expand Down Expand Up @@ -59,6 +67,11 @@ interface OpenStreamGroup<R extends StreamFoldableSessionEventRow> {
runId: SessionRunId | null;
}

interface ClosedFragmentSegment {
content: string;
outputIndex: number;
}

function appendStreamText(current: string, next: string): string {
if (next.length === 0) {
return current;
Expand Down Expand Up @@ -87,6 +100,14 @@ function mergeSnapshotText(current: string, next: string): string {
return `${current}${next}`;
}

function isSnapshotOfFragments(fragments: string, snapshot: string): boolean {
if (fragments.length === 0 || snapshot.length === 0) {
return false;
}

return snapshot.startsWith(fragments) || fragments.startsWith(snapshot);
}

function createOpenStreamGroup<R extends StreamFoldableSessionEventRow>(
row: R,
classification: StreamRowClassification,
Expand Down Expand Up @@ -140,15 +161,60 @@ export function foldStreamedSessionEventRows<R extends StreamFoldableSessionEven
rows: readonly R[],
options: { flushOpenStreams: boolean },
): FoldedStreamedSessionEventRows<R> {
const output: R[] = [];
const output: (R | null)[] = [];
const openGroups = new Map<string, OpenStreamGroup<R>>();
const fragmentSegments = new Map<string, ClosedFragmentSegment[]>();

function closeGroup(group: OpenStreamGroup<R>): void {
function closeGroupAsFragment(key: string, group: OpenStreamGroup<R>): void {
const folded = createFoldedStreamRow(group);

if (folded !== null) {
if (folded === null) {
return;
}

output.push(folded);
const segments = fragmentSegments.get(key) ?? [];
segments.push({ content: folded.content_text, outputIndex: output.length - 1 });
fragmentSegments.set(key, segments);
}

// A snapshot row is the authoritative text of the message it closes. When
// its text extends (or repeats) the fragment rows already emitted for the
// same stream key, those fragments were partial views of this snapshot:
// collapse them into one row at the first fragment's timeline position.
function closeGroupWithSnapshot(key: string, folded: R | null): void {
if (folded === null) {
return;
}

const segments = fragmentSegments.get(key) ?? [];
fragmentSegments.delete(key);
const fragments = segments.map((segment) => segment.content).join("");

if (!isSnapshotOfFragments(fragments, folded.content_text)) {
output.push(folded);
return;
}

const firstSegment = segments[0];
const anchorRow = firstSegment === undefined ? null : output[firstSegment.outputIndex];

if (firstSegment === undefined || anchorRow === null || anchorRow === undefined) {
output.push(folded);
return;
}

for (const segment of segments) {
output[segment.outputIndex] = null;
}

output[firstSegment.outputIndex] = {
...folded,
content_text:
folded.content_text.length >= fragments.length ? folded.content_text : fragments,
occurred_at: anchorRow.occurred_at,
seq: anchorRow.seq,
};
}

for (const row of rows) {
Expand All @@ -162,7 +228,7 @@ export function foldStreamedSessionEventRows<R extends StreamFoldableSessionEven
for (const [key, group] of openGroups) {
if (group.runId === row.run_id) {
openGroups.delete(key);
closeGroup(group);
closeGroupAsFragment(key, group);
}
}
}
Expand All @@ -177,7 +243,7 @@ export function foldStreamedSessionEventRows<R extends StreamFoldableSessionEven
if (classification.phase === "started") {
if (existing !== null) {
openGroups.delete(key);
closeGroup(existing);
closeGroupAsFragment(key, existing);
}

const group = createOpenStreamGroup(row, classification);
Expand All @@ -187,35 +253,44 @@ export function foldStreamedSessionEventRows<R extends StreamFoldableSessionEven
}

if (classification.phase === "added" && existing === null) {
// A standalone message.added row is already a complete message (legacy
// compacted rows and non-streamed messages) — pass it through verbatim.
output.push(row);
// A snapshot with no open stream is either a complete standalone
// message (legacy compacted rows and non-streamed messages) or the
// authoritative copy of fragments that already closed; the supersede
// check keeps the first case verbatim.
closeGroupWithSnapshot(key, row.content_text === classification.placeholder ? null : row);

if (row.content_text === classification.placeholder) {
output.push(row);
}
continue;
}

const group = existing ?? createOpenStreamGroup(row, classification);
mergeStreamRow(group, row, classification);

if (classification.phase === "completed" || classification.phase === "added") {
if (classification.phase === "added") {
openGroups.delete(key);
closeGroupWithSnapshot(key, createFoldedStreamRow(group));
} else if (classification.phase === "completed") {
openGroups.delete(key);
closeGroup(group);
closeGroupAsFragment(key, group);
} else if (existing === null) {
openGroups.set(key, group);
}
}

if (options.flushOpenStreams) {
for (const group of openGroups.values()) {
closeGroup(group);
for (const [key, group] of openGroups) {
closeGroupAsFragment(key, group);
}

return { openStreamRows: [], rows: output };
return { openStreamRows: [], rows: output.filter((row): row is R => row !== null) };
}

return {
openStreamRows: [...openGroups.values()]
.flatMap((group) => group.rows)
.toSorted((a, b) => a.seq - b.seq),
rows: output,
rows: output.filter((row): row is R => row !== null),
};
}
135 changes: 135 additions & 0 deletions apps/api/tests/session-event-stream-fold.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,141 @@ describe("session event stream folding", () => {
expect(secondFold.rows[0]).toMatchObject({ content_text: "流式输出完成", id: "m-end" });
});

test("supersedes a closed stream with its trailing snapshot instead of duplicating it", () => {
// Streamed replies persist message.completed (at message_stop) before the
// aggregated assistant snapshot arrives, so the snapshot lands after its
// stream already closed and must replace the folded row, not repeat it.
const folded = foldStreamedSessionEventRows(
[
row({ content: "Message updated.", eventType: "message.started", id: "m-start", seq: 1 }),
row({ content: "P", eventType: "message.delta", id: "m-1", seq: 2 }),
row({
content: "ong. What would you like to work on?",
eventType: "message.delta",
id: "m-2",
seq: 3,
}),
row({ content: "Message updated.", eventType: "message.completed", id: "m-end", seq: 4 }),
row({
content: "Pong. What would you like to work on?",
eventType: "message.added",
id: "m-added",
seq: 5,
}),
],
{ flushOpenStreams: true },
);

expect(folded.rows).toHaveLength(1);
expect(folded.rows[0]).toMatchObject({
content_text: "Pong. What would you like to work on?",
event_type: "message.added",
id: "m-added",
seq: 1,
});
});

test("collapses a fractured stream whose snapshot matches the concatenated fragments", () => {
// YEF-884: a dropped message_start fractures one reply into per-fragment
// messages, so the rows arrive as started/delta pairs per fragment plus a
// final full snapshot. The snapshot equals the fragment concatenation and
// must fold everything into a single timeline entry.
const folded = foldStreamedSessionEventRows(
[
row({ content: "Message updated.", eventType: "message.started", id: "s-1", seq: 1 }),
row({ content: "P", eventType: "message.delta", id: "m-1", seq: 2 }),
row({ content: "Message updated.", eventType: "message.started", id: "s-2", seq: 3 }),
row({
content: "ong. What would you like to work on?",
eventType: "message.delta",
id: "m-2",
seq: 4,
}),
row({ content: "Message updated.", eventType: "message.started", id: "s-3", seq: 5 }),
row({
content: "Pong. What would you like to work on?",
eventType: "message.delta",
id: "m-3",
seq: 6,
}),
row({
content: "Pong. What would you like to work on?",
eventType: "message.added",
id: "m-added",
seq: 7,
}),
row({ content: "Message updated.", eventType: "message.completed", id: "c-3", seq: 8 }),
row({ content: "Message updated.", eventType: "message.completed", id: "c-1", seq: 9 }),
row({ content: "Message updated.", eventType: "message.completed", id: "c-2", seq: 10 }),
row({
content: "run-1",
eventType: "run.completed",
id: "r-done",
processType: "run.completed",
seq: 11,
}),
],
{ flushOpenStreams: true },
);

expect(folded.rows.map((entry) => entry.content_text)).toEqual([
"Pong. What would you like to work on?",
"run-1",
]);
expect(folded.rows[0]).toMatchObject({ id: "m-added", seq: 1 });
});

test("supersedes a truncated stream with the longer prefix-matching snapshot", () => {
const folded = foldStreamedSessionEventRows(
[
row({ content: "Final ans", eventType: "message.delta", id: "m-1", seq: 1 }),
row({ content: "Message updated.", eventType: "message.completed", id: "m-end", seq: 2 }),
row({ content: "Final answer.", eventType: "message.added", id: "m-added", seq: 3 }),
],
{ flushOpenStreams: true },
);

expect(folded.rows).toHaveLength(1);
expect(folded.rows[0]).toMatchObject({
content_text: "Final answer.",
id: "m-added",
seq: 1,
});
});

test("keeps a snapshot that does not extend the closed stream as its own message", () => {
const folded = foldStreamedSessionEventRows(
[
row({ content: "第一条进度", eventType: "message.delta", id: "m-1", seq: 1 }),
row({ content: "Message updated.", eventType: "message.completed", id: "m-end", seq: 2 }),
row({ content: "另一条最终回复", eventType: "message.added", id: "m-added", seq: 3 }),
],
{ flushOpenStreams: true },
);

expect(folded.rows.map((entry) => entry.content_text)).toEqual([
"第一条进度",
"另一条最终回复",
]);
});

test("consumes fragments per snapshot across a multi-message turn", () => {
const folded = foldStreamedSessionEventRows(
[
row({ content: "进度说明", eventType: "message.delta", id: "m-1", seq: 1 }),
row({ content: "Message updated.", eventType: "message.completed", id: "c-1", seq: 2 }),
row({ content: "进度说明", eventType: "message.added", id: "a-1", seq: 3 }),
row({ content: "最终回复", eventType: "message.delta", id: "m-2", seq: 4 }),
row({ content: "Message updated.", eventType: "message.completed", id: "c-2", seq: 5 }),
row({ content: "最终回复", eventType: "message.added", id: "a-2", seq: 6 }),
],
{ flushOpenStreams: true },
);

expect(folded.rows.map((entry) => entry.content_text)).toEqual(["进度说明", "最终回复"]);
expect(folded.rows.map((entry) => entry.id)).toEqual(["a-1", "a-2"]);
});

test("drops streams that never carried text", () => {
const folded = foldStreamedSessionEventRows(
[
Expand Down
Loading