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
341 changes: 341 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,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";
Expand Down
Loading
Loading