diff --git a/apps/server/src/orchestration-v2/EffectWorker.ts b/apps/server/src/orchestration-v2/EffectWorker.ts index 48dcfa78e96b..1abf070b87ec 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.ts @@ -244,7 +244,12 @@ export const executorLayer: Layer.Layer< text: message.text, ...(message.context ? { context: message.context } : {}), attachments: message.attachments, - modelSelection: run.modelSelection, + // A user's follow-up starts on the thread's saved selection, + // which already holds the steer's choice. A delegated + // completion stays pinned to the run it reports to. + ...(message.delegatedCompletion === undefined + ? {} + : { modelSelection: run.modelSelection }), dispatchMode: { type: message.delegatedCompletion === undefined diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index df9805d504c9..c58dd46609ec 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -3631,17 +3631,25 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ); if (steeringPolicy === "active_steering") { + // The steer's selection becomes the saved next-turn choice, even when it + // matches the running run again. A delegated completion carries the + // run's selection, not a user choice, so it never replaces the saved one. + // The saved choice may name another instance, so it moves with the steer. + const instanceChanged = + input.projection.thread.providerInstanceId !== input.modelSelection.instanceId; if ( - selectionChanged && - !modelSelectionsEqual(input.projection.thread.modelSelection, input.modelSelection) + input.delegatedCompletion === undefined && + (instanceChanged || + !modelSelectionsEqual(input.projection.thread.modelSelection, input.modelSelection)) ) { yield* emitEvent({ - type: "thread.model-selection-updated", + type: instanceChanged ? "thread.provider-switched" : "thread.model-selection-updated", threadId: input.command.threadId, providerInstanceId: input.modelSelection.instanceId, occurredAt: now, payload: { ...input.projection.thread, + providerInstanceId: input.modelSelection.instanceId, modelSelection: input.modelSelection, updatedAt: now, }, diff --git a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts index 9f9f17dedfe2..6a360dd8b6b4 100644 --- a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts +++ b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts @@ -10,6 +10,7 @@ import { ProviderThreadId, ProviderTurnId, ThreadId, + type ModelSelection, type OrchestrationV2DomainEvent, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; @@ -402,183 +403,300 @@ for (const mailbox of [false, true]) { } } -// Claude steers live but cannot interrupt-and-restart. A composer whose model -// options differ from the running run (an explicit `fastMode: false`) must still -// steer, keeping the new selection for the next turn. -it.effect("steers a changed turn-scoped selection into a provider that cannot restart", () => - Effect.scoped( - Effect.gen(function* () { - const cwd = yield* checkpointWorkspace("steering-selection-change"); - const events = yield* Queue.unbounded(); - const started: ProviderAdapterV2TurnInput[] = []; - const steered: string[] = []; - const capabilities = { - ...CodexProviderCapabilitiesV2, - turns: { - ...CodexProviderCapabilitiesV2.turns, - supportsActiveSteering: true, - supportsSteeringByInterruptRestart: false, - }, - }; - const runSelection = { - instanceId, - model: "test-model", - options: [{ id: "effort", value: "xhigh" }], - }; - const composerSelection = { - ...runSelection, - options: [...runSelection.options, { id: "fastMode", value: false }], - }; - const adapter: ProviderAdapterV2Shape = { - instanceId, - driver, - getCapabilities: () => Effect.succeed(capabilities), - planSelectionTransition: () => Effect.succeed({ type: "apply_on_next_turn" }), - openSession: (input) => - Effect.gen(function* () { - const now = yield* DateTime.now; - return { - instanceId, +// Claude and Pi steer live but cannot interrupt-and-restart, so a changed +// selection waits for the next turn as the thread's saved selection. +const runSelection = { + instanceId, + model: "test-model", + options: [{ id: "effort", value: "xhigh" }], +}; +const composerSelection = { + ...runSelection, + options: [...runSelection.options, { id: "fastMode", value: false }], +}; + +const nextTurnSelectionHarness = Effect.fn("nextTurnSelectionHarness")(function* (name: string) { + const cwd = yield* checkpointWorkspace(name); + const events = yield* Queue.unbounded(); + const started: ProviderAdapterV2TurnInput[] = []; + const steered: string[] = []; + const capabilities = { + ...CodexProviderCapabilitiesV2, + turns: { + ...CodexProviderCapabilitiesV2.turns, + supportsActiveSteering: true, + supportsSteeringByInterruptRestart: false, + }, + }; + const adapter: ProviderAdapterV2Shape = { + instanceId, + driver, + getCapabilities: () => Effect.succeed(capabilities), + planSelectionTransition: () => Effect.succeed({ type: "apply_on_next_turn" }), + openSession: (input) => + Effect.gen(function* () { + const now = yield* DateTime.now; + return { + instanceId, + driver, + providerSessionId: input.providerSessionId, + providerSession: { + id: input.providerSessionId, + driver, + providerInstanceId: instanceId, + status: "ready", + cwd, + model: runSelection.model, + capabilities, + createdAt: now, + updatedAt: now, + lastError: null, + }, + events: Stream.fromQueue(events), + ensureThread: ({ threadId }) => + Effect.succeed({ + id: ProviderThreadId.make(`provider-thread:${threadId}`), driver, + providerInstanceId: instanceId, providerSessionId: input.providerSessionId, - providerSession: { - id: input.providerSessionId, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: { driver, nativeId: "native-thread", strength: "strong" }, + nativeConversationHeadRef: null, + status: "idle", + firstRunOrdinal: null, + lastRunOrdinal: null, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }), + resumeThread: ({ providerThread }) => Effect.succeed(providerThread), + startTurn: (turn) => + Effect.gen(function* () { + started.push(turn); + yield* Queue.offer(events, { + type: "provider_turn.updated", driver, - providerInstanceId: instanceId, - status: "ready", - cwd, - model: runSelection.model, - capabilities, - createdAt: now, - updatedAt: now, - lastError: null, - }, - events: Stream.fromQueue(events), - ensureThread: ({ threadId }) => - Effect.succeed({ - id: ProviderThreadId.make(`provider-thread:${threadId}`), - driver, - providerInstanceId: instanceId, - providerSessionId: input.providerSessionId, - appThreadId: threadId, - ownerNodeId: null, - nativeThreadRef: { driver, nativeId: "native-thread", strength: "strong" }, - nativeConversationHeadRef: null, - status: "idle", - firstRunOrdinal: null, - lastRunOrdinal: null, - handoffIds: [], - forkedFrom: null, - createdAt: now, - updatedAt: now, - }), - resumeThread: ({ providerThread }) => Effect.succeed(providerThread), - startTurn: (turn) => - Effect.gen(function* () { - started.push(turn); - yield* Queue.offer(events, { - type: "provider_turn.updated", + providerTurn: { + id: ProviderTurnId.make(`provider-turn:${turn.attemptId}`), + providerThreadId: turn.providerThread.id, + nodeId: turn.rootNodeId, + runAttemptId: turn.attemptId, + nativeTurnRef: { driver, - providerTurn: { - id: ProviderTurnId.make(`provider-turn:${turn.attemptId}`), - providerThreadId: turn.providerThread.id, - nodeId: turn.rootNodeId, - runAttemptId: turn.attemptId, - nativeTurnRef: { - driver, - nativeId: `native:${turn.attemptId}`, - strength: "strong", - }, - ordinal: turn.providerTurnOrdinal, - status: "running", - startedAt: now, - completedAt: null, - }, - }); - }), - steerTurn: (turn) => - Effect.sync(() => { - steered.push(turn.message.text); - }), - interruptTurn: () => Effect.void, - respondToRuntimeRequest: () => Effect.void, - readThreadSnapshot: () => Effect.die("unused"), - rollbackThread: () => Effect.die("unused"), - forkThread: () => Effect.die("unused"), - }; - }), - }; + nativeId: `native:${turn.attemptId}`, + strength: "strong", + }, + ordinal: turn.providerTurnOrdinal, + status: "running", + startedAt: now, + completedAt: null, + }, + }); + }), + steerTurn: (turn) => + Effect.sync(() => { + steered.push(turn.message.text); + }), + interruptTurn: () => Effect.void, + respondToRuntimeRequest: () => Effect.void, + readThreadSnapshot: () => Effect.die("unused"), + rollbackThread: () => Effect.die("unused"), + forkThread: () => Effect.die("unused"), + }; + }), + }; + const layer = makeOrchestratorV2ReplayLayerWithRegistry( + { name }, + ProviderAdapterRegistry.makeSingleLayer(adapter), + { runEffectWorker: false }, + ); + // Creates the thread and starts its first turn on `runSelection`. + const startFirstTurn = Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const threadId = ThreadId.make(`thread:${name}`); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create"), + threadId, + projectId: ProjectId.make(`project:${name}`), + title: "Steer with changed options", + modelSelection: runSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + createdBy: "user", + creationSource: "web", + }); + const running = yield* orchestrator.streamDomainEvents.pipe( + Stream.filter( + (event) => event.type === "provider-turn.updated" && event.payload.status === "running", + ), + Stream.take(1), + Stream.runDrain, + Effect.forkScoped, + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("first"), + threadId, + messageId: MessageId.make("message:first"), + text: "first", + attachments: [], + modelSelection: runSelection, + dispatchMode: { type: "start_immediately" }, + createdBy: "user", + creationSource: "web", + }); + yield* worker.drain(); + yield* Fiber.join(running); + return threadId; + }); + return { events, started, steered, layer, startFirstTurn }; +}); + +it.effect("steers a changed turn-scoped selection into a provider that cannot restart", () => + Effect.scoped( + Effect.gen(function* () { + const { started, steered, layer, startFirstTurn } = yield* nextTurnSelectionHarness( + "steering-selection-change", + ); yield* Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; - const threadId = ThreadId.make("thread:steering-selection-change"); + const threadId = yield* startFirstTurn; + const steer = (id: string, selection: ModelSelection) => + orchestrator + .dispatch({ + type: "message.dispatch", + commandId: CommandId.make(id), + threadId, + messageId: MessageId.make(`message:${id}`), + text: id, + attachments: [], + modelSelection: selection, + dispatchMode: { type: "steer_active", targetRunId: started[0]!.runId }, + createdBy: "user", + creationSource: "web", + }) + .pipe(Effect.andThen(worker.drain())); + + yield* steer("steer-changed", composerSelection); + const changed = yield* orchestrator.getThreadProjection(threadId); + assert.deepEqual(steered, ["steer-changed"]); + assert.equal(started.length, 1); + assert.lengthOf(changed.attempts, 1); + assert.equal(changed.runs[0]?.status, "running"); + assert.deepEqual(changed.runs[0]?.modelSelection, runSelection); + assert.deepEqual(changed.thread.modelSelection, composerSelection); + + // Choosing the running run's selection again replaces the saved choice. + yield* steer("steer-reverted", runSelection); + const reverted = yield* orchestrator.getThreadProjection(threadId); + assert.deepEqual(steered, ["steer-changed", "steer-reverted"]); + assert.lengthOf(reverted.attempts, 1); + assert.deepEqual(reverted.thread.modelSelection, runSelection); + + // The saved choice moves to another instance while the run keeps going. + // Steering with the run's selection brings the thread's instance back too. + const otherSelection = { + instanceId: ProviderInstanceId.make("codex-work"), + model: "other", + }; + const sink = yield* EventSink.EventSinkV2; + yield* sink.write({ + events: [ + { + id: EventId.make("switched-away"), + type: "thread.provider-switched", + threadId, + providerInstanceId: otherSelection.instanceId, + occurredAt: yield* DateTime.now, + payload: { + ...reverted.thread, + providerInstanceId: otherSelection.instanceId, + modelSelection: otherSelection, + }, + }, + ], + }); + yield* steer("steer-back", runSelection); + const back = yield* orchestrator.getThreadProjection(threadId); + assert.deepEqual(steered, ["steer-changed", "steer-reverted", "steer-back"]); + assert.equal(back.thread.providerInstanceId, instanceId); + assert.deepEqual(back.thread.modelSelection, runSelection); + }).pipe(Effect.provide(layer)); + }), + ), +); + +it.effect("starts a steer that missed the turn on the saved next-turn selection", () => + Effect.scoped( + Effect.gen(function* () { + const { events, started, steered, layer, startFirstTurn } = yield* nextTurnSelectionHarness( + "steering-selection-follow-up", + ); + yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const threadId = yield* startFirstTurn; + const first = started[0]!; yield* orchestrator.dispatch({ - type: "thread.create", - commandId: CommandId.make("create"), + type: "message.dispatch", + commandId: CommandId.make("steer"), threadId, - projectId: ProjectId.make("project:steering-selection-change"), - title: "Steer with changed options", - modelSelection: runSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: cwd, + messageId: MessageId.make("message:steer"), + text: "late steer", + attachments: [], + modelSelection: composerSelection, + dispatchMode: { type: "steer_active", targetRunId: first.runId }, createdBy: "user", creationSource: "web", }); - const running = yield* orchestrator.streamDomainEvents.pipe( + // The turn ends before the worker delivers the steer, so the steer + // becomes a follow-up turn. + const settled = yield* orchestrator.streamDomainEvents.pipe( Stream.filter( - (event) => event.type === "provider-turn.updated" && event.payload.status === "running", + (event) => + event.type === "run.updated" && + event.payload.id === first.runId && + event.payload.status === "waiting", ), Stream.take(1), Stream.runDrain, Effect.forkScoped, ); - yield* orchestrator.dispatch({ - type: "message.dispatch", - commandId: CommandId.make("first"), - threadId, - messageId: MessageId.make("message:first"), - text: "first", - attachments: [], - modelSelection: runSelection, - dispatchMode: { type: "start_immediately" }, - createdBy: "user", - creationSource: "web", + const turn = (yield* orchestrator.getThreadProjection(threadId)).providerTurns[0]!; + yield* Queue.offer(events, { + type: "provider_turn.updated", + driver, + providerTurn: { ...turn, status: "completed", completedAt: yield* DateTime.now }, }); - yield* worker.drain(); - yield* Fiber.join(running); - - yield* orchestrator.dispatch({ - type: "message.dispatch", - commandId: CommandId.make("steer"), - threadId, - messageId: MessageId.make("message:steer"), - text: "use the File component", - attachments: [], - modelSelection: composerSelection, - dispatchMode: { type: "steer_active", targetRunId: started[0]!.runId }, - createdBy: "user", - creationSource: "web", + yield* Queue.offer(events, { + type: "turn.terminal", + driver, + providerThreadId: turn.providerThreadId, + providerTurnId: turn.id, + runOrdinal: first.runOrdinal, + status: "completed", + failure: null, + threadDisposition: "reusable", }); + yield* Fiber.join(settled); + yield* worker.drain(); + yield* orchestrator.resumeQueuedRuns; yield* worker.drain(); const projection = yield* orchestrator.getThreadProjection(threadId); - assert.deepEqual(steered, ["use the File component"]); - assert.equal(started.length, 1); - assert.lengthOf(projection.attempts, 1); - assert.equal(projection.runs[0]?.status, "running"); - assert.deepEqual(projection.runs[0]?.modelSelection, runSelection); + assert.deepEqual(steered, []); + assert.equal(started.length, 2); + assert.equal(started[1]?.message.messageId, MessageId.make("message:steer")); + assert.deepEqual(started[1]?.modelSelection, composerSelection); assert.deepEqual(projection.thread.modelSelection, composerSelection); - }).pipe( - Effect.provide( - makeOrchestratorV2ReplayLayerWithRegistry( - { name: "steering-selection-change" }, - ProviderAdapterRegistry.makeSingleLayer(adapter), - { runEffectWorker: false }, - ), - ), - ); + }).pipe(Effect.provide(layer)); }), ), );