Skip to content

Commit 4e716a2

Browse files
feat(director): liven idle-with-fleet flag and decouple host from workflow seam typing (#1048)
* feat(director): liven idle-with-fleet flag and decouple host from workflow seam typing Summary: optionalize the host-to-runtime workflow seam so the host degrades gracefully against directors without workflow support; replace the static idle-with-fleet seed with a live narrow setter driven by fleet-wake publisher transitions. Verification: bun run check passes (7393 tests, 0 fail). * fix(tui): re-sync idle-with-fleet flag on director rebuilds (#1058)
1 parent ed1fd79 commit 4e716a2

10 files changed

Lines changed: 317 additions & 33 deletions

File tree

‎src/agent/director.ts‎

Lines changed: 26 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -447,23 +447,24 @@ export interface ChatDirectorOptions {
447447
* uses a switched model still stamp the previous id — the director learns
448448
* the new id from that inference's completion event.
449449
*
450-
* - getLiveFleetCount → static config (allowIdleWithFleet below). The count
451-
* is genuinely external (subagent lane statuses the reactor never sees —
452-
* its own tasks only carry todo/doing/done/cancelled), so neither
453-
* BaseEnv-derived nor reactor-supplied can reproduce its liveness.
454-
* Idle-with-fleet itself is unchanged (fleet-running TUI sessions allow
455-
* the terminal wait). Accepted loss: a drained fleet no longer resumes
456-
* the open-task nudge — bounded, since the nudge capitulates to terminal
457-
* after its cap anyway.
450+
* - getLiveFleetCount → seeded config + live narrow setter
451+
* (setAllowIdleWithFleet below). The count is genuinely external (subagent
452+
* lane statuses the reactor never sees — its own tasks only carry
453+
* todo/doing/done/cancelled), so neither BaseEnv-derived nor
454+
* reactor-supplied can reproduce its liveness. Idle-with-fleet itself is
455+
* unchanged (fleet-running TUI sessions allow the terminal wait); the
456+
* fleet-wake publisher drives the setter on count transitions, so a
457+
* drained fleet resumes the open-task nudge.
458458
*/
459459
/** Explicit retry policy; when set, skips the default Corbits policy. */
460460
retryPolicy?: RetryPolicy | undefined;
461461
/**
462-
* Static idle-with-fleet allowance (CL-7918 replacement for the former
462+
* Initial idle-with-fleet allowance (CL-7918 replacement for the former
463463
* getLiveFleetCount closure). When true the director allows a terminal
464464
* wait/reply with open tasks; when omitted or false it keeps the open-task
465-
* nudge. The TUI sets this (fleet lanes may appear mid-session); exec
466-
* omits it.
465+
* nudge. The TUI seeds this (fleet lanes may appear mid-session); exec
466+
* omits it. The live fleet-wake publisher then keeps it current through
467+
* setAllowIdleWithFleet, so a drained fleet resumes the nudge.
467468
*/
468469
allowIdleWithFleet?: boolean | undefined;
469470
}
@@ -515,8 +516,8 @@ class ChatDirectorImpl extends DefaultDirector {
515516
// refreshed on every inference completion, so mid-session /model switches
516517
// remap without rebuilding the agent.
517518
private currentSourceId: string | undefined;
518-
/** CL-7918 static replacement for the former getLiveFleetCount closure. */
519-
private readonly allowIdleWithFleet: boolean;
519+
/** CL-7918 live replacement for the former getLiveFleetCount closure. */
520+
private allowIdleWithFleet: boolean;
520521
// Consecutive assistant turns that contain tool calls and no text. Reset on
521522
// any turn with text and on every fresh user message — a weak model that
522523
// spins in place on one thread of tool calls still converges to the
@@ -578,6 +579,13 @@ class ChatDirectorImpl extends DefaultDirector {
578579
this.workflowCoordinator = coordinator;
579580
}
580581

582+
// Narrow live setter for the idle-with-fleet allowance (CL-7972): the
583+
// fleet-wake publisher drives this on fleet-count transitions, so a drained
584+
// fleet resumes the open-task nudge instead of holding the seeded value.
585+
setAllowIdleWithFleet(value: boolean): void {
586+
this.allowIdleWithFleet = value;
587+
}
588+
581589
updateToolDefinitions(toolDefinitions: ToolDefinition[]): void {
582590
const before = toolSetDigest(this._toolDefinitions);
583591
const after = toolSetDigest(toolDefinitions);
@@ -1095,9 +1103,10 @@ class ChatDirectorImpl extends DefaultDirector {
10951103
(a) => a.type === "wait" || a.type === "reply",
10961104
);
10971105
if (hasTerminal) {
1098-
// CL-7918: static idle-with-fleet allowance (replaces the former
1099-
// getLiveFleetCount closure). TUI sessions set allowIdleWithFleet;
1100-
// exec omits it and keeps the nudge.
1106+
// CL-7918 live idle-with-fleet allowance (replaces the former
1107+
// getLiveFleetCount closure): seeded at construction, then kept
1108+
// current by the fleet-wake publisher. TUI seeds true (fleet lanes
1109+
// may appear mid-session); exec omits it and keeps the nudge.
11011110
if (this.allowIdleWithFleet) {
11021111
return base;
11031112
}
@@ -1172,6 +1181,7 @@ export function hydrateTasksFromTurns(turns: ConversationTurn[]): Task[] {
11721181
export interface ChatDirector extends ReactorDirector {
11731182
updateToolDefinitions(toolDefinitions: ToolDefinition[]): void;
11741183
setWorkflowCoordinator(coordinator: WorkflowCoordinator | undefined): void;
1184+
setAllowIdleWithFleet(value: boolean): void;
11751185
getTasks(): Task[];
11761186
restoreTasks(tasks: Task[]): void;
11771187
getContextEstimate(): { tokens: number; isEstimate: boolean };

‎src/director.test.ts‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -442,6 +442,41 @@ describe("open-task termination guard", () => {
442442
).toBe(true);
443443
});
444444

445+
test("setAllowIdleWithFleet tracks fleet transitions off the seeded value", async () => {
446+
const director = createChatDirector("base", [], {
447+
onTasksChange: () => undefined,
448+
allowIdleWithFleet: true,
449+
});
450+
await director.decide(
451+
manageTasksEvent("doing"),
452+
mockState,
453+
mockCapabilities,
454+
);
455+
456+
// Seeded allowance: terminal reply with open tasks, no nudge spent.
457+
const seeded = actionsArray(
458+
await director.decide(textTurn(), mockState, mockCapabilities),
459+
);
460+
expect(hasReply(seeded)).toBe(true);
461+
expect(hasInfer(seeded)).toBe(false);
462+
463+
// Drained fleet resumes the open-task nudge.
464+
director.setAllowIdleWithFleet(false);
465+
const nudged = actionsArray(
466+
await director.decide(textTurn(), mockState, mockCapabilities),
467+
);
468+
expect(hasInfer(nudged)).toBe(true);
469+
expect(hasReply(nudged)).toBe(false);
470+
471+
// Fleet back: terminal allowed again.
472+
director.setAllowIdleWithFleet(true);
473+
const settled = actionsArray(
474+
await director.decide(textTurn(), mockState, mockCapabilities),
475+
);
476+
expect(hasReply(settled)).toBe(true);
477+
expect(hasInfer(settled)).toBe(false);
478+
});
479+
445480
test("empty model turn settles with a valid empty reply", async () => {
446481
// DefaultDirector ends empty responses with bare wait; without a reply,
447482
// agent.send hangs and the TUI Working spinner sticks forever.

‎src/session/assemble-runtime.ts‎

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -474,9 +474,10 @@ export interface ChatAgentWiring {
474474
totalTimeoutMs?: number | undefined;
475475
onTasksChange: (tasks: Task[]) => void;
476476
/**
477-
* Static idle-with-fleet allowance for the chat director (CL-7918: replaces
478-
* the former getLiveFleetCount closure). Omitted in exec (keeps the
479-
* open-task nudge); the TUI sets it (fleet lanes may appear mid-session).
477+
* Seed for the idle-with-fleet allowance for the chat director (CL-7918:
478+
* replaces the former getLiveFleetCount closure; CL-7972 keeps it live via
479+
* the fleet-wake publisher). Omitted in exec (keeps the open-task nudge);
480+
* the TUI seeds it (fleet lanes may appear mid-session).
480481
*/
481482
allowIdleWithFleet?: boolean;
482483
/** Compaction governor re-entry (the reactor emits no event after compact). */
@@ -542,8 +543,9 @@ export function assembleChatAgent(wiring: ChatAgentWiring): AssembledChatAgent {
542543
onTasksChange: wiring.onTasksChange,
543544
requestContinuation: wiring.requestContinuation,
544545
provider: { ...wiring.getProvider() },
545-
// CL-7918: static idle-with-fleet allowance replaces the former
546-
// getLiveFleetCount closure; the live source id for retry stamping
546+
// CL-7918: idle-with-fleet seed replaces the former
547+
// getLiveFleetCount closure (CL-7972 keeps it live via the
548+
// fleet-wake publisher); the live source id for retry stamping
547549
// is reactor-tracked now (no getProviderId).
548550
allowIdleWithFleet: wiring.allowIdleWithFleet,
549551
},

‎src/tui/runner/exit.test.ts‎

Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,14 @@ import { getLogger } from "@intx/log";
55
import type { InferenceSource } from "@intx/types/runtime";
66

77
import * as codexSession from "../../auth/codex/session.js";
8+
import { createChatDirector } from "../../agent/director.js";
9+
import { createSubAgentSessionStore } from "../../subagent/session-store.js";
10+
import type {
11+
ReactorAction,
12+
ReactorCapabilities,
13+
ReactorInboundEvent,
14+
ReactorState,
15+
} from "@intx/types/runtime";
816
import { LOG_NAMESPACE_ROOT } from "../../branding.js";
917
import { defined } from "../../../tests/helpers/defined.js";
1018
import {
@@ -259,3 +267,126 @@ describe("agentProxy.send vs /clear", () => {
259267
}
260268
});
261269
});
270+
271+
const rebuildMockState: ReactorState = {} as unknown as ReactorState;
272+
273+
const rebuildMockCapabilities: ReactorCapabilities = {
274+
infer: (options) =>
275+
({
276+
type: "infer",
277+
...(options !== undefined ? { options } : {}),
278+
}) as ReactorAction,
279+
executeTools: (calls) => ({ type: "execute_tools", calls }),
280+
suspend: (gate) => ({ type: "suspend", gate }),
281+
fork: (mode, forkId) => ({ type: "fork", mode, forkId }),
282+
emit: (eventType, data) => ({ type: "emit", eventType, data }),
283+
reply: (content) => ({ type: "reply", content }),
284+
checkpoint: (message = "") => ({ type: "checkpoint", message }),
285+
compact: (compactor, reason) => ({ type: "compact", compactor, reason }),
286+
wait: () => ({ type: "wait" }),
287+
done: () => ({ type: "done" }),
288+
};
289+
290+
function rebuildManageTasksEvent(): ReactorInboundEvent {
291+
return {
292+
type: "inference.done",
293+
turn: {
294+
role: "assistant",
295+
model: "test",
296+
timestamp: 0,
297+
content: [
298+
{
299+
type: "tool_call",
300+
id: "m",
301+
name: "manage_tasks",
302+
arguments: {
303+
action: "create",
304+
tasks: [{ id: "t1", title: "work", status: "doing" }],
305+
},
306+
},
307+
],
308+
},
309+
usage: { input: 0, output: 1, cacheRead: 0, cacheWrite: 0, thinking: 0 },
310+
source: { model: "test-model" },
311+
} as unknown as ReactorInboundEvent;
312+
}
313+
314+
function rebuildTextTurn(): ReactorInboundEvent {
315+
return {
316+
type: "inference.done",
317+
turn: {
318+
role: "assistant",
319+
model: "test",
320+
timestamp: 0,
321+
content: [{ type: "text", text: "all set" }],
322+
},
323+
usage: { input: 10, output: 1, cacheRead: 0, cacheWrite: 0, thinking: 0 },
324+
source: { model: "test-model" },
325+
} as unknown as ReactorInboundEvent;
326+
}
327+
328+
describe("rebuild re-syncs idle-with-fleet while drained", () => {
329+
test("reload-if-idle and interrupt rebuilds resume the open-task nudge with no fleet transition", async () => {
330+
const store = createSubAgentSessionStore();
331+
const directorHolder: RunnerServices["directorHolder"] = {};
332+
const agent = recordingAgent([]);
333+
const { state, services } = stubSendLifecycle(agent);
334+
services.directorHolder =
335+
directorHolder as unknown as RunnerServices["directorHolder"];
336+
services.subAgentSessions =
337+
store as unknown as RunnerServices["subAgentSessions"];
338+
services.workflowHost = {
339+
reattach: () => undefined,
340+
} as unknown as RunnerServices["workflowHost"];
341+
services.cycleRecorder = {
342+
dispose: async () => "",
343+
reset: () => undefined,
344+
handleEvent: () => undefined,
345+
} as unknown as RunnerServices["cycleRecorder"];
346+
services.buildAgent = (async () => {
347+
// Every rebuild mints a fresh director from the static true seed (fleet
348+
// lanes may appear mid-session), exactly like the TUI session assembly.
349+
directorHolder.instance = createChatDirector("base", [], {
350+
onTasksChange: () => undefined,
351+
allowIdleWithFleet: true,
352+
});
353+
return agent;
354+
}) as unknown as RunnerServices["buildAgent"];
355+
const fleetEvents: unknown[] = [];
356+
services.emitter.on("event", (event: { type: string }) => {
357+
if (event.type === "fleet") fleetEvents.push(event);
358+
});
359+
await createRunLifecycle(state, services);
360+
const expectOpenTaskNudge = async (): Promise<void> => {
361+
const director = defined(
362+
directorHolder.instance,
363+
"directorHolder.instance",
364+
);
365+
await director.decide(
366+
rebuildManageTasksEvent(),
367+
rebuildMockState,
368+
rebuildMockCapabilities,
369+
);
370+
const actions = await director.decide(
371+
rebuildTextTurn(),
372+
rebuildMockState,
373+
rebuildMockCapabilities,
374+
);
375+
const list = Array.isArray(actions) ? actions : [actions];
376+
expect(list.some((action) => action.type === "infer")).toBe(true);
377+
};
378+
// Drained fleet: the idle reload rebuilds onto the static true seed.
379+
state.pendingReload = true;
380+
defined(state.reloadIfIdle, "reloadIfIdle")();
381+
await services.sessionOps.awaitTail();
382+
expect(state.fatalBuildError).toBeNull();
383+
await expectOpenTaskNudge();
384+
// The interrupt rebuild inherits the same seed.
385+
defined(state.interrupt, "interrupt")();
386+
await services.sessionOps.awaitTail();
387+
expect(state.fatalBuildError).toBeNull();
388+
await expectOpenTaskNudge();
389+
expect(store.list()).toEqual([]);
390+
expect(fleetEvents).toEqual([]);
391+
});
392+
});

‎src/tui/runner/exit.ts‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import {
1313
import { getLogger } from "@intx/log";
1414
import type { InferenceSource } from "@intx/types/runtime";
1515
import { consumeStream } from "../../session/stream-consumer.js";
16+
import { liveFleetCount } from "../../subagent/index.js";
1617
import { getTelemetry } from "../../telemetry/singleton.js";
1718
import { onTurnBoundary } from "../../agent/reactor-events.js";
1819
import { setAgentSourceUnlessClosed } from "../agent-source-sync.js";
@@ -145,6 +146,18 @@ export async function closeAgentForRebuild(
145146
}
146147
}
147148

149+
// Rebuilt directors seed allowIdleWithFleet=true (fleet lanes may appear
150+
// mid-session), so a rebuild while drained must re-sync the new director from
151+
// the live fleet count — otherwise the open-task nudge stays suppressed until
152+
// the next fleet transition, which never comes for an already-drained fleet.
153+
function resyncIdleWithFleetFlag(
154+
services: Pick<RunnerServices, "directorHolder" | "subAgentSessions">,
155+
): void {
156+
services.directorHolder.instance?.setAllowIdleWithFleet(
157+
liveFleetCount(services.subAgentSessions.list()) > 0,
158+
);
159+
}
160+
148161
// Every rebuild site funnels its failure (a lock left held by a failed
149162
// close, or any other buildAgent failure) through here so it surfaces as a
150163
// plain-language, caught error rather than an unhandled rejection.
@@ -345,6 +358,7 @@ export async function createRunLifecycle(
345358
throw new AgentContextLockError(state.workdir);
346359
}
347360
state.currentAgent = await services.buildAgent();
361+
resyncIdleWithFleetFlag(services);
348362
state.streamPromise = consumeStream(
349363
liveAgent(state).stream(),
350364
streamSink,
@@ -547,6 +561,7 @@ export async function createRunLifecycle(
547561
throw new AgentContextLockError(state.workdir);
548562
}
549563
state.currentAgent = await services.buildAgent();
564+
resyncIdleWithFleetFlag(services);
550565
services.cycleRecorder.reset();
551566
state.streamPromise = consumeStream(
552567
liveAgent(state).stream(),

‎src/tui/runner/session.ts‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -619,8 +619,9 @@ export async function assembleTUISession(
619619
inactivityTimeoutMs: config.inactivityTimeoutMs ?? 750_000,
620620
totalTimeoutMs: config.totalTimeoutMs,
621621
onTasksChange: (tasks) => emitter.emit("tasks", tasks),
622-
// CL-7918: static idle-with-fleet allowance (fleet lanes may appear
623-
// mid-session); retry stamping tracks the live source id in-reactor now.
622+
// CL-7918 seed for the idle-with-fleet allowance (fleet lanes may appear
623+
// mid-session; CL-7972 keeps it live via the fleet-wake publisher);
624+
// retry stamping tracks the live source id in-reactor now.
624625
allowIdleWithFleet: true,
625626
requestContinuation: () => {
626627
const targetAgent = liveAgent(state);

‎src/tui/runner/wiring.ask-wake.test.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,36 @@ test("failed reset releases publication without flushing partially cancelled wor
6868
}
6969
});
7070

71+
test("fleet-wake publisher reports fleet-count transitions to the idle-with-fleet follower", () => {
72+
const store = createSubAgentSessionStore();
73+
const emitter = new EventEmitter();
74+
const seen: number[] = [];
75+
const publisher = createFleetWakePublisher(store, emitter, (running) => {
76+
seen.push(running);
77+
});
78+
// Steady empty state: no transition, no callback.
79+
publisher.publish();
80+
expect(seen).toEqual([]);
81+
store.start({
82+
id: "worker-one",
83+
agentId: "builder",
84+
description: "worker-one",
85+
brief: "build",
86+
});
87+
// A started lane counts as running: 0 -> 1 transition.
88+
publisher.publish();
89+
expect(seen).toEqual([1]);
90+
store.markRunning("worker-one");
91+
publisher.publish();
92+
expect(seen).toEqual([1]);
93+
// Steady fleet: no repeat callback.
94+
publisher.publish();
95+
expect(seen).toEqual([1]);
96+
store.cancelAll("Session cleared");
97+
publisher.publish();
98+
expect(seen).toEqual([1, 0]);
99+
});
100+
71101
for (const phase of ["settled", "prequeued", "deferred"] as const) {
72102
test(`rotation suppresses old worker snapshots (${phase}), including repeated resets`, async () => {
73103
await withTestRenderer(

0 commit comments

Comments
 (0)