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
30 changes: 29 additions & 1 deletion apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3604,17 +3604,45 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
});
});

// A selection the provider applies on its next turn restarts the run when
// the provider can restart it. Otherwise the steer joins the running turn
// and the selection waits for the next one, rather than failing the steer.
const turnCapabilities = session.providerSession.capabilities.turns;
const selectionMustApplyNow =
selectionChanged &&
(providerInstanceChanged || selectionTransition?.type !== "apply_on_next_turn");
const steeringPolicy = yield* enforceCommandPolicy(input.command)(
commandPolicy.decideSteeringExecution({
commandId: input.command.commandId,
threadId: input.command.threadId,
providerInstanceId: targetRun.providerInstanceId,
capabilities: session.providerSession.capabilities,
forceRestart: input.forceRestart || selectionChanged,
forceRestart:
input.forceRestart ||
selectionMustApplyNow ||
(selectionChanged &&
turnCapabilities.supportsInterrupt &&
turnCapabilities.supportsSteeringByInterruptRestart),
}),
);

if (steeringPolicy === "active_steering") {
if (
selectionChanged &&
!modelSelectionsEqual(input.projection.thread.modelSelection, input.modelSelection)
) {
yield* emitEvent({
type: "thread.model-selection-updated",
threadId: input.command.threadId,
providerInstanceId: input.modelSelection.instanceId,
occurredAt: now,
payload: {
...input.projection.thread,
modelSelection: input.modelSelection,
updatedAt: now,
},
});
}
yield* appendSteeringMessage({
runId: targetRun.id,
nodeId: rootNodeId,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -401,3 +401,184 @@ 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<ProviderAdapterV2Event>();
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,
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,
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,
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"),
};
}),
};
yield* Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const worker = yield* EffectWorker.OrchestrationEffectWorkerV2;
const threadId = ThreadId.make("thread:steering-selection-change");
yield* orchestrator.dispatch({
type: "thread.create",
commandId: CommandId.make("create"),
threadId,
projectId: ProjectId.make("project:steering-selection-change"),
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);

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* 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(projection.thread.modelSelection, composerSelection);
}).pipe(
Effect.provide(
makeOrchestratorV2ReplayLayerWithRegistry(
{ name: "steering-selection-change" },
ProviderAdapterRegistry.makeSingleLayer(adapter),
{ runEffectWorker: false },
),
),
);
}),
),
);
Loading