feat(FN-2920): improve remote tunnel setup and heartbeat scheduling
- Add cloudflared install/detection support in remote settings API, UI, and route tests - Surface Cloudflare tunnel prerequisites in Settings modal with remote access docs updates - Harden heartbeat runtime scheduling by avoiding stale timeout state and simplifying runtime timeout handling - Expand CLI/core/dashboard/engine coverage for task lifecycle, agent health, and runtime heartbeat behavior - Add changesets for heartbeat scheduling fixes and PR approval setting updates Fusion-Task-Id: FN-2920
This commit is contained in:
@@ -862,24 +862,46 @@ describe("InProcessRuntime", () => {
|
||||
expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id);
|
||||
});
|
||||
|
||||
it("re-registers an existing agent when agent:updated event is emitted", async () => {
|
||||
// Create a new agent
|
||||
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);
|
||||
|
||||
// Update the agent
|
||||
await vi.advanceTimersByTimeAsync(400);
|
||||
await store.updateAgent(agent.id, {
|
||||
name: "test-agent-update-renamed",
|
||||
});
|
||||
|
||||
// Verify the agent is still registered (re-registration succeeded)
|
||||
expect(scheduler!.getRegisteredAgents()).toContain(agent.id);
|
||||
await vi.advanceTimersByTimeAsync(599);
|
||||
expect(executeHeartbeatSpy).not.toHaveBeenCalled();
|
||||
|
||||
await vi.advanceTimersByTimeAsync(1);
|
||||
await vi.waitFor(() => {
|
||||
expect(executeHeartbeatSpy).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
expect(executeHeartbeatSpy).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
agentId: agent.id,
|
||||
source: "timer",
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it("unregisters an agent when enabled is set to false in update", async () => {
|
||||
@@ -907,6 +929,41 @@ describe("InProcessRuntime", () => {
|
||||
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 vi.waitFor(() => {
|
||||
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();
|
||||
|
||||
@@ -98,8 +98,6 @@ export class InProcessRuntime
|
||||
private triageProcessor?: TriageProcessor;
|
||||
private messageStore?: MessageStore;
|
||||
private concurrencyChangedListener?: (state: { globalMaxConcurrent: number }) => void;
|
||||
private agentCreatedListener?: (agent: import("@fusion/core").Agent) => void;
|
||||
private agentUpdatedListener?: (agent: import("@fusion/core").Agent, previousState?: import("@fusion/core").AgentState) => void;
|
||||
/** Set of agent IDs with scheduled ephemeral cleanup (prevents duplicate deletion) */
|
||||
private pendingEphemeralDeletions = new Set<string>();
|
||||
/** Map of agent IDs to their cleanup timer IDs */
|
||||
@@ -506,9 +504,8 @@ export class InProcessRuntime
|
||||
);
|
||||
this.triggerScheduler.start();
|
||||
|
||||
// Dynamic registration follows per-agent heartbeat enablement and tickable state.
|
||||
// Non-ephemeral agents are managed unless runtimeConfig.enabled is explicitly false.
|
||||
// Paused/error/terminated states are never timer-armed.
|
||||
// Startup bootstrap for already-persisted agents. Ongoing lifecycle
|
||||
// updates are handled inside HeartbeatTriggerScheduler itself.
|
||||
const isHeartbeatEnabledAgent = (agent: import("@fusion/core").Agent) =>
|
||||
!isEphemeralAgent(agent) && agent.runtimeConfig?.enabled !== false;
|
||||
const isTickableHeartbeatState = (state: import("@fusion/core").AgentState) =>
|
||||
@@ -516,34 +513,6 @@ export class InProcessRuntime
|
||||
const isTimerManagedAgent = (agent: import("@fusion/core").Agent) =>
|
||||
isHeartbeatEnabledAgent(agent) && isTickableHeartbeatState(agent.state);
|
||||
|
||||
this.agentCreatedListener = (agent) => {
|
||||
if (!this.triggerScheduler) return;
|
||||
if (!isTimerManagedAgent(agent)) return;
|
||||
const rc = agent.runtimeConfig;
|
||||
this.triggerScheduler.registerAgent(agent.id, {
|
||||
heartbeatIntervalMs: rc?.heartbeatIntervalMs as number | undefined,
|
||||
maxConcurrentRuns: rc?.maxConcurrentRuns as number | undefined,
|
||||
});
|
||||
runtimeLog.log(`Registered new agent ${agent.id} for heartbeat triggers`);
|
||||
};
|
||||
this.agentStore.on("agent:created", this.agentCreatedListener);
|
||||
|
||||
this.agentUpdatedListener = (agent) => {
|
||||
if (!this.triggerScheduler) return;
|
||||
if (!isTimerManagedAgent(agent)) {
|
||||
this.triggerScheduler.unregisterAgent(agent.id);
|
||||
runtimeLog.log(`Unregistered agent ${agent.id} from heartbeat triggers`);
|
||||
return;
|
||||
}
|
||||
const rc = agent.runtimeConfig;
|
||||
this.triggerScheduler.registerAgent(agent.id, {
|
||||
heartbeatIntervalMs: rc?.heartbeatIntervalMs as number | undefined,
|
||||
maxConcurrentRuns: rc?.maxConcurrentRuns as number | undefined,
|
||||
});
|
||||
runtimeLog.log(`Re-registered agent ${agent.id} for heartbeat triggers`);
|
||||
};
|
||||
this.agentStore.on("agent:updated", this.agentUpdatedListener);
|
||||
|
||||
// Listen for agent state transitions to clean up terminated ephemeral agents.
|
||||
// This catches cases where ephemeral agents (task-workers, spawned children) are
|
||||
// terminated by HeartbeatMonitor or other pathways outside of onComplete/onError callbacks.
|
||||
@@ -595,6 +564,7 @@ export class InProcessRuntime
|
||||
if (!isTimerManagedAgent(agent)) continue;
|
||||
const rc = agent.runtimeConfig;
|
||||
this.triggerScheduler.registerAgent(agent.id, {
|
||||
enabled: rc?.enabled as boolean | undefined,
|
||||
heartbeatIntervalMs: rc?.heartbeatIntervalMs as number | undefined,
|
||||
maxConcurrentRuns: rc?.maxConcurrentRuns as number | undefined,
|
||||
});
|
||||
@@ -792,16 +762,6 @@ export class InProcessRuntime
|
||||
|
||||
// 3. Remove agent event listeners (before stopping trigger scheduler)
|
||||
// Guard on this.agentStore being defined - it may not exist if AgentStore init failed
|
||||
if (this.agentCreatedListener && this.agentStore) {
|
||||
this.agentStore.off("agent:created", this.agentCreatedListener);
|
||||
this.agentCreatedListener = undefined;
|
||||
runtimeLog.log("AgentStore agent:created listener removed");
|
||||
}
|
||||
if (this.agentUpdatedListener && this.agentStore) {
|
||||
this.agentStore.off("agent:updated", this.agentUpdatedListener);
|
||||
this.agentUpdatedListener = undefined;
|
||||
runtimeLog.log("AgentStore agent:updated listener removed");
|
||||
}
|
||||
if (this.ephemeralTerminationListener && this.agentStore) {
|
||||
this.agentStore.off("agent:stateChanged", this.ephemeralTerminationListener);
|
||||
this.ephemeralTerminationListener = undefined;
|
||||
|
||||
Reference in New Issue
Block a user