Skip to content
Open
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
127 changes: 111 additions & 16 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,9 @@ import {
EventId,
MessageId,
ProjectId,
ProviderDriverKind,
ProviderItemId,
RuntimeItemId,
ThreadId,
TurnId,
} from "@t3tools/contracts";
Expand Down Expand Up @@ -674,16 +676,104 @@ describe("ProviderRuntimeIngestion", () => {
const thread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-1" && !message.streaming,
message.id === "assistant:turn-2:item-1" && !message.streaming,
),
);
const message = thread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-1",
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:turn-2:item-1",
);
expect(message?.text).toBe("hello world");
expect(message?.streaming).toBe(false);
});

it("keeps assistant replies distinct when a resumed session reuses an item id", async () => {
const harness = await createHarness();
const firstTurnAt = "2026-07-01T21:06:10.000Z";
const secondTurnAt = "2026-07-02T01:44:30.000Z";
const reusedItemId = RuntimeItemId.make("assistant:session-1:segment:0");

harness.emit({
type: "content.delta",
eventId: asEventId("evt-resume-turn-1"),
provider: ProviderDriverKind.make("cursor"),
createdAt: firstTurnAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-1"),
itemId: reusedItemId,
payload: {
streamKind: "assistant_text",
delta: "first turn reply",
},
});
harness.emit({
type: "item.completed",
eventId: asEventId("evt-resume-turn-1-complete"),
provider: ProviderDriverKind.make("cursor"),
createdAt: firstTurnAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-1"),
itemId: reusedItemId,
payload: {
itemType: "assistant_message",
status: "completed",
},
});

await waitForThread(harness.engine, (thread) =>
thread.messages.some(
(message) =>
message.id === "assistant:turn-1:assistant:session-1:segment:0" &&
message.text === "first turn reply",
),
);

harness.emit({
type: "content.delta",
eventId: asEventId("evt-resume-turn-2"),
provider: ProviderDriverKind.make("cursor"),
createdAt: secondTurnAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-2"),
itemId: reusedItemId,
payload: {
streamKind: "assistant_text",
delta: "second turn reply",
},
});
harness.emit({
type: "item.completed",
eventId: asEventId("evt-resume-turn-2-complete"),
provider: ProviderDriverKind.make("cursor"),
createdAt: secondTurnAt,
threadId: asThreadId("thread-1"),
turnId: asTurnId("turn-2"),
itemId: reusedItemId,
payload: {
itemType: "assistant_message",
status: "completed",
},
});

const thread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message) =>
message.id === "assistant:turn-2:assistant:session-1:segment:0" &&
message.text === "second turn reply",
),
);

expect(
thread.messages.find(
(message) => message.id === "assistant:turn-1:assistant:session-1:segment:0",
)?.text,
).toBe("first turn reply");
expect(
thread.messages.find(
(message) => message.id === "assistant:turn-2:assistant:session-1:segment:0",
)?.text,
).toBe("second turn reply");
});

