diff --git a/apps/mobile/src/components/AppSymbol.tsx b/apps/mobile/src/components/AppSymbol.tsx index 17eaa1cc718b..0126347a3f82 100644 --- a/apps/mobile/src/components/AppSymbol.tsx +++ b/apps/mobile/src/components/AppSymbol.tsx @@ -98,6 +98,7 @@ import IconStar from "@tabler/icons-react-native/IconStar"; import IconStarFilled from "@tabler/icons-react-native/IconStarFilled"; import IconStethoscope from "@tabler/icons-react-native/IconStethoscope"; import IconSun from "@tabler/icons-react-native/IconSun"; +import IconTarget from "@tabler/icons-react-native/IconTarget"; import IconTerminal2 from "@tabler/icons-react-native/IconTerminal2"; import IconTextDecrease from "@tabler/icons-react-native/IconTextDecrease"; import IconTextIncrease from "@tabler/icons-react-native/IconTextIncrease"; @@ -215,6 +216,7 @@ const ANDROID_ICON_BY_SF_SYMBOL = { "star.fill": IconStarFilled, "sun.max": IconSun, "stop.fill": IconPlayerStopFilled, + target: IconTarget, terminal: IconTerminal2, "text.alignleft": IconAlignLeft, "text.bubble": IconMessage, diff --git a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx index 013bf900912d..cd536556b352 100644 --- a/apps/mobile/src/features/threads/ThreadDetailScreen.tsx +++ b/apps/mobile/src/features/threads/ThreadDetailScreen.tsx @@ -30,7 +30,10 @@ import { type CodexArtifactTemplate, } from "@t3tools/client-runtime/codex-artifact-templates"; import type { ThreadUserInputQuestion } from "@t3tools/client-runtime/state/thread-requests"; -import { presentPendingBackgroundWork } from "@t3tools/client-runtime/state/thread-execution"; +import { + presentPendingBackgroundWork, + presentProviderGoal, +} from "@t3tools/client-runtime/state/thread-execution"; import { resolveSubagentPillSegment } from "@t3tools/client-runtime/state/thread-subagents"; import { formatModelSelectionEffort, @@ -484,6 +487,14 @@ export const ThreadDetailScreen = memo(function ThreadDetailScreen(props: Thread waiting: pendingBackgroundWork.waiting, }; } + if (props.selectedThread.goal !== null && contentPresentationKind === "ready") { + const goal = presentProviderGoal(props.selectedThread.goal, false); + return { + kind: "goal", + label: goal.title, + accessibilityLabel: `${goal.title}: ${goal.objective}`, + }; + } return null; })(); const showWorkingControl = floatingStatus !== null; diff --git a/apps/mobile/src/features/threads/floating-working-control.tsx b/apps/mobile/src/features/threads/floating-working-control.tsx index 95ea8ca5534a..27d6757f3158 100644 --- a/apps/mobile/src/features/threads/floating-working-control.tsx +++ b/apps/mobile/src/features/threads/floating-working-control.tsx @@ -407,6 +407,21 @@ function FloatingStatusLabel(props: { ); } + if (props.status.kind === "goal") { + return ( + + + + {props.status.label} + + + ); + } if (props.status.kind === "preparing") { return ( { }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ), ); + + describe("native goals", () => { + const syntheticFrame = (uuid: string, text: string) => + claudeSdkFrame({ + type: "assistant", + message: { + model: "", + id: uuid, + type: "message", + role: "assistant", + content: [{ type: "text", text }], + stop_reason: "end_turn", + stop_sequence: null, + usage: { + input_tokens: 0, + output_tokens: 0, + cache_creation_input_tokens: 0, + cache_read_input_tokens: 0, + }, + }, + parent_tool_use_id: null, + uuid, + session_id: WAKE_NATIVE_SESSION, + }); + const stopHookFeedback = (uuid: string, condition: string, reason: string) => + claudeSdkFrame({ + type: "user", + message: { + role: "user", + content: [{ type: "text", text: `Stop hook feedback:\n[${condition}]: ${reason}` }], + }, + parent_tool_use_id: null, + isSynthetic: true, + uuid, + session_id: WAKE_NATIVE_SESSION, + }); + const goalStatuses = (events: ReadonlyArray) => + events.flatMap((event) => + event.type === "provider_thread.updated" && event.providerThread.goal != null + ? [event.providerThread.goal] + : [], + ); + + it.effect("tracks a /goal through unmet checks until Claude stops on its own", () => + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const condition = "all tests pass"; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-attempt"), + text: `/goal ${condition}`, + attachments: [], + }), + ); + for (const frame of [ + syntheticFrame("goal-set", `Goal set: ${condition}`), + makeAssistantTextFrame({ uuid: "goal-work-1", text: "Fixing the first test." }), + stopHookFeedback("goal-check-1", condition, "One test still fails."), + makeAssistantTextFrame({ uuid: "goal-work-2", text: "All tests pass now." }), + makeResultFrame({ uuid: "goal-result", result: "All tests pass now." }), + ]) { + yield* Queue.offer(harness.sdkMessages, frame); + } + const terminal = yield* Queue.take(harness.terminalReceipts); + assert.equal(terminal.status, "completed"); + assert.deepEqual(goalStatuses(harness.events), [ + { objective: condition, status: "active", checks: 0 }, + { + objective: condition, + status: "active", + checks: 1, + lastCheck: "One test still fails.", + }, + { objective: condition, status: "complete", checks: 1 }, + ]); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("keeps a goal active when a hook stops the turn before the goal passes", () => + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const condition = "the deploy succeeds"; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-hook-stop-attempt"), + text: `/goal ${condition}`, + attachments: [], + }), + ); + for (const frame of [ + syntheticFrame("goal-hook-set", `Goal set: ${condition}`), + makeAssistantTextFrame({ uuid: "goal-hook-work", text: "Deploying." }), + makeResultFrame({ + uuid: "goal-hook-result", + result: "Deploying.", + terminalReason: "hook_stopped", + }), + ]) { + yield* Queue.offer(harness.sdkMessages, frame); + } + yield* Queue.take(harness.terminalReceipts); + assert.deepEqual(goalStatuses(harness.events).at(-1), { + objective: condition, + status: "active", + checks: 0, + }); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("keeps a goal active after a command turn with no model output", () => + Effect.gen(function* () { + const harness = yield* makeWakeHarness; + const condition = "the build is green"; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: { + ...harness.providerThread, + goal: { objective: condition, status: "active", checks: 2 }, + }, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-show-attempt"), + text: "/goal", + attachments: [], + }), + ); + yield* Queue.offer( + harness.sdkMessages, + syntheticFrame("goal-show", `Goal active: ${condition} (2 turns)`), + ); + yield* Queue.offer( + harness.sdkMessages, + makeResultFrame({ uuid: "goal-show-result", result: "", numTurns: 0 }), + ); + yield* Queue.take(harness.terminalReceipts); + assert.deepEqual(goalStatuses(harness.events).at(-1), { + objective: condition, + status: "active", + checks: 2, + }); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + }); }); describe("ClaudeAdapterV2 query message stream", () => { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index 0aca36abe77c..f4e53c395ce4 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -48,6 +48,7 @@ import { type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderCapabilities, type OrchestrationV2ProviderFailure, + OrchestrationV2ProviderGoal, type OrchestrationV2PlanArtifact, type OrchestrationV2PlanStep, type OrchestrationV2PendingBackgroundTask, @@ -1014,6 +1015,64 @@ function textFromClaudeContent(content: SDKAssistantMessage["message"]["content" return content.flatMap((part) => (part.type === "text" ? [part.text] : [])).join(""); } +// In SDK mode Claude reports `/goal` only through the transcript (its +// `active_goal` event is remote-only): synthetic command output names the +// goal, and each unmet evaluator check returns as Stop hook feedback. +const CLAUDE_GOAL_SET_PREFIX = "Goal set: "; +const CLAUDE_GOAL_ACTIVE = /^Goal active: ([\s\S]+?) \((?:not yet evaluated|(\d+) turns?)\)/u; + +/** + * The goal after one root SDK frame, or undefined when the frame says nothing + * about it. Completion has no frame of its own; see finalizeActiveTurn. + */ +function nextClaudeGoal( + current: OrchestrationV2ProviderGoal | null, + message: SDKMessage, +): OrchestrationV2ProviderGoal | null | undefined { + if (message.type === "assistant") { + if (message.parent_tool_use_id !== null || message.message.model !== "") { + return undefined; + } + const text = textFromClaudeContent(message.message.content).trim(); + if (text.startsWith(CLAUDE_GOAL_SET_PREFIX)) { + const objective = text.slice(CLAUDE_GOAL_SET_PREFIX.length).trim(); + return objective.length === 0 ? undefined : { objective, status: "active", checks: 0 }; + } + if (text.startsWith("Goal cleared: ") || text.startsWith("No goal set")) return null; + const active = CLAUDE_GOAL_ACTIVE.exec(text); + if (active?.[1] === undefined) return undefined; + return { + ...(current?.objective === active[1] ? current : {}), + objective: active[1], + status: "active", + checks: Number(active[2] ?? 0), + }; + } + if ( + message.type !== "user" || + message.parent_tool_use_id !== null || + message.isSynthetic !== true || + current === null + ) { + return undefined; + } + const prefix = `Stop hook feedback:\n[${current.objective}]: `; + const content = message.message.content; + const text = + typeof content === "string" + ? content + : content.flatMap((part) => (part.type === "text" ? [part.text] : [])).join(""); + if (!text.startsWith(prefix)) return undefined; + return { + ...current, + status: "active", + checks: (current.checks ?? 0) + 1, + lastCheck: text.slice(prefix.length).trim(), + }; +} + +const providerGoalsEqual = Schema.toEquivalence(Schema.NullOr(OrchestrationV2ProviderGoal)); + /** Compares a subagent result with its routed text regardless of block joins. */ function normalizeClaudeResultText(text: string): string { return text.replace(/\s+/g, " ").trim(); @@ -3035,6 +3094,10 @@ export function makeClaudeAdapterV2( const lastProviderThreadByNativeThread = yield* Ref.make( new Map(), ); + // Native `/goal` per session, seeded from the persisted provider thread. + const goalsByNativeThread = new Map(); + // Turns whose model output met an active goal's Stop hook check at its end. + const goalCheckedTurns = new Set(); // Subagent registry that survives turn settle: a background subagent // (Agent with run_in_background) can complete after the root turn // ended, and its task_notification must both count as wake evidence @@ -3438,6 +3501,9 @@ export function makeClaudeAdapterV2( providerSessionId: session.id, ...(input.status === undefined ? {} : { status: input.status }), pendingBackgroundTasks: claudePendingBackgroundTasksFromRoster(roster), + ...(goalsByNativeThread.has(input.nativeThreadId) + ? { goal: goalsByNativeThread.get(input.nativeThreadId) ?? null } + : {}), updatedAt: now, }; yield* rememberProviderThread(providerThread); @@ -3448,6 +3514,37 @@ export function makeClaudeAdapterV2( }); }); + /** Applies one root frame's goal signal and writes any change onto the provider thread. */ + const trackClaudeGoal = Effect.fnUntraced(function* (input: { + readonly nativeThreadId: string; + readonly context: ActiveClaudeTurnContext; + readonly message: SDKMessage; + }) { + const current = goalsByNativeThread.get(input.nativeThreadId) ?? null; + if ( + current?.status === "active" && + input.message.type === "assistant" && + input.message.parent_tool_use_id === null && + input.message.message.model !== "" + ) { + goalCheckedTurns.add(input.context.providerTurnId); + } + const next = nextClaudeGoal(current, input.message); + if (next === undefined || providerGoalsEqual(current, next)) return; + goalsByNativeThread.set(input.nativeThreadId, next); + const providerThread = (yield* Ref.get(lastProviderThreadByNativeThread)).get( + input.nativeThreadId, + ); + if (providerThread === undefined) return; + const updated = { ...providerThread, goal: next, updatedAt: yield* DateTime.now }; + yield* rememberProviderThread(updated); + yield* emitProviderEvent({ + type: "provider_thread.updated", + driver: CLAUDE_PROVIDER, + providerThread: updated, + }); + }); + const rememberOpaqueTasks = (tasks: ReadonlyArray) => tasks.length === 0 ? Effect.void @@ -4936,10 +5033,36 @@ export function makeClaudeAdapterV2( const clearConversationHead = input.status === "completed" && input.context.input.providerThread.nativeConversationHeadRef !== null; + // While a goal is set, Claude ends a turn on its own only after the + // goal's Stop hook passes. It defers that check while background + // work runs, and a hook that stops the turn reports another reason. + // SDK mode does not report an evaluator timeout or an impossible + // verdict, so those still read as complete. + const goal = + nativeThreadId === null ? undefined : goalsByNativeThread.get(nativeThreadId); + const goalChecked = goalCheckedTurns.delete(input.context.providerTurnId); + const terminalReason = input.result?.terminal_reason; + if ( + nativeThreadId !== null && + goal?.status === "active" && + goalChecked && + input.status === "completed" && + (terminalReason === undefined || terminalReason === "completed") && + roster.size === 0 + ) { + goalsByNativeThread.set(nativeThreadId, { + objective: goal.objective, + status: "complete", + ...(goal.checks === undefined ? {} : { checks: goal.checks }), + }); + } const providerThread: OrchestrationV2ProviderThread = { ...input.context.input.providerThread, providerSessionId: session.id, ...(clearConversationHead ? { nativeConversationHeadRef: null } : {}), + ...(nativeThreadId !== null && goalsByNativeThread.has(nativeThreadId) + ? { goal: goalsByNativeThread.get(nativeThreadId) ?? null } + : {}), firstRunOrdinal: input.context.input.providerThread.firstRunOrdinal ?? input.context.input.runOrdinal, @@ -5513,6 +5636,7 @@ export function makeClaudeAdapterV2( } return; } + yield* trackClaudeGoal({ nativeThreadId: liveQuery.nativeThreadId, context, message }); // Subagent narration belongs to its child thread, never the parent log. if (message.type === "stream_event" && !message.parent_tool_use_id) { @@ -7083,6 +7207,9 @@ export function makeClaudeAdapterV2( return updated; }); yield* rememberProviderThread(turnInput.providerThread); + if (!goalsByNativeThread.has(nativeThreadId)) { + goalsByNativeThread.set(nativeThreadId, turnInput.providerThread.goal ?? null); + } const context: ActiveClaudeTurnContext = { input: turnInput, nativeTurnId, diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index 85748fa056e7..bc93ec1ec0ea 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -1147,6 +1147,11 @@ describe("CodexAdapterV2 rollback mapping", () => { const secondTurn = providerTurn("provider-turn-second", 2, "completed"); const runningTurn = providerTurn("provider-turn-running", 3, "running"); const interruptedTurn = providerTurn("provider-turn-interrupted", 4, "interrupted"); + // `/goal` control turns never reached Codex, so revert must not count them. + const goalCommandTurn = { + ...providerTurn("provider-turn-goal-command", 5, "completed"), + nativeTurnRef: null, + }; const numTurns = yield* CodexAdapterV2.resolveCodexRollbackTurnCount({ providerThread, @@ -1156,7 +1161,7 @@ describe("CodexAdapterV2 rollback mapping", () => { appRunOrdinal: 1, providerTurn: firstTurn, }, - providerThreadTurns: [interruptedTurn, runningTurn, secondTurn, firstTurn], + providerThreadTurns: [goalCommandTurn, interruptedTurn, runningTurn, secondTurn, firstTurn], }); assert.equal(numTurns, 2); @@ -7272,4 +7277,841 @@ describe("CodexAdapterV2 post-settle continuation", () => { assert.include(errorCauseChainText(error), "fork exploded"); }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), ); + + describe("native goals", () => { + const nativeThreadId = "goal-thread"; + const objective = "Ship the feature"; + const codexGoal = (status: string, tokensUsed: number) => ({ + threadId: nativeThreadId, + objective, + status, + tokenBudget: null, + tokensUsed, + timeUsedSeconds: 3, + createdAt: 1782622440, + updatedAt: 1782622450, + }); + const notification = ( + label: string, + method: string, + params: unknown, + ): CodexReplay.CodexAppServerReplayEntry => ({ + type: "emit_inbound", + label, + frame: { method, params }, + }); + const request = (id: number, method: string, params: unknown, result: unknown) => + [ + { type: "expect_outbound", label: method, frame: { id, method, params } }, + { type: "emit_inbound", label: method, frame: { id, result } }, + ] satisfies Array; + const withReplayRequestId = ( + entries: ReadonlyArray, + id: number, + ) => + entries.map((entry) => + "frame" in entry && Predicate.isObject(entry.frame) && "id" in entry.frame + ? { ...entry, frame: { ...entry.frame, id } } + : entry, + ); + const sessionStart = codexReplayPreamble({ + nativeThreadId, + nativeTurnId: "unused", + prompt: "unused", + }).slice(0, 5); + const goalTurnInput = ( + harness: { + readonly threadId: ThreadId; + readonly providerThread: OrchestrationV2ProviderThread; + }, + text: string, + ) => + DateTime.now.pipe( + Effect.map((now) => + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("goal-attempt"), + text, + }), + ), + ); + const providerGoals = (events: ReadonlyArray) => + events.flatMap((event) => + event.type === "provider_thread.updated" ? [event.providerThread.goal?.status ?? null] : [], + ); + + it("parses /goal like the Codex TUI", () => { + assert.deepEqual(CodexAdapterV2.parseCodexGoalCommand("/goal"), { type: "show" }); + assert.deepEqual(CodexAdapterV2.parseCodexGoalCommand(" /goal Pause "), { type: "pause" }); + assert.deepEqual(CodexAdapterV2.parseCodexGoalCommand("/goal resume"), { type: "resume" }); + assert.deepEqual(CodexAdapterV2.parseCodexGoalCommand("/goal clear"), { type: "clear" }); + assert.deepEqual(CodexAdapterV2.parseCodexGoalCommand("/goal fix the\nflaky tests"), { + type: "set", + objective: "fix the\nflaky tests", + }); + assert.isNull(CodexAdapterV2.parseCodexGoalCommand("/goals")); + assert.isNull(CodexAdapterV2.parseCodexGoalCommand("set a /goal")); + }); + + it.effect("keeps a /goal run open across the turns Codex continues on its own", () => + Effect.gen(function* () { + const transcript = makeCodexReplayTranscript({ + scenario: "goal-continuation", + entries: [ + ...sessionStart, + ...request(3, "thread/goal/get", { threadId: nativeThreadId }, { goal: null }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, objective, status: "paused" }, + { goal: codexGoal("paused", 0) }, + ), + // The first goal turn carries this run's turn configuration. + ...withReplayRequestId( + codexReplayPreamble({ + nativeThreadId, + nativeTurnId: "goal-turn-1", + prompt: objective, + }).slice(5), + 5, + ), + ...request( + 6, + "thread/goal/set", + { threadId: nativeThreadId, status: "active" }, + { goal: codexGoal("active", 0) }, + ), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-1", + goal: codexGoal("active", 0), + }), + notification("first turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-1", status: "completed" }), + }), + notification("continuation", "turn/started", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-2", status: "inProgress" }), + }), + notification("goal complete", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-2", + goal: codexGoal("complete", 1200), + }), + notification("continuation done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-2", status: "completed" }), + }), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn(yield* goalTurnInput(harness, `/goal ${objective}`)); + yield* harness.firstTerminal; + + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const providerTurnId = (nativeTurnId: string) => + idAllocator.derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId, + }); + assert.deepEqual( + harness.terminalEvents().map((event) => [event.providerTurnId, event.status]), + [[providerTurnId("goal-turn-2"), "completed"]], + ); + assert.deepEqual( + harness.events.flatMap((event) => + event.type === "provider_turn.updated" && event.providerTurn.status === "running" + ? [ + [ + event.providerTurn.id, + event.providerTurn.runAttemptId, + event.providerTurn.ordinal, + ], + ] + : [], + ), + [ + [providerTurnId("goal-turn-1"), RunAttemptId.make("goal-attempt"), 1], + [providerTurnId("goal-turn-2"), RunAttemptId.make("goal-attempt"), 2], + ], + ); + assert.equal(providerGoals(harness.events).at(-1), "complete"); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("settles /goal pause with a reply instead of a native turn", () => + Effect.gen(function* () { + const transcript = makeCodexReplayTranscript({ + scenario: "goal-pause", + entries: [ + ...sessionStart, + ...request( + 3, + "thread/goal/get", + { threadId: nativeThreadId }, + { goal: codexGoal("active", 10) }, + ), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 10) }, + ), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn(yield* goalTurnInput(harness, "/goal pause")); + yield* harness.firstTerminal; + + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["completed"], + ); + const turns = harness.events.flatMap((event) => + event.type === "provider_turn.updated" ? [event.providerTurn] : [], + ); + assert.isTrue(turns.every((turn) => turn.nativeTurnRef === null)); + assert.include( + harness.events.flatMap((event) => + event.type === "message.updated" && event.message.role === "assistant" + ? [event.message.text] + : [], + ), + "Goal paused. Send /goal resume to continue.", + ); + assert.equal(providerGoals(harness.events).at(-1), "paused"); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("settles a run held for the next goal turn when Stop arrives in between", () => + Effect.gen(function* () { + const prompt = "Keep going"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-between-turns", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn", prompt }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn", + goal: codexGoal("active", 10), + }), + notification("turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn", status: "completed" }), + }), + // Notifications run in order, so seeing this one proves the turn settled. + notification("goal accounted", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: null, + goal: codexGoal("active", 11), + }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 10) }, + ), + notification("goal paused", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: null, + goal: codexGoal("paused", 10), + }), + ], + }); + const requests: Array = []; + const harness = yield* makeCodexReplayHarness( + transcript, + () => Effect.void, + (method) => + Effect.sync(() => { + requests.push(method); + }), + ); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-hold-attempt"), + text: prompt, + }), + ); + const providerTurnId = (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn", + }); + const turnStatuses = () => + harness.events.flatMap((event) => + event.type === "provider_turn.updated" && event.providerTurn.id === providerTurnId + ? [event.providerTurn.status] + : [], + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.tokensUsed === 11, + ), + "the completed turn to be held", + ); + // The held turn keeps reading as running, so Stop still reaches the adapter. + assert.deepEqual(turnStatuses(), ["running"]); + assert.lengthOf(harness.terminalEvents(), 0); + + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId, + }); + yield* harness.firstTerminal; + + assert.notInclude(requests, "turn/interrupt"); + assert.deepEqual(turnStatuses(), ["running", "completed"]); + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("moves a Stop that names an earlier goal turn to the turn Codex continued with", () => + Effect.gen(function* () { + const prompt = "Keep going"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-after-continuation", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn-a", prompt }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-a", + goal: codexGoal("active", 10), + }), + notification("first turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-a", status: "completed" }), + }), + notification("continuation", "turn/started", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "inProgress" }), + }), + notification("goal accounted", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-b", + goal: codexGoal("active", 11), + }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 11) }, + ), + ...request( + 5, + "turn/interrupt", + { threadId: nativeThreadId, turnId: "goal-turn-b" }, + {}, + ), + notification("continuation interrupted", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "interrupted" }), + }), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-successor-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.tokensUsed === 11, + ), + "the continuation to be adopted", + ); + // The projection may still name the first turn when Stop is dispatched. + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn-a", + }), + requestRuntimeRestart: true, + }); + yield* harness.firstTerminal; + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("stops a held goal run's background command when Stop arrives between turns", () => + Effect.gen(function* () { + const prompt = "Start the server"; + const command = "node server.js"; + const commandItem = (status: "inProgress" | "completed") => ({ + type: "commandExecution", + id: "goal-bg-command", + command, + cwd: "/workspace", + processId: "4242", + source: "unifiedExecStartup", + status, + commandActions: [{ type: "unknown", command }], + aggregatedOutput: null, + exitCode: null, + durationMs: null, + }); + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-held-background", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn", prompt }), + notification("command started", "item/started", { + item: commandItem("inProgress"), + threadId: nativeThreadId, + turnId: "goal-turn", + startedAtMs: 1782622440500, + }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn", + goal: codexGoal("active", 10), + }), + notification("turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn", status: "completed" }), + }), + notification("goal accounted", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: null, + goal: codexGoal("active", 11), + }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 11) }, + ), + ...request( + 5, + "thread/backgroundTerminals/terminate", + { threadId: nativeThreadId, processId: "4242" }, + { terminated: true }, + ), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-held-background-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.tokensUsed === 11, + ), + "the completed turn to be held", + ); + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn", + }), + requestRuntimeRestart: true, + }); + yield* harness.firstTerminal; + // The settled-turn path terminates the process (the replay expects that + // request) before it marks the command interrupted. + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "turn_item.updated" && + event.turnItem.type === "command_execution" && + event.turnItem.status === "interrupted", + ), + "the background command to be interrupted", + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("moves Stop to a continuation Codex starts while the goal pause is in flight", () => + Effect.gen(function* () { + const prompt = "Keep going"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-races-continuation", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn-a", prompt }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-a", + goal: codexGoal("active", 10), + }), + { + type: "expect_outbound", + label: "thread/goal/set", + frame: { + id: 4, + method: "thread/goal/set", + params: { threadId: nativeThreadId, status: "paused" }, + }, + }, + notification("first turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-a", status: "completed" }), + }), + notification("continuation", "turn/started", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "inProgress" }), + }), + { + type: "emit_inbound", + label: "thread/goal/set", + frame: { id: 4, result: { goal: codexGoal("paused", 11) } }, + }, + ...request( + 5, + "turn/interrupt", + { threadId: nativeThreadId, turnId: "goal-turn-b" }, + {}, + ), + notification("continuation interrupted", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "interrupted" }), + }), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-stop-race-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.status === "active", + ), + "the goal to be active", + ); + // Stop names the first turn while it still runs. + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn-a", + }), + requestRuntimeRestart: true, + }); + yield* harness.firstTerminal; + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("stops a subagent an earlier goal turn left running when Stop settles the run", () => + Effect.gen(function* () { + const prompt = "Delegate the audit"; + const childThreadId = "goal-child-thread"; + const childTurnId = "goal-child-turn"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-held-subagent", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn", prompt }), + notification("subagent started", "item/completed", { + item: { + type: "subAgentActivity", + id: "goal-subagent-call", + kind: "started", + agentThreadId: childThreadId, + agentPath: "/root/audit", + }, + threadId: nativeThreadId, + turnId: "goal-turn", + completedAtMs: 1782622441000, + }), + notification("child turn", "turn/started", { + threadId: childThreadId, + turn: makeCodexReplayTurn({ id: childTurnId, status: "inProgress" }), + }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn", + goal: codexGoal("active", 10), + }), + notification("turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn", status: "completed" }), + }), + notification("goal accounted", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: null, + goal: codexGoal("active", 11), + }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 11) }, + ), + ...request(5, "turn/interrupt", { threadId: childThreadId, turnId: childTurnId }, {}), + notification("child interrupted", "turn/completed", { + threadId: childThreadId, + turn: makeCodexReplayTurn({ id: childTurnId, status: "interrupted" }), + }), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-held-subagent-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.tokensUsed === 11, + ), + "the completed turn to be held", + ); + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn", + }), + requestRuntimeRestart: true, + }); + yield* harness.firstTerminal; + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_turn.updated" && + event.providerTurn.nativeTurnRef?.nativeId === childTurnId && + event.providerTurn.status === "interrupted", + ), + "the subagent turn to be interrupted", + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("stops an earlier goal turn's subagent when the run settles during the pause", () => + Effect.gen(function* () { + const prompt = "Delegate the audit"; + const childThreadId = "goal-race-child-thread"; + const childTurnId = "goal-race-child-turn"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop-run-settles-during-pause", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn-a", prompt }), + notification("subagent started", "item/completed", { + item: { + type: "subAgentActivity", + id: "goal-race-subagent-call", + kind: "started", + agentThreadId: childThreadId, + agentPath: "/root/audit", + }, + threadId: nativeThreadId, + turnId: "goal-turn-a", + completedAtMs: 1782622441000, + }), + notification("child turn", "turn/started", { + threadId: childThreadId, + turn: makeCodexReplayTurn({ id: childTurnId, status: "inProgress" }), + }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-a", + goal: codexGoal("active", 10), + }), + notification("first turn done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-a", status: "completed" }), + }), + notification("continuation", "turn/started", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "inProgress" }), + }), + notification("goal accounted", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-b", + goal: codexGoal("active", 11), + }), + { + type: "expect_outbound", + label: "thread/goal/set", + frame: { + id: 4, + method: "thread/goal/set", + params: { threadId: nativeThreadId, status: "paused" }, + }, + }, + notification("goal paused", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn-b", + goal: codexGoal("paused", 11), + }), + notification("continuation done", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn-b", status: "completed" }), + }), + { + type: "emit_inbound", + label: "thread/goal/set", + frame: { id: 4, result: { goal: codexGoal("paused", 11) } }, + }, + ...request(5, "turn/interrupt", { threadId: childThreadId, turnId: childTurnId }, {}), + notification("child interrupted", "turn/completed", { + threadId: childThreadId, + turn: makeCodexReplayTurn({ id: childTurnId, status: "interrupted" }), + }), + ], + }); + const harness = yield* makeCodexReplayHarness(transcript); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-settles-during-pause-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_thread.updated" && + event.providerThread.goal?.tokensUsed === 11, + ), + "the continuation to be adopted", + ); + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn-b", + }), + requestRuntimeRestart: true, + }); + yield* awaitUntil( + () => + harness.events.some( + (event) => + event.type === "provider_turn.updated" && + event.providerTurn.nativeTurnRef?.nativeId === childTurnId && + event.providerTurn.status === "interrupted", + ), + "the subagent turn to be interrupted", + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + + it.effect("pauses an active goal before Stop interrupts its turn", () => + Effect.gen(function* () { + const prompt = "Keep going"; + const transcript = makeCodexReplayTranscript({ + scenario: "goal-stop", + entries: [ + ...codexReplayPreamble({ nativeThreadId, nativeTurnId: "goal-turn", prompt }), + notification("goal active", "thread/goal/updated", { + threadId: nativeThreadId, + turnId: "goal-turn", + goal: codexGoal("active", 10), + }), + ...request( + 4, + "thread/goal/set", + { threadId: nativeThreadId, status: "paused" }, + { goal: codexGoal("paused", 10) }, + ), + ...request(5, "turn/interrupt", { threadId: nativeThreadId, turnId: "goal-turn" }, {}), + notification("interrupted", "turn/completed", { + threadId: nativeThreadId, + turn: makeCodexReplayTurn({ id: "goal-turn", status: "interrupted" }), + }), + ], + }); + const requests: Array = []; + const harness = yield* makeCodexReplayHarness( + transcript, + () => Effect.void, + (method) => + Effect.sync(() => { + requests.push(method); + }), + ); + yield* harness.runtime.startTurn( + makeCodexTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now: yield* DateTime.now, + attemptId: RunAttemptId.make("goal-stop-attempt"), + text: prompt, + }), + ); + yield* awaitUntil( + () => providerGoals(harness.events).includes("active"), + "the active goal to reach the provider thread", + ); + yield* harness.runtime.interruptTurn({ + providerThread: harness.providerThread, + providerTurnId: (yield* IdAllocator.IdAllocatorV2).derive.providerTurn({ + driver: CodexAdapterV2.CODEX_DRIVER_KIND, + nativeTurnId: "goal-turn", + }), + }); + yield* harness.firstTerminal; + + assert.deepEqual(requests.slice(-2), ["thread/goal/set", "turn/interrupt"]); + assert.deepEqual( + harness.terminalEvents().map((event) => event.status), + ["interrupted"], + ); + }).pipe(Effect.scoped, Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ); + }); }); diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 8da6ffe4357d..b56e6f045443 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -25,6 +25,7 @@ import { CodexSettings, defaultInstanceIdForDriver, isOrchestrationV2WorkActive, + OrchestrationV2ProviderGoal, ProviderDriverKind, type ProviderSetupError, } from "@t3tools/contracts"; @@ -148,6 +149,7 @@ import { type ProviderAdapterV2RollbackThreadInput, type ProviderAdapterV2RuntimePolicy, type ProviderAdapterV2SessionRuntime, + type ProviderAdapterV2InterruptInput, type ProviderAdapterV2SteerInput, type ProviderAdapterV2TurnInput, } from "../ProviderAdapter.ts"; @@ -846,11 +848,78 @@ function providerThreadFromCodexThread(input: { }; } +// Counts toward native rollback. `/goal` control turns never reach Codex and +// carry no native ref, so they must not consume a native turn on revert. const isTerminalProviderTurn = (turn: OrchestrationV2ProviderTurn): boolean => - turn.status === "completed" || - turn.status === "interrupted" || - turn.status === "failed" || - turn.status === "cancelled"; + turn.nativeTurnRef !== null && + (turn.status === "completed" || + turn.status === "interrupted" || + turn.status === "failed" || + turn.status === "cancelled"); + +export type CodexGoalCommand = + | { readonly type: "show" } + | { readonly type: "clear" } + | { readonly type: "pause" } + | { readonly type: "resume" } + | { readonly type: "set"; readonly objective: string }; + +/** + * Parses a `/goal` message the way the Codex TUI does: `clear`, `pause` and + * `resume` control the current goal, a bare `/goal` shows it, and any other + * text becomes the new objective. Returns null for every other message. + */ +export function parseCodexGoalCommand(text: string): CodexGoalCommand | null { + const match = /^\/goal(?:\s+([\s\S]*))?$/u.exec(text.trim()); + if (match === null) return null; + const argument = (match[1] ?? "").trim(); + switch (argument.toLowerCase()) { + case "": + case "edit": + return { type: "show" }; + case "clear": + return { type: "clear" }; + case "pause": + return { type: "pause" }; + case "resume": + return { type: "resume" }; + default: + return { type: "set", objective: argument }; + } +} + +type CodexThreadGoal = CodexSchema.V2ThreadGoalUpdatedNotification["goal"]; + +const CODEX_GOAL_STATUSES = { + active: "active", + paused: "paused", + blocked: "blocked", + usageLimited: "usage_limited", + budgetLimited: "budget_limited", + complete: "complete", +} as const satisfies Record; + +function providerGoalFromCodex(goal: CodexThreadGoal): OrchestrationV2ProviderGoal | null { + const objective = goal.objective.trim(); + if (objective.length === 0) return null; + return { + objective, + status: CODEX_GOAL_STATUSES[goal.status], + tokensUsed: Math.max(0, goal.tokensUsed), + tokenBudget: goal.tokenBudget ?? null, + timeUsedSeconds: Math.max(0, goal.timeUsedSeconds), + }; +} + +const providerGoalsEqual = Schema.toEquivalence(Schema.NullOr(OrchestrationV2ProviderGoal)); + +function describeCodexGoal(goal: OrchestrationV2ProviderGoal): string { + return `Goal ${goal.status.replace("_", " ")}: ${goal.objective}`; +} + +// Codex starts the next goal turn milliseconds after the last one completes. +// A run waits this long for it before settling, in case Codex declines. +const CODEX_GOAL_CONTINUATION_GRACE = "5 seconds"; const providerTurnsForThread = ( providerTurns: ReadonlyArray, @@ -1069,6 +1138,16 @@ function isPersistentCodexDynamicTool(item: CodexDynamicToolItem): boolean { type CodexRootTerminalEvent = Extract; +/** A completed root turn whose run stays open for the goal turn Codex starts next. */ +interface CodexGoalHold { + readonly context: ActiveCodexTurnContext; + readonly event: Extract; + /** Sent when the hold resolves, so until then Stop and steering still target this turn. */ + readonly completedTurn: OrchestrationV2ProviderTurn; + /** The turn Codex continued with, or undefined when the run settled instead. */ + readonly next: Deferred.Deferred; +} + interface DeferredCodexRootTerminal { readonly context: ActiveCodexTurnContext; readonly event: CodexRootTerminalEvent; @@ -1703,10 +1782,53 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi // Native completion and the interrupt timeout share one finalization // path. Serialize the race so only one can publish terminal events. const turnTerminalizationPermit = yield* Semaphore.make(1); + // Native goals by Codex thread, and the root provider thread snapshot + // each goal change is written onto. + const goalsByNativeThread = new Map(); + const rootProviderThreads = new Map(); + // A completed root turn whose goal is still active. Its run stays open + // for the turn Codex starts next, so goal phases do not read as Done. + const goalHolds = new Map(); + // The root turns of each continued goal run, keyed by each of its turns. + // A Stop or steer that names an earlier turn reaches the newest one, and + // Stop also reaches work an earlier turn left running. Dropped when the + // run settles. + const goalRuns = new Map>(); + // `/goal` runs whose goal turns active after the first turn starts. + const goalActivations = new Map(); + // Stops in progress by native thread. Until a Stop resolves its target, + // a held run must not settle as completed. + const goalStops = new Map(); + const latestGoalTurnId = (providerTurnId: ProviderTurnId) => + goalRuns.get(providerTurnId)?.at(-1)?.providerTurnId ?? providerTurnId; const emitProviderEvent = (event: ProviderAdapterV2Event) => Queue.offer(events, event).pipe(Effect.asVoid); + /** Writes the thread's current native goal onto its root provider thread. */ + const emitGoalUpdate = Effect.fnUntraced(function* (nativeThreadId: string) { + const providerThread = rootProviderThreads.get(nativeThreadId); + if (providerThread === undefined || !goalsByNativeThread.has(nativeThreadId)) return; + const goal = goalsByNativeThread.get(nativeThreadId) ?? null; + if (providerGoalsEqual(providerThread.goal ?? null, goal)) return; + const updated = { ...providerThread, goal, updatedAt: yield* DateTime.now }; + rootProviderThreads.set(nativeThreadId, updated); + yield* emitProviderEvent({ + type: "provider_thread.updated", + driver: CODEX_PROVIDER, + providerThread: updated, + }); + }); + + const rememberRootProviderThread = Effect.fnUntraced(function* ( + providerThread: OrchestrationV2ProviderThread, + ) { + const nativeThreadId = providerThread.nativeThreadRef?.nativeId; + if (nativeThreadId == null) return; + rootProviderThreads.set(nativeThreadId, providerThread); + yield* emitGoalUpdate(nativeThreadId); + }); + // Call only for new model-output activity. A local item/completed can // arrive while the upstream response stream is still retrying. const completeProviderRetry = Effect.fn("CodexAdapterV2.completeProviderRetry")(function* ( @@ -1813,6 +1935,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi completedAt: null, }, }); + yield* rememberRootProviderThread(input.turnInput.providerThread); return context; }); @@ -3902,6 +4025,63 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }); return; } + // Codex continues an active goal on its own. The next turn joins the + // run its previous turn belongs to. + const goalHold = goalHolds.get(payload.threadId); + if (goalHold !== undefined) { + goalHolds.delete(payload.threadId); + yield* emitProviderEvent({ + type: "provider_turn.updated", + driver: CODEX_PROVIDER, + threadId: goalHold.context.projectionThreadId, + providerTurn: goalHold.completedTurn, + }); + const next = yield* registerRootTurn({ + turnInput: { + ...goalHold.context.input, + providerTurnOrdinal: goalHold.context.providerTurnOrdinal + 1, + }, + nativeTurnId: payload.turn.id, + startedAt: codexTimestamp(payload.turn.startedAt), + }); + const goalRun = goalRuns.get(goalHold.context.providerTurnId) ?? [goalHold.context]; + goalRun.push(next); + goalRuns.set(goalHold.context.providerTurnId, goalRun); + goalRuns.set(next.providerTurnId, goalRun); + yield* Deferred.succeed(goalHold.next, next); + return; + } + // A goal turn with no run to own it (Codex continued after the run + // settled, or raced a Stop): stop it and pause the goal rather than + // work out of sight. + if ( + rootProviderThreads.has(payload.threadId) && + (goalsByNativeThread.get(payload.threadId)?.status === "active" || + goalStops.has(payload.threadId)) + ) { + yield* Effect.logWarning("orchestration-v2.codex-goal-turn-without-run", { + nativeThreadId: payload.threadId, + nativeTurnId: payload.turn.id, + }); + yield* client + .request("thread/goal/set", { threadId: payload.threadId, status: "paused" }) + .pipe( + Effect.catch((cause) => + Effect.logWarning("orchestration-v2.codex-goal-pause-failed", { cause }), + ), + Effect.andThen( + client.request("turn/interrupt", { + threadId: payload.threadId, + turnId: payload.turn.id, + }), + ), + Effect.catch((cause) => + Effect.logWarning("orchestration-v2.codex-goal-turn-stop-failed", { cause }), + ), + Effect.forkIn(scope), + ); + return; + } yield* rememberSubagentTurnStarted({ nativeThreadId: payload.threadId, nativeTurnId: payload.turn.id, @@ -3910,6 +4090,22 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }).pipe(Effect.orDie, turnTerminalizationPermit.withPermits(1)), ); + // A goal that stops being active ends any run held open for its next turn. + yield* client.handleServerNotification("thread/goal/updated", (payload) => + Effect.gen(function* () { + goalsByNativeThread.set(payload.threadId, providerGoalFromCodex(payload.goal)); + yield* emitGoalUpdate(payload.threadId); + if (payload.goal.status !== "active") yield* releaseGoalHold(payload.threadId); + }).pipe(turnTerminalizationPermit.withPermits(1)), + ); + yield* client.handleServerNotification("thread/goal/cleared", (payload) => + Effect.gen(function* () { + goalsByNativeThread.set(payload.threadId, null); + yield* emitGoalUpdate(payload.threadId); + yield* releaseGoalHold(payload.threadId); + }).pipe(turnTerminalizationPermit.withPermits(1)), + ); + yield* client.handleServerNotification("error", (payload) => Effect.gen(function* () { const context = yield* awaitActiveTurn(payload.turnId); @@ -5003,6 +5199,16 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi } : event; yield* emitProviderEvent(current); + for (const turn of goalRuns.get(context.providerTurnId) ?? []) { + goalRuns.delete(turn.providerTurnId); + } + // Later goal updates must not mark the settled thread active again. + const nativeThreadId = context.providerThread.nativeThreadRef?.nativeId; + const rootProviderThread = + nativeThreadId == null ? undefined : rootProviderThreads.get(nativeThreadId); + if (nativeThreadId != null && rootProviderThread !== undefined) { + rootProviderThreads.set(nativeThreadId, { ...rootProviderThread, status: "idle" }); + } if (current.status === "failed" && current.failure.class === "usage_limit") { const item = makeProviderFailureTurnItem({ idAllocator, @@ -5030,6 +5236,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi readonly failureMessage?: string; readonly failureCode?: string | null; readonly providerRetry?: ActiveCodexProviderRetry; + readonly goalHoldTurn?: OrchestrationV2ProviderTurn; }) { const event = yield* makeRootTerminalEvent(input); const hasActiveDescendants = Array.from((yield* Ref.get(activeTurns)).values()).some( @@ -5043,10 +5250,60 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }); return; } + const nativeThreadId = input.context.providerThread.nativeThreadRef?.nativeId; + if ( + input.goalHoldTurn !== undefined && + event.status === "completed" && + nativeThreadId != null + ) { + goalHolds.set(nativeThreadId, { + context: input.context, + event, + completedTurn: input.goalHoldTurn, + next: yield* Deferred.make(), + }); + yield* Effect.sleep(CODEX_GOAL_CONTINUATION_GRACE).pipe( + Effect.andThen( + turnTerminalizationPermit.withPermits(1)(releaseGoalHold(nativeThreadId, event)), + ), + Effect.forkIn(scope), + ); + return; + } yield* emitRootTerminal(input.context, event); }, ); + /** + * Settles a held goal run with its last turn. Call with the + * terminalization permit, after removing the hold from `goalHolds`. + */ + const settleGoalHold = Effect.fnUntraced(function* ( + hold: CodexGoalHold, + status: "completed" | "interrupted", + ) { + yield* emitProviderEvent({ + type: "provider_turn.updated", + driver: CODEX_PROVIDER, + threadId: hold.context.projectionThreadId, + providerTurn: hold.completedTurn, + }); + yield* emitRootTerminal(hold.context, { ...hold.event, status }); + yield* Deferred.succeed(hold.next, undefined); + }); + + /** Settles a held goal run unless Stop owns it or a newer hold replaced it. */ + const releaseGoalHold = Effect.fnUntraced(function* ( + nativeThreadId: string, + expected?: CodexRootTerminalEvent, + ) { + const hold = goalHolds.get(nativeThreadId); + if (hold === undefined || goalStops.has(nativeThreadId)) return; + if (expected !== undefined && hold.event !== expected) return; + goalHolds.delete(nativeThreadId); + yield* settleGoalHold(hold, "completed"); + }); + const flushReadyRootTerminals = Effect.fn("CodexAdapterV2.flushReadyRootTerminals")( function* () { const activeTurnContexts = Array.from((yield* Ref.get(activeTurns)).values()); @@ -5133,35 +5390,49 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi for (const [key, part] of reasoningParts) { if (part.turnId === input.nativeTurnId) reasoningParts.delete(key); } - yield* emitProviderEvent({ - type: "provider_turn.updated", - driver: CODEX_PROVIDER, - threadId: input.context.projectionThreadId, - providerTurn: { - id: input.context.providerTurnId, - providerThreadId: input.context.providerThread.id, - nodeId: input.context.providerNodeId, - runAttemptId: - input.context.subagent === null ? input.context.input.attemptId : null, - nativeTurnRef: { - driver: CODEX_PROVIDER, - nativeId: input.nativeTurnId, - strength: "strong", - }, - ordinal: input.context.providerTurnOrdinal, - status: input.status, - startedAt: input.context.startedAt, - completedAt: input.completedAt, - turnTokenUsage: completeCodexTurnTokenUsage( - usageStateForThread( - input.context.providerThread.nativeThreadRef?.nativeId ?? - String(input.context.providerThread.id), - ), - input.nativeTurnId, - input.status === "completed", - ), + const nativeThreadId = input.context.providerThread.nativeThreadRef?.nativeId; + // Codex continues an active goal with another turn, so the run stays + // open and this turn reads as running until that turn starts. + const activation = + nativeThreadId == null ? undefined : goalActivations.get(nativeThreadId); + const holdsForGoal = + input.context.subagent === null && + input.status === "completed" && + nativeThreadId != null && + (goalsByNativeThread.get(nativeThreadId)?.status === "active" || + (activation !== undefined && !activation.stopped)); + const completedTurn: OrchestrationV2ProviderTurn = { + id: input.context.providerTurnId, + providerThreadId: input.context.providerThread.id, + nodeId: input.context.providerNodeId, + runAttemptId: + input.context.subagent === null ? input.context.input.attemptId : null, + nativeTurnRef: { + driver: CODEX_PROVIDER, + nativeId: input.nativeTurnId, + strength: "strong", }, - }); + ordinal: input.context.providerTurnOrdinal, + status: input.status, + startedAt: input.context.startedAt, + completedAt: input.completedAt, + turnTokenUsage: completeCodexTurnTokenUsage( + usageStateForThread( + input.context.providerThread.nativeThreadRef?.nativeId ?? + String(input.context.providerThread.id), + ), + input.nativeTurnId, + input.status === "completed", + ), + }; + if (!holdsForGoal) { + yield* emitProviderEvent({ + type: "provider_turn.updated", + driver: CODEX_PROVIDER, + threadId: input.context.projectionThreadId, + providerTurn: completedTurn, + }); + } if (input.context.subagent !== null) { yield* emitProviderEvent({ type: "node.updated", @@ -5243,6 +5514,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi yield* emitOrDeferRootTerminal({ ...input, ...(providerRetry === undefined ? {} : { providerRetry }), + ...(holdsForGoal ? { goalHoldTurn: completedTurn } : {}), }); } const waiter = (yield* Ref.get(turnWaiters)).get(input.nativeTurnId); @@ -5345,6 +5617,346 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }), ); + /** + * Settles a `/goal` run that started no Codex turn. Its reply is plain + * assistant text, as Claude's own `/goal` output is. + */ + const completeGoalCommandTurn = Effect.fnUntraced(function* ( + turnInput: ProviderAdapterV2TurnInput, + reply: string, + ) { + const now = yield* DateTime.now; + const nativeId = `goal-command:${turnInput.attemptId}`; + const context: ActiveCodexTurnContext = { + input: turnInput, + projectionAppThread: turnInput.appThread, + projectionThreadId: turnInput.threadId, + projectionRunId: turnInput.runId, + nativeTurnId: nativeId, + providerThread: turnInput.providerThread, + providerTurnId: idAllocator.derive.providerTurn({ + driver: CODEX_PROVIDER, + nativeTurnId: nativeId, + }), + providerTurnOrdinal: turnInput.providerTurnOrdinal, + providerNodeId: turnInput.rootNodeId, + providerNodeKind: "root_turn", + providerNodeStartedAt: now, + itemParentNodeId: turnInput.rootNodeId, + rootNodeId: turnInput.rootNodeId, + subagent: null, + startedAt: now, + itemPositions: new Map(), + }; + const providerTurn = { + id: context.providerTurnId, + providerThreadId: turnInput.providerThread.id, + nodeId: turnInput.rootNodeId, + runAttemptId: turnInput.attemptId, + nativeTurnRef: null, + ordinal: turnInput.providerTurnOrdinal, + status: "running", + startedAt: now, + completedAt: null, + } satisfies OrchestrationV2ProviderTurn; + yield* emitProviderEvent({ + type: "provider_turn.updated", + driver: CODEX_PROVIDER, + threadId: turnInput.threadId, + providerTurn, + }); + yield* rememberRootProviderThread(turnInput.providerThread); + const artifacts = yield* buildAgentMessageArtifacts( + context, + { id: nativeId, text: reply }, + true, + ); + yield* emitProviderEvent({ + type: "node.updated", + driver: CODEX_PROVIDER, + node: artifacts.node, + }); + yield* emitProviderEvent({ + type: "message.updated", + driver: CODEX_PROVIDER, + message: artifacts.message, + }); + yield* emitProviderEvent({ + type: "turn_item.updated", + driver: CODEX_PROVIDER, + turnItem: artifacts.turnItem, + }); + yield* emitProviderEvent({ + type: "provider_turn.updated", + driver: CODEX_PROVIDER, + threadId: turnInput.threadId, + providerTurn: { ...providerTurn, status: "completed", completedAt: now }, + }); + yield* emitRootTerminal(context, { + type: "turn.terminal", + driver: CODEX_PROVIDER, + providerThreadId: turnInput.providerThread.id, + providerTurnId: context.providerTurnId, + runOrdinal: turnInput.runOrdinal, + status: "completed", + failure: null, + threadDisposition: "reusable", + }); + }); + + /** The goal hold whose run still reports this provider turn as running. */ + const heldGoalTurn = ( + providerThread: OrchestrationV2ProviderThread, + providerTurnId: ProviderTurnId, + ) => { + const nativeThreadId = providerThread.nativeThreadRef?.nativeId; + const hold = nativeThreadId == null ? undefined : goalHolds.get(nativeThreadId); + return hold?.context.providerTurnId === providerTurnId && nativeThreadId != null + ? { nativeThreadId, hold } + : undefined; + }; + + /** + * Like the Codex TUI, Stop pauses an active goal first so Codex starts + * no further goal turn. Codex can end or continue the turn while the + * pause is in flight, so the target resolves after it: a held run + * settles here and leaves only its retained work to stop, and a + * continued run moves Stop to its newest turn. `goalTurns` lists the + * run's root turns, seen before and after the pause, whose descendants + * Stop also reaches. + */ + const resolveGoalStopTarget = Effect.fnUntraced(function* ( + turnInput: ProviderAdapterV2InterruptInput, + ) { + const goalThreadId = turnInput.providerThread.nativeThreadRef?.nativeId; + if (goalThreadId == null) { + return { turnInput, goalTurns: [] as ReadonlyArray }; + } + goalStops.set(goalThreadId, (goalStops.get(goalThreadId) ?? 0) + 1); + return yield* Effect.gen(function* () { + // The run can settle while the pause is in flight, which drops its record. + const requestedTurnId = latestGoalTurnId(turnInput.providerTurnId); + const heldBefore = goalHolds.get(goalThreadId)?.context; + const turnsBefore = [ + ...(goalRuns.get(requestedTurnId) ?? []), + ...Array.from((yield* Ref.get(activeTurns)).values()).filter( + (context) => context.providerTurnId === requestedTurnId, + ), + ...(heldBefore?.providerTurnId === requestedTurnId ? [heldBefore] : []), + ]; + const activation = goalActivations.get(goalThreadId); + if (activation !== undefined) activation.stopped = true; + if (goalsByNativeThread.get(goalThreadId)?.status === "active") { + yield* client + .request("thread/goal/set", { threadId: goalThreadId, status: "paused" }) + .pipe( + Effect.tap(({ goal }) => + Effect.sync(() => + goalsByNativeThread.set(goalThreadId, providerGoalFromCodex(goal)), + ), + ), + Effect.catch((cause) => + Effect.logWarning("orchestration-v2.codex-goal-pause-failed", { cause }), + ), + ); + } + return yield* turnTerminalizationPermit.withPermits(1)( + Effect.gen(function* () { + const providerTurnId = latestGoalTurnId(turnInput.providerTurnId); + const hold = goalHolds.get(goalThreadId); + const held = hold !== undefined && hold.context.providerTurnId === providerTurnId; + const goalTurns: ReadonlyArray = Array.from( + new Set([ + ...turnsBefore, + ...(goalRuns.get(providerTurnId) ?? []), + ...(held ? [hold.context] : []), + ]), + ); + if (!held) return { turnInput: { ...turnInput, providerTurnId }, goalTurns }; + goalHolds.delete(goalThreadId); + yield* settleGoalHold(hold, "interrupted"); + return { + turnInput: { ...turnInput, providerTurnId, requestRuntimeRestart: true }, + goalTurns, + }; + }), + ); + }).pipe( + Effect.ensuring( + Effect.sync(() => { + const stops = (goalStops.get(goalThreadId) ?? 1) - 1; + if (stops > 0) goalStops.set(goalThreadId, stops); + else goalStops.delete(goalThreadId); + }), + ), + ); + }); + + /** Starts a native turn with this run's full turn configuration. */ + const startNativeTurn = ( + turnInput: ProviderAdapterV2TurnInput, + codexInput: ReadonlyArray, + ) => + Effect.gen(function* () { + const threadId = yield* getNativeThreadId(turnInput.providerThread); + const mcpSession = McpProviderSession.readMcpProviderSession(turnInput.threadId); + const turnStartParams = yield* buildCodexTurnStartParams({ + nativeThreadId: threadId, + codexInput, + runtimePolicy: turnInput.runtimePolicy, + modelSelection: turnInput.modelSelection, + hasT3Mcp: mcpSession !== undefined, + browserToolsAvailable: mcpSession?.browserToolsAvailable ?? true, + deviceToolsAvailable: mcpSession?.capabilities?.has("device") ?? false, + omitServiceTier: adapterOptions.resolveRuntime !== undefined, + }); + yield* Ref.update(pendingRootTurns, (current) => { + const updated = new Map(current); + updated.set(threadId, turnInput); + return updated; + }); + yield* Ref.update(additionalContextByThread, (current) => { + const next = new Map(current); + if (turnStartParams.additionalContext) + next.set(threadId, turnStartParams.additionalContext); + else next.delete(threadId); + return next; + }); + const started = yield* client.request("turn/start", turnStartParams); + yield* registerRootTurn({ + turnInput, + nativeTurnId: started.turn.id, + startedAt: codexTimestamp(started.turn.startedAt), + waitForNativeStart: started.turn.startedAt === null, + }); + }).pipe( + Effect.ensuring( + Effect.flatMap(getNativeThreadId(turnInput.providerThread), (threadId) => + Ref.update(pendingRootTurns, (current) => { + const updated = new Map(current); + updated.delete(threadId); + return updated; + }), + ).pipe(Effect.ignore), + ), + ); + + /** + * Runs `/goal` through Codex's native goal API, as the Codex TUI does. + * Setting or resuming a goal starts the first goal turn with this run's + * turn configuration (approvals, sandbox, model, plan mode), since a + * turn Codex starts on its own reuses the last turn's settings. Codex + * continues later goal turns itself. Other commands settle with a reply. + */ + const runGoalCommand = (turnInput: ProviderAdapterV2TurnInput, command: CodexGoalCommand) => + Effect.gen(function* () { + const threadId = yield* getNativeThreadId(turnInput.providerThread); + // Goal notifications during the command belong to this run's snapshot. + rootProviderThreads.set(threadId, turnInput.providerThread); + const readGoal = client + .request("thread/goal/get", { threadId }) + .pipe( + Effect.map((response) => + response.goal == null ? null : providerGoalFromCodex(response.goal), + ), + ); + const current = yield* readGoal; + goalsByNativeThread.set(threadId, current); + if (command.type === "show") { + return yield* completeGoalCommandTurn( + turnInput, + current === null ? "No goal is set." : describeCodexGoal(current), + ); + } + if (command.type === "clear") { + const { cleared } = yield* client.request("thread/goal/clear", { threadId }); + goalsByNativeThread.set(threadId, null); + return yield* completeGoalCommandTurn( + turnInput, + cleared ? "Goal cleared." : "No goal to clear.", + ); + } + if (current === null && command.type !== "set") { + return yield* completeGoalCommandTurn(turnInput, "No goal is set."); + } + if (command.type === "pause") { + const { goal } = yield* client.request("thread/goal/set", { + threadId, + status: "paused", + }); + goalsByNativeThread.set(threadId, providerGoalFromCodex(goal)); + return yield* completeGoalCommandTurn( + turnInput, + "Goal paused. Send /goal resume to continue.", + ); + } + if (command.type === "set") { + // A new objective replaces the goal and its accounting, like the TUI. + if (current !== null) yield* client.request("thread/goal/clear", { threadId }); + // Paused until our turn runs, so Codex does not start one first. + const { goal } = yield* client.request("thread/goal/set", { + threadId, + objective: command.objective, + status: "paused", + }); + goalsByNativeThread.set(threadId, providerGoalFromCodex(goal)); + } + // Until the goal is active, a fast first turn still holds the run, + // and Stop cancels the activation. + const activation = { stopped: false }; + goalActivations.set(threadId, activation); + yield* Effect.gen(function* () { + // The first goal turn carries the objective; a resume continues from history. + yield* startNativeTurn( + turnInput, + command.type === "set" + ? yield* toCodexInput({ + ...turnInput, + message: { ...turnInput.message, text: command.objective }, + }) + : [], + ); + if (activation.stopped) return; + const { goal } = yield* client.request("thread/goal/set", { + threadId, + status: "active", + }); + goalsByNativeThread.set(threadId, providerGoalFromCodex(goal)); + if (activation.stopped) { + const paused = yield* client.request("thread/goal/set", { + threadId, + status: "paused", + }); + goalsByNativeThread.set(threadId, providerGoalFromCodex(paused.goal)); + } + }).pipe( + Effect.ensuring( + turnTerminalizationPermit.withPermits(1)( + Effect.gen(function* () { + if (goalActivations.get(threadId) === activation) { + goalActivations.delete(threadId); + } + // A run held for an activation that did not happen settles now. + if (goalsByNativeThread.get(threadId)?.status !== "active") { + yield* releaseGoalHold(threadId); + } + }), + ), + ), + ); + }).pipe( + Effect.mapError( + (cause) => + new ProviderAdapterTurnStartError({ + driver: CODEX_PROVIDER, + threadId: turnInput.threadId, + providerThreadId: turnInput.providerThread.id, + runId: turnInput.runId, + cause, + }), + ), + ); + const runtime: ProviderAdapterV2SessionRuntime = { instanceId: adapterOptions.instanceId, driver: CODEX_PROVIDER, @@ -5465,6 +6077,8 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ), Effect.flatMap(decodeCodexResumeMetadata), ); + // Codex follows a resume with a goal snapshot notification; the + // run's first turn writes it if it differs from the stored goal. return { ...threadInput.providerThread, providerSessionId: input.providerSessionId, @@ -5544,61 +6158,21 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }), ), ), - startTurn: (turnInput) => - Effect.gen(function* () { - const threadId = yield* getNativeThreadId(turnInput.providerThread); - - const codexInput = + startTurn: (turnInput) => { + const goalCommand = + turnInput.message.attachments.length === 0 && + turnInput.restartContinuationOfRunId === undefined + ? parseCodexGoalCommand(turnInput.message.text) + : null; + if (goalCommand !== null) return runGoalCommand(turnInput, goalCommand); + return Effect.gen(function* () { + yield* startNativeTurn( + turnInput, turnInput.restartContinuationOfRunId === undefined ? yield* toCodexInput(turnInput) - : []; - const mcpSession = McpProviderSession.readMcpProviderSession(turnInput.threadId); - const turnStartParams = yield* buildCodexTurnStartParams({ - nativeThreadId: threadId, - codexInput, - runtimePolicy: turnInput.runtimePolicy, - modelSelection: turnInput.modelSelection, - hasT3Mcp: mcpSession !== undefined, - browserToolsAvailable: mcpSession?.browserToolsAvailable ?? true, - deviceToolsAvailable: mcpSession?.capabilities?.has("device") ?? false, - omitServiceTier: adapterOptions.resolveRuntime !== undefined, - }); - yield* Ref.update(pendingRootTurns, (current) => { - const updated = new Map(current); - updated.set(threadId, turnInput); - return updated; - }); - yield* Ref.update(additionalContextByThread, (current) => { - const next = new Map(current); - if (turnStartParams.additionalContext) - next.set(threadId, turnStartParams.additionalContext); - else next.delete(threadId); - return next; - }); - const started = yield* client.request("turn/start", turnStartParams); - const nativeTurnId = started.turn.id; - const startedAt = codexTimestamp(started.turn.startedAt); - yield* registerRootTurn({ - turnInput, - nativeTurnId, - startedAt, - waitForNativeStart: started.turn.startedAt === null, - }); - yield* Ref.update(pendingRootTurns, (current) => { - const updated = new Map(current); - updated.delete(threadId); - return updated; - }); + : [], + ); }).pipe( - Effect.ensuring( - Effect.flatMap(getNativeThreadId(turnInput.providerThread), (threadId) => - Ref.update(pendingRootTurns, (current) => { - const updated = new Map(current); - updated.delete(threadId); - return updated; - }), - ).pipe(Effect.ignore), - ), Effect.mapError( (cause) => new ProviderAdapterTurnStartError({ @@ -5609,12 +6183,20 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi cause, }), ), - ), + ); + }, steerTurn: (turnInput) => Effect.gen(function* () { const threadId = yield* getNativeThreadId(turnInput.providerThread); + // Between goal turns, steer the turn Codex continues with. + const latestTurnId = latestGoalTurnId(turnInput.providerTurnId); + const held = heldGoalTurn(turnInput.providerThread, latestTurnId); + const providerTurnId = + held === undefined + ? latestTurnId + : ((yield* Deferred.await(held.hold.next))?.providerTurnId ?? latestTurnId); const activeTurn = Array.from((yield* Ref.get(activeTurns)).values()).find( - (candidate) => candidate.providerTurnId === turnInput.providerTurnId, + (candidate) => candidate.providerTurnId === providerTurnId, ); if (activeTurn === undefined) { return yield* toProtocolError( @@ -5657,8 +6239,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi }), ), ), - interruptTurn: (turnInput) => + interruptTurn: (requestedInput) => Effect.gen(function* () { + const { turnInput, goalTurns } = yield* resolveGoalStopTarget(requestedInput); const [activeTurnContexts, settledTurnContexts] = yield* turnTerminalizationPermit.withPermits(1)( Effect.gen(function* () { @@ -5677,28 +6260,36 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi (candidate) => candidate.providerThread.id === turnInput.providerThread.id, ) : undefined); - if (activeTurn === undefined) { - // Stop on a settled turn this process retains nothing for - // (released, restarted, or every command already reported). - if (turnInput.requestRuntimeRestart === true) return; - return yield* toProtocolError( - `Provider turn ${turnInput.providerTurnId} is not active and cannot be interrupted.`, + // A goal run spans several root turns, and work an earlier one + // started can outlive it. + const lineageRoots = Array.from( + new Set([...(activeTurn === undefined ? [] : [activeTurn]), ...goalTurns]), + ); + const inLineage = (context: ActiveCodexTurnContext) => + lineageRoots.some( + (root) => context === root || isDescendantCodexTurn(context, root), ); - } const interruptTargetContexts = [ - activeTurn, + ...(activeTurn === undefined ? [] : [activeTurn]), ...activeTurnContexts.filter( - (candidate) => - candidate !== activeTurn && isDescendantCodexTurn(candidate, activeTurn), + (candidate) => candidate !== activeTurn && inLineage(candidate), ), ...settledTurnContexts.filter( (candidate) => candidate !== activeTurn && ((turnInput.requestRuntimeRestart === true && candidate.providerThread.id === turnInput.providerThread.id) || - isDescendantCodexTurn(candidate, activeTurn)), + inLineage(candidate)), ), ]; + if (interruptTargetContexts.length === 0) { + // Stop on a settled turn this process retains nothing for + // (released, restarted, or every command already reported). + if (turnInput.requestRuntimeRestart === true) return; + return yield* toProtocolError( + `Provider turn ${turnInput.providerTurnId} is not active and cannot be interrupted.`, + ); + } const interruptTargets: Array<{ readonly context: ActiveCodexTurnContext; readonly completion: Deferred.Deferred; @@ -5761,8 +6352,7 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi const activeLineage = yield* turnTerminalizationPermit.withPermits(1)( Effect.gen(function* () { const lineage = Array.from((yield* Ref.get(activeTurns)).values()).filter( - (context) => - context === activeTurn || isDescendantCodexTurn(context, activeTurn), + inLineage, ); for (const context of lineage) { if (trackedInterruptNativeTurnIds.has(context.nativeTurnId)) { @@ -5839,8 +6429,8 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi Effect.gen(function* () { const lineage = Array.from((yield* Ref.get(activeTurns)).values()).filter( (context) => - context !== activeTurn && - isDescendantCodexTurn(context, activeTurn) && + !lineageRoots.includes(context) && + inLineage(context) && !trackedInterruptNativeTurnIds.has(context.nativeTurnId), ); for (const context of lineage) { @@ -6030,8 +6620,8 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi (cause) => new ProviderAdapterInterruptError({ driver: CODEX_PROVIDER, - providerThreadId: turnInput.providerThread.id, - providerTurnId: turnInput.providerTurnId, + providerThreadId: requestedInput.providerThread.id, + providerTurnId: requestedInput.providerTurnId, cause, }), ), diff --git a/apps/server/src/orchestration-v2/CheckpointRollbackService.ts b/apps/server/src/orchestration-v2/CheckpointRollbackService.ts index 37884d56b70d..1710e7a8cd84 100644 --- a/apps/server/src/orchestration-v2/CheckpointRollbackService.ts +++ b/apps/server/src/orchestration-v2/CheckpointRollbackService.ts @@ -1,6 +1,7 @@ import { CheckpointId, CheckpointScopeId, + latestProviderTurnForAttempt, type OrchestrationV2DomainEvent, ProviderThreadId, ThreadId, @@ -233,11 +234,10 @@ export const layer: Layer.Layer< const targetAttempt = projection.attempts.find( (attempt) => attempt.id === targetRun?.activeAttemptId, ); - const targetTurn = projection.providerTurns.find( - (turn) => - turn.id === targetAttempt?.providerTurnId || - turn.runAttemptId === targetAttempt?.id, - ); + // A goal run can span several native turns; roll back to its last. + const targetTurn = + latestProviderTurnForAttempt(projection.providerTurns, targetAttempt?.id) ?? + projection.providerTurns.find((turn) => turn.id === targetAttempt?.providerTurnId); if (targetTurn === undefined || targetTurn.providerThreadId !== providerThread.id) { return yield* new CheckpointRollbackExecutionError({ reason: "provider-turn-unavailable", diff --git a/apps/server/src/orchestration-v2/NotificationMailbox.ts b/apps/server/src/orchestration-v2/NotificationMailbox.ts index 38845466c142..e7045e28d9e1 100644 --- a/apps/server/src/orchestration-v2/NotificationMailbox.ts +++ b/apps/server/src/orchestration-v2/NotificationMailbox.ts @@ -1,4 +1,8 @@ -import type { MessageId, OrchestrationV2ThreadProjection } from "@t3tools/contracts"; +import { + latestProviderTurnForAttempt, + type MessageId, + type OrchestrationV2ThreadProjection, +} from "@t3tools/contracts"; /** * A persisted steer without an acceptance receipt can be retried as a continuation. @@ -16,8 +20,7 @@ export function isUndeliveredMailboxSteer( run !== undefined && run.userMessageId !== message.id && (["completed", "failed", "interrupted", "cancelled", "rolled_back"].includes(run.status) || - projection.providerTurns.some( - (turn) => turn.runAttemptId === run.activeAttemptId && turn.status === "completed", - )) + latestProviderTurnForAttempt(projection.providerTurns, run.activeAttemptId)?.status === + "completed") ); } diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 42c55371d5f7..3bb5bdd0e7c6 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -43,6 +43,7 @@ import { type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, type OrchestrationV2TurnItem, + latestProviderTurnForAttempt, orchestrationV2RunWorkStartedAt, ProviderInstanceId, type ProviderSessionId, @@ -364,6 +365,14 @@ export function isNativeMaintenanceCommand(message: { ); } +/** A native `/goal` command. It changes the provider's goal, so it never steers a running turn. */ +function isGoalCommand(message: { + readonly text: string; + readonly attachments: ReadonlyArray; +}): boolean { + return message.attachments.length === 0 && /^\/goal(?:\s|$)/u.test(message.text.trim()); +} + const threadPullRequestLinksEqual = Schema.toEquivalence(Schema.NullOr(ThreadLinkedPullRequest)); function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { @@ -594,7 +603,7 @@ function providerTurnForRun( } return ( - projection.providerTurns.find((turn) => turn.runAttemptId === run.activeAttemptId) ?? + latestProviderTurnForAttempt(projection.providerTurns, run.activeAttemptId) ?? projection.providerTurns.find((turn) => { const attempt = projection.attempts.find((candidate) => candidate.id === run.activeAttemptId); return attempt?.providerTurnId === turn.id; @@ -3631,6 +3640,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio : "Signing out must run as a separate turn. Queue it or wait for the active turn to finish.", }); } + if (isGoalCommand(input)) { + return yield* new OrchestratorDispatchError({ + commandId: input.command.commandId, + commandType: input.command.type, + cause: + "Goal commands must run as a separate turn. Queue it or wait for the active turn to finish.", + }); + } const targetMessage = input.projection.messages.find( (message) => message.id === targetRun.userMessageId, ); @@ -4492,10 +4509,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio if (dispatchMode.type === "steer_active") { const targetRunId = dispatchMode.targetRunId; const target = projection.runs.find((run) => run.id === targetRunId); - const turn = projection.providerTurns.find( - (candidate) => - candidate.runAttemptId === target?.activeAttemptId && - candidate.nodeId === target?.rootNodeId, + const turn = latestProviderTurnForAttempt( + projection.providerTurns.filter((candidate) => candidate.nodeId === target?.rootNodeId), + target?.activeAttemptId, ); // The client may still show a running turn while its completion is being // projected. Preserve the submission as a new turn when steering is too late. @@ -5402,9 +5418,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const sourceProviderTurnId = sourceProjection === null || sourceRun === null || sourceRun.activeAttemptId === null ? undefined - : (sourceProjection.providerTurns.find( - (candidate) => candidate.runAttemptId === sourceRun.activeAttemptId, - )?.id ?? + : (latestProviderTurnForAttempt(sourceProjection.providerTurns, sourceRun.activeAttemptId) + ?.id ?? sourceProjection.attempts.find( (candidate) => candidate.id === sourceRun.activeAttemptId, )?.providerTurnId ?? diff --git a/apps/server/src/orchestration-v2/ProjectionStore.test.ts b/apps/server/src/orchestration-v2/ProjectionStore.test.ts index 3e979544c87b..82cdfd996200 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.test.ts @@ -8,6 +8,7 @@ import { CheckpointScopeId, MessageId, type ModelSelection, + type OrchestrationV2ProviderThread, NodeId, ProjectId, ProviderDriverKind, @@ -1596,6 +1597,84 @@ it.layer(TestLayer)("ProjectionStoreV2", (it) => { }), ); + it.effect("shows the native goal of the active provider thread on the shell", () => + Effect.gen(function* () { + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:provider-goal"); + yield* projectionStore.apply({ + id: EventId.make("event:provider-goal:thread"), + type: "thread.created", + threadId, + occurredAt: now, + payload: { + createdBy: "user", + creationSource: "web", + id: threadId, + projectId: ProjectId.make("project:provider-goal"), + title: "Provider goal", + providerInstanceId, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: null, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }, + }); + const applyProviderThread = ( + suffix: string, + goal: OrchestrationV2ProviderThread["goal"], + seconds: number, + ) => + projectionStore.apply({ + id: EventId.make(`event:provider-goal:${suffix}:${seconds}`), + type: "provider-thread.updated", + threadId, + driver, + occurredAt: DateTime.add(now, { seconds }), + payload: { + id: ProviderThreadId.make(`provider-thread:provider-goal:${suffix}`), + driver, + providerInstanceId, + providerSessionId: null, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "idle", + firstRunOrdinal: null, + lastRunOrdinal: null, + handoffIds: [], + forkedFrom: null, + goal, + createdAt: now, + updatedAt: DateTime.add(now, { seconds }), + }, + }); + const shellGoal = Effect.map( + projectionStore.getShellSnapshot(), + (snapshot) => snapshot.threads.find((thread) => thread.id === threadId)?.goal, + ); + const goal = { objective: "Ship it", status: "active" as const, tokensUsed: 10 }; + + yield* applyProviderThread("first", goal, 0); + assert.deepEqual(yield* shellGoal, goal); + // A handoff moves the conversation; the previous provider's goal stays behind. + yield* applyProviderThread("second", null, 1); + assert.isNull(yield* shellGoal); + }), + ); + it.effect("does not treat visited or marked-unread state as thread activity", () => Effect.gen(function* () { const projectionStore = yield* ProjectionStore.ProjectionStoreV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 93f1e1e90de2..f7bac27ef363 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -1405,6 +1405,10 @@ export function threadShellFromProjection( threadId: projection.thread.id, providerThreads: projection.providerThreads, }), + goal: activeProviderGoalForShell( + projection.providerThreads, + projection.thread.activeProviderThreadId, + ), itemCount: activeLocalTurnItems(projection).length, visibleItemCount: projection.visibleTurnItems.length, createdAt: projection.thread.createdAt, @@ -1452,6 +1456,14 @@ function providerInstanceHistoryForShell(input: { return history; } +/** The native goal lives on the provider thread that currently owns the conversation. */ +function activeProviderGoalForShell( + providerThreads: ReadonlyArray, + activeProviderThreadId: OrchestrationV2ThreadProjection["thread"]["activeProviderThreadId"], +): OrchestrationV2ThreadShell["goal"] { + return providerThreads.find((thread) => thread.id === activeProviderThreadId)?.goal ?? null; +} + function isInterruptibleRunForShell(run: OrchestrationV2ThreadProjection["runs"][number]): boolean { return run.status === "preparing" || run.status === "starting" || run.status === "running"; } @@ -1485,6 +1497,7 @@ type ShellThreadState = { readonly hasActionableProposedPlan: boolean; readonly pendingBackgroundTasks: OrchestrationV2ThreadShell["pendingBackgroundTasks"]; readonly providerInstanceHistory: OrchestrationV2ThreadShell["providerInstanceHistory"]; + readonly goal: OrchestrationV2ThreadShell["goal"]; readonly itemCount: number; readonly runlessItemCount: number; readonly updatedAt: OrchestrationV2ThreadProjection["updatedAt"]; @@ -1631,6 +1644,7 @@ function shellFromState(input: { hasActionableProposedPlan: input.state.hasActionableProposedPlan, pendingBackgroundTasks: input.state.pendingBackgroundTasks, providerInstanceHistory: input.state.providerInstanceHistory, + goal: input.state.goal, itemCount: input.state.itemCount, visibleItemCount: input.visibleItemCount, createdAt: input.state.thread.createdAt, @@ -5347,6 +5361,10 @@ export const layer: Layer.Layer = threadId: thread.id, providerThreads: providerThreadsByThreadId.get(thread.id) ?? [], }), + goal: activeProviderGoalForShell( + providerThreadsByThreadId.get(thread.id) ?? [], + thread.activeProviderThreadId, + ), itemCount: row.item_count, runlessItemCount: row.runless_item_count, updatedAt: thread.updatedAt, diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index a3e3b122ca45..19c97ae9b0f1 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -2,6 +2,7 @@ import { modelSelectionsEqual } from "@t3tools/shared/model"; import { projectComposerContextForProvider } from "@t3tools/shared/composerContextReferences"; import { CommandId, + latestProviderTurnForAttempt, type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderThread, @@ -617,11 +618,11 @@ export const layer: Layer.Layer< const sourceAttempt = sourceProjection.attempts.find( (candidate) => candidate.id === sourceRun?.activeAttemptId, ); - const sourceProviderTurn = sourceProjection.providerTurns.find( - (candidate) => - candidate.id === sourceAttempt?.providerTurnId || - candidate.runAttemptId === sourceAttempt?.id, - ); + const sourceProviderTurn = + latestProviderTurnForAttempt(sourceProjection.providerTurns, sourceAttempt?.id) ?? + sourceProjection.providerTurns.find( + (candidate) => candidate.id === sourceAttempt?.providerTurnId, + ); if (sourceRun === undefined || sourceProviderThread === undefined) { return yield* new ProviderTurnStartError({ runId, diff --git a/apps/server/src/provider/Layers/CodexProvider.ts b/apps/server/src/provider/Layers/CodexProvider.ts index dcb606932ab9..5bdc389ff6ed 100644 --- a/apps/server/src/provider/Layers/CodexProvider.ts +++ b/apps/server/src/provider/Layers/CodexProvider.ts @@ -691,6 +691,11 @@ export const checkCodexProviderStatus = Effect.fn("checkCodexProviderStatus")(fu skills: snapshot.skills, slashCommands: [ COMPACT_SLASH_COMMAND, + { + name: "goal", + description: "Set a goal Codex keeps working toward until it is done", + input: { hint: "Objective, or pause, resume, clear" }, + }, { name: "feedback", description: "Send this thread and Codex logs to OpenAI", diff --git a/apps/server/src/provider/Layers/ProviderRegistry.test.ts b/apps/server/src/provider/Layers/ProviderRegistry.test.ts index 332a5e931d29..5f1eae506b05 100644 --- a/apps/server/src/provider/Layers/ProviderRegistry.test.ts +++ b/apps/server/src/provider/Layers/ProviderRegistry.test.ts @@ -446,6 +446,11 @@ it.layer(Layer.mergeAll(TestNodeServices, ServerSettingsModule.layerTest(), Test }, ]); assert.deepStrictEqual(status.slashCommands.slice(1), [ + { + name: "goal", + description: "Set a goal Codex keeps working toward until it is done", + input: { hint: "Objective, or pause, resume, clear" }, + }, { name: "feedback", description: "Send this thread and Codex logs to OpenAI", diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index e001dad666d8..6fb6d4f3fb44 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -97,6 +97,7 @@ import { deriveLatestThreadRun, deriveThreadRuntime, presentPendingBackgroundWork, + presentProviderGoal, } from "@t3tools/client-runtime/state/thread-execution"; import { threadSupportsProviderHandoff } from "@t3tools/client-runtime/state/thread-workflows"; import { @@ -292,6 +293,7 @@ import { ChevronDownIcon, DownloadIcon, GitBranchIcon, + TargetIcon, WifiOffIcon, } from "lucide-react"; import { cn, randomUUID } from "~/lib/utils"; @@ -7106,6 +7108,130 @@ export default function ChatView(props: ChatViewProps) { [environmentId, navigate], ); + // Commands such as /compact and /goal clear run as their own turn. The draft + // and its attachments stay local. + const sendStandaloneCommand = useCallback( + async (text: string, failureMessage: string) => { + if (!activeThread || !clientSettingsHydrated || sendInFlightRef.current) return; + const context = composerRef.current?.getSendContext(); + if (!context?.providerAvailable) return; + + const threadId = activeThread.id; + const messageId = newMessageId(); + const createdAt = new Date().toISOString(); + sendInFlightRef.current = true; + beginLocalDispatch(); + setThreadError(threadId, null); + setOptimisticUserMessages((messages) => [ + ...messages, + { + id: messageId, + role: "user", + text, + runId: null, + createdAt, + updatedAt: createdAt, + streaming: false, + }, + ]); + scrollToEnd(); + try { + const settingsResult = await persistThreadSettingsForNextTurn({ + threadId, + createdAt, + ...(localCheckoutBranchMismatch + ? { branch: localCheckoutBranchMismatch.currentBranch } + : {}), + runtimeMode, + interactionMode: context.interactionMode, + }); + const result = + settingsResult._tag === "Failure" + ? settingsResult + : await startThreadTurn({ + environmentId, + input: { + threadId, + message: { messageId, role: "user", text, attachments: [] }, + modelSelection: context.selectedModelSelection, + runtimeMode, + interactionMode: context.interactionMode, + createdAt, + }, + }); + if (result._tag === "Failure") { + setOptimisticUserMessages((messages) => + messages.filter((message) => message.id !== messageId), + ); + resetLocalDispatch(); + if (!isAtomCommandInterrupted(result)) { + const error = squashAtomCommandFailure(result); + setThreadError(threadId, error instanceof Error ? error.message : failureMessage); + } + } else { + clearUsageLimitsFor(routeThreadKey); + } + } finally { + sendInFlightRef.current = false; + } + }, + [ + activeThread, + beginLocalDispatch, + clientSettingsHydrated, + composerRef, + environmentId, + localCheckoutBranchMismatch, + persistThreadSettingsForNextTurn, + resetLocalDispatch, + routeThreadKey, + runtimeMode, + scrollToEnd, + sendInFlightRef, + setThreadError, + startThreadTurn, + ], + ); + // A native /goal keeps the agent working across turns. Stop pauses a Codex + // goal; once the thread is idle the row offers the native follow-ups. + const activeGoal = activeThreadShell?.goal ?? null; + const goalBannerItem = useMemo(() => { + if (!activeThread || activeGoal === null) return null; + const presentation = presentProviderGoal(activeGoal, isWorking); + return { + id: `goal:${activeThread.id}`, + variant: "default", + priority: "activity", + icon: , + // Usage stays in the title so a long objective cannot clip it. + title: + presentation.usage === null + ? presentation.title + : `${presentation.title} · ${presentation.usage}`, + description: presentation.objective, + actions: isWorking ? undefined : ( + <> + {presentation.canResume ? ( + + ) : null} + + + ), + }; + }, [activeGoal, activeThread, isWorking, sendStandaloneCommand]); + const backgroundWorkBannerItem = useMemo(() => { const presentation = presentPendingBackgroundWork(activeBackgroundTasks); if (presentation === null || !activeThread) { @@ -7372,7 +7498,9 @@ export default function ChatView(props: ChatViewProps) { : null; const composerBannerItems = useMemo(() => { const limitRecoveryItems = limitRecoveryBanner === null ? [] : [limitRecoveryBanner]; - const backgroundWorkItems = backgroundWorkBannerItem === null ? [] : [backgroundWorkBannerItem]; + const backgroundWorkItems = [goalBannerItem, backgroundWorkBannerItem].filter( + (item) => item !== null, + ); const resumeCompactionItems = resumeCompactionBannerItem === null ? [] : [resumeCompactionBannerItem]; const wokeThreadItems = wokeThreadBannerItem === null ? [] : [wokeThreadBannerItem]; @@ -7451,6 +7579,7 @@ export default function ChatView(props: ChatViewProps) { handleRestoreThreadBranch, isRestoringThreadBranch, backgroundWorkBannerItem, + goalBannerItem, localCheckoutBranchMismatch, parkedThreadBannerItem, projectCloneBannerItem, @@ -8168,75 +8297,9 @@ export default function ChatView(props: ChatViewProps) { setThreadError, ], ); - const onCompactContext = async () => { - if (compactDisabled || !activeThread || !clientSettingsHydrated || sendInFlightRef.current) { - return; - } - const context = composerRef.current?.getSendContext(); - if (!context?.providerAvailable) return; - - // Compaction is a standalone command; the draft and its attachments stay local. - const threadId = activeThread.id; - const messageId = newMessageId(); - const createdAt = new Date().toISOString(); - sendInFlightRef.current = true; - beginLocalDispatch(); - setThreadError(threadId, null); - setOptimisticUserMessages((messages) => [ - ...messages, - { - id: messageId, - role: "user", - text: "/compact", - runId: null, - createdAt, - updatedAt: createdAt, - streaming: false, - }, - ]); - scrollToEnd(); - try { - const settingsResult = await persistThreadSettingsForNextTurn({ - threadId, - createdAt, - ...(localCheckoutBranchMismatch - ? { branch: localCheckoutBranchMismatch.currentBranch } - : {}), - runtimeMode, - interactionMode: context.interactionMode, - }); - const result = - settingsResult._tag === "Failure" - ? settingsResult - : await startThreadTurn({ - environmentId, - input: { - threadId, - message: { messageId, role: "user", text: "/compact", attachments: [] }, - modelSelection: context.selectedModelSelection, - runtimeMode, - interactionMode: context.interactionMode, - createdAt, - }, - }); - if (result._tag === "Failure") { - setOptimisticUserMessages((messages) => - messages.filter((message) => message.id !== messageId), - ); - resetLocalDispatch(); - if (!isAtomCommandInterrupted(result)) { - const error = squashAtomCommandFailure(result); - setThreadError( - threadId, - error instanceof Error ? error.message : "Failed to compact context.", - ); - } - } else { - clearUsageLimitsFor(routeThreadKey); - } - } finally { - sendInFlightRef.current = false; - } + const onCompactContext = () => { + if (compactDisabled) return; + void sendStandaloneCommand("/compact", "Failed to compact context."); }; const onResume = async () => { diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index 3ee0f867cae3..eae15550d7b8 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -1258,7 +1258,8 @@ const SidebarThreadRow = memo(function SidebarThreadRow(props: { const topStatus = status === "working" ? { - label: "Working", + // A native /goal keeps the agent going across turns until it is met. + label: thread.goal?.status === "active" ? "Goal" : "Working", icon: "working" as const, // No shimmer: a label that animates forever is noise in a sidebar // full of them (and repaints every vsync on high-refresh displays). diff --git a/docs/internals/providers.md b/docs/internals/providers.md index 52003163e882..a619fbd0dcec 100644 --- a/docs/internals/providers.md +++ b/docs/internals/providers.md @@ -116,6 +116,14 @@ resumes a run, queues behind active work, or steers when the adapter supports it retain the provider's live response path. Do not infer that a request has disappeared merely because it is outside the recent history window. +Native `/goal` state belongs to the provider and is mirrored on the provider thread. Codex starts +the next goal turn on its own milliseconds after the last one completes, so the +[adapter](../../apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts) keeps the run open +and adds that turn to it. One run can therefore own several native turns. `/goal` commands that +start no Codex turn settle on a provider turn without a native ref, which native rollback must +not count. Claude's SDK mode emits no goal events; its adapter reads goal state from the +synthetic command output and Stop hook feedback in the transcript. + Capabilities must describe what the provider can actually do. Antigravity can capture workspace checkpoints but cannot roll back its conversation. The [checkpoint boundary](./overview.md#turn-completion-and-checkpoints) therefore rejects revert before touching files. Native permission and question option IDs must diff --git a/docs/user/composer.md b/docs/user/composer.md index 8cb8d7805ac9..f4c2829019a7 100644 --- a/docs/user/composer.md +++ b/docs/user/composer.md @@ -193,6 +193,19 @@ Provider commands must start the message to run. T3 Code commands such as Send `/compact` in an existing conversation to reduce context usage when the provider supports it. Web and desktop also offer compaction from the context meter. +## Goals + +With Codex and Claude, send `/goal` followed by what "done" means, for example +`/goal all tests in packages/api pass`. The agent keeps working across turns +until it judges the goal met. The thread shows **Goal** while it works, and a +row above the composer shows the goal and its progress. + +- `/goal` alone shows the current goal. `/goal clear` removes it. +- Codex also supports `/goal pause` and `/goal resume`. Stopping a Codex goal + pauses it. +- Stopping Claude ends the current turn, but the goal stays set. Claude checks it + again at the end of your next message. + ## Context in your message Context you attach lands where your cursor is, as a chip inside your text: a terminal excerpt, diff --git a/packages/client-runtime/src/state/models.ts b/packages/client-runtime/src/state/models.ts index 5b319eb10627..55fa5a673616 100644 --- a/packages/client-runtime/src/state/models.ts +++ b/packages/client-runtime/src/state/models.ts @@ -5,6 +5,7 @@ import type { EnvironmentId, MessageId, OrchestrationProjectShell, + OrchestrationV2ProviderGoal, OrchestrationV2RunStatus, OrchestrationV2ProviderFailureClass, OrchestrationV2ThreadProjection, @@ -112,6 +113,8 @@ export interface EnvironmentThreadShell { >; /** Provider instances that have owned the root conversation, oldest first. */ readonly providerInstanceHistory: ReadonlyArray; + /** Native `/goal` on the active provider thread. */ + readonly goal: OrchestrationV2ProviderGoal | null; readonly itemCount: number; readonly visibleItemCount: number; readonly createdAt: string; @@ -254,6 +257,7 @@ export function presentThreadShell( hasActionableProposedPlan: thread.hasActionableProposedPlan, pendingBackgroundTasks: thread.pendingBackgroundTasks ?? [], providerInstanceHistory: thread.providerInstanceHistory ?? [], + goal: thread.goal ?? null, itemCount: thread.itemCount, visibleItemCount: thread.visibleItemCount, createdAt: iso(thread.createdAt), diff --git a/packages/client-runtime/src/state/threadExecution.test.ts b/packages/client-runtime/src/state/threadExecution.test.ts index 6ee64d2eaf66..ba426b8195fe 100644 --- a/packages/client-runtime/src/state/threadExecution.test.ts +++ b/packages/client-runtime/src/state/threadExecution.test.ts @@ -18,6 +18,7 @@ import { describe, expect, it } from "vite-plus/test"; import { v2Projection } from "./orchestrationV2TestFixtures.ts"; import { presentPendingBackgroundWork, + presentProviderGoal, deriveReportedModelSelection, deriveLatestThreadRun, deriveProviderSubagentStatus, @@ -677,3 +678,43 @@ describe("provider-reported model selection", () => { expect(formatModelSelectionEffort(selected, models)).toBe("Unknown"); }); }); + +describe("presentProviderGoal", () => { + it("summarizes Codex accounting and offers resume once the goal stops short", () => { + expect( + presentProviderGoal( + { + objective: "Ship the feature", + status: "paused", + tokensUsed: 12_400, + tokenBudget: 50_000, + timeUsedSeconds: 245, + }, + false, + ), + ).toEqual({ + title: "Goal paused", + objective: "Ship the feature", + usage: "12k / 50k tokens · 4m 5s", + canResume: true, + }); + }); + + it("counts Claude evaluator checks and never offers resume", () => { + expect( + presentProviderGoal({ objective: "All tests pass", status: "active", checks: 1 }, true), + ).toEqual({ + title: "Pursuing goal", + objective: "All tests pass", + usage: "1 check", + canResume: false, + }); + expect( + presentProviderGoal({ objective: "All tests pass", status: "complete", checks: 0 }, false) + .usage, + ).toBeNull(); + expect( + presentProviderGoal({ objective: "All tests pass", status: "active" }, false).title, + ).toBe("Goal set"); + }); +}); diff --git a/packages/client-runtime/src/state/threadExecution.ts b/packages/client-runtime/src/state/threadExecution.ts index a7722af0391b..9139c1004e4e 100644 --- a/packages/client-runtime/src/state/threadExecution.ts +++ b/packages/client-runtime/src/state/threadExecution.ts @@ -10,6 +10,7 @@ import { type ModelSelection, type OrchestrationV2NotificationSource, type OrchestrationV2PendingBackgroundTask, + type OrchestrationV2ProviderGoal, type ServerProviderModel, type OrchestrationV2ExecutionNode, type OrchestrationV2ThreadProjection, @@ -369,6 +370,61 @@ export function presentPendingBackgroundWork( return { title: `${waiting ? "Waiting on" : "Running"} ${joinWithAnd(groups)}`, items, waiting }; } +export interface ProviderGoalPresentation { + /** "Pursuing goal", "Goal paused", "Goal complete"... */ + readonly title: string; + readonly objective: string; + /** "12k / 50k tokens · 4m" for Codex, "2 checks" for Claude; null when nothing is reported. */ + readonly usage: string | null; + /** Only Codex stops a goal short of done, and only a stopped goal can resume. */ + readonly canResume: boolean; +} + +const PROVIDER_GOAL_TITLES: Record = { + active: "Pursuing goal", + paused: "Goal paused", + blocked: "Goal blocked", + usage_limited: "Goal hit a usage limit", + budget_limited: "Goal reached its token budget", + complete: "Goal complete", +}; + +function formatGoalTokens(tokens: number): string { + if (tokens < 1_000) return `${tokens}`; + if (tokens < 1_000_000) return `${Math.round(tokens / 1_000)}k`; + return `${(tokens / 1_000_000).toFixed(1).replace(/\.0$/, "")}m`; +} + +/** + * Status line for a native `/goal`, shared by the web composer row and mobile. + * An active goal on an idle thread (Claude after Stop) is set, not pursued. + */ +export function presentProviderGoal( + goal: OrchestrationV2ProviderGoal, + working: boolean, +): ProviderGoalPresentation { + const usage: Array = []; + if (goal.tokensUsed !== undefined && goal.tokensUsed > 0) { + usage.push( + goal.tokenBudget == null + ? `${formatGoalTokens(goal.tokensUsed)} tokens` + : `${formatGoalTokens(goal.tokensUsed)} / ${formatGoalTokens(goal.tokenBudget)} tokens`, + ); + } + if (goal.timeUsedSeconds !== undefined && goal.timeUsedSeconds >= 60) { + usage.push(formatDuration(goal.timeUsedSeconds * 1_000)); + } + if (goal.checks !== undefined && goal.checks > 0) { + usage.push(`${goal.checks} ${goal.checks === 1 ? "check" : "checks"}`); + } + return { + title: goal.status === "active" && !working ? "Goal set" : PROVIDER_GOAL_TITLES[goal.status], + objective: goal.objective, + usage: usage.length === 0 ? null : usage.join(" · "), + canResume: goal.status !== "active" && goal.status !== "complete", + }; +} + /** The thread a notification row opens: that of the one subagent or delegated task it reports. */ export function notificationChildThreadId( source: OrchestrationV2NotificationSource, diff --git a/packages/contracts/src/orchestrationV2.test.ts b/packages/contracts/src/orchestrationV2.test.ts index 601597fb568a..bb34c3ce998b 100644 --- a/packages/contracts/src/orchestrationV2.test.ts +++ b/packages/contracts/src/orchestrationV2.test.ts @@ -16,12 +16,14 @@ import { ProviderInstanceId, ProviderReplayTranscript, ProviderThreadId, + RunAttemptId, RunId, ThreadId, TrimmedNonEmptyString, TurnItemId, } from "./index.ts"; import { + latestProviderTurnForAttempt, OrchestrationV2Checkpoint, OrchestrationV2CheckpointScope, OrchestrationV2Command, @@ -1248,3 +1250,16 @@ describe("limit recovery choice updates", () => { expect(decode({ ...identity, ...choice })).toEqual({ ...identity, ...choice }); }); }); + +describe("latestProviderTurnForAttempt", () => { + it("returns the attempt's highest-ordinal turn, as a Codex goal run spans several", () => { + const turns = [ + { id: "first", runAttemptId: RunAttemptId.make("goal-attempt"), ordinal: 3 }, + { id: "other", runAttemptId: RunAttemptId.make("other-attempt"), ordinal: 9 }, + { id: "last", runAttemptId: RunAttemptId.make("goal-attempt"), ordinal: 5 }, + { id: "subagent", runAttemptId: null, ordinal: 7 }, + ]; + expect(latestProviderTurnForAttempt(turns, RunAttemptId.make("goal-attempt"))?.id).toBe("last"); + expect(latestProviderTurnForAttempt(turns, null)).toBeUndefined(); + }); +}); diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index c6e7f3ff0140..afbbff67b509 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -832,6 +832,34 @@ export const OrchestrationV2ProviderThreadNativeMetadata = Schema.Struct({ export type OrchestrationV2ProviderThreadNativeMetadata = typeof OrchestrationV2ProviderThreadNativeMetadata.Type; +export const OrchestrationV2ProviderGoalStatus = Schema.Literals([ + "active", + "paused", + "blocked", + "usage_limited", + "budget_limited", + "complete", +]); +export type OrchestrationV2ProviderGoalStatus = typeof OrchestrationV2ProviderGoalStatus.Type; + +/** + * A provider-native goal set with `/goal` (Codex and Claude). The provider + * keeps working until it judges the objective met and owns this state; T3 + * mirrors the latest native report. Usage fields are provider-specific. + */ +export const OrchestrationV2ProviderGoal = Schema.Struct({ + objective: TrimmedNonEmptyString, + status: OrchestrationV2ProviderGoalStatus, + /** Codex accounting for the goal across its turns. */ + tokensUsed: Schema.optional(NonNegativeInt), + tokenBudget: Schema.optional(Schema.NullOr(NonNegativeInt)), + timeUsedSeconds: Schema.optional(NonNegativeInt), + /** Claude: evaluator checks that found the goal unmet, and the latest reason. */ + checks: Schema.optional(NonNegativeInt), + lastCheck: Schema.optional(Schema.String), +}); +export type OrchestrationV2ProviderGoal = typeof OrchestrationV2ProviderGoal.Type; + export const OrchestrationV2ProviderThread = Schema.Struct({ id: ProviderThreadId, driver: ProviderDriverKind, @@ -862,6 +890,10 @@ export const OrchestrationV2ProviderThread = Schema.Struct({ nativeMetadata: Schema.optional(Schema.NullOr(OrchestrationV2ProviderThreadNativeMetadata)).pipe( Schema.withDecodingDefault(Effect.succeed(null)), ), + /** Native goal on this provider thread; rows written before goals decode to null. */ + goal: Schema.optional(Schema.NullOr(OrchestrationV2ProviderGoal)).pipe( + Schema.withDecodingDefault(Effect.succeed(null)), + ), createdAt: Schema.DateTimeUtc, updatedAt: Schema.DateTimeUtc, }); @@ -963,6 +995,25 @@ export const OrchestrationV2ProviderTurn = Schema.Struct({ }); export type OrchestrationV2ProviderTurn = typeof OrchestrationV2ProviderTurn.Type; +/** + * The provider turn a run attempt is on or ended with. A Codex goal keeps one + * run open across several native turns, so an attempt's first turn is not + * always its last. + */ +export function latestProviderTurnForAttempt< + Turn extends Pick, +>( + providerTurns: ReadonlyArray, + attemptId: RunAttemptId | null | undefined, +): Turn | undefined { + let latest: Turn | undefined; + for (const turn of providerTurns) { + if (attemptId == null || turn.runAttemptId !== attemptId) continue; + if (latest === undefined || turn.ordinal > latest.ordinal) latest = turn; + } + return latest; +} + export const OrchestrationV2RuntimeRequest = Schema.Struct({ id: RuntimeRequestId, nodeId: NodeId, @@ -1764,6 +1815,8 @@ export const OrchestrationV2ThreadShell = Schema.Struct({ providerInstanceHistory: Schema.optional(Schema.Array(ProviderInstanceId)).pipe( Schema.withDecodingDefault(Effect.succeed([])), ), + /** Native goal on the active provider thread; omitted by servers without goals. */ + goal: Schema.optional(Schema.NullOr(OrchestrationV2ProviderGoal)), itemCount: NonNegativeInt, visibleItemCount: NonNegativeInt, createdAt: Schema.DateTimeUtc,