diff --git a/.changeset/per-step-review-session-switch.md b/.changeset/per-step-review-session-switch.md new file mode 100644 index 0000000000..1b28d9d030 --- /dev/null +++ b/.changeset/per-step-review-session-switch.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Align per-step review coding with default coding gates and session settings. +category: fix +dev: Removes the extra generic review seam from Coding (per-step review) and makes StepSessionExecutor honor runStepsInNewSessions=false by reusing the primary sequential session while keeping graph step-review boundaries. diff --git a/packages/core/src/__tests__/builtin-workflows.test.ts b/packages/core/src/__tests__/builtin-workflows.test.ts index 066adf2215..8cd0e33c25 100644 --- a/packages/core/src/__tests__/builtin-workflows.test.ts +++ b/packages/core/src/__tests__/builtin-workflows.test.ts @@ -220,7 +220,8 @@ describe("built-in workflows", () => { expect(ir.edges.some((edge) => edge.from === "steps" && edge.to === "browser-verification" && edge.condition === "success")).toBe(true); expect(ir.edges.some((edge) => edge.from === "browser-verification" && edge.to === "code-review" && edge.condition === "success")).toBe(true); expect(ir.edges.some((edge) => edge.from === "code-review" && edge.to === "completion-summary" && edge.condition === "success")).toBe(true); - expect(ir.edges.some((edge) => edge.from === "completion-summary" && edge.to === "review" && edge.condition === "success")).toBe(true); + expect(ir.edges.some((edge) => edge.from === "completion-summary" && edge.to === "merge-gate" && edge.condition === "success")).toBe(true); + expect(ir.nodes.some((node) => node.id === "review")).toBe(false); const foreach = ir.nodes.find((n) => n.kind === "foreach"); expect(foreach).toBeDefined(); const template = ( @@ -883,7 +884,7 @@ describe("built-in workflows", () => { expect(ce.ir.edges.some((edge) => edge.from === "manual-pr-review" && edge.to === "completion-summary")).toBe(true); }); - it("non-default coding built-ins retain their generic review nodes", () => { + it("coding variants only retain generic review nodes where the workflow requires them", () => { const coding = getBuiltinWorkflow("builtin:coding")!; const legacy = getBuiltinWorkflow("builtin:legacy-coding")!; const stepwise = getBuiltinWorkflow("builtin:stepwise-coding")!; @@ -891,7 +892,7 @@ describe("built-in workflows", () => { expect(coding.ir.nodes.some((node) => node.id === "review" && node.config?.seam === "review")).toBe(false); expect(legacy.ir.nodes.some((node) => node.id === "review" && node.config?.seam === "review")).toBe(true); - expect(stepwise.ir.nodes.some((node) => node.id === "review" && node.config?.seam === "review")).toBe(true); + expect(stepwise.ir.nodes.some((node) => node.id === "review" && node.config?.seam === "review")).toBe(false); expect(reviewHeavy.ir.nodes.some((node) => node.id === "review" && node.config?.seam === "review")).toBe(true); }); diff --git a/packages/core/src/builtin-stepwise-coding-workflow-ir.ts b/packages/core/src/builtin-stepwise-coding-workflow-ir.ts index dcf503caf6..8d427771fc 100644 --- a/packages/core/src/builtin-stepwise-coding-workflow-ir.ts +++ b/packages/core/src/builtin-stepwise-coding-workflow-ir.ts @@ -50,6 +50,11 @@ import { * Coding (per-step review) also needs the default-on optional Plan Review gate before * execution. The group sits between `plan` and `parse` so operators can toggle plan * review independently while preserving the per-step code review/rework loop. + * + * FNXC:WorkflowReviewGates 2026-06-29-23:27: + * Per-step review should inherit the regular Coding workflow's graph-native suffix: + * completion summary flows directly to the merge gate, with the final optional Code + * Review group providing the only end-of-task automated review gate. */ const RAW_BUILTIN_STEPWISE_CODING_WORKFLOW_IR: WorkflowIr = { version: "v2", @@ -161,7 +166,6 @@ const RAW_BUILTIN_STEPWISE_CODING_WORKFLOW_IR: WorkflowIr = { codeReviewOptionalGroupNode("in-progress"), codeReviewRemediationNode("in-progress"), completionSummaryNode("in-review"), - { id: "review", kind: "prompt", column: "in-review", config: builtinPromptConfig("review", "Review") }, { id: "merge-gate", kind: "merge-gate", column: "in-review", config: { gate: "auto-merge" } }, { id: "merge-retry", kind: "retry-backoff", column: "in-review", config: { policy: "merge", maxAttempts: 3 } }, { id: "merge-manual-hold", kind: "manual-merge-hold", column: "in-review", config: { release: "manual" } }, @@ -203,18 +207,17 @@ const RAW_BUILTIN_STEPWISE_CODING_WORKFLOW_IR: WorkflowIr = { // KTD-5: bounded rework exhaustion → manual hold; release re-enters the group. { from: "steps", to: "rework-hold", condition: "outcome:rework-exhausted" }, { from: "rework-hold", to: "browser-verification", condition: "success" }, - // browser-verification → code-review → review; each optional-group passes through - // (outcome=success) when disabled, so a task with both off routes straight to review. + // browser-verification → code-review → completion-summary → merge-gate; each + // optional-group passes through (outcome=success) when disabled, so a task with + // both off routes straight to completion summary and merge policy. { from: "browser-verification", to: "code-review", condition: "success" }, { from: "code-review", to: "completion-summary", condition: "success" }, - { from: "completion-summary", to: "review", condition: "success" }, + { from: "completion-summary", to: "merge-gate", condition: "success" }, { from: "browser-verification", to: "browser-verification-remediation", condition: "failure" }, { from: "browser-verification-remediation", to: "browser-verification", condition: "success", kind: "rework" }, { from: "code-review", to: "code-review-remediation", condition: "failure" }, { from: "code-review-remediation", to: "code-review", condition: "success", kind: "rework" }, { from: "steps", to: "end", condition: "failure" }, - { from: "review", to: "merge-gate", condition: "success" }, - { from: "review", to: "end", condition: "failure" }, { from: "merge-gate", to: "branch-group-member-integration", condition: "outcome:auto-on" }, { from: "merge-gate", to: "merge-manual-hold", condition: "outcome:auto-off" }, { from: "merge-retry", to: "merge-attempt", condition: "success", kind: "rework" }, diff --git a/packages/core/src/builtin-workflows.ts b/packages/core/src/builtin-workflows.ts index 39629e03d2..21e6d7db5b 100644 --- a/packages/core/src/builtin-workflows.ts +++ b/packages/core/src/builtin-workflows.ts @@ -566,16 +566,15 @@ export const BUILTIN_WORKFLOWS: WorkflowDefinition[] = [ "code-review": { x: 1080, y: 160 }, "code-review-remediation": { x: 1080, y: 320 }, "completion-summary": { x: 1250, y: 160 }, - review: { x: 1420, y: 160 }, - "merge-gate": { x: 1590, y: 160 }, - "branch-group-member-integration": { x: 1760, y: 80 }, - "branch-group-promotion": { x: 1930, y: 80 }, - "merge-attempt": { x: 2100, y: 160 }, - "merge-retry": { x: 2270, y: 80 }, - "recovery-router": { x: 2270, y: 240 }, - "merge-manual-hold": { x: 1760, y: 240 }, - "post-merge-verification": { x: 2440, y: 160 }, - end: { x: 2610, y: 160 }, + "merge-gate": { x: 1420, y: 160 }, + "branch-group-member-integration": { x: 1590, y: 80 }, + "branch-group-promotion": { x: 1760, y: 80 }, + "merge-attempt": { x: 1930, y: 160 }, + "merge-retry": { x: 2100, y: 80 }, + "recovery-router": { x: 2100, y: 240 }, + "merge-manual-hold": { x: 1590, y: 240 }, + "post-merge-verification": { x: 2270, y: 160 }, + end: { x: 2440, y: 160 }, }, createdAt: BUILTIN_TS, updatedAt: BUILTIN_TS, diff --git a/packages/engine/src/__tests__/step-session-executor.test.ts b/packages/engine/src/__tests__/step-session-executor.test.ts index 6bb71e32e7..ee9af1ccba 100644 --- a/packages/engine/src/__tests__/step-session-executor.test.ts +++ b/packages/engine/src/__tests__/step-session-executor.test.ts @@ -1192,6 +1192,89 @@ describe("StepSessionExecutor", () => { ); }); + it("reuses the primary step session when runStepsInNewSessions is false", async () => { + const prompt = makeStepPrompt("FN-001", 2); + const task = makeTaskDetail({ + prompt, + steps: [ + { name: "Step 0", status: "pending" }, + { name: "Step 1", status: "pending" }, + ], + }); + const settings = makeSettings({ maxParallelSteps: 1, runStepsInNewSessions: false }); + let statsCall = 0; + const session = { + ...makeMockSession(), + getSessionStats: vi.fn(() => { + statsCall++; + return { + tokens: { + input: statsCall * 10, + output: statsCall * 20, + cacheRead: statsCall * 3, + cacheWrite: statsCall, + total: statsCall * 34, + }, + }; + }), + }; + mockedCreateFnAgent.mockResolvedValue({ session } as any); + + const executor = new StepSessionExecutor({ + taskDetail: task, + worktreePath: "/project/.worktrees/main", + rootDir: "/project", + settings, + pluginRunner: undefined, + } as any); + + const result = await executor.executeAll(); + await executor.cleanup(); + + expect(result).toHaveLength(2); + expect(result.every((step) => step.success)).toBe(true); + expect(mockedCreateFnAgent).toHaveBeenCalledTimes(1); + expect(session.prompt).toHaveBeenCalledTimes(2); + expect(session.dispose).toHaveBeenCalledTimes(1); + expect(result[0]?.tokenUsage?.inputTokens).toBe(10); + expect(result[1]?.tokenUsage?.inputTokens).toBe(10); + expect(result[1]?.tokenUsage?.totalTokens).toBe(34); + }); + + it("creates a fresh primary step session for each step when runStepsInNewSessions is true", async () => { + const prompt = makeStepPrompt("FN-001", 2); + const task = makeTaskDetail({ + prompt, + steps: [ + { name: "Step 0", status: "pending" }, + { name: "Step 1", status: "pending" }, + ], + }); + const settings = makeSettings({ maxParallelSteps: 1, runStepsInNewSessions: true }); + const sessions = [makeMockSession(), makeMockSession()]; + mockedCreateFnAgent + .mockResolvedValueOnce({ session: sessions[0] } as any) + .mockResolvedValueOnce({ session: sessions[1] } as any); + + const executor = new StepSessionExecutor({ + taskDetail: task, + worktreePath: "/project/.worktrees/main", + rootDir: "/project", + settings, + pluginRunner: undefined, + } as any); + + const result = await executor.executeAll(); + + expect(result).toHaveLength(2); + expect(result.every((step) => step.success)).toBe(true); + expect(mockedCreateFnAgent).toHaveBeenCalledTimes(2); + expect(sessions[0]?.prompt).toHaveBeenCalledTimes(1); + expect(sessions[1]?.prompt).toHaveBeenCalledTimes(1); + expect(sessions[0]?.dispose).toHaveBeenCalledTimes(1); + expect(sessions[1]?.dispose).toHaveBeenCalledTimes(1); + }); + it("delivers pending steering comments in exactly one subsequent step prompt", async () => { const prompt = makeStepPrompt("FN-001", 2); const task = makeTaskDetail({ diff --git a/packages/engine/src/step-session-executor.ts b/packages/engine/src/step-session-executor.ts index 0dc0433430..bc5e93b493 100644 --- a/packages/engine/src/step-session-executor.ts +++ b/packages/engine/src/step-session-executor.ts @@ -1,9 +1,10 @@ /** - * StepSessionExecutor — runs each task step in its own fresh agent session. + * StepSessionExecutor — runs task steps with graph-controlled step boundaries. * * This module enables per-step error recovery with retry semantics, optional * parallel execution for non-conflicting steps (via git worktree isolation), - * and clean lifecycle management (pause, cleanup). + * optional primary-worktree session reuse, and clean lifecycle management + * (pause, cleanup). * * The class is a standalone engine subsystem with minimal integration surface. * It receives TaskDetail (read-only), a TaskStore for agent logs, and emits @@ -669,10 +670,12 @@ const NOOP_TASK_STORE: Pick = { }; /** - * StepSessionExecutor — runs each task step in its own fresh agent session. + * StepSessionExecutor — runs each task step behind a deterministic boundary. * * This class orchestrates per-step agent sessions with: * - **Sequential execution** (default): steps run one at a time + * - **Session policy**: `runStepsInNewSessions=false` reuses one primary-worktree + * session across sequential steps; `true` creates a fresh session per step * - **Parallel execution**: non-conflicting steps run simultaneously in * separate git worktrees (when `maxParallelSteps > 1`) * - **Per-step retry**: failed steps retry up to 3 times with exponential backoff @@ -712,6 +715,10 @@ export class StepSessionExecutor { private aborted = false; private maxParallel: number; private deliveredSteeringCommentIds = new Set(); + private reusablePrimarySession: AgentSession | null = null; + private reusablePrimaryHandle: SessionHandle | null = null; + private reusableStepTelemetry: { agentLogger: AgentLogger; trackingKey: string } | null = null; + private reusablePrimaryLastTokenUsage: StepResult["tokenUsage"] | undefined; private registerActiveStepSession(stepIndex: number, handle: SessionHandle, worktreePath: string): void { this.activeSessions.set(stepIndex, handle); @@ -744,6 +751,37 @@ export class StepSessionExecutor { this.parallelWorktrees.delete(stepIndex); } + /* + * FNXC:WorkflowStepSessions 2026-06-29-22:58: + * Coding (per-step review) still needs the StepSessionExecutor boundary so the graph can run `step-review` between steps, but the operator's "Each step in a new session" switch must control session freshness. Reuse only the primary sequential worktree when `runStepsInNewSessions` is false; parallel/isolated worktrees always need their own sessions because their cwd differs and they may run concurrently. + */ + private shouldReusePrimarySession(worktreePath: string): boolean { + return this.options.settings.runStepsInNewSessions === false && worktreePath === this.options.worktreePath; + } + + private selectReusableTelemetry(fallback: { agentLogger: AgentLogger; trackingKey: string }): { agentLogger: AgentLogger; trackingKey: string } { + return this.reusableStepTelemetry ?? fallback; + } + + private async disposeReusablePrimarySession(): Promise { + if (!this.reusablePrimarySession) return; + try { + this.reusablePrimaryHandle?.abortBash(); + } catch (err) { + stepExecLog.warn(`Failed to abort reusable primary step session: ${err}`); + } + try { + this.reusablePrimarySession.dispose(); + } catch (err) { + stepExecLog.warn(`Failed to dispose reusable primary step session: ${err}`); + } finally { + this.reusablePrimarySession = null; + this.reusablePrimaryHandle = null; + this.reusableStepTelemetry = null; + this.reusablePrimaryLastTokenUsage = undefined; + } + } + constructor(options: StepSessionExecutorOptions) { this.options = options; this.store = options.store ?? (NOOP_TASK_STORE as TaskStore); @@ -938,6 +976,7 @@ export class StepSessionExecutor { activeSessionRegistry.unregisterPath(worktreePath); } this.activeSessions.clear(); + await this.disposeReusablePrimarySession(); } /** @@ -951,6 +990,7 @@ export class StepSessionExecutor { if (this.activeSessions.size > 0) { await this.terminateAllSessions(); } + await this.disposeReusablePrimarySession(); // Remove parallel worktrees for (const [stepIdx, worktreePath] of this.parallelWorktrees) { @@ -1043,6 +1083,29 @@ export class StepSessionExecutor { } } + private async extractStepTokenUsage(session: AgentSession | null | undefined, reusePrimarySession: boolean): Promise { + const current = await this.extractTokenUsageFromSession(session); + if (!reusePrimarySession || !current) { + return current; + } + + const previous = this.reusablePrimaryLastTokenUsage; + this.reusablePrimaryLastTokenUsage = current; + if (!previous) { + return current; + } + + return { + inputTokens: Math.max(0, current.inputTokens - previous.inputTokens), + outputTokens: Math.max(0, current.outputTokens - previous.outputTokens), + cachedTokens: Math.max(0, current.cachedTokens - previous.cachedTokens), + cacheWriteTokens: Math.max(0, current.cacheWriteTokens - previous.cacheWriteTokens), + totalTokens: Math.max(0, current.totalTokens - previous.totalTokens), + modelProvider: current.modelProvider, + modelId: current.modelId, + }; + } + // ── Internal: Step Execution ──────────────────────────────────────── /** @@ -1073,6 +1136,7 @@ export class StepSessionExecutor { // Build reduced step prompt for context-limit recovery (simpler, shorter) const reducedStepPrompt = buildReducedStepPrompt(promptTaskDetail, stepIndex, this.options.rootDir); + const reusePrimarySession = this.shouldReusePrimarySession(worktreePath); // Acquire semaphore if provided if (semaphore) { @@ -1110,6 +1174,7 @@ export class StepSessionExecutor { persistAgentThinkingLog: resolvePersistAgentThinkingLog(settings, { ephemeral: true }), }); let session: AgentSession | null = null; + const localTelemetry = { agentLogger, trackingKey }; try { // Get plugin tools from plugin runner if available @@ -1150,7 +1215,7 @@ export class StepSessionExecutor { ] : []; - // Create fresh agent session for this attempt + // Create or reuse the agent session for this attempt // Resolve executor model using canonical lane hierarchy: // 1. Task override pair (taskDetail.modelProvider + taskDetail.modelId) // 2. Project execution lane pair (settings.executionProvider + settings.executionModelId) @@ -1164,94 +1229,113 @@ export class StepSessionExecutor { this.options.assignedAgentRuntimeConfig, ); - const createResult = await createResolvedAgentSession({ - sessionPurpose: "executor", - runtimeHint: this.options.runtimeHint, - pluginRunner: this.options.pluginRunner, - cwd: worktreePath, - systemPrompt: `You are an AI agent executing step ${stepIndex} of task ${taskDetail.id}. + if (reusePrimarySession && this.reusablePrimarySession) { + session = this.reusablePrimarySession; + } else { + const createResult = await createResolvedAgentSession({ + sessionPurpose: "executor", + runtimeHint: this.options.runtimeHint, + pluginRunner: this.options.pluginRunner, + cwd: worktreePath, + systemPrompt: `You are an AI agent executing steps for task ${taskDetail.id}. Your role: -- Complete only this step's scoped outcomes. +- Complete only the current step's scoped outcomes. - Read step context before editing. - Reuse existing patterns in nearby code. - Run relevant tests for changes made in this step. - Report blockers clearly instead of guessing. Follow instructions precisely and avoid unrelated changes.`, - defaultProvider: executorProvider, - defaultModelId: executorModelId, - fallbackProvider: settings.fallbackProvider, - fallbackModelId: settings.fallbackModelId, - defaultThinkingLevel: taskDetail.thinkingLevel ?? settings.defaultThinkingLevel, - runAuditor: createRunAuditor(this.store, { - runId: generateSyntheticRunId("workflow-step", taskDetail.id), - // Column-agent attribution (U4): the effective column agent is the - // principal that actually ran when the seam node's column governs; - // fall back to the task's assigned agent (legacy, byte-identical). - agentId: this.options.effectiveAgentId ?? taskDetail.assignedAgentId ?? "executor", - taskId: taskDetail.id, - taskLineageId: taskDetail.lineageId, - phase: "execute", - source: "step-session-executor", - }), - settings, - // FNXC:McpConfig 2026-06-25-23:02: Workflow model-node step sessions receive the same resolved, secret-materialized MCP server set as the parent executor; runtime support is still enforced inside the pi session seam without logging server contents. - mcpServers: this.options.mcpServers, - customTools: [ - ...pluginTools, - ...documentTools, - webFetchTool, - ...memoryTools, - ...taskLogTool, - ...taskCreateTool, - ...delegationTools, - ...messagingTools, - ], - onText: (delta) => { - agentLogger.onText(delta); - stuckTaskDetector?.recordActivity(trackingKey); - }, - onThinking: (delta) => { - agentLogger.onThinking(delta); - }, - onToolStart: (name, args) => { - agentLogger.onToolStart(name, args); - stuckTaskDetector?.recordActivity(trackingKey); - }, - onToolEnd: (name, isError, result) => { - agentLogger.onToolEnd(name, isError, result); - stuckTaskDetector?.recordActivity(trackingKey); - }, - // Skill selection from step-session executor options - ...(this.options.skillSelection ? { skillSelection: this.options.skillSelection } : {}), - actionGateContext: this.options.actionGateContext, - permanentAgentGating: this.options.permanentAgentGating, - taskId: taskDetail.id, - taskTitle: taskDetail.title, - onFallbackModelUsed: createFallbackModelObserver({ - agent: "executor", - label: "workflow step agent", - store: this.store, + defaultProvider: executorProvider, + defaultModelId: executorModelId, + fallbackProvider: settings.fallbackProvider, + fallbackModelId: settings.fallbackModelId, + defaultThinkingLevel: taskDetail.thinkingLevel ?? settings.defaultThinkingLevel, + runAuditor: createRunAuditor(this.store, { + runId: generateSyntheticRunId("workflow-step", taskDetail.id), + // Column-agent attribution (U4): the effective column agent is the + // principal that actually ran when the seam node's column governs; + // fall back to the task's assigned agent (legacy, byte-identical). + agentId: this.options.effectiveAgentId ?? taskDetail.assignedAgentId ?? "executor", + taskId: taskDetail.id, + taskLineageId: taskDetail.lineageId, + phase: "execute", + source: "step-session-executor", + }), + settings, + // FNXC:McpConfig 2026-06-25-23:02: Workflow model-node step sessions receive the same resolved, secret-materialized MCP server set as the parent executor; runtime support is still enforced inside the pi session seam without logging server contents. + mcpServers: this.options.mcpServers, + customTools: [ + ...pluginTools, + ...documentTools, + webFetchTool, + ...memoryTools, + ...taskLogTool, + ...taskCreateTool, + ...delegationTools, + ...messagingTools, + ], + onText: (delta) => { + const telemetry = reusePrimarySession ? this.selectReusableTelemetry(localTelemetry) : localTelemetry; + telemetry.agentLogger.onText(delta); + stuckTaskDetector?.recordActivity(telemetry.trackingKey); + }, + onThinking: (delta) => { + const telemetry = reusePrimarySession ? this.selectReusableTelemetry(localTelemetry) : localTelemetry; + telemetry.agentLogger.onThinking(delta); + }, + onToolStart: (name, args) => { + const telemetry = reusePrimarySession ? this.selectReusableTelemetry(localTelemetry) : localTelemetry; + telemetry.agentLogger.onToolStart(name, args); + stuckTaskDetector?.recordActivity(telemetry.trackingKey); + }, + onToolEnd: (name, isError, result) => { + const telemetry = reusePrimarySession ? this.selectReusableTelemetry(localTelemetry) : localTelemetry; + telemetry.agentLogger.onToolEnd(name, isError, result); + stuckTaskDetector?.recordActivity(telemetry.trackingKey); + }, + // Skill selection from step-session executor options + ...(this.options.skillSelection ? { skillSelection: this.options.skillSelection } : {}), + actionGateContext: this.options.actionGateContext, + permanentAgentGating: this.options.permanentAgentGating, taskId: taskDetail.id, taskTitle: taskDetail.title, - }), - taskEnv: this.options.taskEnv, - }); - session = createResult.session; + onFallbackModelUsed: createFallbackModelObserver({ + agent: "executor", + label: "workflow step agent", + store: this.store, + taskId: taskDetail.id, + taskTitle: taskDetail.title, + }), + taskEnv: this.options.taskEnv, + }); + session = createResult.session; + if (reusePrimarySession) { + this.reusablePrimarySession = session; + } + } // Track session for termination and stuck-task detection. // Pass the canonical task ID (e.g. "FN-1452") as the third argument so // that stuck-kill callbacks (beforeRequeue, onStuck) operate on the real // task rather than the compound step key ("FN-1452-step-1"). - const handle: SessionHandle = { - dispose: () => session?.dispose(), - abortBash: () => session?.abortBash(), - steer: async (message) => { - if (!session) return; - await session.steer(message); - }, - }; + const handle: SessionHandle = reusePrimarySession && this.reusablePrimaryHandle + ? this.reusablePrimaryHandle + : { + dispose: () => session?.dispose(), + abortBash: () => session?.abortBash(), + steer: async (message) => { + if (!session) return; + await session.steer(message); + }, + }; + if (reusePrimarySession && !this.reusablePrimaryHandle) { + this.reusablePrimaryHandle = handle; + } + if (reusePrimarySession) { + this.reusableStepTelemetry = localTelemetry; + } this.registerActiveStepSession(stepIndex, handle, worktreePath); stuckTaskDetector?.trackTask(trackingKey, { dispose: () => session?.dispose() }, taskDetail.id); @@ -1273,7 +1357,7 @@ Follow instructions precisely and avoid unrelated changes.`, stepIndex, success: true, retries, - tokenUsage: await this.extractTokenUsageFromSession(session), + tokenUsage: await this.extractStepTokenUsage(session, reusePrimarySession), }; this.options.onStepComplete?.(stepIndex, result); return result; @@ -1313,7 +1397,7 @@ Follow instructions precisely and avoid unrelated changes.`, stepIndex, success: true, retries, - tokenUsage: await this.extractTokenUsageFromSession(session), + tokenUsage: await this.extractStepTokenUsage(session, reusePrimarySession), }; this.options.onStepComplete?.(stepIndex, result); return result; @@ -1340,11 +1424,15 @@ Follow instructions precisely and avoid unrelated changes.`, success: false, error: errorMessage, retries, - tokenUsage: await this.extractTokenUsageFromSession(session), + tokenUsage: await this.extractStepTokenUsage(session, reusePrimarySession), }; this.options.onStepComplete?.(stepIndex, result); return result; } + if (reusePrimarySession) { + await this.disposeReusablePrimarySession(); + session = null; + } } finally { try { await agentLogger.flush(); @@ -1355,11 +1443,15 @@ Follow instructions precisely and avoid unrelated changes.`, this.unregisterActiveStepSession(stepIndex, worktreePath); stuckTaskDetector?.untrackTask(trackingKey); - try { - session?.dispose(); - } catch (err: unknown) { - const msg = err instanceof Error ? err.message : String(err); - stepExecLog.warn(`Failed to dispose session for step ${stepIndex}: ${msg}`); + if (reusePrimarySession) { + this.reusableStepTelemetry = null; + } else { + try { + session?.dispose(); + } catch (err: unknown) { + const msg = err instanceof Error ? err.message : String(err); + stepExecLog.warn(`Failed to dispose session for step ${stepIndex}: ${msg}`); + } } } }