it("uses assistant item completion detail when no assistant deltas were streamed", async () => {
const harness = await createHarness();
const now = new Date().toISOString();
Expand All @@ -706,11 +796,11 @@ describe("ProviderRuntimeIngestion", () => {
const thread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-no-delta" && !message.streaming,
message.id === "assistant:turn-no-delta:item-no-delta" && !message.streaming,
),
);
const message = thread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-no-delta",
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:turn-no-delta:item-no-delta",
);
expect(message?.text).toBe("assistant-only final text");
expect(message?.streaming).toBe(false);
Expand Down Expand Up @@ -1339,7 +1429,8 @@ describe("ProviderRuntimeIngestion", () => {
const midThread = midReadModel.threads.find((entry) => entry.id === ThreadId.make("thread-1"));
expect(
midThread?.messages.some(
(message: ProviderRuntimeTestMessage) => message.id === "assistant:item-buffered",
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:turn-buffered:item-buffered",
),
).toBe(false);

Expand All @@ -1360,11 +1451,11 @@ describe("ProviderRuntimeIngestion", () => {
const thread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-buffered" && !message.streaming,
message.id === "assistant:turn-buffered:item-buffered" && !message.streaming,
),
);
const message = thread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffered",
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:turn-buffered:item-buffered",
);
expect(message?.text).toBe("buffer me");
expect(message?.streaming).toBe(false);
Expand Down Expand Up @@ -1425,13 +1516,14 @@ describe("ProviderRuntimeIngestion", () => {
const liveThread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-streaming-mode" &&
message.id === "assistant:turn-streaming-mode:item-streaming-mode" &&
message.streaming &&
message.text === "hello live",
),
);
const liveMessage = liveThread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-streaming-mode",
(entry: ProviderRuntimeTestMessage) =>
entry.id === "assistant:turn-streaming-mode:item-streaming-mode",
);
expect(liveMessage?.streaming).toBe(true);

Expand All @@ -1453,11 +1545,12 @@ describe("ProviderRuntimeIngestion", () => {
const finalThread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-streaming-mode" && !message.streaming,
message.id === "assistant:turn-streaming-mode:item-streaming-mode" && !message.streaming,
),
);
const finalMessage = finalThread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-streaming-mode",
(entry: ProviderRuntimeTestMessage) =>
entry.id === "assistant:turn-streaming-mode:item-streaming-mode",
);
expect(finalMessage?.text).toBe("hello live");
expect(finalMessage?.streaming).toBe(false);
Expand Down Expand Up @@ -1513,11 +1606,12 @@ describe("ProviderRuntimeIngestion", () => {
const thread = await waitForThread(harness.engine, (entry) =>
entry.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-buffer-spill" && !message.streaming,
message.id === "assistant:turn-buffer-spill:item-buffer-spill" && !message.streaming,
),
);
const message = thread.messages.find(
(entry: ProviderRuntimeTestMessage) => entry.id === "assistant:item-buffer-spill",
(entry: ProviderRuntimeTestMessage) =>
entry.id === "assistant:turn-buffer-spill:item-buffer-spill",
);
expect(message?.text.length).toBe(oversizedText.length);
expect(message?.text).toBe(oversizedText);
Expand Down Expand Up @@ -1589,7 +1683,8 @@ describe("ProviderRuntimeIngestion", () => {
thread.session?.activeTurnId === null &&
thread.messages.some(
(message: ProviderRuntimeTestMessage) =>
message.id === "assistant:item-complete-dedup" && !message.streaming,
message.id === "assistant:turn-complete-dedup:item-complete-dedup" &&
!message.streaming,
),
);

Expand All @@ -1603,7 +1698,7 @@ describe("ProviderRuntimeIngestion", () => {
return false;
}
return (
event.payload.messageId === "assistant:item-complete-dedup" &&
event.payload.messageId === "assistant:turn-complete-dedup:item-complete-dedup" &&
event.payload.streaming === false
);
});
Expand Down Expand Up @@ -1973,7 +2068,7 @@ describe("ProviderRuntimeIngestion", () => {
(entry: ProviderRuntimeTestCheckpoint) => entry.turnId === "turn-p1",
);
expect(checkpoint?.status).toBe("missing");
expect(checkpoint?.assistantMessageId).toBe("assistant:item-p1-assistant");
expect(checkpoint?.assistantMessageId).toBe("assistant:turn-p1:item-p1-assistant");
expect(checkpoint?.checkpointRef).toBe("provider-diff:evt-turn-diff-updated");
});

Expand Down
20 changes: 11 additions & 9 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import {
ProviderRuntimeIngestionService,
type ProviderRuntimeIngestionShape,
} from "../Services/ProviderRuntimeIngestion.ts";
import { assistantMessageBaseKey, assistantMessageId } from "../assistantMessageIds.ts";

