diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 585354f7baf2..bddc39b24b86 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -1727,6 +1727,15 @@ function renderFeedEntry( })} {showAssistantMeta ? ( + {message.actualModel ? ( + + Model: {message.actualModel} + + ) : null} { `; yield* sql` INSERT INTO projection_thread_messages ( - message_id, thread_id, role, text, attachments_json, context_json, is_streaming, created_at, updated_at + message_id, thread_id, role, text, actual_model, attachments_json, context_json, + is_streaming, created_at, updated_at ) VALUES (${messageId}, ${threadId}, 'user', 'Read these notes', - ${attachmentsJson}, ${contextJson}, 0, ${createdAt}, ${createdAt}) + 'openai/gpt-5.6-sol', ${attachmentsJson}, ${contextJson}, 0, ${createdAt}, ${createdAt}) `; yield* sql` INSERT INTO projection_thread_messages ( @@ -840,6 +841,7 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { text: "Read these notes", turnId: null, streaming: false, + actualModel: "openai/gpt-5.6-sol", createdAt, updatedAt: createdAt, attachments, diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index 5bfe411fd845..fc000c6d6a70 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -34,6 +34,7 @@ import { ThreadPullRequestSnapshot, ThreadPullRequestStack, type ThreadPullRequestLink, + TrimmedNonEmptyString, } from "@t3tools/contracts"; import { legacyLinkedPullRequestOf } from "@t3tools/shared/threadPullRequests"; import * as Arr from "effect/Array"; @@ -112,6 +113,7 @@ const ProjectionProjectDbRowSchema = ProjectionProject.mapFields( const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( Struct.assign({ isStreaming: Schema.Number, + actualModel: Schema.NullOr(TrimmedNonEmptyString), attachments: Schema.NullOr(Schema.fromJsonString(Schema.Array(ChatAttachment))), context: Schema.NullOr(Schema.fromJsonString(OrchestrationMessageContext)), }), @@ -691,6 +693,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -1294,6 +1297,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -1327,6 +1331,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -1740,6 +1745,7 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -2150,6 +2156,7 @@ pending_approval_requests AS ( id: row.messageId, role: row.role, text: row.text, + ...(row.actualModel !== null ? { actualModel: row.actualModel } : {}), ...(row.attachments !== null ? { attachments: row.attachments } : {}), ...(row.context !== null ? { context: row.context } : {}), turnId: row.turnId, @@ -3277,6 +3284,7 @@ pending_approval_requests AS ( text: row.text, turnId: row.turnId, streaming: row.isStreaming === 1, + ...(row.actualModel !== null ? { actualModel: row.actualModel } : {}), createdAt: row.createdAt, updatedAt: row.updatedAt, ...(row.attachments !== null ? { attachments: row.attachments } : {}), @@ -3535,6 +3543,7 @@ pending_approval_requests AS ( id: row.messageId, role: row.role, text: row.text, + ...(row.actualModel !== null ? { actualModel: row.actualModel } : {}), turnId: row.turnId, streaming: row.isStreaming === 1, createdAt: row.createdAt, diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index d61739f72c21..21be313da71b 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3838,6 +3838,427 @@ describe("ProviderRuntimeIngestion", () => { expect(completionEvents).toHaveLength(1); }); + it("enriches an already completed assistant message with the actual turn model", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-actual-model"); + const itemId = asItemId("item-actual-model"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-message-delta-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId: asItemId("native-message-actual-model") }, + payload: { streamKind: "assistant_text", delta: "done" }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-message-completed-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + payload: { itemType: "assistant_message", status: "completed" }, + }); + await waitForThread(harness.readModel, (thread) => + thread.messages.some( + (message) => message.id === "assistant:item-actual-model" && !message.streaming, + ), + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + providerRefs: { providerItemId: asItemId("native-message-actual-model") }, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.id === "assistant:item-actual-model" && message.actualModel === "gpt-5.6-luna", + ), + ); + const messages = thread.messages.filter( + (message) => message.id === "assistant:item-actual-model", + ); + expect(messages).toHaveLength(1); + expect(messages[0]?.actualModel).toBe("gpt-5.6-luna"); + }); + + it("targets a non-streamed assistant item when model metadata arrives later", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-non-streamed-model"); + const itemId = asItemId("item-non-streamed-model"); + const providerMessageId = asItemId("native-non-streamed-model"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-non-streamed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-item-completed-non-streamed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId: providerMessageId }, + payload: { + itemType: "assistant_message", + status: "completed", + detail: "non-streamed response", + }, + }); + await waitForThread(harness.readModel, (thread) => + thread.messages.some( + (message) => message.id === "assistant:item-non-streamed-model" && !message.streaming, + ), + ); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-non-streamed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + providerRefs: { providerItemId: providerMessageId }, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.id === "assistant:item-non-streamed-model" && + message.actualModel === "gpt-5.6-luna", + ), + ); + const messages = thread.messages.filter( + (message) => message.id === "assistant:item-non-streamed-model", + ); + expect(messages).toHaveLength(1); + expect(messages[0]?.text).toBe("non-streamed response"); + expect(messages[0]?.actualModel).toBe("gpt-5.6-luna"); + }); + + it("retains the actual model when turn completion precedes the assistant item", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-completed-before-item"); + const itemId = asItemId("item-completed-before-item"); + const providerMessageId = asItemId("native-completed-before-item"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-before-item"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-before-item"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + providerRefs: { providerItemId: providerMessageId }, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === null); + + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-message-completed-after-turn"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId: providerMessageId }, + payload: { + itemType: "assistant_message", + status: "completed", + detail: "completed after the turn event", + }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.id === "assistant:item-completed-before-item" && + message.actualModel === "gpt-5.6-luna", + ), + ); + const messages = thread.messages.filter( + (message) => message.id === "assistant:item-completed-before-item", + ); + expect(messages).toHaveLength(1); + expect(messages[0]?.text).toBe("completed after the turn event"); + expect(messages[0]?.actualModel).toBe("gpt-5.6-luna"); + }); + + it("does not defer unkeyed terminal model metadata to a later assistant item", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-unkeyed-model"); + const itemId = asItemId("item-after-unkeyed-model"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-unkeyed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-unkeyed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === null); + + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-item-completed-after-unkeyed-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId: asItemId("native-after-unkeyed-model") }, + payload: { + itemType: "assistant_message", + status: "completed", + detail: "unrelated later response", + }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => message.id === "assistant:item-after-unkeyed-model" && !message.streaming, + ), + ); + const message = thread.messages.find( + (entry) => entry.id === "assistant:item-after-unkeyed-model", + ); + expect(message?.text).toBe("unrelated later response"); + expect(message?.actualModel).toBeUndefined(); + }); + + it("applies late actual model metadata only to its originating assistant message", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-targeted-actual-model"); + const firstItemId = asItemId("item-targeted-model-first"); + const secondItemId = asItemId("item-targeted-model-second"); + const firstProviderMessageId = asItemId("native-targeted-model-first"); + const secondProviderMessageId = asItemId("native-targeted-model-second"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-targeted-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + for (const [index, itemId, providerItemId, delta] of [ + ["first", firstItemId, firstProviderMessageId, "first response"], + ["second", secondItemId, secondProviderMessageId, "second response"], + ] as const) { + harness.emit({ + type: "content.delta", + eventId: asEventId(`evt-message-delta-targeted-${index}`), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId }, + payload: { streamKind: "assistant_text", delta }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId(`evt-message-completed-targeted-${index}`), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId, + providerRefs: { providerItemId }, + payload: { itemType: "assistant_message", status: "completed" }, + }); + await waitForThread(harness.readModel, (thread) => + thread.messages.some( + (message) => message.id === `assistant:${itemId}` && !message.streaming, + ), + ); + } + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-targeted-actual-model"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + providerRefs: { providerItemId: firstProviderMessageId }, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.id === "assistant:item-targeted-model-first" && + message.actualModel === "gpt-5.6-luna", + ), + ); + const firstMessage = thread.messages.find( + (message) => message.id === "assistant:item-targeted-model-first", + ); + const secondMessage = thread.messages.find( + (message) => message.id === "assistant:item-targeted-model-second", + ); + expect(firstMessage?.actualModel).toBe("gpt-5.6-luna"); + expect(secondMessage?.actualModel).toBeUndefined(); + }); + + it("updates a completed provider-targeted message while another message is streaming", async () => { + const harness = await createHarness(); + const now = "2026-08-14T00:00:00.000Z"; + const threadId = asThreadId("thread-1"); + const turnId = asTurnId("turn-targeted-with-streaming-message"); + const completedItemId = asItemId("item-targeted-completed"); + const streamingItemId = asItemId("item-unrelated-streaming"); + const completedProviderMessageId = asItemId("native-targeted-completed"); + + harness.emit({ + type: "turn.started", + eventId: asEventId("evt-turn-started-targeted-with-streaming"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + }); + await waitForThread(harness.readModel, (thread) => thread.session?.activeTurnId === turnId); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-delta-targeted-completed"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId: completedItemId, + providerRefs: { providerItemId: completedProviderMessageId }, + payload: { streamKind: "assistant_text", delta: "completed response" }, + }); + harness.emit({ + type: "item.completed", + eventId: asEventId("evt-item-targeted-completed"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId: completedItemId, + providerRefs: { providerItemId: completedProviderMessageId }, + payload: { itemType: "assistant_message", status: "completed" }, + }); + await waitForThread(harness.readModel, (thread) => + thread.messages.some( + (message) => message.id === "assistant:item-targeted-completed" && !message.streaming, + ), + ); + + harness.emit({ + type: "content.delta", + eventId: asEventId("evt-delta-unrelated-streaming"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + itemId: streamingItemId, + providerRefs: { providerItemId: asItemId("native-unrelated-streaming") }, + payload: { streamKind: "assistant_text", delta: "still streaming" }, + }); + await harness.drain(); + + harness.emit({ + type: "turn.completed", + eventId: asEventId("evt-turn-completed-targeted-with-streaming"), + provider: ProviderDriverKind.make("opencode"), + createdAt: now, + threadId, + turnId, + providerRefs: { providerItemId: completedProviderMessageId }, + payload: { state: "completed", actualModel: "gpt-5.6-luna" }, + }); + + const thread = await waitForThread(harness.readModel, (entry) => + entry.messages.some( + (message) => + message.id === "assistant:item-targeted-completed" && + message.actualModel === "gpt-5.6-luna", + ), + ); + const completedMessage = thread.messages.find( + (message) => message.id === "assistant:item-targeted-completed", + ); + const streamingMessage = thread.messages.find( + (message) => message.id === "assistant:item-unrelated-streaming", + ); + expect(completedMessage?.actualModel).toBe("gpt-5.6-luna"); + expect(streamingMessage?.actualModel).toBeUndefined(); + expect(streamingMessage?.text).toBe("still streaming"); + expect(streamingMessage?.streaming).toBe(false); + }); + it("maps canonical request events into approval activities with requestKind", 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 84cb34a4e783..076ac65e3c8d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -62,6 +62,8 @@ const segmentStateKey = (threadId: ThreadId, turnId: TurnId, role: MessageStream role === "reasoning" ? `${providerTurnKey(threadId, turnId)}:reasoning` : providerTurnKey(threadId, turnId); +const providerMessageKey = (threadId: ThreadId, turnId: TurnId, providerItemId: string) => + `${threadId}:${turnId}:${providerItemId}`; const providerTaskKey = (threadId: ThreadId, taskId: string) => `${threadId}:${taskId}`; // Fallback when the in-memory description cache no longer has the task name @@ -105,6 +107,11 @@ interface AssistantSegmentState { activeMessageId: MessageId | null; } +interface PendingActualModel { + readonly actualModel: string; + readonly providerItemId?: string | undefined; +} + const TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY = 10_000; const TURN_MESSAGE_IDS_BY_TURN_TTL = Duration.minutes(120); const BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY = 20_000; @@ -1035,6 +1042,18 @@ const make = Effect.gen(function* () { lookup: () => Effect.succeed(new Set()), }); + const assistantMessageIdsByProviderMessageKey = yield* Cache.make>({ + capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY, + timeToLive: TURN_MESSAGE_IDS_BY_TURN_TTL, + lookup: () => Effect.succeed(new Set()), + }); + + const pendingActualModelByTurnKey = yield* Cache.make({ + capacity: TURN_MESSAGE_IDS_BY_TURN_CACHE_CAPACITY, + timeToLive: TURN_MESSAGE_IDS_BY_TURN_TTL, + lookup: () => Effect.die(new Error("pending actual model should be read through getOption")), + }); + const bufferedAssistantTextByMessageId = yield* Cache.make({ capacity: BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY, timeToLive: BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_TTL, @@ -1163,6 +1182,43 @@ const make = Effect.gen(function* () { const clearAssistantMessageIdsForTurn = (threadId: ThreadId, turnId: TurnId) => Cache.invalidate(turnMessageIdsByTurnKey, providerTurnKey(threadId, turnId)); + const rememberProviderAssistantMessageId = ( + threadId: ThreadId, + turnId: TurnId, + providerItemId: string, + messageId: MessageId, + ) => + Cache.getOption( + assistantMessageIdsByProviderMessageKey, + providerMessageKey(threadId, turnId, providerItemId), + ).pipe( + Effect.flatMap((existingIds) => { + const nextIds = Option.match(existingIds, { + onNone: () => new Set([messageId]), + onSome: (ids) => new Set(ids).add(messageId), + }); + return Cache.set( + assistantMessageIdsByProviderMessageKey, + providerMessageKey(threadId, turnId, providerItemId), + nextIds, + ); + }), + ); + + const getProviderAssistantMessageIds = ( + threadId: ThreadId, + turnId: TurnId, + providerItemId: string, + ) => + Cache.getOption( + assistantMessageIdsByProviderMessageKey, + providerMessageKey(threadId, turnId, providerItemId), + ).pipe( + Effect.map((existingIds) => + Option.getOrElse(existingIds, (): Set => new Set()), + ), + ); + const getAssistantSegmentStateForTurn = ( threadId: ThreadId, turnId: TurnId, @@ -1475,6 +1531,7 @@ const make = Effect.gen(function* () { finalDeltaCommandTag: string; fallbackText?: string; hasProjectedMessage?: boolean; + actualModel?: string; }) => Effect.gen(function* () { const bufferedText = yield* takeBufferedAssistantText(input.messageId); @@ -1511,6 +1568,7 @@ const make = Effect.gen(function* () { threadId: input.threadId, messageId: input.messageId, ...(input.turnId ? { turnId: input.turnId } : {}), + ...(input.actualModel ? { actualModel: input.actualModel } : {}), createdAt: input.createdAt, }); } @@ -1620,6 +1678,10 @@ const make = Effect.gen(function* () { const prefix = `${threadId}:`; const proposedPlanPrefix = `plan:${threadId}:`; const turnKeys = Array.from(yield* Cache.keys(turnMessageIdsByTurnKey)); + const providerMessageKeys = Array.from( + yield* Cache.keys(assistantMessageIdsByProviderMessageKey), + ); + const pendingActualModelKeys = Array.from(yield* Cache.keys(pendingActualModelByTurnKey)); const assistantSegmentKeys = Array.from(yield* Cache.keys(assistantSegmentStateByTurnKey)); const proposedPlanKeys = Array.from(yield* Cache.keys(bufferedProposedPlanById)); const taskDescriptionKeys = Array.from(yield* Cache.keys(taskDescriptionByTaskKey)); @@ -1642,6 +1704,20 @@ const make = Effect.gen(function* () { }), { concurrency: 1 }, ).pipe(Effect.asVoid); + yield* Effect.forEach( + providerMessageKeys, + (key) => + key.startsWith(prefix) + ? Cache.invalidate(assistantMessageIdsByProviderMessageKey, key) + : Effect.void, + { concurrency: 1 }, + ).pipe(Effect.asVoid); + yield* Effect.forEach( + pendingActualModelKeys, + (key) => + key.startsWith(prefix) ? Cache.invalidate(pendingActualModelByTurnKey, key) : Effect.void, + { concurrency: 1 }, + ).pipe(Effect.asVoid); yield* Effect.forEach( assistantSegmentKeys, (key) => @@ -1768,6 +1844,8 @@ const make = Effect.gen(function* () { const eventTurnId = toTurnId(event.turnId); const activeTurnId = thread.session?.activeTurnId ?? null; const isTerminalTurn = event.type === "turn.completed" || event.type === "turn.aborted"; + const terminalActualModel = + event.type === "turn.completed" ? event.payload.actualModel : undefined; const isCompactedThreadState = event.type === "thread.state.changed" && event.payload.state === "compacted"; const pendingTurnStart = @@ -2032,6 +2110,15 @@ const make = Effect.gen(function* () { }); if (turnId) { yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId); + const providerItemId = event.providerRefs?.providerItemId; + if (providerItemId) { + yield* rememberProviderAssistantMessageId( + thread.id, + turnId, + providerItemId, + assistantMessageId, + ); + } } const streamingMode = yield* resolveResponseStreamingMode(thread.projectId); @@ -2252,6 +2339,21 @@ const make = Effect.gen(function* () { role: "reasoning", }); } + const pendingActualModel = turnId + ? Option.getOrUndefined( + yield* Cache.getOption( + pendingActualModelByTurnKey, + providerTurnKey(thread.id, turnId), + ), + ) + : undefined; + const completionProviderItemId = event.providerRefs?.providerItemId; + const matchingPendingActualModel = + pendingActualModel && + (pendingActualModel.providerItemId === undefined || + pendingActualModel.providerItemId === completionProviderItemId) + ? pendingActualModel.actualModel + : undefined; const activeAssistantMessageId = turnId ? yield* getActiveAssistantMessageIdForTurn(thread.id, turnId) : Option.none(); @@ -2276,12 +2378,21 @@ const make = Effect.gen(function* () { Option.isNone(activeAssistantMessageId) && turnId !== undefined && hasAssistantMessagesForTurn && + matchingPendingActualModel === undefined && (assistantCompletion.fallbackText?.trim().length ?? 0) === 0; if (!shouldSkipRedundantCompletion) { if (turnId && Option.isNone(activeAssistantMessageId)) { yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId); } + if (turnId && completionProviderItemId) { + yield* rememberProviderAssistantMessageId( + thread.id, + turnId, + completionProviderItemId, + assistantMessageId, + ); + } yield* finalizeAssistantMessage({ event, @@ -2292,6 +2403,7 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete", finalDeltaCommandTag: "assistant-delta-finalize", hasProjectedMessage: existingAssistantMessage !== undefined, + ...(matchingPendingActualModel ? { actualModel: matchingPendingActualModel } : {}), ...(assistantCompletion.fallbackText !== undefined && shouldApplyFallbackCompletionText ? { fallbackText: assistantCompletion.fallbackText } : {}), @@ -2299,6 +2411,12 @@ const make = Effect.gen(function* () { if (turnId) { yield* forgetAssistantMessageId(thread.id, turnId, assistantMessageId); + if (matchingPendingActualModel) { + yield* Cache.invalidate( + pendingActualModelByTurnKey, + providerTurnKey(thread.id, turnId), + ); + } } } @@ -2362,7 +2480,41 @@ const make = Effect.gen(function* () { createdAt: now, }); } - const assistantMessageIds = yield* getAssistantMessageIdsForTurn(thread.id, turnId); + const trackedAssistantMessageIds = yield* getAssistantMessageIdsForTurn( + thread.id, + turnId, + ); + const terminalProviderItemId = + event.type === "turn.completed" ? event.providerRefs?.providerItemId : undefined; + const modelTargetMessageIds = + terminalActualModel && terminalProviderItemId + ? yield* getProviderAssistantMessageIds(thread.id, turnId, terminalProviderItemId) + : new Set(); + const completedAssistantMessageId = + trackedAssistantMessageIds.size === 0 && + modelTargetMessageIds.size === 0 && + terminalActualModel && + !terminalProviderItemId + ? Option.getOrUndefined( + yield* projectionThreadMessages.getLatestAssistantMessageIdForTurn({ + threadId: thread.id, + turnId, + }), + ) + : undefined; + const assistantMessageIds = + trackedAssistantMessageIds.size > 0 || modelTargetMessageIds.size > 0 + ? Array.from(new Set([...trackedAssistantMessageIds, ...modelTargetMessageIds])) + : completedAssistantMessageId + ? [completedAssistantMessageId] + : []; + if (terminalActualModel && terminalProviderItemId && assistantMessageIds.length === 0) { + yield* Cache.set(pendingActualModelByTurnKey, providerTurnKey(thread.id, turnId), { + actualModel: terminalActualModel, + providerItemId: terminalProviderItemId, + }); + } + const terminalAssistantMessageId = assistantMessageIds.at(-1); yield* Effect.forEach( assistantMessageIds, (assistantMessageId) => @@ -2377,11 +2529,24 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete-finalize", finalDeltaCommandTag: "assistant-delta-finalize-fallback", hasProjectedMessage: existingMessage !== undefined, + ...(terminalActualModel && + (modelTargetMessageIds.size > 0 + ? modelTargetMessageIds.has(assistantMessageId) + : !terminalProviderItemId && + assistantMessageId === terminalAssistantMessageId) + ? { actualModel: terminalActualModel } + : {}), }), ), ), { concurrency: 1 }, ).pipe(Effect.asVoid); + if (terminalActualModel && assistantMessageIds.length > 0) { + yield* Cache.invalidate( + pendingActualModelByTurnKey, + providerTurnKey(thread.id, turnId), + ); + } yield* clearAssistantMessageIdsForTurn(thread.id, turnId); yield* clearAssistantSegmentStateForTurn(thread.id, turnId); yield* clearAssistantSegmentStateForTurn(thread.id, turnId, "reasoning"); diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0119c0e8599a..5c10d1241142 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1972,6 +1972,9 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" messageId: command.messageId, role: command.type === "thread.message.reasoning.complete" ? "reasoning" : "assistant", text: "", + ...(command.type === "thread.message.assistant.complete" && command.actualModel + ? { actualModel: command.actualModel } + : {}), turnId: command.turnId ?? null, streaming: false, createdAt: command.createdAt, diff --git a/apps/server/src/orchestration/projector.ts b/apps/server/src/orchestration/projector.ts index 53013770b15b..259a81e4d607 100644 --- a/apps/server/src/orchestration/projector.ts +++ b/apps/server/src/orchestration/projector.ts @@ -765,6 +765,7 @@ export function projectEvent( id: payload.messageId, role: payload.role, text: payload.text, + ...(payload.actualModel !== undefined ? { actualModel: payload.actualModel } : {}), ...(payload.attachments !== undefined ? { attachments: payload.attachments } : {}), ...(payload.context !== undefined ? { context: payload.context } : {}), turnId: payload.turnId, @@ -790,6 +791,9 @@ export function projectEvent( streaming: message.streaming, updatedAt: message.updatedAt, turnId: message.turnId, + ...(message.actualModel !== undefined + ? { actualModel: message.actualModel } + : {}), ...(message.attachments !== undefined ? { attachments: message.attachments } : {}), diff --git a/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts b/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts index 87a15b95e413..93d29d6b2703 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts @@ -2,6 +2,7 @@ import { MessageId, ThreadId, TurnId } from "@t3tools/contracts"; import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as Option from "effect/Option"; import { ProjectionThreadMessageRepository } from "../Services/ProjectionThreadMessages.ts"; import { ProjectionThreadMessageRepositoryLive } from "./ProjectionThreadMessages.ts"; @@ -180,6 +181,43 @@ layer("ProjectionThreadMessageRepository", (it) => { }), ); + it.effect("persists actual model metadata across later partial upserts", () => + Effect.gen(function* () { + const repository = yield* ProjectionThreadMessageRepository; + const threadId = ThreadId.make("thread-actual-model"); + const messageId = MessageId.make("message-actual-model"); + const createdAt = "2026-08-14T00:00:00.000Z"; + + yield* repository.upsert({ + messageId, + threadId, + turnId: null, + role: "assistant", + text: "complete", + actualModel: "gpt-5.6-luna", + isStreaming: false, + createdAt, + updatedAt: createdAt, + }); + yield* repository.upsert({ + messageId, + threadId, + turnId: null, + role: "assistant", + text: "complete", + isStreaming: false, + createdAt, + updatedAt: "2026-08-14T00:00:01.000Z", + }); + + const row = yield* repository.getByMessageId({ messageId }); + assert.equal(row._tag, "Some"); + if (row._tag === "Some") { + assert.equal(row.value.actualModel, "gpt-5.6-luna"); + } + }), + ); + it.effect("preserves existing attachments when upsert omits attachments", () => Effect.gen(function* () { const repository = yield* ProjectionThreadMessageRepository; @@ -324,4 +362,36 @@ layer("ProjectionThreadMessageRepository", (it) => { ); }), ); + + it.effect("finds the latest assistant message id for a turn without hydrating message text", () => + Effect.gen(function* () { + const repository = yield* ProjectionThreadMessageRepository; + const threadId = ThreadId.make("thread-latest-assistant-message"); + const turnId = TurnId.make("turn-latest-assistant-message"); + + assert.isTrue( + Option.isNone(yield* repository.getLatestAssistantMessageIdForTurn({ threadId, turnId })), + ); + for (const [index, createdAt] of [ + "2026-03-01T00:00:00.000Z", + "2026-03-01T00:00:01.000Z", + ].entries()) { + yield* repository.upsert({ + messageId: MessageId.make(`message-latest-assistant-${index}`), + threadId, + turnId, + role: "assistant", + text: "large text that the id query must not select", + isStreaming: false, + createdAt, + updatedAt: createdAt, + }); + } + + assert.deepEqual( + yield* repository.getLatestAssistantMessageIdForTurn({ threadId, turnId }), + Option.some(MessageId.make("message-latest-assistant-1")), + ); + }), + ); }); diff --git a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts index 28aeb6d794e9..8c0086d3575b 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts @@ -5,11 +5,16 @@ import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Schema from "effect/Schema"; import * as Struct from "effect/Struct"; -import { ChatAttachment, OrchestrationMessageContext } from "@t3tools/contracts"; +import { + ChatAttachment, + OrchestrationMessageContext, + TrimmedNonEmptyString, +} from "@t3tools/contracts"; import { toPersistenceSqlError } from "../Errors.ts"; import { AppendStreamingProjectionThreadMessage, + GetLatestProjectionThreadAssistantMessageInput, GetProjectionThreadMessageInput, HasProjectionThreadAssistantMessageInput, ProjectionThreadMessageRepository, @@ -22,11 +27,15 @@ import { const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( Struct.assign({ isStreaming: Schema.Number, + actualModel: Schema.NullOr(TrimmedNonEmptyString), attachments: Schema.NullOr(Schema.fromJsonString(Schema.Array(ChatAttachment))), context: Schema.NullOr(Schema.fromJsonString(OrchestrationMessageContext)), }), ); const ProjectionThreadMessageExistsDbRowSchema = Schema.Struct({ exists: Schema.Number }); +const ProjectionThreadMessageIdDbRowSchema = Schema.Struct({ + messageId: ProjectionThreadMessage.fields.messageId, +}); function toProjectionThreadMessage( row: Schema.Schema.Type, @@ -37,6 +46,7 @@ function toProjectionThreadMessage( turnId: row.turnId, role: row.role, text: row.text, + ...(row.actualModel !== null ? { actualModel: row.actualModel } : {}), isStreaming: row.isStreaming === 1, createdAt: row.createdAt, updatedAt: row.updatedAt, @@ -61,6 +71,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { turn_id, role, text, + actual_model, attachments_json, context_json, is_streaming, @@ -73,6 +84,14 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { ${row.turnId}, ${row.role}, ${row.text}, + COALESCE( + ${row.actualModel ?? null}, + ( + SELECT actual_model + FROM projection_thread_messages + WHERE message_id = ${row.messageId} + ) + ), COALESCE( ${nextAttachmentsJson}, ( @@ -99,6 +118,10 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { turn_id = excluded.turn_id, role = excluded.role, text = excluded.text, + actual_model = COALESCE( + excluded.actual_model, + projection_thread_messages.actual_model + ), attachments_json = COALESCE( excluded.attachments_json, projection_thread_messages.attachments_json @@ -176,6 +199,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -204,6 +228,20 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { `, }); + const getLatestProjectionThreadAssistantMessageIdRow = SqlSchema.findOneOption({ + Request: GetLatestProjectionThreadAssistantMessageInput, + Result: ProjectionThreadMessageIdDbRowSchema, + execute: ({ threadId, turnId }) => sql` + SELECT message_id AS "messageId" + FROM projection_thread_messages + WHERE thread_id = ${threadId} + AND turn_id = ${turnId} + AND role = 'assistant' + ORDER BY created_at DESC, message_id DESC + LIMIT 1 + `, + }); + const listProjectionThreadMessageRows = SqlSchema.findAll({ Request: ListProjectionThreadMessagesInput, Result: ProjectionThreadMessageDbRowSchema, @@ -215,6 +253,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { turn_id AS "turnId", role, text, + actual_model AS "actualModel", attachments_json AS "attachments", context_json AS "context", is_streaming AS "isStreaming", @@ -279,6 +318,17 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { Effect.map((row) => row.exists === 1), ); + const getLatestAssistantMessageIdForTurn: ProjectionThreadMessageRepositoryShape["getLatestAssistantMessageIdForTurn"] = + (input) => + getLatestProjectionThreadAssistantMessageIdRow(input).pipe( + Effect.mapError( + toPersistenceSqlError( + "ProjectionThreadMessageRepository.getLatestAssistantMessageIdForTurn:query", + ), + ), + Effect.map(Option.map((row) => row.messageId)), + ); + const listByThreadId: ProjectionThreadMessageRepositoryShape["listByThreadId"] = (input) => listProjectionThreadMessageRows(input).pipe( Effect.mapError( @@ -309,6 +359,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { appendStreaming, getByMessageId, hasAssistantMessageForTurn, + getLatestAssistantMessageIdForTurn, listByThreadId, getLatestUserMessageAt, deleteByThreadId, diff --git a/apps/server/src/persistence/Migrations.ts b/apps/server/src/persistence/Migrations.ts index ad015534ea01..7536a358e216 100644 --- a/apps/server/src/persistence/Migrations.ts +++ b/apps/server/src/persistence/Migrations.ts @@ -64,6 +64,7 @@ import Migration0049 from "./Migrations/049_ProjectionThreadsActiveOrderKey.ts"; import Migration0050 from "./Migrations/050_ProjectionThreadPullRequests.ts"; import Migration0051 from "./Migrations/051_ProjectionThreadMessageContext.ts"; import Migration0052 from "./Migrations/052_ProjectionThreadTitleState.ts"; +import Migration0053 from "./Migrations/053_ProjectionThreadMessageActualModel.ts"; /** * Migration loader with all migrations defined inline. @@ -128,6 +129,7 @@ const migrationEntries = [ [50, "ProjectionThreadPullRequests", Migration0050], [51, "ProjectionThreadMessageContext", Migration0051], [52, "ProjectionThreadTitleState", Migration0052], + [53, "ProjectionThreadMessageActualModel", Migration0053], ] as const; export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const); diff --git a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts new file mode 100644 index 000000000000..c9bea56a6310 --- /dev/null +++ b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts @@ -0,0 +1,27 @@ +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; + +import { runMigrations } from "../Migrations.ts"; + +const layer = it.layer(NodeSqliteClient.layerMemory()); + +layer("053_ProjectionThreadMessageActualModel", (it) => { + it.effect("adds the nullable actual model to message projections", () => + Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + + yield* runMigrations({ toMigrationInclusive: 49 }); + yield* runMigrations({ toMigrationInclusive: 53 }); + + const columns = yield* sql<{ readonly name: string; readonly notnull: number }>` + PRAGMA table_info(projection_thread_messages) + `; + const actualModel = columns.find((column) => column.name === "actual_model"); + + assert.equal(actualModel?.name, "actual_model"); + assert.equal(actualModel?.notnull, 0); + }), + ); +}); diff --git a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.ts b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.ts new file mode 100644 index 000000000000..4e363e1c4dcd --- /dev/null +++ b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.ts @@ -0,0 +1,16 @@ +import * as Effect from "effect/Effect"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +export default Effect.gen(function* () { + const sql = yield* SqlClient.SqlClient; + const columns = yield* sql<{ readonly name: string }>` + PRAGMA table_info(projection_thread_messages) + `; + + if (!columns.some((column) => column.name === "actual_model")) { + yield* sql` + ALTER TABLE projection_thread_messages + ADD COLUMN actual_model TEXT + `; + } +}); diff --git a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts index e3e5b6e5151d..c3460cd693dd 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts @@ -14,6 +14,7 @@ import { ThreadId, TurnId, IsoDateTime, + TrimmedNonEmptyString, } from "@t3tools/contracts"; import * as Schema from "effect/Schema"; import * as Context from "effect/Context"; @@ -29,6 +30,7 @@ export const ProjectionThreadMessage = Schema.Struct({ turnId: Schema.NullOr(TurnId), role: OrchestrationMessageRole, text: Schema.String, + actualModel: Schema.optional(TrimmedNonEmptyString), attachments: Schema.optional(Schema.Array(ChatAttachment)), context: Schema.optional(OrchestrationMessageContext), isStreaming: Schema.Boolean, @@ -61,6 +63,13 @@ export const HasProjectionThreadAssistantMessageInput = Schema.Struct({ export type HasProjectionThreadAssistantMessageInput = typeof HasProjectionThreadAssistantMessageInput.Type; +export const GetLatestProjectionThreadAssistantMessageInput = Schema.Struct({ + threadId: ThreadId, + turnId: TurnId, +}); +export type GetLatestProjectionThreadAssistantMessageInput = + typeof GetLatestProjectionThreadAssistantMessageInput.Type; + export const DeleteProjectionThreadMessagesInput = Schema.Struct({ threadId: ThreadId, }); @@ -98,6 +107,11 @@ export interface ProjectionThreadMessageRepositoryShape { input: HasProjectionThreadAssistantMessageInput, ) => Effect.Effect; + /** Read the latest assistant message id for a turn without hydrating its text. */ + readonly getLatestAssistantMessageIdForTurn: ( + input: GetLatestProjectionThreadAssistantMessageInput, + ) => Effect.Effect, ProjectionRepositoryError>; + /** * List projected thread messages for a thread. * diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index 96d6b10d3839..54e273a4f67d 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -30,6 +30,7 @@ import { ProviderDriverKind, ProviderInstanceId, ThreadId, + TurnId, } from "@t3tools/contracts"; import { createModelSelection } from "@t3tools/shared/model"; import { ServerConfig } from "../../config.ts"; @@ -47,6 +48,8 @@ import { isSameOpenCodeDirectory, makeOpenCodeAdapter, mergeOpenCodeAssistantText, + rememberCompletedTurnWithoutModel, + rememberMessageTurn, } from "./OpenCodeAdapter.ts"; import { symlinksSupported } from "@t3tools/shared/testing/symlinks"; @@ -6892,6 +6895,373 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { }), ); + it("bounds completed turns waiting for actual model metadata", () => { + const pendingTurns = new Map, "completed" | "failed">(); + for (let index = 0; index < 20; index += 1) { + rememberCompletedTurnWithoutModel( + pendingTurns, + TurnId.make(`turn-pending-model-${index}`), + "completed", + ); + } + + NodeAssert.equal(pendingTurns.size, 16); + NodeAssert.equal(pendingTurns.has(TurnId.make("turn-pending-model-0")), false); + NodeAssert.equal(pendingTurns.has(TurnId.make("turn-pending-model-3")), false); + NodeAssert.equal(pendingTurns.has(TurnId.make("turn-pending-model-4")), true); + NodeAssert.equal(pendingTurns.has(TurnId.make("turn-pending-model-19")), true); + }); + + it("does not rebind a message to a newer active turn", () => { + const messageTurns = new Map>(); + const originalTurnId = TurnId.make("turn-original-message"); + rememberMessageTurn(messageTurns, "msg-late-role", originalTurnId); + rememberMessageTurn(messageTurns, "msg-late-role", TurnId.make("turn-new-active")); + + NodeAssert.equal(messageTurns.size, 1); + NodeAssert.equal(messageTurns.get("msg-late-role"), originalTurnId); + }); + + it.effect("enriches turn completion when response model metadata arrives after idle", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-late-actual-model"); + const pushEvent = makeOpenCodeEventQueue(); + + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(2), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "identify the late model", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "kedvai/auto"), + }); + const promptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { id: "msg-late-actual-model", role: "assistant", parentID: promptMessageId }, + }, + }); + pushEvent({ + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + pushEvent({ + type: "message.part.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + part: { + id: "part-late-step-finish", + sessionID: "http://127.0.0.1:9999/session", + messageID: "msg-late-actual-model", + type: "step-finish", + reason: "stop", + modelID: "gpt-5.6-luna", + cost: 0, + tokens: { input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }, + }, + }, + }); + + const completions = Array.from( + yield* Fiber.join(eventsFiber).pipe(Effect.timeout("1 second")), + ); + const actualModels = completions.map((event) => + event.type === "turn.completed" ? event.payload.actualModel : undefined, + ); + NodeAssert.equal(completions.length, 2); + NodeAssert.equal(actualModels[0], undefined); + NodeAssert.equal(actualModels[1], "gpt-5.6-luna"); + NodeAssert.equal(completions[1]?.providerRefs?.providerItemId, "msg-late-actual-model"); + }), + ); + + it.effect("includes response model metadata on normal turn completion", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-current-actual-model"); + const pushEvent = makeOpenCodeEventQueue(); + const completedFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.runHead, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "identify the current model", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "kedvai/auto"), + }); + const promptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { id: "msg-current-actual-model", role: "assistant", parentID: promptMessageId }, + }, + }); + pushEvent({ + type: "message.part.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + part: { + id: "part-current-step-finish", + sessionID: "http://127.0.0.1:9999/session", + messageID: "msg-current-actual-model", + type: "step-finish", + reason: "stop", + modelID: "gpt-5.6-sol", + cost: 0, + tokens: { input: 2, output: 3, reasoning: 1, cache: { read: 0, write: 0 } }, + }, + }, + }); + pushEvent({ + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const completed = Option.getOrUndefined( + yield* Fiber.join(completedFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(completed?.type, "turn.completed"); + if (completed?.type === "turn.completed") { + NodeAssert.equal(completed.payload.actualModel, "gpt-5.6-sol"); + NodeAssert.equal(completed.providerRefs?.providerItemId, "msg-current-actual-model"); + } + }), + ); + + it.effect("captures response model metadata received before assistant message metadata", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-reordered-actual-model"); + const pushEvent = makeOpenCodeEventQueue(); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(2), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "identify the reordered model", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "kedvai/auto"), + }); + const promptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + pushEvent({ + type: "message.part.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + part: { + id: "part-reordered-step-finish", + sessionID: "http://127.0.0.1:9999/session", + messageID: "msg-reordered-actual-model", + type: "step-finish", + reason: "stop", + modelID: "gpt-5.6-terra", + cost: 0, + tokens: { input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }, + }, + }, + }); + pushEvent({ + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg-reordered-actual-model", + role: "assistant", + parentID: promptMessageId, + }, + }, + }); + + const completions = Array.from( + yield* Fiber.join(eventsFiber).pipe(Effect.timeout("1 second")), + ); + const actualModels = completions.map((event) => + event.type === "turn.completed" ? event.payload.actualModel : undefined, + ); + NodeAssert.equal(completions.length, 2); + NodeAssert.equal(actualModels[0], undefined); + NodeAssert.equal(actualModels[1], "gpt-5.6-terra"); + NodeAssert.equal(completions[1]?.providerRefs?.providerItemId, "msg-reordered-actual-model"); + }), + ); + + it.effect("does not attach delayed metadata from an older response to the active turn", () => + Effect.gen(function* () { + const adapter = yield* OpenCodeAdapter; + const threadId = asThreadId("thread-opencode-delayed-previous-model"); + const pushEvent = makeOpenCodeEventQueue(); + const eventsFiber = yield* adapter.streamEvents.pipe( + Stream.filter((event) => event.threadId === threadId && event.type === "turn.completed"), + Stream.take(3), + Stream.runCollect, + Effect.forkChild, + ); + + yield* adapter.startSession({ + provider: ProviderDriverKind.make("opencode"), + threadId, + runtimeMode: "full-access", + }); + yield* adapter.sendTurn({ + threadId, + input: "first turn", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "kedvai/auto"), + }); + const firstPromptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + pushEvent({ + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + for (let attempt = 0; attempt < 100; attempt += 1) { + const session = (yield* adapter.listSessions()).find( + (candidate) => candidate.threadId === threadId, + ); + if (session?.status === "ready") break; + yield* Effect.yieldNow; + } + NodeAssert.equal( + (yield* adapter.listSessions()).find((candidate) => candidate.threadId === threadId) + ?.status, + "ready", + ); + + yield* adapter.sendTurn({ + threadId, + input: "second turn", + modelSelection: createModelSelection(ProviderInstanceId.make("opencode"), "kedvai/auto"), + }); + const secondPromptMessageId = (runtimeMock.state.promptCalls.at(-1) as { messageID: string }) + .messageID; + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg-current-second-turn", + role: "assistant", + parentID: secondPromptMessageId, + }, + }, + }); + pushEvent({ + type: "message.part.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + part: { + id: "part-current-second-turn", + sessionID: "http://127.0.0.1:9999/session", + messageID: "msg-current-second-turn", + type: "step-finish", + reason: "stop", + modelID: "gpt-5.6-sol", + cost: 0, + tokens: { input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }, + }, + }, + }); + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { + id: "msg-delayed-first-turn", + role: "assistant", + parentID: firstPromptMessageId, + }, + }, + }); + pushEvent({ + type: "message.part.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + part: { + id: "part-delayed-first-turn", + sessionID: "http://127.0.0.1:9999/session", + messageID: "msg-delayed-first-turn", + type: "step-finish", + reason: "stop", + modelID: "wrong-old-model", + cost: 0, + tokens: { input: 1, output: 1, reasoning: 0, cache: { read: 0, write: 0 } }, + }, + }, + }); + pushEvent({ + type: "session.status", + properties: { + sessionID: "http://127.0.0.1:9999/session", + status: { type: "idle" }, + }, + }); + + const completions = Array.from( + yield* Fiber.join(eventsFiber).pipe(Effect.timeout("1 second")), + ); + NodeAssert.equal(completions.length, 3); + NodeAssert.equal(completions[0]?.type, "turn.completed"); + NodeAssert.equal(completions[1]?.type, "turn.completed"); + NodeAssert.equal(completions[2]?.type, "turn.completed"); + if ( + completions[0]?.type === "turn.completed" && + completions[1]?.type === "turn.completed" && + completions[2]?.type === "turn.completed" + ) { + NodeAssert.equal(completions[0].payload.actualModel, undefined); + NodeAssert.equal(completions[1].payload.actualModel, "wrong-old-model"); + NodeAssert.equal(completions[1].providerRefs?.providerItemId, "msg-delayed-first-turn"); + NodeAssert.equal(completions[2].payload.actualModel, "gpt-5.6-sol"); + NodeAssert.equal(completions[2].providerRefs?.providerItemId, "msg-current-second-turn"); + } + }), + ); + it.effect("does not strip coincidental prefix overlap from OpenCode part deltas", () => Effect.gen(function* () { const adapter = yield* OpenCodeAdapter; diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 41bf634c0d3b..78fe3fd1352f 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -2,6 +2,7 @@ import { EventId, type OpenCodeSettings, ProviderDriverKind, + ProviderItemId, ProviderInstanceId, type ProviderRuntimeEvent, type ProviderSendTurnInput, @@ -69,6 +70,7 @@ const PROVIDER = ProviderDriverKind.make("opencode"); * rather than misread (mirrors GROK_RESUME_VERSION / CURSOR_RESUME_VERSION). */ const OPENCODE_RESUME_VERSION = 1 as const; +const MAX_PENDING_ACTUAL_MODEL_TURNS = 16; /** * Decode a persisted resume cursor into the upstream `ses_…` id. Anything @@ -326,6 +328,32 @@ function isOpenCodeDefaultTitle(title: string): boolean { return OPENCODE_DEFAULT_TITLE_PATTERN.test(title); } +function actualModelFromPart(part: Part): string | undefined { + if (part.type !== "step-finish") return undefined; + const modelID = (part as Part & { readonly modelID?: unknown }).modelID; + return typeof modelID === "string" ? trimText(modelID) : undefined; +} + +export function rememberCompletedTurnWithoutModel( + pendingTurns: Map, + turnId: TurnId, + state: "completed" | "failed", +): void { + pendingTurns.delete(turnId); + pendingTurns.set(turnId, state); + if (pendingTurns.size <= MAX_PENDING_ACTUAL_MODEL_TURNS) return; + const oldestTurnId = pendingTurns.keys().next().value; + if (oldestTurnId !== undefined) pendingTurns.delete(oldestTurnId); +} + +export function rememberMessageTurn( + messageTurns: Map, + messageId: string, + turnId: TurnId, +): void { + if (!messageTurns.has(messageId)) messageTurns.set(messageId, turnId); +} + type OpenCodeTextPart = Extract; type OpenCodeTextPartState = Pick & { @@ -350,11 +378,16 @@ interface OpenCodeSessionContext { readonly pendingPermissions: Map; readonly pendingQuestions: Map; readonly messageRoleById: Map; + readonly turnIdByPromptMessageId: Map; + readonly turnIdByMessageId: Map; + readonly actualModelByMessageId: Map; // OpenCode permits edits to completed parts. Keep text for snapshot comparison // until native removal or session teardown, but do not retain other part payloads. readonly textPartsByMessageId: Map>; turnTokenUsage: OpenCodeTurnTokenUsageAccumulator | undefined; activeTurnId: TurnId | undefined; + activeActualModel: { readonly messageId: string; readonly model: string } | undefined; + readonly completedTurnsWithoutModel: Map; activeAgent: string | undefined; activeVariant: string | undefined; cancellation: OpenCodeCancellation | undefined; @@ -498,6 +531,7 @@ type EventBaseInput = { readonly threadId: ThreadId; readonly turnId?: TurnId | undefined; readonly itemId?: string | undefined; + readonly providerItemId?: string | undefined; readonly requestId?: string | undefined; readonly createdAt?: string | undefined; readonly raw?: unknown; @@ -1036,6 +1070,9 @@ export function makeOpenCodeAdapter( createdAt, ...(input.turnId ? { turnId: input.turnId } : {}), ...(input.itemId ? { itemId: RuntimeItemId.make(input.itemId) } : {}), + ...(input.providerItemId + ? { providerRefs: { providerItemId: ProviderItemId.make(input.providerItemId) } } + : {}), ...(input.requestId ? { requestId: RuntimeRequestId.make(input.requestId) } : {}), ...(input.raw !== undefined ? { @@ -1133,12 +1170,19 @@ export function makeOpenCodeAdapter( context.pendingIdleReconciliation = undefined; } const tokenUsage = takeOpenCodeTurnTokenUsage(context, true); + const actualModel = context.activeActualModel; context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.interruptedTurnId = undefined; context.awaitingBusyAfterInterruption = false; context.reconcileIdleStatus = false; + if (actualModel) { + context.completedTurnsWithoutModel.delete(turnId); + } else { + rememberCompletedTurnWithoutModel(context.completedTurnsWithoutModel, turnId, "completed"); + } for (const requestId of context.autoRepliedRequestIds) { context.emittedTerminalRequestIds.add(requestId); } @@ -1157,12 +1201,14 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId, + providerItemId: actualModel?.messageId, raw, })), type: "turn.completed", payload: { state: "completed", tokenUsage, + ...(actualModel ? { actualModel: actualModel.model } : {}), }, }); }); @@ -1301,12 +1347,23 @@ export function makeOpenCodeAdapter( return; } const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); + const actualModel = context.activeActualModel; context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; context.reconcileIdleStatus = false; + if (actualModel) { + context.completedTurnsWithoutModel.delete(promptAdmission.turnId); + } else { + rememberCompletedTurnWithoutModel( + context.completedTurnsWithoutModel, + promptAdmission.turnId, + "failed", + ); + } yield* updateProviderSession( context, { status: "error", lastError: detail }, @@ -1316,6 +1373,7 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId: promptAdmission.turnId, + providerItemId: actualModel?.messageId, raw: promptAdmission.recoveryRaw, })), type: "turn.completed", @@ -1323,6 +1381,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: detail, tokenUsage, + ...(actualModel ? { actualModel: actualModel.model } : {}), }, }); yield* emit({ @@ -1527,6 +1586,7 @@ export function makeOpenCodeAdapter( if (context.activeTurnId === turnId) { tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; yield* updateProviderSession( @@ -1632,6 +1692,7 @@ export function makeOpenCodeAdapter( threadId: context.session.threadId, turnId, itemId: part.id, + providerItemId: part.messageID, createdAt: part.time !== undefined ? isoFromEpochMs(part.time.start) : undefined, raw, })), @@ -1650,6 +1711,7 @@ export function makeOpenCodeAdapter( threadId: context.session.threadId, turnId, itemId: part.id, + providerItemId: part.messageID, createdAt: isoFromEpochMs(part.time.end), raw, })), @@ -2175,6 +2237,37 @@ export function makeOpenCodeAdapter( yield* run.pipe(Effect.forkIn(context.sessionScope)); }); + const captureActualModel = Effect.fn("captureActualModel")(function* ( + context: OpenCodeSessionContext, + messageId: string, + actualModel: string, + raw: unknown, + ) { + const partTurnId = context.turnIdByMessageId.get(messageId); + if (!partTurnId) return false; + if (context.activeTurnId === partTurnId) { + context.activeActualModel = { messageId, model: actualModel }; + return true; + } + const completedState = context.completedTurnsWithoutModel.get(partTurnId); + if (!completedState) return false; + context.completedTurnsWithoutModel.delete(partTurnId); + yield* emit({ + ...(yield* buildEventBase({ + threadId: context.session.threadId, + turnId: partTurnId, + providerItemId: messageId, + raw, + })), + type: "turn.completed", + payload: { + state: completedState, + actualModel, + }, + }); + return true; + }); + const handleSubscribedEvent = Effect.fn("handleSubscribedEvent")(function* ( context: OpenCodeSessionContext, event: OpenCodeSubscribedEvent, @@ -2358,6 +2451,9 @@ export function makeOpenCodeAdapter( event.properties.info.parentID.trim().length > 0 ? event.properties.info.parentID : undefined; + const messageTurnId = parentMessageId + ? context.turnIdByPromptMessageId.get(parentMessageId) + : undefined; const observedOwnership = parentMessageId === undefined ? "unknown" @@ -2383,6 +2479,20 @@ export function makeOpenCodeAdapter( usage.unresolvedStepsByMessageId.delete(event.properties.info.id); } } + if (messageTurnId) { + rememberMessageTurn( + context.turnIdByMessageId, + event.properties.info.id, + messageTurnId, + ); + } + const actualModel = context.actualModelByMessageId.get(event.properties.info.id); + if ( + actualModel && + (yield* captureActualModel(context, event.properties.info.id, actualModel, event)) + ) { + context.actualModelByMessageId.delete(event.properties.info.id); + } for (const part of context.textPartsByMessageId .get(event.properties.info.id) ?.values() ?? []) { @@ -2394,6 +2504,9 @@ export function makeOpenCodeAdapter( case "message.removed": { context.messageRoleById.delete(event.properties.messageID); + context.turnIdByPromptMessageId.delete(event.properties.messageID); + context.turnIdByMessageId.delete(event.properties.messageID); + context.actualModelByMessageId.delete(event.properties.messageID); context.textPartsByMessageId.delete(event.properties.messageID); break; } @@ -2435,6 +2548,7 @@ export function makeOpenCodeAdapter( threadId: context.session.threadId, turnId, itemId: event.properties.partID, + providerItemId: event.properties.messageID, raw: event, })), type: "content.delta", @@ -2450,6 +2564,24 @@ export function makeOpenCodeAdapter( const part = event.properties.part; const messageRole = messageRoleForPart(context, part); + const actualModel = actualModelFromPart(part); + if (actualModel) { + context.actualModelByMessageId.delete(part.messageID); + context.actualModelByMessageId.set(part.messageID, actualModel); + if (context.actualModelByMessageId.size > MAX_PENDING_ACTUAL_MODEL_TURNS) { + const oldestMessageId = context.actualModelByMessageId.keys().next().value; + if (oldestMessageId !== undefined) { + context.actualModelByMessageId.delete(oldestMessageId); + } + } + if ( + messageRole === "assistant" && + (yield* captureActualModel(context, part.messageID, actualModel, event)) + ) { + context.actualModelByMessageId.delete(part.messageID); + } + } + if (turnId && part.type === "step-finish" && context.turnTokenUsage) { const usage = context.turnTokenUsage; const ownership = usage.assistantOwnershipByMessageId.get(part.messageID); @@ -2646,6 +2778,7 @@ export function makeOpenCodeAdapter( case "session.error": { const message = sessionErrorMessage(event.properties.error); const activeTurnId = context.activeTurnId; + const actualModel = context.activeActualModel; const cancellation = context.cancellation; if (isOpenCodeAbortError(event.properties.error)) { if (cancellation !== undefined && cancellation.turnId === undefined) { @@ -2673,9 +2806,21 @@ export function makeOpenCodeAdapter( } const tokenUsage = activeTurnId ? takeOpenCodeTurnTokenUsage(context, false) : undefined; context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.reconcileIdleStatus = false; + if (activeTurnId) { + if (actualModel) { + context.completedTurnsWithoutModel.delete(activeTurnId); + } else { + rememberCompletedTurnWithoutModel( + context.completedTurnsWithoutModel, + activeTurnId, + "failed", + ); + } + } yield* schedulePendingRequestRecovery(context); yield* updateProviderSession( context, @@ -2690,6 +2835,7 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId: activeTurnId, + providerItemId: actualModel?.messageId, raw: event, })), type: "turn.completed", @@ -2697,6 +2843,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: message, tokenUsage, + ...(actualModel ? { actualModel: actualModel.model } : {}), }, }); } @@ -3017,8 +3164,13 @@ export function makeOpenCodeAdapter( pendingQuestions: new Map(), textPartsByMessageId: new Map(), messageRoleById: new Map(), + turnIdByPromptMessageId: new Map(), + turnIdByMessageId: new Map(), + actualModelByMessageId: new Map(), turnTokenUsage: undefined, activeTurnId: undefined, + activeActualModel: undefined, + completedTurnsWithoutModel: new Map(), activeAgent: undefined, activeVariant: undefined, cancellation: undefined, @@ -3201,8 +3353,17 @@ export function makeOpenCodeAdapter( context.activeTurnId = turnId; if (steeringTurnId === undefined) { + context.activeActualModel = undefined; context.turnTokenUsage = makeOpenCodeTurnTokenUsageAccumulator(); } + context.turnIdByPromptMessageId.delete(messageId); + context.turnIdByPromptMessageId.set(messageId, turnId); + if (context.turnIdByPromptMessageId.size > MAX_PENDING_ACTUAL_MODEL_TURNS) { + const oldestPromptMessageId = context.turnIdByPromptMessageId.keys().next().value; + if (oldestPromptMessageId !== undefined) { + context.turnIdByPromptMessageId.delete(oldestPromptMessageId); + } + } context.turnTokenUsage?.promptMessageIds.add(messageId); context.activeAgent = agent ?? (input.interactionMode === "plan" ? "plan" : undefined); context.activeVariant = variant; @@ -3354,6 +3515,7 @@ export function makeOpenCodeAdapter( const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; yield* updateProviderSession( @@ -3402,6 +3564,7 @@ export function makeOpenCodeAdapter( const tokenUsage = takeOpenCodeTurnTokenUsage(context, false); context.promptAdmission = undefined; context.activeTurnId = undefined; + context.activeActualModel = undefined; context.activeAgent = undefined; context.activeVariant = undefined; context.awaitingBusyAfterInterruption = false; diff --git a/apps/web/src/components/chat/MessagesTimeline.test.tsx b/apps/web/src/components/chat/MessagesTimeline.test.tsx index 3bf5f482948d..27040db6d0a9 100644 --- a/apps/web/src/components/chat/MessagesTimeline.test.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.test.tsx @@ -304,6 +304,28 @@ describe("MessagesTimeline", () => { expect(markup).toContain('aria-label="Next turn"'); }); + it("shows the actual model on a completed assistant turn", () => { + const entry = buildAssistantTimelineEntry("Hello"); + const markup = renderToStaticMarkup( + , + ); + + expect(markup).toContain("data-assistant-actual-model"); + expect(markup).toContain("Model: gpt-5.6-luna"); + }); + // Expanding history uses this suite's existing test renderer, deprecated in // React 19. Migrate these interaction tests together when a DOM test setup is added. it.each([{}, { text: "Text-only answer", file: "Answer with a file" }])( diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 117e2d51c276..4a89ca6e2a18 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -2431,30 +2431,46 @@ function AssistantMessageMeta({ const ctx = use(TimelineRowCtx); return ( -
- - {!message.streaming && ( +
+ {message.actualModel ? ( - }> - {formatDayAwareTimestamp(message.updatedAt, ctx.timestampFormat)} + + } + > + Model: {message.actualModel} - - {formatChatTimestampTooltip(message.updatedAt, ctx.timestampFormat)} - + Actual model: {message.actualModel} - )} + ) : null} +
+ + {!message.streaming && ( + + }> + {formatDayAwareTimestamp(message.updatedAt, ctx.timestampFormat)} + + + {formatChatTimestampTooltip(message.updatedAt, ctx.timestampFormat)} + + + )} +
); } diff --git a/packages/client-runtime/src/state/threadReducer.ts b/packages/client-runtime/src/state/threadReducer.ts index 101bb34fba91..78f28b95f6f1 100644 --- a/packages/client-runtime/src/state/threadReducer.ts +++ b/packages/client-runtime/src/state/threadReducer.ts @@ -380,6 +380,9 @@ export function applyThreadDetailEvent( id: event.payload.messageId, role: event.payload.role, text: event.payload.text, + ...(event.payload.actualModel !== undefined + ? { actualModel: event.payload.actualModel } + : {}), ...(event.payload.attachments !== undefined ? { attachments: event.payload.attachments } : {}), @@ -403,6 +406,7 @@ export function applyThreadDetailEvent( : entry.text, streaming: message.streaming, ...(message.turnId !== undefined ? { turnId: message.turnId } : {}), + ...(message.actualModel !== undefined ? { actualModel: message.actualModel } : {}), ...(message.streaming ? {} : { updatedAt: message.updatedAt }), ...(message.attachments !== undefined ? { attachments: message.attachments } : {}), ...(message.context !== undefined ? { context: message.context } : {}), diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index cdcc5b7e437a..56e73d09565f 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -554,6 +554,7 @@ export const OrchestrationMessage = Schema.Struct({ id: MessageId, role: OrchestrationMessageRole, text: Schema.String, + actualModel: Schema.optional(TrimmedNonEmptyString), attachments: Schema.optional(Schema.Array(ChatAttachment)), context: Schema.optional(OrchestrationMessageContext), turnId: Schema.NullOr(TurnId), @@ -1478,6 +1479,7 @@ const ThreadMessageAssistantCompleteCommand = Schema.Struct({ threadId: ThreadId, messageId: MessageId, turnId: Schema.optional(TurnId), + actualModel: Schema.optional(TrimmedNonEmptyString), createdAt: IsoDateTime, }); @@ -1869,6 +1871,7 @@ export const ThreadMessageSentPayload = Schema.Struct({ messageId: MessageId, role: OrchestrationMessageRole, text: Schema.String, + actualModel: Schema.optional(TrimmedNonEmptyString), attachments: Schema.optional(Schema.Array(ChatAttachment)), context: Schema.optional(OrchestrationMessageContext), turnId: Schema.NullOr(TurnId), diff --git a/packages/contracts/src/providerRuntime.ts b/packages/contracts/src/providerRuntime.ts index af1baac74f9d..929876bfaa6e 100644 --- a/packages/contracts/src/providerRuntime.ts +++ b/packages/contracts/src/providerRuntime.ts @@ -399,6 +399,7 @@ export type TurnTokenUsage = typeof TurnTokenUsage.Type; const TurnCompletedPayload = Schema.Struct({ state: RuntimeTurnState, + actualModel: Schema.optional(TrimmedNonEmptyStringSchema), stopReason: Schema.optional(Schema.NullOr(TrimmedNonEmptyStringSchema)), usage: Schema.optional(Schema.Unknown), modelUsage: Schema.optional(UnknownRecordSchema),