feat(FN-3481): clean up ephemeral runtimes and spawned agents immediately o
Merges three major changesets: ephemeral agent cleanup for FN-3481 (runtime and spawned agent teardown), a fix for planning-mode refine continuation flow (FN-3209) plus a new local startup script, and chat SSE broadcast isolation with QuickChat backend unification. Key components affected include th Fusion-Task-Id: FN-3481
This commit is contained in:
@@ -3854,7 +3854,6 @@ describe("swallowed async store failure observability", () => {
|
||||
});
|
||||
|
||||
await (executor as any).terminateChildAgent("child-007");
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
await Promise.resolve();
|
||||
|
||||
expect(warnSpy).toHaveBeenCalledWith(
|
||||
@@ -10809,71 +10808,38 @@ describe("Agent Spawning - Child Termination", () => {
|
||||
expect(internals.totalSpawnedCount).toBe(0);
|
||||
});
|
||||
|
||||
it("terminateChildAgent auto-deletes agent after 5 second delay", async () => {
|
||||
vi.useFakeTimers();
|
||||
it("terminateChildAgent auto-deletes agent immediately", async () => {
|
||||
const agentStore = createMockAgentStore() as any;
|
||||
agentStore.deleteAgent = vi.fn().mockResolvedValue(undefined);
|
||||
const store = createMockStore();
|
||||
|
||||
try {
|
||||
const agentStore = createMockAgentStore() as any;
|
||||
// Add deleteAgent mock to the agent store
|
||||
agentStore.deleteAgent = vi.fn().mockResolvedValue(undefined);
|
||||
const store = createMockStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore } as any);
|
||||
const internals = executor as any;
|
||||
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore } as any);
|
||||
const internals = executor as any;
|
||||
const mockSession = { dispose: vi.fn() };
|
||||
const childId = "agent-auto-delete-test";
|
||||
internals.childSessions.set(childId, mockSession);
|
||||
internals.totalSpawnedCount = 1;
|
||||
|
||||
const mockSession = { dispose: vi.fn() };
|
||||
const childId = "agent-auto-delete-test";
|
||||
internals.childSessions.set(childId, mockSession);
|
||||
internals.totalSpawnedCount = 1;
|
||||
await internals.terminateChildAgent(childId);
|
||||
|
||||
// Terminate the child
|
||||
const terminatePromise = internals.terminateChildAgent(childId);
|
||||
await terminatePromise;
|
||||
|
||||
// Session should be disposed immediately
|
||||
expect(mockSession.dispose).toHaveBeenCalled();
|
||||
expect(internals.pendingEphemeralDeletions.has(childId)).toBe(true);
|
||||
|
||||
// deleteAgent should not be called yet (before 5 seconds)
|
||||
expect(agentStore.deleteAgent).not.toHaveBeenCalled();
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// Now deleteAgent should have been called
|
||||
expect(agentStore.deleteAgent).toHaveBeenCalledTimes(1);
|
||||
expect(agentStore.deleteAgent).toHaveBeenCalledWith(childId);
|
||||
expect(internals.pendingEphemeralDeletions.has(childId)).toBe(false);
|
||||
|
||||
// Should not throw even when delete fails
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
expect(mockSession.dispose).toHaveBeenCalled();
|
||||
expect(agentStore.deleteAgent).toHaveBeenCalledTimes(1);
|
||||
expect(agentStore.deleteAgent).toHaveBeenCalledWith(childId);
|
||||
expect(internals.pendingEphemeralDeletions.has(childId)).toBe(false);
|
||||
});
|
||||
|
||||
it("disposeEphemeralTimers clears pending spawned cleanup timers", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
const agentStore = createMockAgentStore() as any;
|
||||
agentStore.deleteAgent = vi.fn().mockResolvedValue(undefined);
|
||||
const store = createMockStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore } as any);
|
||||
const internals = executor as any;
|
||||
it("disposeEphemeralTimers clears pending deletion bookkeeping", async () => {
|
||||
const agentStore = createMockAgentStore() as any;
|
||||
agentStore.deleteAgent = vi.fn().mockResolvedValue(undefined);
|
||||
const store = createMockStore();
|
||||
const executor = new TaskExecutor(store, "/tmp/test", { agentStore } as any);
|
||||
const internals = executor as any;
|
||||
|
||||
internals.childSessions.set("agent-dispose-test", { dispose: vi.fn() });
|
||||
internals.totalSpawnedCount = 1;
|
||||
await internals.terminateChildAgent("agent-dispose-test");
|
||||
expect(internals.pendingEphemeralDeletions.has("agent-dispose-test")).toBe(true);
|
||||
internals.pendingEphemeralDeletions.add("agent-dispose-test");
|
||||
executor.disposeEphemeralTimers();
|
||||
|
||||
executor.disposeEphemeralTimers();
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
expect(agentStore.deleteAgent).not.toHaveBeenCalled();
|
||||
expect(internals.pendingEphemeralDeletions.size).toBe(0);
|
||||
expect(internals.ephemeralCleanupTimers.size).toBe(0);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
expect(internals.pendingEphemeralDeletions.size).toBe(0);
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -603,10 +603,8 @@ export class TaskExecutor {
|
||||
private completedTaskWatchdogs = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
/** One-shot watchdogs for workflow reruns that should have bounced back to in-progress. */
|
||||
private workflowRerunWatchdogs = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
/** Set of ephemeral spawned agent IDs with scheduled cleanup (prevents duplicate deletion attempts). */
|
||||
/** Set of ephemeral spawned agent IDs with in-flight cleanup (prevents duplicate deletion attempts). */
|
||||
private pendingEphemeralDeletions = new Set<string>();
|
||||
/** Map of spawned agent IDs to scheduled cleanup timer handles for shutdown disposal. */
|
||||
private ephemeralCleanupTimers = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
|
||||
private async finalizeAlreadyReviewedTask(taskId: string): Promise<"merged" | "blocked" | "missing"> {
|
||||
const latestTask = await this.store.getTask(taskId);
|
||||
@@ -736,15 +734,7 @@ export class TaskExecutor {
|
||||
}
|
||||
|
||||
disposeEphemeralTimers(): void {
|
||||
const timerCount = this.ephemeralCleanupTimers.size;
|
||||
for (const timerId of this.ephemeralCleanupTimers.values()) {
|
||||
clearTimeout(timerId);
|
||||
}
|
||||
this.ephemeralCleanupTimers.clear();
|
||||
this.pendingEphemeralDeletions.clear();
|
||||
if (timerCount > 0) {
|
||||
executorLog.log(`Cleared ${timerCount} pending spawned-agent cleanup timer(s)`);
|
||||
}
|
||||
}
|
||||
|
||||
private isBenignEphemeralDeleteRaceError(agentId: string, err: unknown): boolean {
|
||||
@@ -6656,23 +6646,17 @@ and show an appropriate message to the user.\`
|
||||
executorLog.warn(`Failed to update spawned child ${childId} state to 'terminated' during cleanup: ${msg}`);
|
||||
}
|
||||
|
||||
// Auto-delete the child agent after a short delay so the UI can observe
|
||||
// the terminal state before the agent is removed.
|
||||
this.pendingEphemeralDeletions.add(childId);
|
||||
const timerId = setTimeout(async () => {
|
||||
this.ephemeralCleanupTimers.delete(childId);
|
||||
this.pendingEphemeralDeletions.delete(childId);
|
||||
try {
|
||||
await this.options.agentStore?.deleteAgent(childId);
|
||||
} catch (err: unknown) {
|
||||
if (this.isBenignEphemeralDeleteRaceError(childId, err)) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await this.options.agentStore?.deleteAgent(childId);
|
||||
} catch (err: unknown) {
|
||||
if (!this.isBenignEphemeralDeleteRaceError(childId, err)) {
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
executorLog.warn(`Failed to delete spawned agent ${childId}: ${msg}`);
|
||||
}
|
||||
}, 5000);
|
||||
this.ephemeralCleanupTimers.set(childId, timerId);
|
||||
} finally {
|
||||
this.pendingEphemeralDeletions.delete(childId);
|
||||
}
|
||||
|
||||
this.totalSpawnedCount = Math.max(0, this.totalSpawnedCount - 1);
|
||||
}
|
||||
|
||||
@@ -746,7 +746,7 @@ describe("InProcessRuntime", () => {
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
it("auto-deletes task-worker agent on task completion after 5 second delay", async () => {
|
||||
it("auto-deletes task-worker agent on task completion immediately", async () => {
|
||||
vi.useFakeTimers();
|
||||
|
||||
try {
|
||||
@@ -773,14 +773,9 @@ describe("InProcessRuntime", () => {
|
||||
deleteAgentSpy.mockClear();
|
||||
executorOptions.onComplete?.({ id: "FN-AUTO1" } as Task);
|
||||
|
||||
// Verify deleteAgent was not called immediately (before 5 seconds)
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// Now deleteAgent should have been called
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
@@ -819,7 +814,7 @@ describe("InProcessRuntime", () => {
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
it("auto-deletes task-worker agent on task error after 5 second delay", async () => {
|
||||
it("auto-deletes task-worker agent on task error immediately", async () => {
|
||||
vi.useFakeTimers();
|
||||
|
||||
try {
|
||||
@@ -848,14 +843,9 @@ describe("InProcessRuntime", () => {
|
||||
deleteAgentSpy.mockClear();
|
||||
executorOptions.onError?.({ id: "FN-AUTO2" } as Task, new Error("Task failed"));
|
||||
|
||||
// Verify deleteAgent was not called immediately (before 5 seconds)
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// Now deleteAgent should have been called
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
@@ -1274,14 +1264,9 @@ describe("InProcessRuntime", () => {
|
||||
// Wait for async handler
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Verify deleteAgent was NOT called immediately (needs 5s delay)
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// Now deleteAgent should have been called
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id);
|
||||
|
||||
// Note: We verified deleteAgent was called, which is the key behavior.
|
||||
@@ -1318,8 +1303,7 @@ describe("InProcessRuntime", () => {
|
||||
// Wait for async handler
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Advance timers to ensure cleanup would have run
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// deleteAgent should NOT have been called for non-ephemeral agent
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
@@ -1360,11 +1344,9 @@ describe("InProcessRuntime", () => {
|
||||
// Wait for async handlers
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// deleteAgent should have been called only once (deduplicated)
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
@@ -1397,9 +1379,8 @@ describe("InProcessRuntime", () => {
|
||||
// Emit termination event
|
||||
store.emit("agent:stateChanged", agent.id, "running", "terminated");
|
||||
|
||||
// Wait for async handler, then fire delayed cleanup
|
||||
// Wait for async handler
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// Cleanup should still be attempted
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
@@ -1442,8 +1423,7 @@ describe("InProcessRuntime", () => {
|
||||
// Wait for async handler
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Advance timers to trigger deletion
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Should have attempted deletion
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
@@ -1462,7 +1442,7 @@ describe("InProcessRuntime", () => {
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
it("clears pending timers on runtime stop", async () => {
|
||||
it("handles runtime stop racing with in-flight cleanup", async () => {
|
||||
vi.useFakeTimers();
|
||||
|
||||
try {
|
||||
@@ -1487,14 +1467,9 @@ describe("InProcessRuntime", () => {
|
||||
// Wait for async handler
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Stop runtime before timer fires
|
||||
await runtime.stop();
|
||||
|
||||
// Advance timers - deletion should NOT happen because timer was cleared
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
|
||||
// deleteAgent should NOT have been called (timer was cleared)
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
expect(deleteAgentSpy.mock.calls.length).toBeLessThanOrEqual(1);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
@@ -1527,10 +1502,11 @@ describe("InProcessRuntime", () => {
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
// Advance timers by 5 seconds
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
// deleteAgent should have been called for spawned ephemeral agent
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
@@ -1560,14 +1536,15 @@ describe("InProcessRuntime", () => {
|
||||
executorOptions.onComplete?.({ id: "FN-DUP-COMPLETE" } as Task);
|
||||
store.emit("agent:stateChanged", worker!.id, "running", "terminated");
|
||||
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(deleteAgentSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
it("clears onComplete cleanup timer on stop", async () => {
|
||||
it("handles onComplete cleanup racing with runtime stop", async () => {
|
||||
vi.useFakeTimers();
|
||||
try {
|
||||
await runtime.start();
|
||||
@@ -1582,8 +1559,7 @@ describe("InProcessRuntime", () => {
|
||||
executorOptions.onComplete?.({ id: "FN-STOP-COMPLETE" } as Task);
|
||||
|
||||
await runtime.stop();
|
||||
await vi.advanceTimersByTimeAsync(5000);
|
||||
expect(deleteAgentSpy).not.toHaveBeenCalled();
|
||||
expect(deleteAgentSpy.mock.calls.length).toBeLessThanOrEqual(1);
|
||||
} finally {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
|
||||
@@ -106,10 +106,8 @@ export class InProcessRuntime
|
||||
private triageProcessor?: TriageProcessor;
|
||||
private messageStore?: MessageStore;
|
||||
private concurrencyChangedListener?: (state: { globalMaxConcurrent: number }) => void;
|
||||
/** Set of agent IDs with scheduled ephemeral cleanup (prevents duplicate deletion) */
|
||||
/** Set of agent IDs with in-flight ephemeral cleanup (prevents duplicate deletion) */
|
||||
private pendingEphemeralDeletions = new Set<string>();
|
||||
/** Map of agent IDs to their cleanup timer IDs */
|
||||
private ephemeralCleanupTimers = new Map<string, ReturnType<typeof setTimeout>>();
|
||||
/** Listener for agent:stateChanged events to clean up terminated ephemeral agents */
|
||||
private ephemeralTerminationListener?: (agentId: string, from: import("@fusion/core").AgentState, to: import("@fusion/core").AgentState) => void;
|
||||
/**
|
||||
@@ -459,11 +457,7 @@ export class InProcessRuntime
|
||||
});
|
||||
this.taskAgentMap.delete(task.id);
|
||||
if (!ephemeral) return;
|
||||
// Auto-delete the task-worker agent after a short delay so the UI
|
||||
// can observe the terminal state before the agent is removed.
|
||||
const timerId = setTimeout(async () => {
|
||||
this.ephemeralCleanupTimers.delete(agentId);
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
void (async () => {
|
||||
try {
|
||||
await this.agentStore?.deleteAgent(agentId);
|
||||
} catch (err: unknown) {
|
||||
@@ -472,9 +466,10 @@ export class InProcessRuntime
|
||||
}
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
runtimeLog.warn(`Failed to delete agent ${agentId} after completion: ${msg}`);
|
||||
} finally {
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
}
|
||||
}, 5000);
|
||||
this.ephemeralCleanupTimers.set(agentId, timerId);
|
||||
})();
|
||||
}
|
||||
},
|
||||
onError: (task, error) => {
|
||||
@@ -514,11 +509,7 @@ export class InProcessRuntime
|
||||
});
|
||||
this.taskAgentMap.delete(task.id);
|
||||
if (!ephemeral) return;
|
||||
// Auto-delete the task-worker agent after a short delay so the UI
|
||||
// can observe the terminal state before the agent is removed.
|
||||
const timerId = setTimeout(async () => {
|
||||
this.ephemeralCleanupTimers.delete(agentId);
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
void (async () => {
|
||||
try {
|
||||
await this.agentStore?.deleteAgent(agentId);
|
||||
} catch (err: unknown) {
|
||||
@@ -527,9 +518,10 @@ export class InProcessRuntime
|
||||
}
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
runtimeLog.warn(`Failed to delete agent ${agentId} after error: ${msg}`);
|
||||
} finally {
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
}
|
||||
}, 5000);
|
||||
this.ephemeralCleanupTimers.set(agentId, timerId);
|
||||
})();
|
||||
}
|
||||
},
|
||||
};
|
||||
@@ -622,22 +614,18 @@ export class InProcessRuntime
|
||||
if (!agent) return;
|
||||
if (!isEphemeralAgent(agent)) return;
|
||||
|
||||
// Schedule deletion after delay so UI can observe terminal state
|
||||
this.pendingEphemeralDeletions.add(agentId);
|
||||
const timerId = setTimeout(async () => {
|
||||
this.ephemeralCleanupTimers.delete(agentId);
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
try {
|
||||
await this.agentStore?.deleteAgent(agentId);
|
||||
} catch (err: unknown) {
|
||||
if (this.isBenignEphemeralDeleteRaceError(agentId, err)) {
|
||||
return;
|
||||
}
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
runtimeLog.warn(`Failed to delete ephemeral agent ${agentId} after termination: ${msg}`);
|
||||
try {
|
||||
await this.agentStore?.deleteAgent(agentId);
|
||||
} catch (err: unknown) {
|
||||
if (this.isBenignEphemeralDeleteRaceError(agentId, err)) {
|
||||
return;
|
||||
}
|
||||
}, 5000);
|
||||
this.ephemeralCleanupTimers.set(agentId, timerId);
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
runtimeLog.warn(`Failed to delete ephemeral agent ${agentId} after termination: ${msg}`);
|
||||
} finally {
|
||||
this.pendingEphemeralDeletions.delete(agentId);
|
||||
}
|
||||
} catch (err: unknown) {
|
||||
const msg = err instanceof Error ? err.message : String(err);
|
||||
runtimeLog.warn(`Failed to process termination event for agent ${agentId}: ${msg}`);
|
||||
@@ -912,12 +900,6 @@ export class InProcessRuntime
|
||||
this.ephemeralTerminationListener = undefined;
|
||||
runtimeLog.log("AgentStore agent:stateChanged listener removed");
|
||||
}
|
||||
// Clear any pending ephemeral cleanup timers to prevent leaks during shutdown
|
||||
for (const [agentId, timerId] of this.ephemeralCleanupTimers) {
|
||||
clearTimeout(timerId);
|
||||
runtimeLog.log(`Cleared pending cleanup timer for ephemeral agent ${agentId}`);
|
||||
}
|
||||
this.ephemeralCleanupTimers.clear();
|
||||
this.pendingEphemeralDeletions.clear();
|
||||
this.executor?.disposeEphemeralTimers();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user