const providerTurnKey = (threadId: ThreadId, turnId: TurnId) => `${threadId}:${turnId}`;
const providerCommandId = (event: ProviderRuntimeEvent, tag: string): CommandId =>
Expand Down Expand Up @@ -60,6 +61,13 @@ function toTurnId(value: TurnId | string | undefined): TurnId | undefined {
return value === undefined ? undefined : TurnId.make(String(value));
}

function assistantMessageIdForEvent(event: ProviderRuntimeEvent): MessageId {
return assistantMessageId(
assistantMessageBaseKey(event.itemId, event.turnId, event.eventId),
toTurnId(event.turnId),
);
}

function toApprovalRequestId(value: string | undefined): ApprovalRequestId | undefined {
return value === undefined ? undefined : ApprovalRequestId.make(value);
}
Expand Down Expand Up @@ -1026,9 +1034,7 @@ const make = Effect.gen(function* () {
event.type === "turn.proposed.delta" ? event.payload.delta : undefined;

if (assistantDelta && assistantDelta.length > 0) {
const assistantMessageId = MessageId.make(
`assistant:${event.itemId ?? event.turnId ?? event.eventId}`,
);
const assistantMessageId = assistantMessageIdForEvent(event);
const turnId = toTurnId(event.turnId);
if (turnId) {
yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId);
Expand Down Expand Up @@ -1069,9 +1075,7 @@ const make = Effect.gen(function* () {
const assistantCompletion =
event.type === "item.completed" && event.payload.itemType === "assistant_message"
? {
messageId: MessageId.make(
`assistant:${event.itemId ?? event.turnId ?? event.eventId}`,
),
messageId: assistantMessageIdForEvent(event),
fallbackText: event.payload.detail,
}
: undefined;
Expand Down Expand Up @@ -1209,9 +1213,7 @@ const make = Effect.gen(function* () {
if (thread.checkpoints.some((c) => c.turnId === turnId)) {
// Already tracked; no-op.
} else {
const assistantMessageId = MessageId.make(
`assistant:${event.itemId ?? event.turnId ?? event.eventId}`,
);
const assistantMessageId = assistantMessageIdForEvent(event);
const maxTurnCount = thread.checkpoints.reduce(
(max, c) => Math.max(max, c.checkpointTurnCount),
0,
Expand Down
26 changes: 26 additions & 0 deletions apps/server/src/orchestration/assistantMessageIds.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,26 @@
import { describe, expect, it } from "vitest";
import { TurnId } from "@t3tools/contracts";

import { assistantMessageBaseKey, assistantMessageId } from "./assistantMessageIds.ts";

describe("assistant message ids", () => {
it("scopes reused provider item ids to their turn", () => {
const baseKey = assistantMessageBaseKey("assistant:session-1:segment:0", "turn-1", "event-1");

const firstTurnId = assistantMessageId(baseKey, TurnId.make("turn-1"));
const secondTurnId = assistantMessageId(baseKey, TurnId.make("turn-2"));

expect(firstTurnId).toBe("assistant:turn-1:assistant:session-1:segment:0");
expect(secondTurnId).toBe("assistant:turn-2:assistant:session-1:segment:0");
expect(firstTurnId).not.toBe(secondTurnId);
});

it("falls back from item id to turn id and event id", () => {
expect(assistantMessageBaseKey(undefined, "turn-1", "event-1")).toBe("turn-1");
expect(assistantMessageBaseKey(undefined, undefined, "event-1")).toBe("event-1");
});

it("preserves the legacy shape when no turn id is available", () => {
expect(assistantMessageId("item-1")).toBe("assistant:item-1");
});
});
13 changes: 13 additions & 0 deletions apps/server/src/orchestration/assistantMessageIds.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
import { MessageId, type TurnId } from "@t3tools/contracts";

export function assistantMessageBaseKey(
itemId: string | undefined,
turnId: string | undefined,
eventId: string,
): string {
return String(itemId ?? turnId ?? eventId);
}

export function assistantMessageId(baseKey: string, turnId?: TurnId): MessageId {
return MessageId.make(`assistant:${turnId ? `${turnId}:` : ""}${baseKey}`);
}