Merge branch 'main' into timothyjlaurent/stuck-spinners
This commit is contained in:
@@ -747,7 +747,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 {
|
||||
@@ -774,14 +774,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();
|
||||
}
|
||||
@@ -820,7 +815,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 {
|
||||
@@ -849,14 +844,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();
|
||||
}
|
||||
@@ -1275,14 +1265,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.
|
||||
@@ -1319,8 +1304,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();
|
||||
@@ -1361,11 +1345,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();
|
||||
@@ -1398,9 +1380,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);
|
||||
@@ -1443,8 +1424,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);
|
||||
@@ -1463,7 +1443,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 {
|
||||
@@ -1488,14 +1468,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();
|
||||
}
|
||||
@@ -1528,10 +1503,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();
|
||||
@@ -1561,14 +1537,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();
|
||||
@@ -1583,8 +1560,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