diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 4c6e3703ba82..e0e3bb8d52ae 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -1797,6 +1797,347 @@ describe("ProviderRuntimeIngestion", () => { expect(message?.streaming).toBe(false); }); + it("appends the missing suffix when assistant completion extends a buffered prefix", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const full = "Hi! I'm Codex, ready to help."; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-buffered-prefix"); + const itemId = asItemId("item-buffered-prefix"); + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-buffered-prefix-started"), + provider: codex, + createdAt: now, + threadId, + turnId, + }, + { + type: "content.delta", + eventId: asEventId("evt-buffered-prefix-delta"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { streamKind: "assistant_text", delta: "Hi! I" }, + }, + ]); + expect( + (await harness.readModel()).threads + .find((thread) => thread.id === threadId) + ?.messages.find((message) => message.id === `assistant:${itemId}`), + ).toBeUndefined(); + + await harness.emitAndDrain([ + { + type: "item.completed", + eventId: asEventId("evt-buffered-prefix-completed"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed", detail: full }, + }, + ]); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && !message.streaming, + ), + ); + const matches = thread.messages.filter( + (message: ProviderRuntimeTestMessage) => message.id === `assistant:${itemId}`, + ); + expect(matches).toHaveLength(1); + expect(matches[0]?.text).toBe(full); + }); + + it("appends the missing suffix when assistant completion extends an already projected prefix", async () => { + const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } }); + const now = "2026-01-01T00:00:00.000Z"; + const full = "Hi! I'm Codex, ready to help."; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-projected-prefix"); + const itemId = asItemId("item-projected-prefix"); + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-projected-prefix-started"), + provider: codex, + createdAt: now, + threadId, + turnId, + }, + { + type: "content.delta", + eventId: asEventId("evt-projected-prefix-delta"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { streamKind: "assistant_text", delta: "Hi! I" }, + }, + ]); + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && message.streaming && message.text === "Hi! I", + ), + ); + + await harness.emitAndDrain([ + { + type: "item.completed", + eventId: asEventId("evt-projected-prefix-completed"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed", detail: full }, + }, + ]); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && !message.streaming && message.text === full, + ), + ); + const matches = thread.messages.filter( + (message: ProviderRuntimeTestMessage) => message.id === `assistant:${itemId}`, + ); + expect(matches).toHaveLength(1); + expect(matches[0]?.text).toBe(full); + }); + + it("appends the missing suffix when completion extends a flushed prefix and buffered remainder", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-split-prefix"); + const itemId = asItemId("item-split-prefix"); + const full = "First paragraph.\n\nHi! I'm Codex, ready to help."; + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-split-prefix-started"), + provider: codex, + createdAt: now, + threadId, + turnId, + }, + { + type: "content.delta", + eventId: asEventId("evt-split-prefix-delta"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { streamKind: "assistant_text", delta: "First paragraph.\n\nHi! I" }, + }, + ]); + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && + message.streaming && + message.text === "First paragraph.\n\n", + ), + ); + + await harness.emitAndDrain([ + { + type: "item.completed", + eventId: asEventId("evt-split-prefix-completed"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed", detail: full }, + }, + ]); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && !message.streaming && message.text === full, + ), + ); + const matches = thread.messages.filter( + (message: ProviderRuntimeTestMessage) => message.id === `assistant:${itemId}`, + ); + expect(matches).toHaveLength(1); + expect(matches[0]?.text).toBe(full); + }); + + it("persists a whitespace-only suffix when assistant completion extends projected text", async () => { + const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } }); + const now = "2026-01-01T00:00:00.000Z"; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-whitespace-suffix"); + const itemId = asItemId("item-whitespace-suffix"); + const full = "Hello\n"; + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-whitespace-suffix-started"), + provider: codex, + createdAt: now, + threadId, + turnId, + }, + { + type: "content.delta", + eventId: asEventId("evt-whitespace-suffix-delta"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { streamKind: "assistant_text", delta: "Hello" }, + }, + ]); + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && message.streaming && message.text === "Hello", + ), + ); + + await harness.emitAndDrain([ + { + type: "item.completed", + eventId: asEventId("evt-whitespace-suffix-completed"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed", detail: full }, + }, + ]); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && !message.streaming && message.text === full, + ), + ); + const matches = thread.messages.filter( + (message: ProviderRuntimeTestMessage) => message.id === `assistant:${itemId}`, + ); + expect(matches).toHaveLength(1); + expect(matches[0]?.text).toBe(full); + }); + + it("does not create an assistant message from a whitespace-only completion snapshot", async () => { + const harness = await createHarness(); + const now = "2026-01-01T00:00:00.000Z"; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-whitespace-only-started"), + provider: codex, + createdAt: now, + threadId, + turnId: asTurnId("turn-whitespace-only"), + }, + { + type: "item.completed", + eventId: asEventId("evt-whitespace-only-completed"), + provider: codex, + createdAt: now, + threadId, + turnId: asTurnId("turn-whitespace-only"), + itemId: asItemId("item-whitespace-only"), + payload: { itemType: "assistant_message", status: "completed", detail: "\n" }, + }, + ]); + + const thread = (await harness.readModel()).threads.find((entry) => entry.id === threadId); + expect(thread?.messages.filter((message) => message.role === "assistant")).toEqual([]); + }); + + it("keeps streamed assistant text when the completion snapshot diverges", async () => { + const harness = await createHarness({ serverSettings: { responseStreamingMode: "token" } }); + const now = "2026-01-01T00:00:00.000Z"; + const codex = ProviderDriverKind.make("codex"); + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-divergent-snapshot"); + const itemId = asItemId("item-divergent-snapshot"); + + await harness.emitAndDrain([ + { + type: "turn.started", + eventId: asEventId("evt-divergent-snapshot-started"), + provider: codex, + createdAt: now, + threadId, + turnId, + }, + { + type: "content.delta", + eventId: asEventId("evt-divergent-snapshot-delta"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { streamKind: "assistant_text", delta: "Hello" }, + }, + ]); + await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && message.streaming && message.text === "Hello", + ), + ); + + await harness.emitAndDrain([ + { + type: "item.completed", + eventId: asEventId("evt-divergent-snapshot-completed"), + provider: codex, + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed", detail: "Goodbye" }, + }, + ]); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message: ProviderRuntimeTestMessage) => + message.id === `assistant:${itemId}` && !message.streaming, + ), + ); + const matches = thread.messages.filter( + (message: ProviderRuntimeTestMessage) => message.id === `assistant:${itemId}`, + ); + expect(matches).toHaveLength(1); + expect(matches[0]?.text).toBe("Hello"); + }); + it("preserves completed tool metadata on projected tool activities", async () => { const harness = await createHarness(); const now = "2026-01-01T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 2584a7578200..f958aabbe008 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -200,6 +200,37 @@ function hasRenderableAssistantText(text: string | undefined): boolean { return (text?.trim().length ?? 0) > 0; } +/** + * `detail` on assistant item completion is a full snapshot, not a delta. + * When it strictly extends text already accumulated (projected plus still + * buffered), return only the part not yet projected. The buffer has been + * taken but not written, so that remainder includes the buffered text and + * any missing suffix, including one that is only whitespace. Equal, empty, + * and divergent snapshots keep the streamed text. + */ +function assistantCompletionText(input: { + readonly projectedText: string; + readonly bufferedText: string; + readonly detail: string | undefined; +}): string { + const accumulated = `${input.projectedText}${input.bufferedText}`; + const detail = input.detail; + if ( + detail !== undefined && + detail.startsWith(accumulated) && + detail.length > accumulated.length + ) { + return detail.slice(input.projectedText.length); + } + if (input.bufferedText.length > 0) { + return input.bufferedText; + } + if (input.projectedText.length === 0 && detail !== undefined && detail.trim().length > 0) { + return detail; + } + return ""; +} + // An opening fence may sit at any indentation, since fences inside list // items are indented past the marker. A closing fence may be indented at most // three spaces more than its opener. Deeper lines are content in the block. @@ -1503,21 +1534,25 @@ const make = Effect.gen(function* () { commandTag: string; finalDeltaCommandTag: string; fallbackText?: string; + projectedText?: string; hasProjectedMessage?: boolean; }) => Effect.gen(function* () { const bufferedText = yield* takeBufferedAssistantText(input.messageId); - const text = - bufferedText.length > 0 - ? bufferedText - : (input.fallbackText?.trim().length ?? 0) > 0 - ? input.fallbackText! - : ""; + const projectedText = input.projectedText ?? ""; + const text = assistantCompletionText({ + projectedText, + bufferedText, + detail: input.fallbackText, + }); const hasRenderableText = hasRenderableAssistantText(text); + // A whitespace-only suffix is still part of a reply that is already + // visible. It must not be the only content of a new message. + const shouldPersistText = hasRenderableText || (projectedText.length > 0 && text.length > 0); const isReasoning = messageStreamRoleOf(input.messageId) === "reasoning"; - if (hasRenderableText) { + if (shouldPersistText) { yield* orchestrationEngine.dispatch({ type: isReasoning ? "thread.message.reasoning.delta" : "thread.message.assistant.delta", commandId: yield* providerCommandId(input.event, input.finalDeltaCommandTag), @@ -1531,7 +1566,7 @@ const make = Effect.gen(function* () { }); } - if (input.hasProjectedMessage || hasRenderableText) { + if (input.hasProjectedMessage || shouldPersistText) { yield* orchestrationEngine.dispatch({ type: isReasoning ? "thread.message.reasoning.complete" @@ -2298,9 +2333,6 @@ const make = Effect.gen(function* () { streamingOnly: false, }), ]); - const shouldApplyFallbackCompletionText = - !existingAssistantMessage || existingAssistantMessage.text.length === 0; - const shouldSkipRedundantCompletion = Option.isNone(activeAssistantMessageId) && turnId !== undefined && @@ -2321,7 +2353,8 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete", finalDeltaCommandTag: "assistant-delta-finalize", hasProjectedMessage: existingAssistantMessage !== undefined, - ...(assistantCompletion.fallbackText !== undefined && shouldApplyFallbackCompletionText + projectedText: existingAssistantMessage?.text ?? "", + ...(assistantCompletion.fallbackText !== undefined ? { fallbackText: assistantCompletion.fallbackText } : {}), });