From ec4b247a86dc3c8fc06425625270cc3f02fa6026 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Fri, 12 Jun 2026 03:00:45 -0700 Subject: [PATCH] FN-6284: refire deferred agent assignments Ensure assigned agents resume task work even when their heartbeat loop is idle. - add deferred assignment refiring when an agent is assigned work without an active heartbeat run - track in-process runtime activity so scheduled work is not double-started - cover idle, active, and concurrent heartbeat assignment paths with scheduler and runtime tests - document the deferred assignment wake-up behavior and add a changeset Files changed: .changeset/fn-6284-deferred-assignment-refire.md | 5 + docs/agents.md | 2 + docs/architecture.md | 1 + .../src/__tests__/heartbeat-scheduler.test.ts | 174 ++++++++++++++++++++- packages/engine/src/agent-heartbeat.ts | 109 ++++++++++++- .../runtimes/__tests__/in-process-runtime.test.ts | 24 +++ packages/engine/src/runtimes/in-process-runtime.ts | 3 + 7 files changed, 311 insertions(+), 7 deletions(-) Fusion-Task-Id: FN-6284 Fusion-Task-Lineage: cf8a9abe-163e-45c6-961a-095891e45ea5 --- .../fn-6284-deferred-assignment-refire.md | 5 + docs/agents.md | 2 + docs/architecture.md | 1 + .../src/__tests__/heartbeat-scheduler.test.ts | 176 +++++++++++++++++- packages/engine/src/agent-heartbeat.ts | 109 ++++++++++- .../__tests__/in-process-runtime.test.ts | 24 +++ .../engine/src/runtimes/in-process-runtime.ts | 3 + 7 files changed, 312 insertions(+), 8 deletions(-) create mode 100644 .changeset/fn-6284-deferred-assignment-refire.md diff --git a/.changeset/fn-6284-deferred-assignment-refire.md b/.changeset/fn-6284-deferred-assignment-refire.md new file mode 100644 index 0000000000..5dd5ffda5b --- /dev/null +++ b/.changeset/fn-6284-deferred-assignment-refire.md @@ -0,0 +1,5 @@ +--- +"@runfusion/fusion": patch +--- + +Re-fire durable-agent assignment wakes that were skipped because the agent was mid-heartbeat, so newly assigned tasks are worked when the active run completes instead of waiting for the next timer tick. diff --git a/docs/agents.md b/docs/agents.md index 40f133dbe0..16612f9d31 100644 --- a/docs/agents.md +++ b/docs/agents.md @@ -527,6 +527,8 @@ The `runtimeConfig` field on agents supports the following options: | `modelId` | `string` | — | AI model ID override for heartbeat session | | `budgetConfig` | `AgentBudgetConfig` | — | Token budget governance settings | +Assignment-triggered heartbeats are completion-resilient: if an `agent:assigned` wake is skipped only because the durable agent already has an active heartbeat run, Fusion records the latest assigned task as a pending assignment and re-fires that assignment wake once the active run completes. This prevents assigned work from being stranded by long heartbeat intervals or `skipHeartbeatWhenIdle`; disabled agents (`enabled === false`) and budget-exhausted agents still do not defer assignment wakes. + Heartbeat values are validated and minimum-clamped to 5 minutes (300,000 ms). Project setting `heartbeatMultiplier` (default `1`) scales resolved heartbeat timing globally: both the heartbeat interval (`pollIntervalMs`) and unresponsive timeout base (`heartbeatTimeoutMs`) are multiplied. Per-agent `heartbeatIntervalMs`/`heartbeatTimeoutMs` remain base values before multiplier scaling. This setting is configured from the **Agents** screen's **Controls** popup under "Heartbeat Speed". diff --git a/docs/architecture.md b/docs/architecture.md index 71ac7cf571..357e47027a 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -1305,6 +1305,7 @@ Limits are controlled by project settings (`maxSpawnedAgentsPerParent`, `maxSpaw - timer - task assignment - on-demand runs +- Assignment triggers skipped because a heartbeat run is already active are deferred and re-fired from `HeartbeatMonitor.onRunCompleted`, preserving the existing completion recovery path while avoiding timer-dependent stalls. ### Custom instructions `packages/engine/src/agent-instructions.ts` resolves per-agent instruction text/path with path-traversal and extension validation. diff --git a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts index f9857710c8..a293f22971 100644 --- a/packages/engine/src/__tests__/heartbeat-scheduler.test.ts +++ b/packages/engine/src/__tests__/heartbeat-scheduler.test.ts @@ -1206,6 +1206,7 @@ describe("HeartbeatTriggerScheduler", () => { describe("assignment watching", () => { let eventStore: EventEmitter & { + getAgent: ReturnType; getActiveHeartbeatRun: ReturnType; getBudgetStatus: ReturnType; getRecentRuns: ReturnType; @@ -1215,6 +1216,7 @@ describe("HeartbeatTriggerScheduler", () => { vi.useRealTimers(); // Ensure real timers for these tests eventStore = Object.assign(new EventEmitter(), { + getAgent: vi.fn().mockResolvedValue({ id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} }), getActiveHeartbeatRun: vi.fn().mockResolvedValue(null), getBudgetStatus: vi.fn().mockRejectedValue(new Error("budget status unavailable")), getRecentRuns: vi.fn().mockResolvedValue([]), @@ -1276,16 +1278,176 @@ describe("HeartbeatTriggerScheduler", () => { expect(eventStore.getActiveHeartbeatRun).not.toHaveBeenCalled(); }); - it("skips trigger when agent has active run", async () => { - (eventStore.getActiveHeartbeatRun as ReturnType).mockResolvedValue({ - id: "run-active", - status: "active", - }); + // Regression surface checklist for deferred assignments: + // - active-run assignment skip records pending work; no-active-run control remains immediate + // - run-completion drain re-fires once, latest rapid re-assignment wins + // - transient global/engine pause, new active run, and parallel-execution guards preserve pending work + // - terminal missing/disabled/budget-exhausted states and unregister clear pending work + // - skipHeartbeatWhenIdle/long timer stalls are avoided because drain is completion-driven, not timer-driven + it("defers an active-run assignment and re-fires it exactly once on drain", async () => { + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); - const agent = { id: "agent-test", name: "Test" } as import("@fusion/core").Agent; - eventStore.emit("agent:assigned", agent, "FN-003"); + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-001"); await new Promise((resolve) => setTimeout(resolve, 10)); + expect(callback).not.toHaveBeenCalled(); + + await scheduler.drainPendingAssignment("agent-test"); + + expect(callback).toHaveBeenCalledOnce(); + expect(callback).toHaveBeenCalledWith("agent-test", "assignment", expect.objectContaining({ + taskId: "FN-001", + wakeReason: "assignment", + triggerDetail: "task-assigned", + })); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).toHaveBeenCalledOnce(); + }); + + it("does not record pending work when assignment fires immediately", async () => { + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-002"); + + await vi.waitFor(() => { + expect(callback).toHaveBeenCalledOnce(); + }, { timeout: 1000 }); + + callback.mockClear(); + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + }); + + it("keeps only the latest task when assignments are repeated during an active run", async () => { + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); + + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-OLD"); + eventStore.emit("agent:assigned", agent, "FN-LATEST"); + + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(callback).not.toHaveBeenCalled(); + + await scheduler.drainPendingAssignment("agent-test"); + + expect(callback).toHaveBeenCalledOnce(); + expect(callback).toHaveBeenCalledWith("agent-test", "assignment", expect.objectContaining({ taskId: "FN-LATEST" })); + }); + + it.each([ + ["globalPause", { globalPause: true }], + ["enginePaused", { enginePaused: true }], + ])("preserves pending assignment while %s blocks drain", async (_name, settings) => { + scheduler.stop(); + const pausedTaskStore = { getSettings: vi.fn().mockResolvedValue(settings) } as unknown as TaskStore; + scheduler = new HeartbeatTriggerScheduler(eventStore as unknown as AgentStore, callback, pausedTaskStore); + scheduler.start(); + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); + + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-PAUSED"); + await new Promise((resolve) => setTimeout(resolve, 10)); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + + (pausedTaskStore.getSettings as ReturnType).mockResolvedValue({}); + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).toHaveBeenCalledOnce(); + expect(callback).toHaveBeenCalledWith("agent-test", "assignment", expect.objectContaining({ taskId: "FN-PAUSED" })); + }); + + it("preserves pending assignment when a new active run exists at drain time", async () => { + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValueOnce({ id: "run-new", status: "active" }) + .mockResolvedValue(null); + + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-ACTIVE"); + await new Promise((resolve) => setTimeout(resolve, 10)); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).toHaveBeenCalledOnce(); + expect(callback).toHaveBeenCalledWith("agent-test", "assignment", expect.objectContaining({ taskId: "FN-ACTIVE" })); + }); + + it.each([ + ["missing agent", async () => { + eventStore.getAgent.mockResolvedValue(null); + }], + ["disabled agent", async () => { + eventStore.getAgent.mockResolvedValue({ id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {}, runtimeConfig: { enabled: false } }); + }], + ["budget exhausted", async () => { + eventStore.getBudgetStatus.mockResolvedValue(createBudgetStatus({ + agentId: "agent-test", + isOverBudget: true, + isOverThreshold: true, + usagePercent: 100, + budgetLimit: 1000, + thresholdPercent: 80, + })); + }], + ])("clears pending assignment without re-fire for %s", async (_name, configureTerminal) => { + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-CLEAR"); + await new Promise((resolve) => setTimeout(resolve, 10)); + await configureTerminal(); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + + eventStore.getAgent.mockResolvedValue({ id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} }); + eventStore.getBudgetStatus.mockRejectedValue(new Error("budget status unavailable")); + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + }); + + it("preserves pending assignment while parallel execution guard blocks drain", async () => { + scheduler.stop(); + scheduler = new HeartbeatTriggerScheduler(eventStore as unknown as AgentStore, callback, undefined, { + isTaskExecuting: (taskId) => taskId === "FN-EXECUTING", + }); + scheduler.start(); + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); + + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {}, runtimeConfig: { allowParallelExecution: false } } as import("@fusion/core").Agent; + eventStore.getAgent.mockResolvedValue(agent); + eventStore.emit("agent:assigned", agent, "FN-EXECUTING"); + await new Promise((resolve) => setTimeout(resolve, 10)); + + await scheduler.drainPendingAssignment("agent-test"); + expect(callback).not.toHaveBeenCalled(); + }); + + it("clears pending assignment when unregistering an agent", async () => { + (eventStore.getActiveHeartbeatRun as ReturnType) + .mockResolvedValueOnce({ id: "run-active", status: "active" }) + .mockResolvedValue(null); + + const agent = { id: "agent-test", name: "Test", role: "executor", state: "active", metadata: {} } as import("@fusion/core").Agent; + eventStore.emit("agent:assigned", agent, "FN-UNREGISTER"); + await new Promise((resolve) => setTimeout(resolve, 10)); + + scheduler.unregisterAgent("agent-test"); + await scheduler.drainPendingAssignment("agent-test"); expect(callback).not.toHaveBeenCalled(); }); diff --git a/packages/engine/src/agent-heartbeat.ts b/packages/engine/src/agent-heartbeat.ts index 8184148e66..5baa6a41e2 100644 --- a/packages/engine/src/agent-heartbeat.ts +++ b/packages/engine/src/agent-heartbeat.ts @@ -3627,11 +3627,19 @@ function readHeartbeatTimerRepairMetadata(agent: Agent): HeartbeatTimerRepairMet }; } +type PendingAssignment = { + taskId: string; + triggeringCommentIds?: string[]; + triggeringCommentType?: "steering" | "task" | "pr"; + budgetStatus?: AgentBudgetStatus; +}; + export class HeartbeatTriggerScheduler { private store: AgentStore; private callback: TriggerCallback; private taskStore?: TaskStore; private timers: Map = new Map(); + private pendingAssignments: Map = new Map(); private registrationEpochs: Map = new Map(); private running = false; private assignedListener: ((agent: import("@fusion/core").Agent, taskId: string) => void) | null = null; @@ -3911,6 +3919,7 @@ export class HeartbeatTriggerScheduler { */ unregisterAgent(agentId: string): void { this.registrationEpochs.set(agentId, (this.registrationEpochs.get(agentId) ?? 0) + 1); + this.pendingAssignments.delete(agentId); if (this.timers.has(agentId)) { this.clearAgentTimer(agentId); heartbeatLog.log(`Unregistered timer for ${agentId}`); @@ -3948,9 +3957,12 @@ export class HeartbeatTriggerScheduler { return; } - // Guard: skip if agent already has an active run + // Guard: skip if agent already has an active run. Preserve this + // assignment for completion-driven re-fire so it is not stranded by + // long/idle-skipped timer intervals. const activeRun = await this.store.getActiveHeartbeatRun(agent.id); if (activeRun) { + this.pendingAssignments.set(agent.id, { taskId }); heartbeatLog.log(`Assignment trigger skipped for ${agent.id} (active run)`); return; } @@ -4024,6 +4036,101 @@ export class HeartbeatTriggerScheduler { heartbeatLog.log("Watching agent:assigned events"); } + /** + * Re-evaluate and re-fire an assignment trigger that was deferred because + * the agent already had an active heartbeat run. Transient ineligibility + * keeps the pending entry so a later completion can retry; terminal + * ineligibility clears it. + */ + async drainPendingAssignment(agentId: string): Promise { + if (!this.running) return; + + const pending = this.pendingAssignments.get(agentId); + if (!pending) { + return; + } + + try { + const agent = await this.store.getAgent(agentId); + if (!agent) { + this.pendingAssignments.delete(agentId); + heartbeatLog.log(`Deferred assignment cleared for ${agentId} (agent missing)`); + return; + } + + if (!isHeartbeatManaged(agent)) { + this.pendingAssignments.delete(agentId); + heartbeatLog.log(`Deferred assignment cleared for ${agentId} (ephemeral/internal)`); + return; + } + + const runtimeConfig = (agent.runtimeConfig ?? {}) as { enabled?: boolean; allowParallelExecution?: boolean }; + if (runtimeConfig.enabled === false) { + this.pendingAssignments.delete(agentId); + heartbeatLog.log(`Deferred assignment cleared for ${agentId} (disabled)`); + return; + } + + if (!isTickableState(agent.state)) { + heartbeatLog.log(`Deferred assignment preserved for ${agentId} (state=${agent.state})`); + return; + } + + const settings = this.taskStore ? await this.taskStore.getSettings() : null; + if (settings?.globalPause) { + heartbeatLog.log(`Deferred assignment preserved for ${agentId} (global pause active)`); + return; + } + if (settings?.enginePaused) { + heartbeatLog.log(`Deferred assignment preserved for ${agentId} (engine paused)`); + return; + } + + const activeRun = await this.store.getActiveHeartbeatRun(agentId); + if (activeRun) { + heartbeatLog.log(`Deferred assignment preserved for ${agentId} (active run)`); + return; + } + + if ( + runtimeConfig.allowParallelExecution === false + && (this.isTaskExecuting?.(pending.taskId) || this.isAgentEffectivelyExecuting?.(agentId)) + ) { + heartbeatLog.log(`Deferred assignment preserved for ${agentId} (parallel execution disabled, task ${pending.taskId} or column-bound session executing)`); + return; + } + + let budgetStatus: AgentBudgetStatus | undefined = pending.budgetStatus; + try { + budgetStatus = await this.store.getBudgetStatus(agentId); + if (budgetStatus.isOverBudget) { + this.pendingAssignments.delete(agentId); + heartbeatLog.log(`Deferred assignment cleared for ${agentId} (budget exhausted)`); + return; + } + } catch (budgetErr) { + heartbeatLog.warn(`Deferred assignment budget check failed for ${agentId}: ${budgetErr instanceof Error ? budgetErr.message : String(budgetErr)} — proceeding without budget check`); + } + + this.pendingAssignments.delete(agentId); + heartbeatLog.log(`Deferred assignment re-fired for ${agentId} (task: ${pending.taskId})`); + await this.callback(agentId, "assignment", { + taskId: pending.taskId, + wakeReason: "assignment", + triggerDetail: "task-assigned", + ...(pending.triggeringCommentIds?.length + ? { + triggeringCommentIds: pending.triggeringCommentIds, + triggeringCommentType: pending.triggeringCommentType ?? "steering", + } + : {}), + ...(budgetStatus && { budgetStatus }), + }); + } catch (err) { + heartbeatLog.error(`Deferred assignment drain error for ${agentId}: ${err instanceof Error ? err.message : err}`); + } + } + /** * Unsubscribe from agent:assigned events. */ 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 0dcbefca34..f4c3aaf39a 100644 --- a/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts +++ b/packages/engine/src/runtimes/__tests__/in-process-runtime.test.ts @@ -18,6 +18,7 @@ const { mockRecoverInterruptedRuns, mockExecutorCtor, mockResumeOrphaned, + mockResumeTaskForAgent, mockTaskStoreSettings, mockTaskStoreGetTask, mockTaskStoreUpdateSettings, @@ -36,6 +37,7 @@ const { mockRecoverInterruptedRuns: vi.fn().mockResolvedValue(undefined), mockExecutorCtor: vi.fn(), mockResumeOrphaned: vi.fn().mockResolvedValue(undefined), + mockResumeTaskForAgent: vi.fn().mockResolvedValue(undefined), mockTaskStoreSettings: {} as Record, mockTaskStoreGetTask: vi.fn().mockResolvedValue(null), mockTaskStoreUpdateSettings: vi.fn().mockResolvedValue(undefined), @@ -193,6 +195,7 @@ vi.mock("../../executor.js", async () => { mockExecutorCtor(options); const self = {} as Record; self.resumeOrphaned = mockResumeOrphaned; + self.resumeTaskForAgent = mockResumeTaskForAgent; self.recoverCompletedTask = vi.fn().mockResolvedValue(true); self.getExecutingTaskIds = vi.fn().mockReturnValue(new Set()); self.handleLoopDetected = vi.fn().mockResolvedValue(false); @@ -244,6 +247,8 @@ describe("InProcessRuntime", () => { } mockTaskStoreGetTask.mockReset(); mockTaskStoreGetTask.mockResolvedValue(null); + mockResumeTaskForAgent.mockReset(); + mockResumeTaskForAgent.mockResolvedValue(undefined); mockIsGitRepository.mockReset(); mockIsGitRepository.mockResolvedValue(true); mockReapOrphanWorktrees.mockReset(); @@ -699,6 +704,25 @@ describe("InProcessRuntime", () => { }); describe("trigger scheduler wiring", () => { + it("composes run-completion resume with deferred assignment drain", async () => { + await runtime.start(); + const store = getAgentStore(runtime); + const agent = await store.createAgent({ name: "completion-wiring", role: "executor" }); + const monitor = runtime.getHeartbeatMonitor(); + const triggerScheduler = runtime.getTriggerScheduler(); + expect(monitor).toBeDefined(); + expect(triggerScheduler).toBeDefined(); + const drainSpy = vi.spyOn(triggerScheduler!, "drainPendingAssignment").mockResolvedValue(undefined); + + const run = await monitor!.startRun(agent.id, { source: "timer" }); + await monitor!.completeRun(agent.id, run.id, { status: "completed" }); + + await vi.waitFor(() => { + expect(mockResumeTaskForAgent).toHaveBeenCalledWith(agent.id); + expect(drainSpy).toHaveBeenCalledWith(agent.id); + }); + }, 30000); + it("creates trigger scheduler on start", async () => { await runtime.start(); expect(runtime.getTriggerScheduler()).toBeDefined(); diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 40ed13d67f..ec32551b4b 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -610,6 +610,9 @@ export class InProcessRuntime runtimeLog.warn(`resumeTaskForAgent failed for ${agentId}: ${err instanceof Error ? err.message : String(err)}`); }); } + void this.triggerScheduler?.drainPendingAssignment(agentId).catch((err) => { + runtimeLog.warn(`drainPendingAssignment failed for ${agentId}: ${err instanceof Error ? err.message : String(err)}`); + }); }, }); this.heartbeatMonitor.start();