Skip to content

Commit 9f5920b

Browse files
fix(tui): re-sync idle-with-fleet flag on director rebuilds (#1058)
1 parent 07ce93e commit 9f5920b

2 files changed

Lines changed: 146 additions & 0 deletions

File tree

‎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(),

0 commit comments

Comments
 (0)