diff --git a/src/agent/director.ts b/src/agent/director.ts index 9c6f5077a..48add989c 100644 --- a/src/agent/director.ts +++ b/src/agent/director.ts @@ -447,23 +447,24 @@ export interface ChatDirectorOptions { * uses a switched model still stamp the previous id — the director learns * the new id from that inference's completion event. * - * - getLiveFleetCount → static config (allowIdleWithFleet below). The count - * is genuinely external (subagent lane statuses the reactor never sees — - * its own tasks only carry todo/doing/done/cancelled), so neither - * BaseEnv-derived nor reactor-supplied can reproduce its liveness. - * Idle-with-fleet itself is unchanged (fleet-running TUI sessions allow - * the terminal wait). Accepted loss: a drained fleet no longer resumes - * the open-task nudge — bounded, since the nudge capitulates to terminal - * after its cap anyway. + * - getLiveFleetCount → seeded config + live narrow setter + * (setAllowIdleWithFleet below). The count is genuinely external (subagent + * lane statuses the reactor never sees — its own tasks only carry + * todo/doing/done/cancelled), so neither BaseEnv-derived nor + * reactor-supplied can reproduce its liveness. Idle-with-fleet itself is + * unchanged (fleet-running TUI sessions allow the terminal wait); the + * fleet-wake publisher drives the setter on count transitions, so a + * drained fleet resumes the open-task nudge. */ /** Explicit retry policy; when set, skips the default Corbits policy. */ retryPolicy?: RetryPolicy | undefined; /** - * Static idle-with-fleet allowance (CL-7918 replacement for the former + * Initial idle-with-fleet allowance (CL-7918 replacement for the former * getLiveFleetCount closure). When true the director allows a terminal * wait/reply with open tasks; when omitted or false it keeps the open-task - * nudge. The TUI sets this (fleet lanes may appear mid-session); exec - * omits it. + * nudge. The TUI seeds this (fleet lanes may appear mid-session); exec + * omits it. The live fleet-wake publisher then keeps it current through + * setAllowIdleWithFleet, so a drained fleet resumes the nudge. */ allowIdleWithFleet?: boolean | undefined; } @@ -515,8 +516,8 @@ class ChatDirectorImpl extends DefaultDirector { // refreshed on every inference completion, so mid-session /model switches // remap without rebuilding the agent. private currentSourceId: string | undefined; - /** CL-7918 static replacement for the former getLiveFleetCount closure. */ - private readonly allowIdleWithFleet: boolean; + /** CL-7918 live replacement for the former getLiveFleetCount closure. */ + private allowIdleWithFleet: boolean; // Consecutive assistant turns that contain tool calls and no text. Reset on // any turn with text and on every fresh user message — a weak model that // spins in place on one thread of tool calls still converges to the @@ -578,6 +579,13 @@ class ChatDirectorImpl extends DefaultDirector { this.workflowCoordinator = coordinator; } + // Narrow live setter for the idle-with-fleet allowance (CL-7972): the + // fleet-wake publisher drives this on fleet-count transitions, so a drained + // fleet resumes the open-task nudge instead of holding the seeded value. + setAllowIdleWithFleet(value: boolean): void { + this.allowIdleWithFleet = value; + } + updateToolDefinitions(toolDefinitions: ToolDefinition[]): void { const before = toolSetDigest(this._toolDefinitions); const after = toolSetDigest(toolDefinitions); @@ -1095,9 +1103,10 @@ class ChatDirectorImpl extends DefaultDirector { (a) => a.type === "wait" || a.type === "reply", ); if (hasTerminal) { - // CL-7918: static idle-with-fleet allowance (replaces the former - // getLiveFleetCount closure). TUI sessions set allowIdleWithFleet; - // exec omits it and keeps the nudge. + // CL-7918 live idle-with-fleet allowance (replaces the former + // getLiveFleetCount closure): seeded at construction, then kept + // current by the fleet-wake publisher. TUI seeds true (fleet lanes + // may appear mid-session); exec omits it and keeps the nudge. if (this.allowIdleWithFleet) { return base; } @@ -1172,6 +1181,7 @@ export function hydrateTasksFromTurns(turns: ConversationTurn[]): Task[] { export interface ChatDirector extends ReactorDirector { updateToolDefinitions(toolDefinitions: ToolDefinition[]): void; setWorkflowCoordinator(coordinator: WorkflowCoordinator | undefined): void; + setAllowIdleWithFleet(value: boolean): void; getTasks(): Task[]; restoreTasks(tasks: Task[]): void; getContextEstimate(): { tokens: number; isEstimate: boolean }; diff --git a/src/director.test.ts b/src/director.test.ts index 731f3f360..e17faac8b 100644 --- a/src/director.test.ts +++ b/src/director.test.ts @@ -442,6 +442,41 @@ describe("open-task termination guard", () => { ).toBe(true); }); + test("setAllowIdleWithFleet tracks fleet transitions off the seeded value", async () => { + const director = createChatDirector("base", [], { + onTasksChange: () => undefined, + allowIdleWithFleet: true, + }); + await director.decide( + manageTasksEvent("doing"), + mockState, + mockCapabilities, + ); + + // Seeded allowance: terminal reply with open tasks, no nudge spent. + const seeded = actionsArray( + await director.decide(textTurn(), mockState, mockCapabilities), + ); + expect(hasReply(seeded)).toBe(true); + expect(hasInfer(seeded)).toBe(false); + + // Drained fleet resumes the open-task nudge. + director.setAllowIdleWithFleet(false); + const nudged = actionsArray( + await director.decide(textTurn(), mockState, mockCapabilities), + ); + expect(hasInfer(nudged)).toBe(true); + expect(hasReply(nudged)).toBe(false); + + // Fleet back: terminal allowed again. + director.setAllowIdleWithFleet(true); + const settled = actionsArray( + await director.decide(textTurn(), mockState, mockCapabilities), + ); + expect(hasReply(settled)).toBe(true); + expect(hasInfer(settled)).toBe(false); + }); + test("empty model turn settles with a valid empty reply", async () => { // DefaultDirector ends empty responses with bare wait; without a reply, // agent.send hangs and the TUI Working spinner sticks forever. diff --git a/src/session/assemble-runtime.ts b/src/session/assemble-runtime.ts index 93ee172af..f7cdb1f68 100644 --- a/src/session/assemble-runtime.ts +++ b/src/session/assemble-runtime.ts @@ -474,9 +474,10 @@ export interface ChatAgentWiring { totalTimeoutMs?: number | undefined; onTasksChange: (tasks: Task[]) => void; /** - * Static idle-with-fleet allowance for the chat director (CL-7918: replaces - * the former getLiveFleetCount closure). Omitted in exec (keeps the - * open-task nudge); the TUI sets it (fleet lanes may appear mid-session). + * Seed for the idle-with-fleet allowance for the chat director (CL-7918: + * replaces the former getLiveFleetCount closure; CL-7972 keeps it live via + * the fleet-wake publisher). Omitted in exec (keeps the open-task nudge); + * the TUI seeds it (fleet lanes may appear mid-session). */ allowIdleWithFleet?: boolean; /** Compaction governor re-entry (the reactor emits no event after compact). */ @@ -542,8 +543,9 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent { onTasksChange: wiring.onTasksChange, requestContinuation: wiring.requestContinuation, provider: { ...wiring.getProvider() }, - // CL-7918: static idle-with-fleet allowance replaces the former - // getLiveFleetCount closure; the live source id for retry stamping + // CL-7918: idle-with-fleet seed replaces the former + // getLiveFleetCount closure (CL-7972 keeps it live via the + // fleet-wake publisher); the live source id for retry stamping // is reactor-tracked now (no getProviderId). allowIdleWithFleet: wiring.allowIdleWithFleet, }, 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(), diff --git a/src/tui/runner/session.ts b/src/tui/runner/session.ts index 880fc7f45..a7fb84b86 100644 --- a/src/tui/runner/session.ts +++ b/src/tui/runner/session.ts @@ -619,8 +619,9 @@ export async function assembleTUISession( inactivityTimeoutMs: config.inactivityTimeoutMs ?? 750_000, totalTimeoutMs: config.totalTimeoutMs, onTasksChange: (tasks) => emitter.emit("tasks", tasks), - // CL-7918: static idle-with-fleet allowance (fleet lanes may appear - // mid-session); retry stamping tracks the live source id in-reactor now. + // CL-7918 seed for the idle-with-fleet allowance (fleet lanes may appear + // mid-session; CL-7972 keeps it live via the fleet-wake publisher); + // retry stamping tracks the live source id in-reactor now. allowIdleWithFleet: true, requestContinuation: () => { const targetAgent = liveAgent(state); diff --git a/src/tui/runner/wiring.ask-wake.test.ts b/src/tui/runner/wiring.ask-wake.test.ts index dac583f9c..fefaa12d6 100644 --- a/src/tui/runner/wiring.ask-wake.test.ts +++ b/src/tui/runner/wiring.ask-wake.test.ts @@ -68,6 +68,36 @@ test("failed reset releases publication without flushing partially cancelled wor } }); +test("fleet-wake publisher reports fleet-count transitions to the idle-with-fleet follower", () => { + const store = createSubAgentSessionStore(); + const emitter = new EventEmitter(); + const seen: number[] = []; + const publisher = createFleetWakePublisher(store, emitter, (running) => { + seen.push(running); + }); + // Steady empty state: no transition, no callback. + publisher.publish(); + expect(seen).toEqual([]); + store.start({ + id: "worker-one", + agentId: "builder", + description: "worker-one", + brief: "build", + }); + // A started lane counts as running: 0 -> 1 transition. + publisher.publish(); + expect(seen).toEqual([1]); + store.markRunning("worker-one"); + publisher.publish(); + expect(seen).toEqual([1]); + // Steady fleet: no repeat callback. + publisher.publish(); + expect(seen).toEqual([1]); + store.cancelAll("Session cleared"); + publisher.publish(); + expect(seen).toEqual([1, 0]); +}); + for (const phase of ["settled", "prequeued", "deferred"] as const) { test(`rotation suppresses old worker snapshots (${phase}), including repeated resets`, async () => { await withTestRenderer( diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 21eee1af6..fe922fb79 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -105,6 +105,10 @@ export function createFleetStallPollTick( export function createFleetWakePublisher( sessions: RunnerServices["subAgentSessions"], emitter: RunnerServices["emitter"], + // Live idle-with-fleet flag (CL-7972): fired on fleet-count transitions so + // the chat director's seeded allowance tracks the live fleet instead of + // holding its construction value. Omitted in tests that only assert events. + onFleetCount?: (running: number) => void, ) { let lastLiveFleet = 0; let suspended = false; @@ -120,6 +124,7 @@ export function createFleetWakePublisher( if (fleet !== lastLiveFleet) { lastLiveFleet = fleet; emitter.emit("event", { type: "fleet", running: fleet }); + onFleetCount?.(fleet); } return { previousRunning, running: fleet }; }; @@ -261,6 +266,11 @@ export function wirePostStartup( const fleetWakePublisher = createFleetWakePublisher( services.subAgentSessions, services.emitter, + // Liven the seeded idle-with-fleet allowance (CL-7972): the director + // reads the holder live so rebuilds stay tracked, and degrades gracefully + // while the director is not yet built. + (running) => + services.directorHolder.instance?.setAllowIdleWithFleet(running > 0), ); state.withFleetPublicationSuspended = fleetWakePublisher.withSuspended; sessionBridge.setDryOpenTaskDriver(() => { diff --git a/src/workflows/host.ts b/src/workflows/host.ts index 4b2c30d8f..05e3d1a86 100644 --- a/src/workflows/host.ts +++ b/src/workflows/host.ts @@ -59,8 +59,10 @@ export interface WorkflowHostArgs { getSessionId: () => string; getToolDefinitions: () => ToolDefinition[]; // The live chat director; the workflow coordinator is attached to it when a - // workflow starts. Returns undefined before the director is built. - getDirector: () => { setWorkflowCoordinator: SetCoordinator } | undefined; + // workflow starts. Returns undefined before the director is built. The seam + // is optional so the host degrades gracefully against directors without + // workflow support — every attach point presence-checks before calling. + getDirector: () => { setWorkflowCoordinator?: SetCoordinator } | undefined; // Overrides the state-tree home (defaults to the real user home). Tests // pass a sandboxed dir here so persist()/resume() never touch ~/.corbits. home?: string; @@ -84,7 +86,7 @@ export class WorkflowHost { // Re-attach the active coordinator to a freshly rebuilt director. Safe to // call with no active workflow — it just clears any stale coordinator. reattach(): void { - this.args.getDirector()?.setWorkflowCoordinator(this.coordinator); + this.args.getDirector()?.setWorkflowCoordinator?.(this.coordinator); } // Drop the active workflow and history (e.g. on /clear). @@ -94,7 +96,7 @@ export class WorkflowHost { this.pendingReplace = undefined; this.completedWorkflows = []; this.lastActiveStatus = undefined; - this.args.getDirector()?.setWorkflowCoordinator(undefined); + this.args.getDirector()?.setWorkflowCoordinator?.(undefined); this.notify(); } @@ -166,7 +168,7 @@ export class WorkflowHost { this.listen(runtime); this.runtime = runtime; this.coordinator = coordinator; - this.args.getDirector()?.setWorkflowCoordinator(coordinator); + this.args.getDirector()?.setWorkflowCoordinator?.(coordinator); if (restore !== true) { runtime.start(workflow); this.persist(); diff --git a/tests/unit/workflow-host.test.ts b/tests/unit/workflow-host.test.ts index d0247aab0..e7043fdb9 100644 --- a/tests/unit/workflow-host.test.ts +++ b/tests/unit/workflow-host.test.ts @@ -43,6 +43,12 @@ async function withHost( home: string, ) => void | Promise, onChange?: () => void, + // Overrides the director seam: defaults to a tracking director, pass + // () => undefined (not yet built) or () => ({}) (no workflow support) to + // exercise graceful degradation. + getDirector?: () => + | { setWorkflowCoordinator?: (c: WorkflowCoordinator | undefined) => void } + | undefined, ): Promise { const cwd = await mkdtemp(join(tmpdir(), "wf-host-")); const home = await mkdtemp(join(tmpdir(), "wf-host-home-")); @@ -54,11 +60,13 @@ async function withHost( cwd, getSessionId: () => "session-1", getToolDefinitions: () => tools, - getDirector: () => ({ - setWorkflowCoordinator: (c) => { - director.coordinator = c; - }, - }), + getDirector: + getDirector ?? + (() => ({ + setWorkflowCoordinator: (c) => { + director.coordinator = c; + }, + })), home, ...(onChange !== undefined ? { onChange } : {}), }); @@ -201,3 +209,43 @@ test("resume() restores an on-disk workflow snapshot for the session", async () expect(director.coordinator).toBeInstanceOf(WorkflowCoordinator); }); }); + +test("host degrades gracefully when no director is built yet", async () => { + await withHost( + [], + async (host) => { + expect(host.start("review")).toBe("Started review workflow."); + expect(host.isActive()).toBe(true); + host.reattach(); + expect(host.isActive()).toBe(true); + host.reset(); + expect(host.isActive()).toBe(false); + }, + undefined, + () => undefined, + ); +}); + +test("host degrades gracefully against a director without workflow support", async () => { + await withHost( + [], + async (host, _director, cwd, home) => { + const workflow = findWorkflow("review"); + expect(workflow).toBeDefined(); + const runtime = new WorkflowRuntime(new Map()); + runtime.start(defined(workflow, "workflow")); + await saveWorkflowState(cwd, "session-1", runtime.state(), home); + + await host.resume(); + expect(host.isActive()).toBe(true); + host.reattach(); + expect(host.isActive()).toBe(true); + expect(host.start("build")).toContain("again to replace"); + expect(host.start("build")).toBe("Started build workflow."); + host.reset(); + expect(host.isActive()).toBe(false); + }, + undefined, + () => ({}), + ); +});