From 9346b7a834195a6181d0bb6e238f4ece9272be47 Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Fri, 14 Aug 2026 01:29:06 +0200 Subject: [PATCH 1/9] feat(opencode): show actual model per turn --- .../src/features/threads/ThreadFeed.tsx | 9 ++ .../Layers/ProjectionPipeline.ts | 5 + .../Layers/ProjectionSnapshotQuery.ts | 7 + .../Layers/ProviderRuntimeIngestion.test.ts | 66 ++++++++ .../Layers/ProviderRuntimeIngestion.ts | 24 ++- apps/server/src/orchestration/decider.ts | 1 + apps/server/src/orchestration/projector.ts | 4 + .../Layers/ProjectionThreadMessages.test.ts | 37 +++++ .../Layers/ProjectionThreadMessages.ts | 23 ++- apps/server/src/persistence/Migrations.ts | 2 + ...ProjectionThreadMessageActualModel.test.ts | 28 ++++ .../053_ProjectionThreadMessageActualModel.ts | 16 ++ .../Services/ProjectionThreadMessages.ts | 2 + .../src/provider/Layers/OpenCodeAdapter.ts | 145 ++++++++++++++++++ .../components/chat/MessagesTimeline.test.tsx | 22 +++ .../src/components/chat/MessagesTimeline.tsx | 56 ++++--- .../client-runtime/src/state/threadReducer.ts | 4 + packages/contracts/src/orchestration.ts | 3 + packages/contracts/src/providerRuntime.ts | 1 + 19 files changed, 434 insertions(+), 21 deletions(-) create mode 100644 apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts create mode 100644 apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.ts diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 585354f7baf2..00808373e8cd 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} { 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, + 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, + 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("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..e55a2f372be1 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1475,6 +1475,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 +1512,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, }); } @@ -2362,7 +2364,23 @@ const make = Effect.gen(function* () { createdAt: now, }); } - const assistantMessageIds = yield* getAssistantMessageIdsForTurn(thread.id, turnId); + const trackedAssistantMessageIds = yield* getAssistantMessageIdsForTurn( + thread.id, + turnId, + ); + const completedAssistantMessageId = + trackedAssistantMessageIds.size === 0 && event.payload.actualModel + ? (yield* getLoadedThreadDetail())?.messages.findLast( + (message) => message.role === "assistant" && message.turnId === turnId, + )?.id + : undefined; + const assistantMessageIds = + trackedAssistantMessageIds.size > 0 + ? Array.from(trackedAssistantMessageIds) + : completedAssistantMessageId + ? [completedAssistantMessageId] + : []; + const terminalAssistantMessageId = assistantMessageIds.at(-1); yield* Effect.forEach( assistantMessageIds, (assistantMessageId) => @@ -2377,6 +2395,10 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete-finalize", finalDeltaCommandTag: "assistant-delta-finalize-fallback", hasProjectedMessage: existingMessage !== undefined, + ...(event.payload.actualModel && + assistantMessageId === terminalAssistantMessageId + ? { actualModel: event.payload.actualModel } + : {}), }), ), ), diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0119c0e8599a..2a88a79987c2 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1972,6 +1972,7 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" messageId: command.messageId, role: command.type === "thread.message.reasoning.complete" ? "reasoning" : "assistant", text: "", + ...(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..48a2748d37a1 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts @@ -180,6 +180,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; diff --git a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts index 28aeb6d794e9..beace4628d43 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts @@ -5,7 +5,11 @@ 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 { @@ -22,6 +26,7 @@ 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)), }), @@ -37,6 +42,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 +67,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { turn_id, role, text, + actual_model, attachments_json, context_json, is_streaming, @@ -73,6 +80,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 +114,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 +195,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", @@ -215,6 +235,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", 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..df24e6e5f630 --- /dev/null +++ b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts @@ -0,0 +1,28 @@ +import { assert, it } from "@effect/vitest"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; +import * as SqlClient from "effect/unstable/sql/SqlClient"; + +import { runMigrations } from "../Migrations.ts"; +import * as NodeSqliteClient from "../NodeSqliteClient.ts"; + +const layer = it.layer(Layer.mergeAll(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: 52 }); + + 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..a8f581021318 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, diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 41bf634c0d3b..deb518741c24 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -69,6 +69,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 +327,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 +377,15 @@ interface OpenCodeSessionContext { readonly pendingPermissions: Map; readonly pendingQuestions: Map; readonly messageRoleById: 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: string | undefined; + readonly completedTurnsWithoutModel: Map; activeAgent: string | undefined; activeVariant: string | undefined; cancellation: OpenCodeCancellation | undefined; @@ -1133,12 +1164,23 @@ 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); } @@ -1163,6 +1205,7 @@ export function makeOpenCodeAdapter( payload: { state: "completed", tokenUsage, + ...(actualModel ? { actualModel } : {}), }, }); }); @@ -1301,12 +1344,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 }, @@ -1323,6 +1377,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: detail, tokenUsage, + ...(actualModel ? { actualModel } : {}), }, }); yield* emit({ @@ -1527,6 +1582,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( @@ -2175,6 +2231,36 @@ 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 = 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, + raw, + })), + type: "turn.completed", + payload: { + state: completedState, + actualModel, + }, + }); + return true; + }); + const handleSubscribedEvent = Effect.fn("handleSubscribedEvent")(function* ( context: OpenCodeSessionContext, event: OpenCodeSubscribedEvent, @@ -2352,6 +2438,21 @@ export function makeOpenCodeAdapter( context.textPartsByMessageId.delete(event.properties.info.id); } if (event.properties.info.role === "assistant") { + if (turnId) { + rememberMessageTurn(context.turnIdByMessageId, event.properties.info.id, turnId); + } + 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); + } const usage = context.turnTokenUsage; const parentMessageId = typeof event.properties.info.parentID === "string" && @@ -2394,6 +2495,8 @@ export function makeOpenCodeAdapter( case "message.removed": { context.messageRoleById.delete(event.properties.messageID); + context.turnIdByMessageId.delete(event.properties.messageID); + context.actualModelByMessageId.delete(event.properties.messageID); context.textPartsByMessageId.delete(event.properties.messageID); break; } @@ -2450,6 +2553,27 @@ export function makeOpenCodeAdapter( const part = event.properties.part; const messageRole = messageRoleForPart(context, part); + if (turnId) { + rememberMessageTurn(context.turnIdByMessageId, part.messageID, turnId); + } + 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 +2770,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 +2798,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, @@ -2697,6 +2834,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: message, tokenUsage, + ...(actualModel ? { actualModel } : {}), }, }); } @@ -3017,8 +3155,12 @@ export function makeOpenCodeAdapter( pendingQuestions: new Map(), textPartsByMessageId: new Map(), messageRoleById: new Map(), + turnIdByMessageId: new Map(), + actualModelByMessageId: new Map(), turnTokenUsage: undefined, activeTurnId: undefined, + activeActualModel: undefined, + completedTurnsWithoutModel: new Map(), activeAgent: undefined, activeVariant: undefined, cancellation: undefined, @@ -3201,6 +3343,7 @@ export function makeOpenCodeAdapter( context.activeTurnId = turnId; if (steeringTurnId === undefined) { + context.activeActualModel = undefined; context.turnTokenUsage = makeOpenCodeTurnTokenUsageAccumulator(); } context.turnTokenUsage?.promptMessageIds.add(messageId); @@ -3354,6 +3497,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 +3546,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..57de49e01526 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -2432,29 +2432,47 @@ function AssistantMessageMeta({ 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), From a606062647dbc6eac3d5eab3248e04cd855780c4 Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Fri, 14 Aug 2026 02:19:51 +0200 Subject: [PATCH 2/9] fix(opencode): adapt model provenance to current runtime --- .../src/features/threads/ThreadFeed.tsx | 2 +- .../Layers/ProviderRuntimeIngestion.ts | 18 +- .../Layers/ProjectionThreadMessages.test.ts | 33 +++ .../Layers/ProjectionThreadMessages.ts | 30 +++ ...ProjectionThreadMessageActualModel.test.ts | 5 +- .../Services/ProjectionThreadMessages.ts | 12 + .../provider/Layers/OpenCodeAdapter.test.ts | 223 ++++++++++++++++++ .../src/provider/Layers/OpenCodeAdapter.ts | 13 +- .../src/components/chat/MessagesTimeline.tsx | 4 +- 9 files changed, 315 insertions(+), 25 deletions(-) diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index 00808373e8cd..bddc39b24b86 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -1730,7 +1730,7 @@ function renderFeedEntry( {message.actualModel ? ( Model: {message.actualModel} diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index e55a2f372be1..ff553c8b82cf 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -1770,6 +1770,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 = @@ -2369,10 +2371,13 @@ const make = Effect.gen(function* () { turnId, ); const completedAssistantMessageId = - trackedAssistantMessageIds.size === 0 && event.payload.actualModel - ? (yield* getLoadedThreadDetail())?.messages.findLast( - (message) => message.role === "assistant" && message.turnId === turnId, - )?.id + trackedAssistantMessageIds.size === 0 && terminalActualModel + ? Option.getOrUndefined( + yield* projectionThreadMessages.getLatestAssistantMessageIdForTurn({ + threadId: thread.id, + turnId, + }), + ) : undefined; const assistantMessageIds = trackedAssistantMessageIds.size > 0 @@ -2395,9 +2400,8 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete-finalize", finalDeltaCommandTag: "assistant-delta-finalize-fallback", hasProjectedMessage: existingMessage !== undefined, - ...(event.payload.actualModel && - assistantMessageId === terminalAssistantMessageId - ? { actualModel: event.payload.actualModel } + ...(terminalActualModel && assistantMessageId === terminalAssistantMessageId + ? { actualModel: terminalActualModel } : {}), }), ), diff --git a/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts b/apps/server/src/persistence/Layers/ProjectionThreadMessages.test.ts index 48a2748d37a1..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"; @@ -361,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 beace4628d43..8c0086d3575b 100644 --- a/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Layers/ProjectionThreadMessages.ts @@ -14,6 +14,7 @@ import { import { toPersistenceSqlError } from "../Errors.ts"; import { AppendStreamingProjectionThreadMessage, + GetLatestProjectionThreadAssistantMessageInput, GetProjectionThreadMessageInput, HasProjectionThreadAssistantMessageInput, ProjectionThreadMessageRepository, @@ -32,6 +33,9 @@ const ProjectionThreadMessageDbRowSchema = ProjectionThreadMessage.mapFields( }), ); const ProjectionThreadMessageExistsDbRowSchema = Schema.Struct({ exists: Schema.Number }); +const ProjectionThreadMessageIdDbRowSchema = Schema.Struct({ + messageId: ProjectionThreadMessage.fields.messageId, +}); function toProjectionThreadMessage( row: Schema.Schema.Type, @@ -224,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, @@ -300,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( @@ -330,6 +359,7 @@ const makeProjectionThreadMessageRepository = Effect.gen(function* () { appendStreaming, getByMessageId, hasAssistantMessageForTurn, + getLatestAssistantMessageIdForTurn, listByThreadId, getLatestUserMessageAt, deleteByThreadId, diff --git a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts index df24e6e5f630..df0352415903 100644 --- a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts +++ b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts @@ -1,12 +1,11 @@ import { assert, it } from "@effect/vitest"; import * as Effect from "effect/Effect"; -import * as Layer from "effect/Layer"; import * as SqlClient from "effect/unstable/sql/SqlClient"; +import * as NodeSqliteClient from "@t3tools/shared/nodeSqliteClient"; import { runMigrations } from "../Migrations.ts"; -import * as NodeSqliteClient from "../NodeSqliteClient.ts"; -const layer = it.layer(Layer.mergeAll(NodeSqliteClient.layerMemory())); +const layer = it.layer(NodeSqliteClient.layerMemory()); layer("053_ProjectionThreadMessageActualModel", (it) => { it.effect("adds the nullable actual model to message projections", () => diff --git a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts index a8f581021318..c3460cd693dd 100644 --- a/apps/server/src/persistence/Services/ProjectionThreadMessages.ts +++ b/apps/server/src/persistence/Services/ProjectionThreadMessages.ts @@ -63,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, }); @@ -100,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..c49d58b7bc46 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,226 @@ 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"), + }); + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { id: "msg-late-actual-model", role: "assistant" }, + }, + }); + 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"); + }), + ); + + 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"), + }); + pushEvent({ + type: "message.updated", + properties: { + sessionID: "http://127.0.0.1:9999/session", + info: { id: "msg-current-actual-model", role: "assistant" }, + }, + }); + 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"); + } + }), + ); + + 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"), + }); + 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" }, + }, + }); + + 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"); + }), + ); + 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 deb518741c24..2572b2d2c5d5 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.ts @@ -1175,11 +1175,7 @@ export function makeOpenCodeAdapter( if (actualModel) { context.completedTurnsWithoutModel.delete(turnId); } else { - rememberCompletedTurnWithoutModel( - context.completedTurnsWithoutModel, - turnId, - "completed", - ); + rememberCompletedTurnWithoutModel(context.completedTurnsWithoutModel, turnId, "completed"); } for (const requestId of context.autoRepliedRequestIds) { context.emittedTerminalRequestIds.add(requestId); @@ -2444,12 +2440,7 @@ export function makeOpenCodeAdapter( const actualModel = context.actualModelByMessageId.get(event.properties.info.id); if ( actualModel && - (yield* captureActualModel( - context, - event.properties.info.id, - actualModel, - event, - )) + (yield* captureActualModel(context, event.properties.info.id, actualModel, event)) ) { context.actualModelByMessageId.delete(event.properties.info.id); } diff --git a/apps/web/src/components/chat/MessagesTimeline.tsx b/apps/web/src/components/chat/MessagesTimeline.tsx index 57de49e01526..4a89ca6e2a18 100644 --- a/apps/web/src/components/chat/MessagesTimeline.tsx +++ b/apps/web/src/components/chat/MessagesTimeline.tsx @@ -2431,9 +2431,7 @@ function AssistantMessageMeta({ const ctx = use(TimelineRowCtx); return ( -
+
{message.actualModel ? ( Date: Tue, 8 Sep 2026 14:39:23 +0200 Subject: [PATCH 3/9] fix(opencode): target model metadata by response --- .../Layers/ProviderRuntimeIngestion.test.ts | 83 ++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 90 ++++++++++- .../provider/Layers/OpenCodeAdapter.test.ts | 153 +++++++++++++++++- .../src/provider/Layers/OpenCodeAdapter.ts | 63 +++++--- 4 files changed, 363 insertions(+), 26 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index bdcf87656c50..2b409db26fa5 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3863,6 +3863,7 @@ describe("ProviderRuntimeIngestion", () => { threadId, turnId, itemId, + providerRefs: { providerItemId: asItemId("native-message-actual-model") }, payload: { streamKind: "assistant_text", delta: "done" }, }); harness.emit({ @@ -3888,6 +3889,7 @@ describe("ProviderRuntimeIngestion", () => { createdAt: now, threadId, turnId, + providerRefs: { providerItemId: asItemId("native-message-actual-model") }, payload: { state: "completed", actualModel: "gpt-5.6-luna" }, }); @@ -3904,6 +3906,87 @@ describe("ProviderRuntimeIngestion", () => { expect(messages[0]?.actualModel).toBe("gpt-5.6-luna"); }); + 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("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 ff553c8b82cf..e76692013d34 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 @@ -1035,6 +1037,12 @@ 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 bufferedAssistantTextByMessageId = yield* Cache.make({ capacity: BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_CACHE_CAPACITY, timeToLive: BUFFERED_MESSAGE_TEXT_BY_MESSAGE_ID_TTL, @@ -1163,6 +1171,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, @@ -1622,6 +1667,9 @@ 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 assistantSegmentKeys = Array.from(yield* Cache.keys(assistantSegmentStateByTurnKey)); const proposedPlanKeys = Array.from(yield* Cache.keys(bufferedProposedPlanById)); const taskDescriptionKeys = Array.from(yield* Cache.keys(taskDescriptionByTaskKey)); @@ -1644,6 +1692,14 @@ 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( assistantSegmentKeys, (key) => @@ -2036,6 +2092,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); @@ -2370,8 +2435,17 @@ const make = Effect.gen(function* () { 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 && terminalActualModel + trackedAssistantMessageIds.size === 0 && + modelTargetMessageIds.size === 0 && + terminalActualModel && + !terminalProviderItemId ? Option.getOrUndefined( yield* projectionThreadMessages.getLatestAssistantMessageIdForTurn({ threadId: thread.id, @@ -2382,9 +2456,11 @@ const make = Effect.gen(function* () { const assistantMessageIds = trackedAssistantMessageIds.size > 0 ? Array.from(trackedAssistantMessageIds) - : completedAssistantMessageId - ? [completedAssistantMessageId] - : []; + : modelTargetMessageIds.size > 0 + ? Array.from(modelTargetMessageIds) + : completedAssistantMessageId + ? [completedAssistantMessageId] + : []; const terminalAssistantMessageId = assistantMessageIds.at(-1); yield* Effect.forEach( assistantMessageIds, @@ -2400,7 +2476,11 @@ const make = Effect.gen(function* () { commandTag: "assistant-complete-finalize", finalDeltaCommandTag: "assistant-delta-finalize-fallback", hasProjectedMessage: existingMessage !== undefined, - ...(terminalActualModel && assistantMessageId === terminalAssistantMessageId + ...(terminalActualModel && + (modelTargetMessageIds.size > 0 + ? modelTargetMessageIds.has(assistantMessageId) + : !terminalProviderItemId && + assistantMessageId === terminalAssistantMessageId) ? { actualModel: terminalActualModel } : {}), }), diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts index c49d58b7bc46..54e273a4f67d 100644 --- a/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts +++ b/apps/server/src/provider/Layers/OpenCodeAdapter.test.ts @@ -6945,11 +6945,13 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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" }, + info: { id: "msg-late-actual-model", role: "assistant", parentID: promptMessageId }, }, }); pushEvent({ @@ -6985,6 +6987,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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"); }), ); @@ -7009,11 +7012,13 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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" }, + info: { id: "msg-current-actual-model", role: "assistant", parentID: promptMessageId }, }, }); pushEvent({ @@ -7046,6 +7051,7 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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"); } }), ); @@ -7072,6 +7078,8 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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: { @@ -7099,7 +7107,11 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { type: "message.updated", properties: { sessionID: "http://127.0.0.1:9999/session", - info: { id: "msg-reordered-actual-model", role: "assistant" }, + info: { + id: "msg-reordered-actual-model", + role: "assistant", + parentID: promptMessageId, + }, }, }); @@ -7112,6 +7124,141 @@ it.layer(OpenCodeAdapterTestLayer)("OpenCodeAdapterLive", (it) => { 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"); + } }), ); diff --git a/apps/server/src/provider/Layers/OpenCodeAdapter.ts b/apps/server/src/provider/Layers/OpenCodeAdapter.ts index 2572b2d2c5d5..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, @@ -377,6 +378,7 @@ 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 @@ -384,7 +386,7 @@ interface OpenCodeSessionContext { readonly textPartsByMessageId: Map>; turnTokenUsage: OpenCodeTurnTokenUsageAccumulator | undefined; activeTurnId: TurnId | undefined; - activeActualModel: string | undefined; + activeActualModel: { readonly messageId: string; readonly model: string } | undefined; readonly completedTurnsWithoutModel: Map; activeAgent: string | undefined; activeVariant: string | undefined; @@ -529,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; @@ -1067,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 ? { @@ -1195,13 +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 ? { actualModel: actualModel.model } : {}), }, }); }); @@ -1366,6 +1373,7 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId: promptAdmission.turnId, + providerItemId: actualModel?.messageId, raw: promptAdmission.recoveryRaw, })), type: "turn.completed", @@ -1373,7 +1381,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: detail, tokenUsage, - ...(actualModel ? { actualModel } : {}), + ...(actualModel ? { actualModel: actualModel.model } : {}), }, }); yield* emit({ @@ -1684,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, })), @@ -1702,6 +1711,7 @@ export function makeOpenCodeAdapter( threadId: context.session.threadId, turnId, itemId: part.id, + providerItemId: part.messageID, createdAt: isoFromEpochMs(part.time.end), raw, })), @@ -2236,7 +2246,7 @@ export function makeOpenCodeAdapter( const partTurnId = context.turnIdByMessageId.get(messageId); if (!partTurnId) return false; if (context.activeTurnId === partTurnId) { - context.activeActualModel = actualModel; + context.activeActualModel = { messageId, model: actualModel }; return true; } const completedState = context.completedTurnsWithoutModel.get(partTurnId); @@ -2246,6 +2256,7 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId: partTurnId, + providerItemId: messageId, raw, })), type: "turn.completed", @@ -2434,22 +2445,15 @@ export function makeOpenCodeAdapter( context.textPartsByMessageId.delete(event.properties.info.id); } if (event.properties.info.role === "assistant") { - if (turnId) { - rememberMessageTurn(context.turnIdByMessageId, event.properties.info.id, turnId); - } - 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); - } const usage = context.turnTokenUsage; const parentMessageId = typeof event.properties.info.parentID === "string" && 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" @@ -2475,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() ?? []) { @@ -2486,6 +2504,7 @@ 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); @@ -2529,6 +2548,7 @@ export function makeOpenCodeAdapter( threadId: context.session.threadId, turnId, itemId: event.properties.partID, + providerItemId: event.properties.messageID, raw: event, })), type: "content.delta", @@ -2544,9 +2564,6 @@ export function makeOpenCodeAdapter( const part = event.properties.part; const messageRole = messageRoleForPart(context, part); - if (turnId) { - rememberMessageTurn(context.turnIdByMessageId, part.messageID, turnId); - } const actualModel = actualModelFromPart(part); if (actualModel) { context.actualModelByMessageId.delete(part.messageID); @@ -2818,6 +2835,7 @@ export function makeOpenCodeAdapter( ...(yield* buildEventBase({ threadId: context.session.threadId, turnId: activeTurnId, + providerItemId: actualModel?.messageId, raw: event, })), type: "turn.completed", @@ -2825,7 +2843,7 @@ export function makeOpenCodeAdapter( state: "failed", errorMessage: message, tokenUsage, - ...(actualModel ? { actualModel } : {}), + ...(actualModel ? { actualModel: actualModel.model } : {}), }, }); } @@ -3146,6 +3164,7 @@ export function makeOpenCodeAdapter( pendingQuestions: new Map(), textPartsByMessageId: new Map(), messageRoleById: new Map(), + turnIdByPromptMessageId: new Map(), turnIdByMessageId: new Map(), actualModelByMessageId: new Map(), turnTokenUsage: undefined, @@ -3337,6 +3356,14 @@ export function makeOpenCodeAdapter( 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; From bba98aa704ed190bc2a5a3d603ad9a6e30dd495b Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Tue, 8 Sep 2026 14:57:14 +0200 Subject: [PATCH 4/9] fix(opencode): retain reordered model metadata --- .../Layers/ProjectionSnapshotQuery.test.ts | 6 +- .../Layers/ProjectionSnapshotQuery.ts | 2 + .../Layers/ProviderRuntimeIngestion.test.ts | 61 +++++++++++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 53 ++++++++++++++++ 4 files changed, 120 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts index 6183dbc66168..200c6e6bcb81 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.test.ts @@ -815,9 +815,10 @@ projectionSnapshotLayer("ProjectionSnapshotQuery", (it) => { `; 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 23ebfd583d3d..fc000c6d6a70 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -1297,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", @@ -3283,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 } : {}), diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 2b409db26fa5..08970ab4019d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3906,6 +3906,67 @@ describe("ProviderRuntimeIngestion", () => { 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("applies late actual model metadata only to its originating assistant message", async () => { const harness = await createHarness(); const now = "2026-08-14T00:00:00.000Z"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index e76692013d34..22159a477e4f 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -107,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; @@ -1043,6 +1048,12 @@ const make = Effect.gen(function* () { 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, @@ -1670,6 +1681,7 @@ const make = Effect.gen(function* () { 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)); @@ -1700,6 +1712,12 @@ const make = Effect.gen(function* () { : 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) => @@ -2321,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(); @@ -2345,6 +2378,7 @@ const make = Effect.gen(function* () { Option.isNone(activeAssistantMessageId) && turnId !== undefined && hasAssistantMessagesForTurn && + matchingPendingActualModel === undefined && (assistantCompletion.fallbackText?.trim().length ?? 0) === 0; if (!shouldSkipRedundantCompletion) { @@ -2361,6 +2395,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 } : {}), @@ -2368,6 +2403,12 @@ const make = Effect.gen(function* () { if (turnId) { yield* forgetAssistantMessageId(thread.id, turnId, assistantMessageId); + if (matchingPendingActualModel) { + yield* Cache.invalidate( + pendingActualModelByTurnKey, + providerTurnKey(thread.id, turnId), + ); + } } } @@ -2461,6 +2502,12 @@ const make = Effect.gen(function* () { : completedAssistantMessageId ? [completedAssistantMessageId] : []; + if (terminalActualModel && assistantMessageIds.length === 0) { + yield* Cache.set(pendingActualModelByTurnKey, providerTurnKey(thread.id, turnId), { + actualModel: terminalActualModel, + ...(terminalProviderItemId ? { providerItemId: terminalProviderItemId } : {}), + }); + } const terminalAssistantMessageId = assistantMessageIds.at(-1); yield* Effect.forEach( assistantMessageIds, @@ -2488,6 +2535,12 @@ const make = Effect.gen(function* () { ), { 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"); From a0e272cde7f65f44ee7fe49f3eeb4bb5c70941c6 Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Tue, 8 Sep 2026 18:49:38 +0200 Subject: [PATCH 5/9] fix(opencode): map non-streamed responses --- .../Layers/ProviderRuntimeIngestion.test.ts | 65 +++++++++++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 8 +++ 2 files changed, 73 insertions(+) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 08970ab4019d..82c0cd65dce0 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -3906,6 +3906,71 @@ describe("ProviderRuntimeIngestion", () => { 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"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 22159a477e4f..1b4819e8cef2 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -2385,6 +2385,14 @@ const make = Effect.gen(function* () { if (turnId && Option.isNone(activeAssistantMessageId)) { yield* rememberAssistantMessageId(thread.id, turnId, assistantMessageId); } + if (turnId && completionProviderItemId) { + yield* rememberProviderAssistantMessageId( + thread.id, + turnId, + completionProviderItemId, + assistantMessageId, + ); + } yield* finalizeAssistantMessage({ event, From 7fceedc80aded8b4b2d86ea06f057aeee07c64de Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Tue, 8 Sep 2026 19:13:29 +0200 Subject: [PATCH 6/9] fix(opencode): require keyed model deferral --- .../Layers/ProviderRuntimeIngestion.test.ts | 56 +++++++++++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 4 +- 2 files changed, 58 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 82c0cd65dce0..9b994131f9af 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -4032,6 +4032,62 @@ describe("ProviderRuntimeIngestion", () => { 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"; diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts index 1b4819e8cef2..b292b612ceab 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -2510,10 +2510,10 @@ const make = Effect.gen(function* () { : completedAssistantMessageId ? [completedAssistantMessageId] : []; - if (terminalActualModel && assistantMessageIds.length === 0) { + if (terminalActualModel && terminalProviderItemId && assistantMessageIds.length === 0) { yield* Cache.set(pendingActualModelByTurnKey, providerTurnKey(thread.id, turnId), { actualModel: terminalActualModel, - ...(terminalProviderItemId ? { providerItemId: terminalProviderItemId } : {}), + providerItemId: terminalProviderItemId, }); } const terminalAssistantMessageId = assistantMessageIds.at(-1); From 9c2ad48128e7c746dc82d6656c944ddc70177f94 Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Tue, 8 Sep 2026 19:59:04 +0200 Subject: [PATCH 7/9] fix(opencode): merge completion targets --- .../Layers/ProviderRuntimeIngestion.test.ts | 90 +++++++++++++++++++ .../Layers/ProviderRuntimeIngestion.ts | 12 ++- 2 files changed, 95 insertions(+), 7 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts index 9b994131f9af..21be313da71b 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.test.ts @@ -4169,6 +4169,96 @@ describe("ProviderRuntimeIngestion", () => { 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 b292b612ceab..076ac65e3c8d 100644 --- a/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts +++ b/apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts @@ -2503,13 +2503,11 @@ const make = Effect.gen(function* () { ) : undefined; const assistantMessageIds = - trackedAssistantMessageIds.size > 0 - ? Array.from(trackedAssistantMessageIds) - : modelTargetMessageIds.size > 0 - ? Array.from(modelTargetMessageIds) - : completedAssistantMessageId - ? [completedAssistantMessageId] - : []; + 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, From 0d9cfa2e78162990a7e8c62cb4b8cba933d6b11b Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Tue, 15 Sep 2026 09:21:38 +0200 Subject: [PATCH 8/9] test(opencode): target renumbered model migration --- .../Migrations/053_ProjectionThreadMessageActualModel.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts index df0352415903..c9bea56a6310 100644 --- a/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts +++ b/apps/server/src/persistence/Migrations/053_ProjectionThreadMessageActualModel.test.ts @@ -13,7 +13,7 @@ layer("053_ProjectionThreadMessageActualModel", (it) => { const sql = yield* SqlClient.SqlClient; yield* runMigrations({ toMigrationInclusive: 49 }); - yield* runMigrations({ toMigrationInclusive: 52 }); + yield* runMigrations({ toMigrationInclusive: 53 }); const columns = yield* sql<{ readonly name: string; readonly notnull: number }>` PRAGMA table_info(projection_thread_messages) From bdf218769d65c4256c1f85629e62931d1af545b5 Mon Sep 17 00:00:00 2001 From: Giuseppe Crescimanno Date: Thu, 17 Sep 2026 05:23:44 +0200 Subject: [PATCH 9/9] fix(opencode): narrow completion model metadata --- apps/server/src/orchestration/decider.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 2a88a79987c2..5c10d1241142 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -1972,7 +1972,9 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" messageId: command.messageId, role: command.type === "thread.message.reasoning.complete" ? "reasoning" : "assistant", text: "", - ...(command.actualModel ? { actualModel: command.actualModel } : {}), + ...(command.type === "thread.message.assistant.complete" && command.actualModel + ? { actualModel: command.actualModel } + : {}), turnId: command.turnId ?? null, streaming: false, createdAt: command.createdAt,