diff --git a/packages/ui/src/sync/event-pipeline.test.ts b/packages/ui/src/sync/event-pipeline.test.ts index 6ebcc41d..f7d690e9 100644 --- a/packages/ui/src/sync/event-pipeline.test.ts +++ b/packages/ui/src/sync/event-pipeline.test.ts @@ -97,6 +97,51 @@ describe("createEventPipeline", () => { })).toEqual(["updated:a", "delta:b", "updated:ab"]) }) + test("does not merge deltas across an intervening part snapshot", async () => { + let resolveStreamFinished!: () => void + const streamFinished = new Promise((resolve) => { + resolveStreamFinished = resolve + }) + let resolveDelivered!: () => void + const deliveredAll = new Promise((resolve) => { + resolveDelivered = resolve + }) + const delivered: Event[] = [] + const pipeline = createEventPipeline({ + sdk: createSdk([ + partUpdatedEvent("a"), + deltaEvent("b"), + partUpdatedEvent("ab"), + deltaEvent("c"), + ], resolveStreamFinished), + onEvent: (_directory, payload) => { + delivered.push(payload) + if (delivered.length === 4) { + resolveDelivered() + } + }, + transport: "sse", + heartbeatTimeoutMs: 1_000, + }) + + try { + await streamFinished + await Promise.race([deliveredAll, new Promise((resolve) => setTimeout(resolve, 300))]) + } finally { + pipeline.cleanup() + } + + // The "ab" snapshot is a coalescing barrier: the trailing "c" delta must + // stay a separate event after it, not merge into the "b" delta queued + // before the snapshot (which the snapshot would then overwrite). + expect(delivered.map((event) => { + if (event.type === "message.part.delta") { + return `delta:${(event.properties as { delta: string }).delta}` + } + return `updated:${((event.properties as { part: { text: string } }).part).text}` + })).toEqual(["updated:a", "delta:b", "updated:ab", "delta:c"]) + }) + test("normalizes openchamber session status events", async () => { let resolveStreamFinished!: () => void const streamFinished = new Promise((resolve) => { diff --git a/packages/ui/src/sync/event-pipeline.ts b/packages/ui/src/sync/event-pipeline.ts index aaa6efbd..701903a4 100644 --- a/packages/ui/src/sync/event-pipeline.ts +++ b/packages/ui/src/sync/event-pipeline.ts @@ -447,6 +447,26 @@ export function createEventPipeline(input: EventPipelineInput): EventPipeline { const normalizedPayload = normalizeEventType(payload) const routedDirectory = routeDirectory?.(directory, normalizedPayload) || directory const d = getOrCreateDir(routedDirectory) + + // A full part snapshot is a coalescing barrier for that part's deltas: + // drop its pending delta coalescing keys so a delta arriving after the + // snapshot starts a fresh queue entry instead of merging into a delta + // queued before the snapshot, which the snapshot would then overwrite and + // drop the later delta's text. The already-queued delta event stays. + if (normalizedPayload.type === "message.part.updated") { + const part = (normalizedPayload.properties as { part?: { id?: unknown; messageID?: unknown } }).part + const messageID = typeof part?.messageID === "string" ? part.messageID : undefined + const partID = typeof part?.id === "string" ? part.id : undefined + if (messageID && partID) { + const deltaPrefix = `message.part.delta:${messageID}:${partID}:` + for (const coalesceKey of d.coalesced.keys()) { + if (coalesceKey.startsWith(deltaPrefix)) { + d.coalesced.delete(coalesceKey) + } + } + } + } + const k = key(normalizedPayload) if (k) { const i = d.coalesced.get(k)