diff --git a/packages/engine/src/__tests__/ephemeral-worker-manager.test.ts b/packages/engine/src/__tests__/ephemeral-worker-manager.test.ts new file mode 100644 index 0000000000..3697b1aaf2 --- /dev/null +++ b/packages/engine/src/__tests__/ephemeral-worker-manager.test.ts @@ -0,0 +1,438 @@ +import { EventEmitter } from "node:events"; +import { describe, it, expect, vi, beforeEach } from "vitest"; +import type { Agent, AgentStore, Task, TaskStore } from "@fusion/core"; +import { EphemeralWorkerManager } from "../ephemeral-worker-manager.js"; + +const BASE_TIME = "2026-07-03T00:00:00.000Z"; + +type AgentPatch = Partial & { metadata?: Record; runtimeConfig?: Record }; + +type FakeAgentStore = AgentStore & EventEmitter & { + agents: Map; + createAgent: ReturnType; + deleteAgent: ReturnType; + assignTask: ReturnType; + syncExecutionTaskLink: ReturnType; + updateAgentState: ReturnType; + findAgentByName: ReturnType; + listAgents: ReturnType; + getAgent: ReturnType; +}; + +type Harness = { + agentStore: FakeAgentStore; + taskStore: TaskStore & { tasks: Map; getTask: ReturnType }; + logger: { log: ReturnType; warn: ReturnType }; + externalPending: ReturnType; + getSettings: ReturnType; + manager: EphemeralWorkerManager; +}; + +function makeAgent(id: string, patch: AgentPatch = {}): Agent { + return { + id, + name: patch.name ?? id, + role: patch.role ?? "executor", + state: patch.state ?? "idle", + taskId: patch.taskId, + createdAt: patch.createdAt ?? BASE_TIME, + updatedAt: patch.updatedAt ?? BASE_TIME, + metadata: patch.metadata ?? {}, + runtimeConfig: patch.runtimeConfig, + } as Agent; +} + +function makeTask(id: string, patch: Partial = {}): Task { + return { + id, + title: id, + description: "test task", + column: "in-progress", + steps: [], + createdAt: BASE_TIME, + updatedAt: BASE_TIME, + ...patch, + } as Task; +} + +function createAgentStore(initialAgents: Agent[] = []): FakeAgentStore { + const emitter = new EventEmitter() as FakeAgentStore; + emitter.agents = new Map(initialAgents.map((agent) => [agent.id, structuredClone(agent)])); + + emitter.getAgent = vi.fn(async (agentId: string) => emitter.agents.get(agentId) ?? null); + emitter.listAgents = vi.fn(async () => Array.from(emitter.agents.values())); + emitter.findAgentByName = vi.fn(async (name: string) => Array.from(emitter.agents.values()).find((agent) => agent.name === name) ?? null); + emitter.createAgent = vi.fn(async (input: Partial) => { + const id = `agent-${emitter.agents.size + 1}`; + const agent = makeAgent(id, { + name: input.name ?? id, + role: input.role ?? "executor", + state: input.state ?? "idle", + metadata: input.metadata as Record | undefined, + runtimeConfig: input.runtimeConfig as Record | undefined, + }); + emitter.agents.set(agent.id, agent); + emitter.emit("agent:created", agent); + return agent; + }); + emitter.assignTask = vi.fn(async (agentId: string, taskId: string) => { + const agent = emitter.agents.get(agentId); + if (agent) { + agent.taskId = taskId; + agent.updatedAt = BASE_TIME; + emitter.emit("agent:assigned", agent, taskId); + } + return agent ?? null; + }); + emitter.syncExecutionTaskLink = vi.fn(async (agentId: string, taskId?: string) => { + const agent = emitter.agents.get(agentId); + if (agent) { + if (taskId) agent.taskId = taskId; + else delete (agent as { taskId?: string }).taskId; + agent.updatedAt = BASE_TIME; + } + return agent ?? null; + }); + emitter.updateAgentState = vi.fn(async (agentId: string, state: Agent["state"]) => { + const agent = emitter.agents.get(agentId); + if (!agent) throw new Error(`Agent ${agentId} not found`); + const from = agent.state; + agent.state = state; + agent.updatedAt = BASE_TIME; + emitter.emit("agent:stateChanged", agentId, from, state); + return agent; + }); + emitter.deleteAgent = vi.fn(async (agentId: string) => { + if (!emitter.agents.has(agentId)) throw new Error(`Agent ${agentId} not found`); + emitter.agents.delete(agentId); + emitter.emit("agent:deleted", agentId); + }); + return emitter; +} + +function createHarness(initialAgents: Agent[] = []): Harness { + const agentStore = createAgentStore(initialAgents); + const tasks = new Map(); + const taskStore = { + tasks, + getTask: vi.fn(async (taskId: string) => tasks.get(taskId) ?? null), + } as unknown as Harness["taskStore"]; + const logger = { log: vi.fn(), warn: vi.fn() }; + const externalPending = vi.fn(() => false); + const getSettings = vi.fn(async () => ({ ephemeralAgentsEnabled: true })); + const manager = new EphemeralWorkerManager({ + agentStore, + taskStore, + logger, + isDeletionPendingExternal: externalPending, + getSettings, + }); + return { agentStore, taskStore, logger, externalPending, getSettings, manager }; +} + +async function flushMicrotasks(turns = 6): Promise { + for (let i = 0; i < turns; i += 1) await Promise.resolve(); +} + +describe("EphemeralWorkerManager", () => { + let harness: Harness; + + beforeEach(() => { + vi.clearAllMocks(); + harness = createHarness(); + }); + + describe("task ownership", () => { + it("uses durable assigned agents without creating task workers", async () => { + const durable = makeAgent("durable-1", { name: "Durable", state: "idle", metadata: {} }); + harness.agentStore.agents.set(durable.id, durable); + + const owner = await harness.manager.onTaskStart(makeTask("FN-DURABLE", { assignedAgentId: durable.id })); + + expect(owner).toEqual({ agentId: durable.id, ephemeral: false }); + expect(harness.agentStore.syncExecutionTaskLink).toHaveBeenCalledWith(durable.id, "FN-DURABLE"); + expect(harness.agentStore.createAgent).not.toHaveBeenCalled(); + expect(harness.agentStore.updateAgentState).toHaveBeenNthCalledWith(1, durable.id, "active"); + expect(harness.agentStore.updateAgentState).toHaveBeenNthCalledWith(2, durable.id, "running"); + expect(harness.manager.getOwner("FN-DURABLE")).toEqual(owner); + }); + + it("creates, assigns, and runs an ephemeral worker for unassigned tasks", async () => { + const owner = await harness.manager.onTaskStart(makeTask("FN-EPHEMERAL")); + + expect(owner).toEqual({ agentId: "agent-1", ephemeral: true }); + expect(harness.agentStore.createAgent).toHaveBeenCalledWith(expect.objectContaining({ + name: "executor-FN-EPHEMERAL", + role: "executor", + metadata: expect.objectContaining({ agentKind: "task-worker", taskWorker: true }), + runtimeConfig: { enabled: false }, + })); + expect(harness.agentStore.assignTask).toHaveBeenCalledWith("agent-1", "FN-EPHEMERAL"); + expect(harness.agentStore.updateAgentState).toHaveBeenNthCalledWith(1, "agent-1", "active"); + expect(harness.agentStore.updateAgentState).toHaveBeenNthCalledWith(2, "agent-1", "running"); + }); + + it("reuses an existing cross-restart ephemeral worker for the same task", async () => { + const existing = makeAgent("worker-1", { + name: "executor-FN-REUSE", + taskId: "FN-REUSE", + metadata: { agentKind: "task-worker" }, + runtimeConfig: { enabled: false }, + }); + harness = createHarness([existing]); + + const owner = await harness.manager.onTaskStart(makeTask("FN-REUSE")); + + expect(owner).toEqual({ agentId: existing.id, ephemeral: true }); + expect(harness.agentStore.createAgent).not.toHaveBeenCalled(); + expect(harness.logger.log).toHaveBeenCalledWith(expect.stringContaining("Reusing existing ephemeral worker")); + }); + + it("deletes stale same-name workers before respawning", async () => { + const stale = makeAgent("worker-stale", { + name: "executor-FN-RESPAWN", + taskId: "FN-OLD", + metadata: { agentKind: "task-worker" }, + runtimeConfig: { enabled: false }, + }); + harness = createHarness([stale]); + + const owner = await harness.manager.onTaskStart(makeTask("FN-RESPAWN")); + + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith(stale.id); + expect(owner).toEqual({ agentId: "agent-1", ephemeral: true }); + expect(harness.agentStore.agents.has(stale.id)).toBe(false); + }); + + it("refuses to spawn when ephemeral agents are disabled", async () => { + harness.getSettings.mockResolvedValueOnce({ ephemeralAgentsEnabled: false }); + + const owner = await harness.manager.onTaskStart(makeTask("FN-DISABLED")); + + expect(owner).toBeNull(); + expect(harness.agentStore.createAgent).not.toHaveBeenCalled(); + expect(harness.logger.warn).toHaveBeenCalledWith(expect.stringContaining("ephemeralAgentsEnabled=false")); + }); + + it("falls back to task-worker ownership when assignedAgentId points to an ephemeral", async () => { + const ephemeral = makeAgent("child-1", { + name: "child", + metadata: { agentKind: "task-worker" }, + runtimeConfig: { enabled: false }, + }); + harness.agentStore.agents.set(ephemeral.id, ephemeral); + + const owner = await harness.manager.onTaskStart(makeTask("FN-ASSIGNED-EPHEMERAL", { assignedAgentId: ephemeral.id })); + + expect(owner).toEqual({ agentId: "agent-2", ephemeral: true }); + expect(harness.agentStore.syncExecutionTaskLink).not.toHaveBeenCalledWith(ephemeral.id, "FN-ASSIGNED-EPHEMERAL"); + expect(harness.agentStore.createAgent).toHaveBeenCalledWith(expect.objectContaining({ name: "executor-FN-ASSIGNED-EPHEMERAL" })); + }); + }); + + describe("completion and error cleanup", () => { + it("returns durable owners to active and does not delete them on completion or error", async () => { + const durable = makeAgent("durable-cleanup", { name: "Durable Cleanup", state: "active" }); + harness.agentStore.agents.set(durable.id, durable); + + await harness.manager.onTaskStart(makeTask("FN-DURABLE-COMPLETE", { assignedAgentId: durable.id })); + await harness.manager.onTaskComplete("FN-DURABLE-COMPLETE"); + expect(harness.agentStore.syncExecutionTaskLink).toHaveBeenLastCalledWith(durable.id, undefined); + expect(harness.agentStore.deleteAgent).not.toHaveBeenCalledWith(durable.id); + expect((await harness.agentStore.getAgent(durable.id))?.state).toBe("active"); + expect((await harness.agentStore.getAgent(durable.id))?.taskId).toBeUndefined(); + + await harness.manager.onTaskStart(makeTask("FN-DURABLE-ERROR", { assignedAgentId: durable.id })); + await harness.manager.onTaskError("FN-DURABLE-ERROR"); + expect(harness.agentStore.deleteAgent).not.toHaveBeenCalledWith(durable.id); + expect((await harness.agentStore.getAgent(durable.id))?.state).toBe("active"); + }); + + it("deletes ephemeral owners on completion and error", async () => { + await harness.manager.onTaskStart(makeTask("FN-COMPLETE")); + await harness.manager.onTaskComplete("FN-COMPLETE"); + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith("agent-1"); + expect(harness.manager.isDeletionPending("agent-1")).toBe(false); + + await harness.manager.onTaskStart(makeTask("FN-ERROR")); + await harness.manager.onTaskError("FN-ERROR"); + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith("agent-1"); + }); + + it("recovers a cross-restart owner by name during completion cleanup", async () => { + const existing = makeAgent("worker-disk", { + name: "executor-FN-DISK", + taskId: "FN-DISK", + metadata: { agentKind: "task-worker" }, + runtimeConfig: { enabled: false }, + }); + harness = createHarness([existing]); + + await harness.manager.onTaskComplete("FN-DISK"); + + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith(existing.id); + expect(harness.logger.log).toHaveBeenCalledWith(expect.stringContaining("Recovered ephemeral owner")); + }); + + it("logs genuine cleanup warnings and suppresses benign delete races", async () => { + await harness.manager.onTaskStart(makeTask("FN-WARN")); + harness.agentStore.deleteAgent.mockRejectedValueOnce(new Error("delete failed")); + await harness.manager.onTaskError("FN-WARN"); + expect(harness.logger.warn).toHaveBeenCalledWith(expect.stringContaining("Failed to delete agent agent-1 after error: delete failed")); + + harness.logger.warn.mockClear(); + const benignOwner = await harness.manager.onTaskStart(makeTask("FN-BENIGN")); + expect(benignOwner).toBeDefined(); + harness.agentStore.deleteAgent.mockRejectedValueOnce(new Error(`Agent ${benignOwner!.agentId} not found`)); + await harness.manager.onTaskComplete("FN-BENIGN"); + expect(harness.logger.warn).not.toHaveBeenCalledWith(expect.stringContaining("Failed to delete agent")); + }); + }); + + describe("halt listener cleanup", () => { + it("deletes task-worker and spawned ephemerals that enter halted states", async () => { + const taskWorker = makeAgent("worker-paused", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + const spawned = makeAgent("spawned-error", { metadata: { type: "spawned" }, runtimeConfig: { enabled: false } }); + harness.agentStore.agents.set(taskWorker.id, taskWorker); + harness.agentStore.agents.set(spawned.id, spawned); + harness.manager.attachStateChangeListener(); + + harness.agentStore.emit("agent:stateChanged", taskWorker.id, "running", "paused"); + harness.agentStore.emit("agent:stateChanged", spawned.id, "running", "error"); + await flushMicrotasks(); + + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith(taskWorker.id); + expect(harness.agentStore.deleteAgent).toHaveBeenCalledWith(spawned.id); + }); + + it("ignores non-ephemeral agents, unchanged states, and externally pending deletes", async () => { + const durable = makeAgent("durable-paused", { metadata: {}, runtimeConfig: { enabled: true } }); + const pending = makeAgent("pending-paused", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + harness.agentStore.agents.set(durable.id, durable); + harness.agentStore.agents.set(pending.id, pending); + harness.externalPending.mockImplementation((agentId: string) => agentId === pending.id); + harness.manager.attachStateChangeListener(); + + harness.agentStore.emit("agent:stateChanged", durable.id, "active", "paused"); + harness.agentStore.emit("agent:stateChanged", pending.id, "running", "paused"); + harness.agentStore.emit("agent:stateChanged", pending.id, "paused", "paused"); + await flushMicrotasks(); + + expect(harness.agentStore.deleteAgent).not.toHaveBeenCalled(); + }); + + it("prevents duplicate listener deletes and detaches cleanly", async () => { + let resolveDelete: (() => void) | undefined; + const worker = makeAgent("worker-dup", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + harness.agentStore.agents.set(worker.id, worker); + harness.agentStore.deleteAgent.mockImplementationOnce(async (agentId: string) => { + await new Promise((resolve) => { resolveDelete = resolve; }); + harness.agentStore.agents.delete(agentId); + }); + const listener = harness.manager.attachStateChangeListener(); + expect(harness.manager.attachStateChangeListener()).toBe(listener); + + harness.agentStore.emit("agent:stateChanged", worker.id, "running", "paused"); + harness.agentStore.emit("agent:stateChanged", worker.id, "running", "error"); + await flushMicrotasks(); + expect(harness.agentStore.deleteAgent).toHaveBeenCalledTimes(1); + resolveDelete?.(); + await flushMicrotasks(); + + harness.manager.detachStateChangeListener(); + const afterDetach = makeAgent("worker-detached", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + harness.agentStore.agents.set(afterDetach.id, afterDetach); + harness.agentStore.emit("agent:stateChanged", afterDetach.id, "running", "paused"); + await flushMicrotasks(); + expect(harness.agentStore.deleteAgent).toHaveBeenCalledTimes(1); + }); + + it("suppresses benign halt-delete races but logs genuine failures", async () => { + const benign = makeAgent("worker-benign", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + const genuine = makeAgent("worker-genuine", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); + harness.agentStore.agents.set(benign.id, benign); + harness.agentStore.agents.set(genuine.id, genuine); + harness.manager.attachStateChangeListener(); + harness.agentStore.deleteAgent + .mockRejectedValueOnce(new Error(`Agent ${benign.id} not found`)) + .mockRejectedValueOnce(new Error("delete failed")); + + harness.agentStore.emit("agent:stateChanged", benign.id, "running", "paused"); + harness.agentStore.emit("agent:stateChanged", genuine.id, "running", "paused"); + await flushMicrotasks(); + + expect(harness.logger.warn).toHaveBeenCalledTimes(1); + expect(harness.logger.warn).toHaveBeenCalledWith(expect.stringContaining(`Failed to delete ephemeral agent ${genuine.id}`)); + }); + }); + + describe("startup reconciliation", () => { + it("returns zero for empty or all-durable agent lists", async () => { + expect(await harness.manager.reconcileOrphaned()).toBe(0); + + harness.agentStore.agents.set("durable", makeAgent("durable", { metadata: {}, taskId: "FN-1" })); + expect(await harness.manager.reconcileOrphaned()).toBe(0); + expect(harness.agentStore.deleteAgent).not.toHaveBeenCalled(); + }); + + it("keeps populated in-progress ephemeral workers and deletes stale task bindings", async () => { + const live = makeAgent("live-worker", { metadata: { agentKind: "task-worker" }, taskId: "FN-LIVE" }); + const done = makeAgent("done-worker", { metadata: { agentKind: "task-worker" }, taskId: "FN-DONE" }); + const todo = makeAgent("todo-worker", { metadata: { agentKind: "task-worker" }, taskId: "FN-TODO" }); + harness.agentStore.agents.set(live.id, live); + harness.agentStore.agents.set(done.id, done); + harness.agentStore.agents.set(todo.id, todo); + harness.taskStore.tasks.set("FN-LIVE", makeTask("FN-LIVE", { column: "in-progress" })); + harness.taskStore.tasks.set("FN-DONE", makeTask("FN-DONE", { column: "done" })); + harness.taskStore.tasks.set("FN-TODO", makeTask("FN-TODO", { column: "todo" })); + + expect(await harness.manager.reconcileOrphaned()).toBe(2); + + expect(harness.agentStore.agents.has(live.id)).toBe(true); + expect(harness.agentStore.agents.has(done.id)).toBe(false); + expect(harness.agentStore.agents.has(todo.id)).toBe(false); + }); + + it("deletes no-task, missing-task, paused, and error ephemerals", async () => { + const agents = [ + makeAgent("no-task", { metadata: { agentKind: "task-worker" } }), + makeAgent("missing-task", { metadata: { agentKind: "task-worker" }, taskId: "FN-MISSING" }), + makeAgent("paused", { state: "paused", metadata: { agentKind: "task-worker" }, taskId: "FN-LIVE" }), + makeAgent("error", { state: "error", metadata: { agentKind: "task-worker" }, taskId: "FN-LIVE" }), + ]; + for (const agent of agents) harness.agentStore.agents.set(agent.id, agent); + harness.taskStore.tasks.set("FN-LIVE", makeTask("FN-LIVE", { column: "in-progress" })); + + expect(await harness.manager.reconcileOrphaned()).toBe(4); + for (const agent of agents) expect(harness.agentStore.agents.has(agent.id)).toBe(false); + }); + + it("counts benign delete races and continues past genuine delete failures", async () => { + const benign = makeAgent("sweep-benign", { metadata: { agentKind: "task-worker" } }); + const failing = makeAgent("sweep-failing", { metadata: { agentKind: "task-worker" } }); + const next = makeAgent("sweep-next", { metadata: { agentKind: "task-worker" } }); + harness.agentStore.agents.set(benign.id, benign); + harness.agentStore.agents.set(failing.id, failing); + harness.agentStore.agents.set(next.id, next); + harness.agentStore.deleteAgent + .mockRejectedValueOnce(new Error(`Agent ${benign.id} not found`)) + .mockRejectedValueOnce(new Error("delete failed")) + .mockImplementationOnce(async (agentId: string) => { harness.agentStore.agents.delete(agentId); }); + + expect(await harness.manager.reconcileOrphaned()).toBe(2); + + expect(harness.logger.warn).toHaveBeenCalledWith(expect.stringContaining(`Startup sweep failed to delete ephemeral agent ${failing.id}`)); + expect(harness.agentStore.agents.has(next.id)).toBe(false); + expect(harness.logger.log).toHaveBeenCalledWith(expect.stringContaining("Startup ephemeral sweep cleaned 2 orphaned agent(s)")); + }); + }); + + it("reset clears owners and pending deletions", async () => { + await harness.manager.onTaskStart(makeTask("FN-RESET")); + expect(harness.manager.getOwner("FN-RESET")).toBeDefined(); + + harness.manager.reset(); + + expect(harness.manager.getOwner("FN-RESET")).toBeUndefined(); + }); +}); diff --git a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts index a293f22971..d57dccefb3 100644 --- a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts +++ b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts @@ -83,6 +83,153 @@ describe("HeartbeatTriggerScheduler", () => { }); }); + describe("agent lifecycle seam registration", () => { + type LifecycleStore = EventEmitter & Pick & { + agents: Map; + }; + + const baseAgent = (id: string, patch: Partial = {}): Agent => ({ + id, + name: patch.name ?? id, + role: patch.role ?? "executor", + state: patch.state ?? "active", + createdAt: patch.createdAt ?? "2026-01-01T00:00:00.000Z", + updatedAt: patch.updatedAt ?? "2026-01-01T00:00:00.000Z", + metadata: patch.metadata ?? {}, + runtimeConfig: patch.runtimeConfig, + lastHeartbeatAt: patch.lastHeartbeatAt, + taskId: patch.taskId, + }) as Agent; + + function createLifecycleStore(initialAgents: Agent[] = []): LifecycleStore { + const eventStore = Object.assign(new EventEmitter(), { + agents: new Map(initialAgents.map((agent) => [agent.id, agent])), + getAgent: vi.fn(async function(this: LifecycleStore, agentId: string) { + return this.agents.get(agentId) ?? null; + }), + getActiveHeartbeatRun: vi.fn().mockResolvedValue(null), + getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()), + listAgents: vi.fn(async function(this: LifecycleStore) { + return Array.from(this.agents.values()); + }), + getRecentRuns: vi.fn().mockResolvedValue([]), + updateAgent: vi.fn().mockImplementation(async function(this: LifecycleStore, agentId: string, patch: Partial) { + const before = this.agents.get(agentId) ?? baseAgent(agentId); + const after = { ...before, ...patch, runtimeConfig: patch.runtimeConfig ?? before.runtimeConfig } as Agent; + this.agents.set(agentId, after); + this.emit("agent:configRevision", agentId, { before, after }); + this.emit("agent:updated", after); + return after; + }), + }) as LifecycleStore; + return eventStore; + } + + beforeEach(() => { + vi.useFakeTimers(); + }); + + afterEach(() => { + scheduler?.stop(); + vi.useRealTimers(); + }); + + it("registers created heartbeat agents and excludes disabled or internal workers", async () => { + /* + FNXC:TestInfrastructure 2026-07-03-11:06: + Scheduler lifecycle behavior should be exercised at the EventEmitter seam instead of paying InProcessRuntime startup cost for created/default/explicit/disabled registration variants. + */ + const eventStore = createLifecycleStore(); + scheduler = new HeartbeatTriggerScheduler(eventStore as unknown as AgentStore, callback); + scheduler.start(); + + const defaultAgent = baseAgent("agent-default"); + const explicitAgent = baseAgent("agent-explicit", { runtimeConfig: { enabled: true, heartbeatIntervalMs: 15_000 } }); + const disabledAgent = baseAgent("agent-disabled", { runtimeConfig: { enabled: false } }); + const internalWorker = baseAgent("agent-worker", { metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: true, heartbeatIntervalMs: 1_000 } }); + for (const agent of [defaultAgent, explicitAgent, disabledAgent, internalWorker]) { + eventStore.agents.set(agent.id, agent); + eventStore.emit("agent:created", agent); + } + + expect(scheduler.getRegisteredAgents()).toContain(defaultAgent.id); + expect(scheduler.getRegisteredAgents()).toContain(explicitAgent.id); + expect(scheduler.getRegisteredAgents()).not.toContain(disabledAgent.id); + expect(scheduler.getRegisteredAgents()).not.toContain(internalWorker.id); + + await vi.advanceTimersByTimeAsync(14_999); + expect(callback).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(callback).toHaveBeenCalledWith(explicitAgent.id, "timer", expect.objectContaining({ intervalMs: 15_000 })); + }); + + it("keeps unrelated updates stable, re-arms interval changes, and clears paused timers", async () => { + const agent = baseAgent("agent-lifecycle", { runtimeConfig: { enabled: true, heartbeatIntervalMs: 1_000 } }); + const eventStore = createLifecycleStore([agent]); + scheduler = new HeartbeatTriggerScheduler(eventStore as unknown as AgentStore, callback); + scheduler.start(); + eventStore.emit("agent:created", agent); + + await vi.advanceTimersByTimeAsync(400); + const renamed = { ...agent, name: "renamed" } as Agent; + eventStore.agents.set(agent.id, renamed); + eventStore.emit("agent:updated", renamed); + await vi.advanceTimersByTimeAsync(599); + expect(callback).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(callback).toHaveBeenCalledTimes(1); + + callback.mockClear(); + await eventStore.updateAgent(agent.id, { runtimeConfig: { enabled: true, heartbeatIntervalMs: 2_000 } }); + await Promise.resolve(); + await vi.advanceTimersByTimeAsync(1_999); + expect(callback).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(callback).toHaveBeenCalledWith(agent.id, "timer", expect.objectContaining({ intervalMs: 2_000 })); + + callback.mockClear(); + const paused = { ...eventStore.agents.get(agent.id)!, state: "paused" as const }; + eventStore.agents.set(agent.id, paused); + eventStore.emit("agent:updated", paused); + expect(scheduler.getRegisteredAgents()).not.toContain(agent.id); + await vi.advanceTimersByTimeAsync(4_000); + expect(callback).not.toHaveBeenCalled(); + + const resumed = { ...paused, state: "active" as const }; + eventStore.agents.set(agent.id, resumed); + eventStore.emit("agent:updated", resumed); + expect(scheduler.getRegisteredAgents()).toContain(agent.id); + await vi.advanceTimersByTimeAsync(1_999); + expect(callback).not.toHaveBeenCalled(); + await vi.advanceTimersByTimeAsync(1); + expect(callback).toHaveBeenCalledTimes(1); + }); + + it("unregisters deleted agents and removes lifecycle listeners on stop", async () => { + const agent = baseAgent("agent-cleanup", { runtimeConfig: { enabled: true, heartbeatIntervalMs: 1_000 } }); + const eventStore = createLifecycleStore([agent]); + scheduler = new HeartbeatTriggerScheduler(eventStore as unknown as AgentStore, callback); + scheduler.start(); + eventStore.emit("agent:created", agent); + expect(scheduler.getRegisteredAgents()).toContain(agent.id); + + eventStore.emit("agent:deleted", agent.id); + expect(scheduler.getRegisteredAgents()).not.toContain(agent.id); + + eventStore.emit("agent:created", agent); + expect(scheduler.getRegisteredAgents()).toContain(agent.id); + scheduler.stop(); + expect(eventStore.listenerCount("agent:created")).toBe(0); + expect(eventStore.listenerCount("agent:updated")).toBe(0); + expect(eventStore.listenerCount("agent:configRevision")).toBe(0); + expect(eventStore.listenerCount("agent:deleted")).toBe(0); + + const afterStop = baseAgent("agent-after-stop", { runtimeConfig: { enabled: true, heartbeatIntervalMs: 1_000 } }); + eventStore.emit("agent:created", afterStop); + expect(scheduler.getRegisteredAgents()).not.toContain(afterStop.id); + }); + }); + describe("scheduler timer audit", () => { it("re-arms a tickable durable agent when timer entry is missing and no lifecycle event fires", async () => { vi.useFakeTimers(); diff --git a/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts b/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts index 13b8108f93..1fdeb1edd5 100644 --- a/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts +++ b/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts @@ -833,27 +833,6 @@ describe("InProcessRuntime", () => { expect(registeredAgents).toContain(createdAgent.id); }); - it("does not register paused agents on startup", async () => { - await runtime.start(); - - const store = getAgentStore(runtime); - const pausedAgent = await store.createAgent({ - name: "Paused Agent", - role: "executor", - runtimeConfig: { heartbeatIntervalMs: 30000, enabled: true }, - }); - await store.updateAgentState(pausedAgent.id, "active"); - await store.updateAgentState(pausedAgent.id, "paused"); - - await runtime.stop(); - runtime = new InProcessRuntime(buildTestConfig(testDir), mockCentralCore); - await runtime.start(); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler).toBeDefined(); - expect(scheduler!.getRegisteredAgents()).not.toContain(pausedAgent.id); - }); - it("routes assignment triggers through executeHeartbeat", async () => { await runtime.start(); @@ -888,370 +867,34 @@ describe("InProcessRuntime", () => { ); }, 30000); - it("reuses assigned durable agent as execution owner without creating a task-worker", async () => { - await runtime.start(); - - const store = getAgentStore(runtime); - const durable = await store.createAgent({ name: "Durable Exec", role: "executor" }); - const createAgentSpy = vi.spyOn(store, "createAgent"); - const assignTaskSpy = vi.spyOn(store, "assignTask"); - const syncLinkSpy = vi.spyOn(store, "syncExecutionTaskLink"); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - executorOptions.onStart?.({ id: "FN-1661", assignedAgentId: durable.id } as Task, join(testDir, "worktree-FN-1661")); - - await vi.waitFor(async () => { - const updated = await store.getAgent(durable.id); - expect(updated?.taskId).toBe("FN-1661"); - expect(updated?.state).toBe("running"); - }); - - expect(syncLinkSpy).toHaveBeenCalledWith(durable.id, "FN-1661"); - expect(assignTaskSpy).not.toHaveBeenCalledWith(durable.id, "FN-1661"); - expect(createAgentSpy).not.toHaveBeenCalledWith(expect.objectContaining({ name: "executor-FN-1661" })); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((agent: Agent) => agent.name === "executor-FN-1661")).toBe(false); - }, 30000); - - it("falls back to runtime task-worker agents for unassigned tasks", async () => { - await runtime.start(); - - const store = getAgentStore(runtime); - - const assignTaskSpy = vi.spyOn(store, "assignTask"); - const updateStateSpy = vi.spyOn(store, "updateAgentState"); - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - expect(executorOptions.onStart).toBeTypeOf("function"); - - executorOptions.onStart?.({ id: "FN-1661" } as Task, join(testDir, "worktree-FN-1661")); - - await vi.waitFor(async () => { - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents).toHaveLength(1); - expect(agents[0]).toMatchObject({ - name: "executor-FN-1661", - role: "executor", - state: "running", - taskId: "FN-1661", - metadata: { - agentKind: "task-worker", - taskWorker: true, - managedBy: "task-executor", - }, - runtimeConfig: { - enabled: false, - }, - }); - }); - - expect(assignTaskSpy).toHaveBeenCalledWith(expect.any(String), "FN-1661"); - expect(updateStateSpy).toHaveBeenNthCalledWith(1, expect.any(String), "active"); - expect(updateStateSpy).toHaveBeenNthCalledWith(2, expect.any(String), "running"); - expect(assignTaskSpy.mock.invocationCallOrder[0]).toBeLessThan(updateStateSpy.mock.invocationCallOrder[0]); - }, 30000); - - it("does not spawn runtime task-worker agents when ephemeral agents are disabled", async () => { - mockTaskStoreSettings.ephemeralAgentsEnabled = false; + it("wires executor ownership callbacks through the worker manager", async () => { await runtime.start(); const store = getAgentStore(runtime); const createAgentSpy = vi.spyOn(store, "createAgent"); - const warnSpy = vi.spyOn(runtimeLog, "warn"); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - executorOptions.onStart?.({ id: "FN-1663" } as Task, join(testDir, "worktree-FN-1663")); - await flushRuntimeCallbackMicrotasks(); - - expect(warnSpy).toHaveBeenCalledWith( - expect.stringContaining("Task FN-1663 has no permanent agent assignment; ephemeralAgentsEnabled=false"), - ); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((agent: Agent) => agent.name === "executor-FN-1663")).toBe(false); - expect(createAgentSpy).not.toHaveBeenCalledWith(expect.objectContaining({ name: "executor-FN-1663" })); - }, 30000); - - it("falls back to runtime task-worker when assignedAgentId points to ephemeral agent", async () => { - await runtime.start(); - - const store = getAgentStore(runtime); - const ephemeral = await store.createAgent({ - name: "Spawned Child", - role: "executor", - metadata: { agentKind: "task-worker", managedBy: "task-executor" }, - runtimeConfig: { enabled: false }, - }); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - executorOptions.onStart?.({ id: "FN-1662", assignedAgentId: ephemeral.id } as Task, join(testDir, "worktree-FN-1662")); - await flushRuntimeCallbackMicrotasks(); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((agent: Agent) => agent.name === "executor-FN-1662")).toBe(true); - }, 30000); - - it("does not wake executeHeartbeat for runtime ownership sync of durable assigned agents", async () => { - /* - FNXC:TestInfrastructure 2026-06-26-21:52: - This negative-assertion test must verify executeHeartbeat is NOT woken by runtime ownership sync. - Previously it paid a real `await new Promise(r => setTimeout(r, 25))` wall-clock sleep to let any - erroneously-scheduled executeHeartbeat fire before asserting it did not — pure dead time on every run. - Per FN-5048 (prefer fake timers over real polling/time waits) we run under fake timers and advance the - window deterministically with advanceTimersByTimeAsync. The inflated 30000ms per-test timeout is removed - now that no real wait remains, and callback microtasks are flushed directly instead of poll-looped. - */ - vi.useFakeTimers(); - try { - await runtime.start(); - - const monitor = runtime.getHeartbeatMonitor(); - expect(monitor).toBeDefined(); - const heartbeatMonitor = monitor!; - const executeResult = { id: "run-task-worker" } as Awaited>; - const executeSpy = vi - .spyOn(heartbeatMonitor, "executeHeartbeat") - .mockResolvedValue(executeResult); - - const store = getAgentStore(runtime); - const durable = await store.createAgent({ name: "Owned Exec", role: "executor" }); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - executorOptions.onStart?.({ id: "FN-2001", assignedAgentId: durable.id } as Task, join(testDir, "worktree-FN-2001")); - await flushRuntimeCallbackMicrotasks(); - - const updated = await store.getAgent(durable.id); - expect(updated?.taskId).toBe("FN-2001"); - - // Drive the negative-assertion window deterministically instead of sleeping 25ms of real time. - await vi.advanceTimersByTimeAsync(25); - expect(executeSpy).not.toHaveBeenCalled(); - } finally { - vi.useRealTimers(); - } - }); - - it("cleans up durable execution owner on completion without deleting agent", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const durable = await store.createAgent({ name: "Durable Cleanup", role: "executor" }); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onComplete?: (task: Task) => void; - }; - - executorOptions.onStart?.({ id: "FN-DURABLE-1", assignedAgentId: durable.id } as Task, join(testDir, "worktree-FN-DURABLE-1")); - await flushRuntimeCallbackMicrotasks(); - expect((await store.getAgent(durable.id))?.taskId).toBe("FN-DURABLE-1"); - - executorOptions.onComplete?.({ id: "FN-DURABLE-1" } as Task); - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - const updated = await store.getAgent(durable.id); - expect(updated?.state).toBe("active"); - expect(updated?.taskId).toBeUndefined(); - expect(deleteAgentSpy).not.toHaveBeenCalledWith(durable.id); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("auto-deletes task-worker agent on task completion immediately", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onComplete?: (task: Task) => void; - }; - expect(executorOptions.onComplete).toBeTypeOf("function"); - - // Create a task-worker agent first via onStart - executorOptions.onStart?.({ id: "FN-AUTO1" } as Task, join(testDir, "worktree-FN-AUTO1")); - await flushRuntimeCallbackMicrotasks(); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((a: Agent) => a.name === "executor-FN-AUTO1")).toBe(true); - - // Clear previous calls and trigger onComplete - deleteAgentSpy.mockClear(); - executorOptions.onComplete?.({ id: "FN-AUTO1" } as Task); - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("cleans up durable execution owner on error without deleting agent", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const durable = await store.createAgent({ name: "Durable Error", role: "executor" }); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onError?: (task: Task, error: Error) => void; - }; - - executorOptions.onStart?.({ id: "FN-DURABLE-2", assignedAgentId: durable.id } as Task, join(testDir, "worktree-FN-DURABLE-2")); - await flushRuntimeCallbackMicrotasks(); - expect((await store.getAgent(durable.id))?.taskId).toBe("FN-DURABLE-2"); - - executorOptions.onError?.({ id: "FN-DURABLE-2" } as Task, new Error("boom")); - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - const updated = await store.getAgent(durable.id); - expect(updated?.state).toBe("active"); - expect(updated?.taskId).toBeUndefined(); - expect(deleteAgentSpy).not.toHaveBeenCalledWith(durable.id); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("auto-deletes task-worker agent on task error immediately", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onError?: (task: Task, error: Error) => void; - }; - expect(executorOptions.onError).toBeTypeOf("function"); - - // Create a task-worker agent first via onStart - const onStartOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - }; - onStartOptions.onStart?.({ id: "FN-AUTO2" } as Task, join(testDir, "worktree-FN-AUTO2")); - await flushRuntimeCallbackMicrotasks(); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((a: Agent) => a.name === "executor-FN-AUTO2")).toBe(true); - - // Clear previous calls and trigger onError - deleteAgentSpy.mockClear(); - executorOptions.onError?.({ id: "FN-AUTO2" } as Task, new Error("Task failed")); - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }, 30000); - }); - - describe("agent cleanup failure diagnostics", () => { - it("logs warning when agent state update fails on task completion", async () => { - const warnSpy = vi.spyOn(runtimeLog, "warn"); - await runtime.start(); - - const store = getAgentStore(runtime); - const updateStateSpy = vi.spyOn(store, "updateAgentState").mockImplementation(async (_agentId, state) => { - if (state === "active") { - throw new Error("state update failed"); - } - return {} as Agent; - }); - + const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { onStart?: (task: Task, worktreePath: string) => void; onComplete?: (task: Task) => void; }; - executorOptions.onStart?.({ id: "FN-DIAG-1" } as Task, join(testDir, "worktree-FN-DIAG-1")); + expect(executorOptions.onStart).toBeTypeOf("function"); + expect(executorOptions.onComplete).toBeTypeOf("function"); + + executorOptions.onStart?.({ id: "FN-WIRING" } as Task, join(testDir, "worktree-FN-WIRING")); await flushRuntimeCallbackMicrotasks(); + expect(createAgentSpy).toHaveBeenCalledWith(expect.objectContaining({ name: "executor-FN-WIRING" })); - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((a: Agent) => a.name === "executor-FN-DIAG-1")).toBe(true); + const worker = (await store.listAgents({ includeEphemeral: true })) + .find((agent: Agent) => agent.name === "executor-FN-WIRING"); + expect(worker).toBeDefined(); - updateStateSpy.mockClear(); - executorOptions.onComplete?.({ id: "FN-DIAG-1" } as Task); - await flushRuntimeCallbackMicrotasks(); - - expect(warnSpy).toHaveBeenCalledWith( - expect.stringContaining("Failed to update agent"), - ); - expect(warnSpy).toHaveBeenCalledWith( - expect.stringContaining("active (completion)"), - ); - - warnSpy.mockRestore(); - }, 30000); - - it("logs warning when agent deletion fails after task error", async () => { - vi.useFakeTimers(); - const warnSpy = vi.spyOn(runtimeLog, "warn"); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockRejectedValue(new Error("delete failed")); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onError?: (task: Task, error: Error) => void; - }; - executorOptions.onStart?.({ id: "FN-DIAG-2" } as Task, join(testDir, "worktree-FN-DIAG-2")); - await flushRuntimeCallbackMicrotasks(); - - const agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((a: Agent) => a.name === "executor-FN-DIAG-2")).toBe(true); - - deleteAgentSpy.mockClear(); - executorOptions.onError?.({ id: "FN-DIAG-2" } as Task, new Error("Task failed")); - - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(warnSpy).toHaveBeenCalledWith( - expect.stringContaining("Failed to delete agent"), - ); - expect(warnSpy).toHaveBeenCalledWith( - expect.stringContaining("after error"), - ); - } finally { - warnSpy.mockRestore(); - vi.useRealTimers(); - } + executorOptions.onComplete?.({ id: "FN-WIRING" } as Task); + await flushRuntimeCallbackMicrotasks(40); + expect(deleteAgentSpy).toHaveBeenCalledWith(worker!.id); }, 30000); }); + describe("configuration", () => { it("should store projectId in config", () => { // Access via the constructor params - runtime is created with testDir @@ -1290,654 +933,5 @@ describe("InProcessRuntime", () => { }); }); - describe("dynamic agent registration with HeartbeatTriggerScheduler", () => { - beforeEach(async () => { - vi.useFakeTimers(); - await runtime.start(); - }); - afterEach(async () => { - await runtime.stop(); - vi.useRealTimers(); - }); - - it("registers heartbeat-managed created agents and skips disabled agents", async () => { - /* - FNXC:TestInfrastructure 2026-07-03-10:45: - Dynamic registration variants share one started scheduler; assert default, explicit, and disabled configs together to preserve lifecycle coverage while reducing repeated runtime setup. - */ - const store = getAgentStore(runtime); - const defaultAgent = await store.createAgent({ - name: "test-agent-dynamic", - role: "executor", - }); - const enabledAgent = await store.createAgent({ - name: "test-agent-default-interval", - role: "executor", - runtimeConfig: { enabled: true }, - }); - const explicitAgent = await store.createAgent({ - name: "test-agent-explicit", - role: "executor", - runtimeConfig: { - heartbeatIntervalMs: 15000, - enabled: true, - }, - }); - const disabledAgent = await store.createAgent({ - name: "test-agent-disabled", - role: "executor", - runtimeConfig: { - enabled: false, - }, - }); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler).toBeDefined(); - expect(scheduler!.getRegisteredAgents()).toContain(defaultAgent.id); - expect(scheduler!.getRegisteredAgents()).toContain(enabledAgent.id); - expect(scheduler!.getRegisteredAgents()).toContain(explicitAgent.id); - expect(scheduler!.getRegisteredAgents()).not.toContain(disabledAgent.id); - }); - - it("does not reset an armed timer on unrelated agent updates", async () => { - const store = getAgentStore(runtime); - const monitor = runtime.getHeartbeatMonitor(); - expect(monitor).toBeDefined(); - - const executeHeartbeatSpy = vi - .spyOn(monitor!, "executeHeartbeat") - .mockResolvedValue({ id: "run-update-timer-stability" } as any); - - const agent = await store.createAgent({ - name: "test-agent-update", - role: "executor", - runtimeConfig: { - enabled: true, - heartbeatIntervalMs: 1000, - }, - }); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - - await vi.advanceTimersByTimeAsync(400); - await store.updateAgent(agent.id, { - name: "test-agent-update-renamed", - }); - - await vi.advanceTimersByTimeAsync(599); - expect(executeHeartbeatSpy).not.toHaveBeenCalled(); - - await vi.advanceTimersByTimeAsync(1); - await flushRuntimeCallbackMicrotasks(); - expect(executeHeartbeatSpy).toHaveBeenCalledTimes(1); - - expect(executeHeartbeatSpy).toHaveBeenCalledWith( - expect.objectContaining({ - agentId: agent.id, - source: "timer", - }), - ); - }); - - it("reconciles a missing timer for a tickable durable agent without state changes", async () => { - const store = getAgentStore(runtime); - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler).toBeDefined(); - - const agent = await store.createAgent({ - name: "audit-rearm-agent", - role: "executor", - runtimeConfig: { - enabled: true, - heartbeatIntervalMs: 1_000, - }, - }); - - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - scheduler!.unregisterAgent(agent.id); - expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id); - - await vi.advanceTimersByTimeAsync(60_000); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - }); - - it("unregisters an agent when enabled is set to false in update", async () => { - // Create a new agent with heartbeat enabled - const store = getAgentStore(runtime); - const agent = await store.createAgent({ - name: "test-agent-toggle", - role: "executor", - runtimeConfig: { - enabled: true, - }, - }); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - - // Update the agent to disable heartbeat - await store.updateAgent(agent.id, { - runtimeConfig: { - enabled: false, - }, - }); - - // Verify the agent was unregistered - expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id); - }); - - it("re-arms the timer when heartbeat interval changes", async () => { - const store = getAgentStore(runtime); - const monitor = runtime.getHeartbeatMonitor(); - expect(monitor).toBeDefined(); - - const executeHeartbeatSpy = vi - .spyOn(monitor!, "executeHeartbeat") - .mockResolvedValue({ id: "run-interval-change" } as any); - - const agent = await store.createAgent({ - name: "interval-change-agent", - role: "executor", - runtimeConfig: { - enabled: true, - heartbeatIntervalMs: 1000, - }, - }); - - await vi.advanceTimersByTimeAsync(400); - await store.updateAgent(agent.id, { - runtimeConfig: { - enabled: true, - heartbeatIntervalMs: 2000, - }, - }); - - await vi.advanceTimersByTimeAsync(1599); - expect(executeHeartbeatSpy).not.toHaveBeenCalled(); - - await vi.advanceTimersByTimeAsync(401); - await flushRuntimeCallbackMicrotasks(); - expect(executeHeartbeatSpy).toHaveBeenCalledTimes(1); - }); - - it("clears timers on pause and re-arms from resume without stale pre-pause firing", async () => { - const store = getAgentStore(runtime); - const monitor = runtime.getHeartbeatMonitor(); - expect(monitor).toBeDefined(); - - const executeHeartbeatSpy = vi - .spyOn(monitor!, "executeHeartbeat") - .mockResolvedValue({ id: "run-resume-test" } as any); - - const agent = await store.createAgent({ - name: "resume-timer-agent", - role: "executor", - runtimeConfig: { - enabled: true, - heartbeatIntervalMs: 1000, - }, - }); - await store.updateAgentState(agent.id, "active"); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler).toBeDefined(); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - - await vi.advanceTimersByTimeAsync(400); - expect(executeHeartbeatSpy).not.toHaveBeenCalled(); - - await store.updateAgentState(agent.id, "paused"); - expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id); - - // Advance beyond the original tick window; stale pre-pause timer must not fire. - await vi.advanceTimersByTimeAsync(800); - expect(executeHeartbeatSpy).not.toHaveBeenCalled(); - - await store.updateAgentState(agent.id, "active"); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - - // Resume should start a fresh interval from now, not from pre-pause start. - await vi.advanceTimersByTimeAsync(900); - expect(executeHeartbeatSpy).not.toHaveBeenCalled(); - - await vi.advanceTimersByTimeAsync(100); - await flushRuntimeCallbackMicrotasks(); - expect(executeHeartbeatSpy).toHaveBeenCalledTimes(1); - - expect(executeHeartbeatSpy).toHaveBeenCalledWith( - expect.objectContaining({ - agentId: agent.id, - source: "timer", - }), - ); - }); - - it("removes event listeners when runtime is stopped", async () => { - // Create a new agent before stopping - const store = getAgentStore(runtime); - const agent = await store.createAgent({ - name: "test-agent-cleanup", - role: "executor", - }); - - const scheduler = runtime.getTriggerScheduler(); - expect(scheduler!.getRegisteredAgents()).toContain(agent.id); - - // Stop the runtime - await runtime.stop(); - - // The agent should still be registered (unregister is internal to scheduler) - // But the listeners should be removed - verify by checking they don't fire - // Create another agent - it won't be registered since runtime is stopped - const agent2 = await store.createAgent({ - name: "test-agent-after-stop", - role: "executor", - }); - - // Since runtime is stopped, trigger scheduler is stopped - // The agent won't be in registered list - expect(scheduler!.getRegisteredAgents()).not.toContain(agent2.id); - }); - }); - - describe("ephemeral paused-state cleanup", () => { - it("auto-deletes ephemeral agent when it transitions to paused via agent:stateChanged", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - // Create an ephemeral task-worker agent - const agent = await store.createAgent({ - name: "executor-FN-TERM-1", - role: "executor", - metadata: { - agentKind: "task-worker", - taskWorker: true, - managedBy: "task-executor", - }, - runtimeConfig: { enabled: false }, - }); - - // Verify agent exists - let agents = await store.listAgents({ includeEphemeral: true }); - expect(agents.some((a: Agent) => a.id === agent.id)).toBe(true); - - // Emit agent:stateChanged event to trigger cleanup - store.emit("agent:stateChanged", agent.id, "running", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id); - - // Note: We verified deleteAgent was called, which is the key behavior. - // The actual removal from listAgents depends on the real AgentStore implementation. - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("does not auto-delete non-ephemeral agent when it transitions to paused", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - // Create a non-ephemeral user-managed agent - const agent = await store.createAgent({ - name: "user-managed-agent", - role: "executor", - // No ephemeral metadata - runtimeConfig: { enabled: true }, - }); - - // Verify agent exists - let agents = await store.listAgents(); - expect(agents.some((a: Agent) => a.id === agent.id)).toBe(true); - - // Emit agent:stateChanged event to trigger cleanup - store.emit("agent:stateChanged", agent.id, "active", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - - await vi.advanceTimersByTimeAsync(0); - - // deleteAgent should NOT have been called for non-ephemeral agent - expect(deleteAgentSpy).not.toHaveBeenCalled(); - - // Agent should still exist - agents = await store.listAgents(); - expect(agents.some((a: Agent) => a.id === agent.id)).toBe(true); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("does not schedule duplicate deletion when termination event fires multiple times", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - // Create an ephemeral task-worker agent - const agent = await store.createAgent({ - name: "executor-FN-DUP-1", - role: "executor", - metadata: { - agentKind: "task-worker", - taskWorker: true, - managedBy: "task-executor", - }, - runtimeConfig: { enabled: false }, - }); - - // Emit termination event multiple times - store.emit("agent:stateChanged", agent.id, "running", "paused"); - store.emit("agent:stateChanged", agent.id, "paused", "paused"); // Already halted - - // Wait for async handlers - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("does not warn when cleanup delete fails only because agent is already gone", async () => { - vi.useFakeTimers(); - const warnSpy = vi.spyOn(runtimeLog, "warn"); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - - // Create an ephemeral agent - const agent = await store.createAgent({ - name: "executor-FN-BENIGN-1", - role: "executor", - metadata: { - agentKind: "task-worker", - }, - runtimeConfig: { enabled: false }, - }); - - const deleteAgentSpy = vi - .spyOn(store, "deleteAgent") - .mockRejectedValueOnce(new Error(`Agent ${agent.id} not found`)); - - // Emit termination event - store.emit("agent:stateChanged", agent.id, "running", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - - // Cleanup should still be attempted - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id); - - // Benign not-found races should not produce warning-level noise - const emittedCleanupWarning = warnSpy.mock.calls.some(([msg]) => - typeof msg === "string" && msg.includes("Failed to delete ephemeral agent"), - ); - expect(emittedCleanupWarning).toBe(false); - } finally { - warnSpy.mockRestore(); - vi.useRealTimers(); - } - }, 30000); - - it("warns on genuine cleanup failure but does not throw", async () => { - vi.useFakeTimers(); - const warnSpy = vi.spyOn(runtimeLog, "warn"); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockRejectedValue(new Error("delete failed")); - - // Create an ephemeral agent - const agent = await store.createAgent({ - name: "executor-FN-WARN-1", - role: "executor", - metadata: { - agentKind: "task-worker", - }, - runtimeConfig: { enabled: false }, - }); - - // Emit termination event - store.emit("agent:stateChanged", agent.id, "running", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - - await vi.advanceTimersByTimeAsync(0); - - // Should have attempted deletion - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id); - - // Genuine failures still log warning-level context - const cleanupWarnings = warnSpy.mock.calls.filter(([msg]) => - typeof msg === "string" && msg.includes("Failed to delete ephemeral agent"), - ); - expect(cleanupWarnings).toHaveLength(1); - expect(cleanupWarnings[0]?.[0]).toContain(agent.id); - expect(cleanupWarnings[0]?.[0]).toContain("delete failed"); - } finally { - warnSpy.mockRestore(); - vi.useRealTimers(); - } - }, 30000); - - it("handles runtime stop racing with in-flight cleanup", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - // Create an ephemeral agent - const agent = await store.createAgent({ - name: "executor-FN-STOP-1", - role: "executor", - metadata: { - taskWorker: true, - }, - runtimeConfig: { enabled: false }, - }); - - // Emit termination event - store.emit("agent:stateChanged", agent.id, "running", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - - await runtime.stop(); - - expect(deleteAgentSpy.mock.calls.length).toBeLessThanOrEqual(1); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("handles spawned ephemeral agents (type=spawned) correctly", async () => { - vi.useFakeTimers(); - - try { - await runtime.start(); - - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - // Create a spawned child agent (type=spawned is ephemeral) - const agent = await store.createAgent({ - name: "child-agent-001", - role: "executor", - metadata: { - type: "spawned", - parentTaskId: "FN-PARENT", - }, - runtimeConfig: { enabled: false }, - }); - - // Emit termination event - store.emit("agent:stateChanged", agent.id, "running", "paused"); - - // Wait for async handler - await vi.advanceTimersByTimeAsync(0); - - await flushRuntimeCallbackMicrotasks(); - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - - // deleteAgent should have been called for spawned ephemeral agent - expect(deleteAgentSpy).toHaveBeenCalledWith(agent.id); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("does not double-delete when onComplete already scheduled cleanup", async () => { - vi.useFakeTimers(); - try { - await runtime.start(); - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onComplete?: (task: Task) => void; - }; - executorOptions.onStart?.({ id: "FN-DUP-COMPLETE" } as Task, join(testDir, "worktree-FN-DUP-COMPLETE")); - - await flushRuntimeCallbackMicrotasks(); - const worker = (await store.listAgents({ includeEphemeral: true })) - .find((a: Agent) => a.name === "executor-FN-DUP-COMPLETE"); - expect(worker).toBeDefined(); - - executorOptions.onComplete?.({ id: "FN-DUP-COMPLETE" } as Task); - store.emit("agent:stateChanged", worker!.id, "running", "paused"); - await vi.advanceTimersByTimeAsync(0); - await flushRuntimeCallbackMicrotasks(); - - expect(deleteAgentSpy).toHaveBeenCalledTimes(1); - } finally { - vi.useRealTimers(); - } - }, 30000); - - it("handles onComplete cleanup racing with runtime stop", async () => { - vi.useFakeTimers(); - try { - await runtime.start(); - const store = getAgentStore(runtime); - const deleteAgentSpy = vi.spyOn(store, "deleteAgent").mockResolvedValue(undefined); - - const executorOptions = mockExecutorCtor.mock.calls.at(-1)?.[0] as { - onStart?: (task: Task, worktreePath: string) => void; - onComplete?: (task: Task) => void; - }; - executorOptions.onStart?.({ id: "FN-STOP-COMPLETE" } as Task, join(testDir, "worktree-FN-STOP-COMPLETE")); - executorOptions.onComplete?.({ id: "FN-STOP-COMPLETE" } as Task); - - await runtime.stop(); - expect(deleteAgentSpy.mock.calls.length).toBeLessThanOrEqual(1); - } finally { - vi.useRealTimers(); - } - }, 30000); - }); - - describe("startup ephemeral sweep", () => { - it("cleans paused ephemeral agents on startup", async () => { - const { AgentStore } = await import("@fusion/core"); - const preStore = new AgentStore({ rootDir: join(testDir, ".fusion") }); - await preStore.init(); - const orphan = await preStore.createAgent({ - name: "orphan-paused", - role: "executor", - metadata: { agentKind: "task-worker" }, - runtimeConfig: { enabled: false }, - }); - await preStore.updateAgentState(orphan.id, "active"); - await preStore.updateAgentState(orphan.id, "paused"); - - await runtime.start(); - const store = getAgentStore(runtime); - expect(await store.getAgent(orphan.id)).toBeNull(); - }, 30000); - - it("cleans ephemeral agents assigned to non-in-progress tasks", async () => { - mockTaskStoreGetTask.mockResolvedValue({ id: "FN-DONE", column: "done" }); - const { AgentStore } = await import("@fusion/core"); - const preStore = new AgentStore({ rootDir: join(testDir, ".fusion") }); - await preStore.init(); - const orphan = await preStore.createAgent({ - name: "orphan-stale-task", - role: "executor", - metadata: { agentKind: "task-worker" }, - runtimeConfig: { enabled: false }, - }); - await preStore.assignTask(orphan.id, "FN-DONE"); - - await runtime.start(); - const store = getAgentStore(runtime); - expect(await store.getAgent(orphan.id)).toBeNull(); - }, 30000); - - it("continues startup sweep when one delete fails", async () => { - const warnSpy = vi.spyOn(runtimeLog, "warn"); - const { AgentStore } = await import("@fusion/core"); - const originalDeleteAgent = AgentStore.prototype.deleteAgent; - const deleteProtoSpy = vi - .spyOn(AgentStore.prototype, "deleteAgent") - .mockRejectedValueOnce(new Error("delete failed")) - .mockImplementation(async function(this: AgentStore, agentId: string) { - return originalDeleteAgent.call(this, agentId); - }); - - try { - const preStore = new AgentStore({ rootDir: join(testDir, ".fusion") }); - await preStore.init(); - const a1 = await preStore.createAgent({ name: "orphan-a1", role: "executor", metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); - const a2 = await preStore.createAgent({ name: "orphan-a2", role: "executor", metadata: { agentKind: "task-worker" }, runtimeConfig: { enabled: false } }); - await preStore.updateAgentState(a1.id, "active"); - await preStore.updateAgentState(a1.id, "paused"); - await preStore.updateAgentState(a2.id, "active"); - await preStore.updateAgentState(a2.id, "paused"); - - await runtime.start(); - const store = getAgentStore(runtime); - expect(runtime.getStatus()).toBe("active"); - const remaining = await store.listAgents({ includeEphemeral: true }); - expect(remaining.filter((a: Agent) => a.id === a1.id || a.id === a2.id)).toHaveLength(1); - expect(warnSpy).toHaveBeenCalled(); - } finally { - warnSpy.mockRestore(); - deleteProtoSpy.mockRestore(); - } - }, 30000); - }); });