From b94aabc58cdc3f290196f302d822860d0ed2c380 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 14 Sep 2026 14:01:23 -0700 Subject: [PATCH] fix(tui): re-sync idle-with-fleet flag on director rebuilds --- src/tui/runner/exit.test.ts | 131 ++++++++++++++++++++++++++++++++++++ src/tui/runner/exit.ts | 15 +++++ 2 files changed, 146 insertions(+) diff --git a/src/tui/runner/exit.test.ts b/src/tui/runner/exit.test.ts index b768e90c9..00c3d2c0f 100644 --- a/src/tui/runner/exit.test.ts +++ b/src/tui/runner/exit.test.ts @@ -5,6 +5,14 @@ import { getLogger } from "@intx/log"; import type { InferenceSource } from "@intx/types/runtime"; import * as codexSession from "../../auth/codex/session.js"; +import { createChatDirector } from "../../agent/director.js"; +import { createSubAgentSessionStore } from "../../subagent/session-store.js"; +import type { + ReactorAction, + ReactorCapabilities, + ReactorInboundEvent, + ReactorState, +} from "@intx/types/runtime"; import { LOG_NAMESPACE_ROOT } from "../../branding.js"; import { defined } from "../../../tests/helpers/defined.js"; import { @@ -259,3 +267,126 @@ describe("agentProxy.send vs /clear", () => { } }); }); + +const rebuildMockState: ReactorState = {} as unknown as ReactorState; + +const rebuildMockCapabilities: ReactorCapabilities = { + infer: (options) => + ({ + type: "infer", + ...(options !== undefined ? { options } : {}), + }) as ReactorAction, + executeTools: (calls) => ({ type: "execute_tools", calls }), + suspend: (gate) => ({ type: "suspend", gate }), + fork: (mode, forkId) => ({ type: "fork", mode, forkId }), + emit: (eventType, data) => ({ type: "emit", eventType, data }), + reply: (content) => ({ type: "reply", content }), + checkpoint: (message = "") => ({ type: "checkpoint", message }), + compact: (compactor, reason) => ({ type: "compact", compactor, reason }), + wait: () => ({ type: "wait" }), + done: () => ({ type: "done" }), +}; + +function rebuildManageTasksEvent(): ReactorInboundEvent { + return { + type: "inference.done", + turn: { + role: "assistant", + model: "test", + timestamp: 0, + content: [ + { + type: "tool_call", + id: "m", + name: "manage_tasks", + arguments: { + action: "create", + tasks: [{ id: "t1", title: "work", status: "doing" }], + }, + }, + ], + }, + usage: { input: 0, output: 1, cacheRead: 0, cacheWrite: 0, thinking: 0 }, + source: { model: "test-model" }, + } as unknown as ReactorInboundEvent; +} + +function rebuildTextTurn(): ReactorInboundEvent { + return { + type: "inference.done", + turn: { + role: "assistant", + model: "test", + timestamp: 0, + content: [{ type: "text", text: "all set" }], + }, + usage: { input: 10, output: 1, cacheRead: 0, cacheWrite: 0, thinking: 0 }, + source: { model: "test-model" }, + } as unknown as ReactorInboundEvent; +} + +describe("rebuild re-syncs idle-with-fleet while drained", () => { + test("reload-if-idle and interrupt rebuilds resume the open-task nudge with no fleet transition", async () => { + const store = createSubAgentSessionStore(); + const directorHolder: RunnerServices["directorHolder"] = {}; + const agent = recordingAgent([]); + const { state, services } = stubSendLifecycle(agent); + services.directorHolder = + directorHolder as unknown as RunnerServices["directorHolder"]; + services.subAgentSessions = + store as unknown as RunnerServices["subAgentSessions"]; + services.workflowHost = { + reattach: () => undefined, + } as unknown as RunnerServices["workflowHost"]; + services.cycleRecorder = { + dispose: async () => "", + reset: () => undefined, + handleEvent: () => undefined, + } as unknown as RunnerServices["cycleRecorder"]; + services.buildAgent = (async () => { + // Every rebuild mints a fresh director from the static true seed (fleet + // lanes may appear mid-session), exactly like the TUI session assembly. + directorHolder.instance = createChatDirector("base", [], { + onTasksChange: () => undefined, + allowIdleWithFleet: true, + }); + return agent; + }) as unknown as RunnerServices["buildAgent"]; + const fleetEvents: unknown[] = []; + services.emitter.on("event", (event: { type: string }) => { + if (event.type === "fleet") fleetEvents.push(event); + }); + await createRunLifecycle(state, services); + const expectOpenTaskNudge = async (): Promise => { + const director = defined( + directorHolder.instance, + "directorHolder.instance", + ); + await director.decide( + rebuildManageTasksEvent(), + rebuildMockState, + rebuildMockCapabilities, + ); + const actions = await director.decide( + rebuildTextTurn(), + rebuildMockState, + rebuildMockCapabilities, + ); + const list = Array.isArray(actions) ? actions : [actions]; + expect(list.some((action) => action.type === "infer")).toBe(true); + }; + // Drained fleet: the idle reload rebuilds onto the static true seed. + state.pendingReload = true; + defined(state.reloadIfIdle, "reloadIfIdle")(); + await services.sessionOps.awaitTail(); + expect(state.fatalBuildError).toBeNull(); + await expectOpenTaskNudge(); + // The interrupt rebuild inherits the same seed. + defined(state.interrupt, "interrupt")(); + await services.sessionOps.awaitTail(); + expect(state.fatalBuildError).toBeNull(); + await expectOpenTaskNudge(); + expect(store.list()).toEqual([]); + expect(fleetEvents).toEqual([]); + }); +}); diff --git a/src/tui/runner/exit.ts b/src/tui/runner/exit.ts index 795e02cca..a313d5e45 100644 --- a/src/tui/runner/exit.ts +++ b/src/tui/runner/exit.ts @@ -13,6 +13,7 @@ import { import { getLogger } from "@intx/log"; import type { InferenceSource } from "@intx/types/runtime"; import { consumeStream } from "../../session/stream-consumer.js"; +import { liveFleetCount } from "../../subagent/index.js"; import { getTelemetry } from "../../telemetry/singleton.js"; import { onTurnBoundary } from "../../agent/reactor-events.js"; import { setAgentSourceUnlessClosed } from "../agent-source-sync.js"; @@ -145,6 +146,18 @@ export async function closeAgentForRebuild( } } +// Rebuilt directors seed allowIdleWithFleet=true (fleet lanes may appear +// mid-session), so a rebuild while drained must re-sync the new director from +// the live fleet count — otherwise the open-task nudge stays suppressed until +// the next fleet transition, which never comes for an already-drained fleet. +function resyncIdleWithFleetFlag( + services: Pick, +): void { + services.directorHolder.instance?.setAllowIdleWithFleet( + liveFleetCount(services.subAgentSessions.list()) > 0, + ); +} + // Every rebuild site funnels its failure (a lock left held by a failed // close, or any other buildAgent failure) through here so it surfaces as a // plain-language, caught error rather than an unhandled rejection. @@ -345,6 +358,7 @@ export async function createRunLifecycle( throw new AgentContextLockError(state.workdir); } state.currentAgent = await services.buildAgent(); + resyncIdleWithFleetFlag(services); state.streamPromise = consumeStream( liveAgent(state).stream(), streamSink, @@ -547,6 +561,7 @@ export async function createRunLifecycle( throw new AgentContextLockError(state.workdir); } state.currentAgent = await services.buildAgent(); + resyncIdleWithFleetFlag(services); services.cycleRecorder.reset(); state.streamPromise = consumeStream( liveAgent(state).stream(),