Skip to content
Closed
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
123 changes: 123 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1797,6 +1797,129 @@ 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("preserves completed tool metadata on projected tool activities", async () => {
const harness = await createHarness();
const now = "2026-01-01T00:00:00.000Z";
Expand Down
49 changes: 39 additions & 10 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,35 @@ function hasRenderableAssistantText(text: string | undefined): boolean {
return (text?.trim().length ?? 0) > 0;
}

/**
* `detail` on item completion is a full snapshot. When it strictly extends
* text already accumulated (projected plus still buffered), only the part
* not yet projected is safe to persist. Equal, empty, and divergent
* snapshots keep the streamed text.
*/
function assistantCompletionDelta(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?.trim().length ?? 0) > 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.
Expand Down Expand Up @@ -1502,16 +1531,16 @@ 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 text = assistantCompletionDelta({
projectedText: input.projectedText ?? "",
bufferedText,
detail: input.fallbackText,
});
const hasRenderableText = hasRenderableAssistantText(text);

const isReasoning = messageStreamRoleOf(input.messageId) === "reasoning";
Expand Down Expand Up @@ -2297,9 +2326,6 @@ const make = Effect.gen(function* () {
streamingOnly: false,
}),
]);
const shouldApplyFallbackCompletionText =
!existingAssistantMessage || existingAssistantMessage.text.length === 0;

const shouldSkipRedundantCompletion =
Option.isNone(activeAssistantMessageId) &&
turnId !== undefined &&
Expand All @@ -2320,7 +2346,10 @@ const make = Effect.gen(function* () {
commandTag: "assistant-complete",
finalDeltaCommandTag: "assistant-delta-finalize",
hasProjectedMessage: existingAssistantMessage !== undefined,
...(assistantCompletion.fallbackText !== undefined && shouldApplyFallbackCompletionText
...(existingAssistantMessage !== undefined
? { projectedText: existingAssistantMessage.text }
: {}),
...(assistantCompletion.fallbackText !== undefined
? { fallbackText: assistantCompletion.fallbackText }
: {}),
});
Expand Down
Loading