Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions apps/mobile/src/features/threads/ThreadDetailScreen.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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.
Expand Down
106 changes: 105 additions & 1 deletion apps/server/src/orchestration-v2/ProviderEventIngestor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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 = {
Expand Down Expand Up @@ -1234,4 +1243,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",
});
}),
);
});
65 changes: 64 additions & 1 deletion apps/server/src/orchestration-v2/ProviderEventIngestor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import {
type OrchestrationV2PlanArtifact,
type OrchestrationV2Run,
type OrchestrationV2ProviderTurn,
type OrchestrationV2Subagent,
type ModelSelection,
type RuntimeMode,
type ProviderInteractionMode,
Expand All @@ -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>()(
"ProviderEventNormalizeError",
Expand Down Expand Up @@ -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<string>();

Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -543,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;
Expand Down
10 changes: 9 additions & 1 deletion apps/server/src/orchestration-v2/ProviderSessionManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
10 changes: 9 additions & 1 deletion apps/server/src/orchestration-v2/runtimeLayer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -301,7 +302,14 @@ export function makeOrchestratorV2ReplayLayerWithRegistry<Error>(
);
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),
Expand Down
11 changes: 10 additions & 1 deletion apps/web/src/components/ChatView.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,7 @@ import {
createModelSelection,
formatModelSlugName,
resolvePromptInjectedEffort,
resolveSelectableModel,
} from "@t3tools/shared/model";
import {
projectScriptCwd,
Expand Down Expand Up @@ -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)
Expand Down
Loading