From d3731441e97a7e799072dd9a7a6a669bce50961c Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 2 Oct 2026 20:28:42 -0700 Subject: [PATCH 1/2] fix(server): Claude subagent threads show the model their agent file picks Co-Authored-By: Claude Opus 5.5 (1M context) --- .../features/threads/ThreadDetailScreen.tsx | 12 ++- .../ProviderEventIngestor.test.ts | 95 +++++++++++++++++++ .../orchestration-v2/ProviderEventIngestor.ts | 48 ++++++++-- apps/web/src/components/ChatView.tsx | 11 ++- 4 files changed, 157 insertions(+), 9 deletions(-) diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index d1470315a566..19dea22e9e66 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -35,7 +35,7 @@ import { formatModelSelectionEffort, type ProviderSubagentStatus, } from "@t3tools/client-runtime/state/thread-execution"; -import { formatModelSlugName } from "@t3tools/shared/model"; +import { formatModelSlugName, resolveSelectableModel } from "@t3tools/shared/model"; import { isProviderNativeSubagentThread } from "@t3tools/contracts"; import type { QueuedRunEdit } from "../../state/queued-run-edit"; import type { FollowUpBehavior } from "../../lib/followUpBehavior"; @@ -760,8 +760,16 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread const providerSubagentProvider = props.serverConfig?.providers.find( (provider) => provider.instanceId === props.selectedThread.modelSelection.instanceId, ); + // Providers can report a dated id or alias (claude-haiku-4-5-20251001). + const providerSubagentModelSlug = providerSubagentProvider + ? resolveSelectableModel( + providerSubagentProvider.driver, + props.selectedThread.modelSelection.model, + providerSubagentProvider.models, + ) + : null; const providerSubagentCatalogModel = providerSubagentProvider?.models.find( - (model) => model.slug === props.selectedThread.modelSelection.model, + (model) => model.slug === providerSubagentModelSlug, ); const workspaceContentWidth = useWorkspaceContentWidth(); // Clearing animated width can retain the unfolded width after Android resumes folded. diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts index 7b73b7d7f1d9..8d07d537f890 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts @@ -1234,4 +1234,99 @@ layer("ProviderEventIngestorV2", (it) => { assert.equal(messageEvents[0]?.threadId, childThreadId); }), ); + + it.effect("moves a native subagent's thread to the model its provider reports later", () => + Effect.gen(function* () { + const now = yield* DateTime.now; + const eventSink = yield* EventSink.EventSinkV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const ingestor = yield* ProviderEventIngestor.ProviderEventIngestorV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const rootEvent = yield* threadCreatedEvent(now); + if (rootEvent.type !== "thread.created") { + throw new Error("Expected a thread.created fixture event"); + } + const childThreadId = idAllocator.derive.threadFromProviderThread({ + driver: CODEX_DRIVER, + nativeThreadId: "native-late-model-subagent", + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId: rootEvent.threadId, + }); + const ingest = (event: ProviderEventIngestor.ProviderEventIngestInput["event"]) => + ingestor.ingestNormalized({ + providerSessionId, + providerInstanceId: modelSelection.instanceId, + threadId: rootEvent.threadId, + event, + }); + yield* eventSink.write({ events: [rootEvent] }); + // The subagent's thread starts on the parent's model and options. + yield* ingest({ + type: "app_thread.created", + driver: CODEX_DRIVER, + appThread: { + ...rootEvent.payload, + id: childThreadId, + title: "review design", + modelSelection: { + ...modelSelection, + options: [{ id: "reasoningEffort", value: "xhigh" }], + }, + activeProviderThreadId: null, + lineage: { + parentThreadId: rootEvent.threadId, + relationshipToParent: "subagent", + rootThreadId: rootEvent.threadId, + }, + }, + }); + const subagentUpdated = { + type: "subagent.updated", + driver: CODEX_DRIVER, + subagent: { + id: NodeId.make("node:late-model-subagent"), + threadId: rootEvent.threadId, + runId: null, + parentNodeId: NodeId.make("node:root"), + origin: "provider_native", + createdBy: "agent", + driver: CODEX_DRIVER, + providerInstanceId: modelSelection.instanceId, + providerThreadId: null, + childThreadId, + nativeTaskRef: null, + prompt: "Review the design", + title: "review design", + model: "gpt-6.1-sol", + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + } satisfies ProviderEventIngestor.ProviderEventIngestInput["event"]; + + const first = yield* ingest(subagentUpdated); + const repeated = yield* ingest(subagentUpdated); + const childThread = yield* projectionStore.getThread(childThreadId); + + assert.deepEqual( + first.map((stored) => [stored.event.type, stored.event.threadId]), + [ + ["subagent.updated", rootEvent.threadId], + ["thread.model-selection-updated", childThreadId], + ], + ); + assert.deepEqual( + repeated.map((stored) => stored.event.type), + ["subagent.updated"], + ); + assert.deepEqual(childThread.modelSelection, { + instanceId: modelSelection.instanceId, + model: "gpt-6.1-sol", + }); + }), + ); }); diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts index 8082962b86de..c6f2e8cc1182 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts @@ -392,16 +392,52 @@ export const layer: Layer.Layer< nodeId: input.event.node.id, }), ]; - case "subagent.updated": + case "subagent.updated": { + const subagent = input.event.subagent; + const subagentEvent = yield* makeDomainEvent(input, { + type: "subagent.updated", + threadId: subagent.threadId, + payload: subagent, + runId: subagent.runId, + nodeId: subagent.id, + }); + // A native subagent's thread starts on the parent's model when the + // provider names the real one later (a Claude agent file's model + // arrives with its first reply). Clients read the thread's model. + if ( + subagent.origin !== "provider_native" || + subagent.childThreadId === null || + subagent.model === null + ) { + return [subagentEvent]; + } + const childThread = yield* projections + .getThread(subagent.childThreadId) + .pipe( + Effect.catchTag("ProjectionStoreThreadNotFoundError", () => Effect.succeed(null)), + ); + if (childThread === null || childThread.modelSelection.model === subagent.model) { + return [subagentEvent]; + } + const now = yield* DateTime.now; return [ + subagentEvent, yield* makeDomainEvent(input, { - type: "subagent.updated", - threadId: input.event.subagent.threadId, - payload: input.event.subagent, - runId: input.event.subagent.runId, - nodeId: input.event.subagent.id, + type: "thread.model-selection-updated", + threadId: childThread.id, + // The parent's options belong to the parent's model. + payload: { + ...childThread, + modelSelection: { + instanceId: childThread.modelSelection.instanceId, + model: subagent.model, + }, + updatedAt: now, + }, + occurredAt: now, }), ]; + } case "message.updated": return [ yield* makeDomainEvent(input, { diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 8fbe528b34b4..b5a96518892c 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -116,6 +116,7 @@ import { createModelSelection, formatModelSlugName, resolvePromptInjectedEffort, + resolveSelectableModel, } from "@t3tools/shared/model"; import { projectScriptCwd, @@ -4055,8 +4056,16 @@ export default function ChatView(props: ChatViewProps) { const showProviderSubagentBar = isProviderSubagent; const composerMounted = !showProviderSubagentBar; const providerSubagentModels = selectedProviderEntry?.models ?? EMPTY_PROVIDER_MODELS; + // Providers can report a dated id or alias (claude-haiku-4-5-20251001). + const providerSubagentModelSlug = selectedProviderEntry + ? resolveSelectableModel( + selectedProviderEntry.driverKind, + activeThread?.modelSelection.model, + providerSubagentModels, + ) + : null; const providerSubagentCatalogModel = providerSubagentModels.find( - (model) => model.slug === activeThread?.modelSelection.model, + (model) => model.slug === providerSubagentModelSlug, ); const providerSubagentModelLabel = providerSubagentCatalogModel ? getTriggerDisplayModelName(providerSubagentCatalogModel) From 37572be2df61885ef9c4ea4822ba2be8f644da33 Mon Sep 17 00:00:00 2001 From: Theo Browne Date: Fri, 2 Oct 2026 20:37:50 -0700 Subject: [PATCH 2/2] fix(server): sync a subagent thread's model under its thread lock Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ProviderEventIngestor.test.ts | 11 +- .../orchestration-v2/ProviderEventIngestor.ts | 113 +++++++++++------- .../ProviderSessionManager.test.ts | 10 +- .../src/orchestration-v2/runtimeLayer.ts | 10 +- .../testkit/ProviderReplayHarness.ts | 10 +- 5 files changed, 107 insertions(+), 47 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts index 8d07d537f890..6bbfecf31d85 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts @@ -33,6 +33,7 @@ import * as EventStore from "./EventStore.ts"; import * as IdAllocator from "./IdAllocator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; +import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; import { makeProviderEventRoutingState, @@ -54,8 +55,16 @@ const TestLayer = Layer.mergeAll( TestStoresLayer, TestEventSinkLayer, IdAllocator.layer, + ThreadCommandExecutor.layer, ProviderEventIngestor.layer.pipe( - Layer.provide(Layer.mergeAll(TestStoresLayer, TestEventSinkLayer, IdAllocator.layer)), + Layer.provide( + Layer.mergeAll( + TestStoresLayer, + TestEventSinkLayer, + IdAllocator.layer, + ThreadCommandExecutor.layer, + ), + ), ), ); const modelSelection = { diff --git a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts index c6f2e8cc1182..01a34547a237 100644 --- a/apps/server/src/orchestration-v2/ProviderEventIngestor.ts +++ b/apps/server/src/orchestration-v2/ProviderEventIngestor.ts @@ -6,6 +6,7 @@ import { type OrchestrationV2PlanArtifact, type OrchestrationV2Run, type OrchestrationV2ProviderTurn, + type OrchestrationV2Subagent, type ModelSelection, type RuntimeMode, type ProviderInteractionMode, @@ -31,6 +32,7 @@ import * as ProjectionStore from "./ProjectionStore.ts"; import * as IdAllocator from "./IdAllocator.ts"; import { ProviderAdapterV2Event } from "./ProviderAdapter.ts"; import { makeProviderFailureTurnItem } from "./ProviderFailure.ts"; +import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; export class ProviderEventNormalizeError extends Schema.TaggedError()( "ProviderEventNormalizeError", @@ -249,13 +251,17 @@ const decodeDomainEvent = Schema.decodeUnknownEffect(OrchestrationV2DomainEvent) export const layer: Layer.Layer< ProviderEventIngestorV2, never, - EventSink.EventSinkV2 | IdAllocator.IdAllocatorV2 | ProjectionStore.ProjectionStoreV2 + | EventSink.EventSinkV2 + | IdAllocator.IdAllocatorV2 + | ProjectionStore.ProjectionStoreV2 + | ThreadCommandExecutor.ThreadCommandExecutor > = Layer.effect( ProviderEventIngestorV2, Effect.gen(function* () { const eventSink = yield* EventSink.EventSinkV2; const projections = yield* ProjectionStore.ProjectionStoreV2; const idAllocator = yield* IdAllocator.IdAllocatorV2; + const threadCommands = yield* ThreadCommandExecutor.ThreadCommandExecutor; const analytics = yield* ProviderTurnAnalytics; const completedTurnAnalytics = new Set(); @@ -338,6 +344,48 @@ export const layer: Layer.Layer< }, ); + /** + * A native subagent's thread starts on the parent's model when the + * provider names the real one later (a Claude agent file's model arrives + * with the subagent's first reply). Clients read the thread's model, so + * move the thread to the reported one. Thread commands rewrite the whole + * thread row under the thread's lock, so this read and write take it too. + */ + const syncSubagentThreadModel = Effect.fn("ProviderEventIngestor.syncSubagentThreadModel")( + function* (input: ProviderEventIngestInput, subagent: OrchestrationV2Subagent) { + const { childThreadId, model } = subagent; + if (subagent.origin !== "provider_native" || childThreadId === null || model === null) { + return []; + } + const staleThread = projections.getThread(childThreadId).pipe( + Effect.map((thread) => (thread.modelSelection.model === model ? null : thread)), + Effect.catchTags({ ProjectionStoreThreadNotFoundError: () => Effect.succeed(null) }), + ); + // Nearly every update already matches; only a mismatch takes the lock. + if ((yield* staleThread) === null) return []; + return yield* threadCommands.withLock( + childThreadId, + Effect.gen(function* () { + const thread = yield* staleThread; + if (thread === null) return []; + const now = yield* DateTime.now; + const event = yield* makeDomainEvent(input, { + type: "thread.model-selection-updated", + threadId: thread.id, + // The parent's options belong to the parent's model. + payload: { + ...thread, + modelSelection: { instanceId: thread.modelSelection.instanceId, model }, + updatedAt: now, + }, + occurredAt: now, + }); + return yield* eventSink.write({ events: [event] }); + }), + ); + }, + ); + const normalize: ProviderEventIngestorV2Shape["normalize"] = (input) => Effect.gen(function* () { switch (input.event.type) { @@ -392,52 +440,16 @@ export const layer: Layer.Layer< nodeId: input.event.node.id, }), ]; - case "subagent.updated": { - const subagent = input.event.subagent; - const subagentEvent = yield* makeDomainEvent(input, { - type: "subagent.updated", - threadId: subagent.threadId, - payload: subagent, - runId: subagent.runId, - nodeId: subagent.id, - }); - // A native subagent's thread starts on the parent's model when the - // provider names the real one later (a Claude agent file's model - // arrives with its first reply). Clients read the thread's model. - if ( - subagent.origin !== "provider_native" || - subagent.childThreadId === null || - subagent.model === null - ) { - return [subagentEvent]; - } - const childThread = yield* projections - .getThread(subagent.childThreadId) - .pipe( - Effect.catchTag("ProjectionStoreThreadNotFoundError", () => Effect.succeed(null)), - ); - if (childThread === null || childThread.modelSelection.model === subagent.model) { - return [subagentEvent]; - } - const now = yield* DateTime.now; + case "subagent.updated": return [ - subagentEvent, yield* makeDomainEvent(input, { - type: "thread.model-selection-updated", - threadId: childThread.id, - // The parent's options belong to the parent's model. - payload: { - ...childThread, - modelSelection: { - instanceId: childThread.modelSelection.instanceId, - model: subagent.model, - }, - updatedAt: now, - }, - occurredAt: now, + type: "subagent.updated", + threadId: input.event.subagent.threadId, + payload: input.event.subagent, + runId: input.event.subagent.runId, + nodeId: input.event.subagent.id, }), ]; - } case "message.updated": return [ yield* makeDomainEvent(input, { @@ -579,6 +591,21 @@ export const layer: Layer.Layer< .pipe(Effect.mapError(mapWriteError)); return result.storedEvents; }).pipe( + Effect.flatMap((storedEvents) => + storedEvents.length === 0 || input.event.type !== "subagent.updated" + ? Effect.succeed(storedEvents) + : syncSubagentThreadModel(input, input.event.subagent).pipe( + Effect.map((synced) => [...storedEvents, ...synced]), + Effect.mapError( + (cause) => + new ProviderEventPublishError({ + providerSessionId: input.providerSessionId, + eventCount: 1, + cause, + }), + ), + ), + ), Effect.tap((storedEvents) => Effect.gen(function* () { if (storedEvents.length === 0 || input.event.type !== "provider_turn.updated") return; diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 46c4695bad1d..a567a4ab4841 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -52,6 +52,7 @@ import { } from "./ProviderAdapter.ts"; import * as ProviderAdapterRegistry from "./ProviderAdapterRegistry.ts"; import * as ProviderEventIngestor from "./ProviderEventIngestor.ts"; +import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; import * as ProviderSessionManager from "./ProviderSessionManager.ts"; const TestDatabaseLayer = SqlitePersistenceMemory; @@ -383,7 +384,14 @@ function makeTestLayer(input: { }), ); const providerEventIngestorTestLayer = ProviderEventIngestor.layer.pipe( - Layer.provide(Layer.mergeAll(configuredEventSinkLayer, IdAllocator.layer, TestStoresLayer)), + Layer.provide( + Layer.mergeAll( + configuredEventSinkLayer, + IdAllocator.layer, + TestStoresLayer, + ThreadCommandExecutor.layer, + ), + ), ); return Layer.mergeAll( TestStoresLayer, diff --git a/apps/server/src/orchestration-v2/runtimeLayer.ts b/apps/server/src/orchestration-v2/runtimeLayer.ts index 7ecbf57f204a..8573c6570795 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.ts @@ -35,6 +35,7 @@ import { layer as providerContinuationRequestsLayer } from "./ProviderContinuati import { workerLive as providerContinuationWorkerLive } from "./ProviderContinuationService.ts"; import { layer as threadTitleRegenerationServiceLayer } from "./ThreadTitleRegenerationService.ts"; import { layer as providerEventIngestorLayer } from "./ProviderEventIngestor.ts"; +import * as ThreadCommandExecutor from "./ThreadCommandExecutor.ts"; import { layer as providerSessionManagerLayer } from "./ProviderSessionManager.ts"; import { layer as providerRuntimeRecoveryLayer } from "./ProviderRuntimeRecoveryService.ts"; import { layer as providerSwitchServiceLayer } from "./ProviderSwitchService.ts"; @@ -98,7 +99,14 @@ export const ProjectServiceLayerLive = projectServiceLayer.pipe( ); const providerEventIngestorProvided = providerEventIngestorLayer.pipe( - Layer.provide(Layer.mergeAll(eventSinkProvided, idAllocatorLayer, projectionStoreLayer)), + Layer.provide( + Layer.mergeAll( + eventSinkProvided, + idAllocatorLayer, + projectionStoreLayer, + ThreadCommandExecutor.layer, + ), + ), ); const checkpointServiceProvided = checkpointServiceLayer.pipe(Layer.provide(idAllocatorLayer)); diff --git a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts index 2fa1345528f9..f4a1201f1f3e 100644 --- a/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts +++ b/apps/server/src/orchestration-v2/testkit/ProviderReplayHarness.ts @@ -50,6 +50,7 @@ import * as ThreadTitleRegenerationService from "../ThreadTitleRegenerationServi import * as RuntimePolicy from "../RuntimePolicy.ts"; import * as TurnItemPositionStore from "../TurnItemPositionStore.ts"; import * as RuntimeRequestService from "../RuntimeRequestService.ts"; +import * as ThreadCommandExecutor from "../ThreadCommandExecutor.ts"; import * as ThreadForkService from "../ThreadForkService.ts"; import { runOrchestratorV2Scenario, @@ -301,7 +302,14 @@ export function makeOrchestratorV2ReplayLayerWithRegistry( ); const commandReceiptStoreProvided = CommandReceiptStore.layer.pipe(Layer.provide(databaseLayer)); const providerEventIngestorProvided = ProviderEventIngestor.layer.pipe( - Layer.provide(Layer.mergeAll(storesLayer, eventSinkProvided, IdAllocator.layer)), + Layer.provide( + Layer.mergeAll( + storesLayer, + eventSinkProvided, + IdAllocator.layer, + ThreadCommandExecutor.layer, + ), + ), ); const vcsDriverRegistryLayer = VcsDriverRegistry.layer.pipe( Layer.provide(VcsProcess.layer),