feat(FN-4119): add stale heartbeat run reaper to prevent engine blockage
Implements a stale heartbeat reaper in the engine that detects and cleans up heartbeat runs exceeding their timeout, preventing blocked heartbeat cycles from stalling agent activity. Includes comprehensive tests for the reaper gates and documentation of the feature in agents.md. Fusion-Task-Id: FN-4119
This commit is contained in:
5
.changeset/FN-4119-stale-heartbeat-reaper.md
Normal file
5
.changeset/FN-4119-stale-heartbeat-reaper.md
Normal file
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@runfusion/fusion": patch
|
||||
---
|
||||
|
||||
Heartbeat scheduling now auto-reaps stale active heartbeat runs so durable agents recover regular timer ticks without requiring a manual stop/start.
|
||||
@@ -1093,9 +1093,32 @@ Effects:
|
||||
|
||||
This covers the untracked timer-loss failure mode where no `agent:updated` event fires after a timer entry disappears. Manual stop/start is no longer required to re-arm the timer in that case.
|
||||
|
||||
### Stale Active-Run Reaper (FN-4119)
|
||||
|
||||
`HeartbeatTriggerScheduler` also reaps **stale persisted `status="active"` heartbeat runs** before they can block future timer progress forever.
|
||||
|
||||
When it fires:
|
||||
- `onTimerTick()` finds an active run row for a durable, tickable agent
|
||||
- or `auditTimerRegistrations()` finds a missing timer plus an active run row for that same durable agent
|
||||
- the persisted run has no fresh heartbeat for longer than **`heartbeatTimeoutMs × heartbeatRepairStaleMultiplier`**
|
||||
- the engine is not globally paused and not timer-paused via `enginePaused`
|
||||
|
||||
Threshold semantics:
|
||||
- The reaper reuses the same `heartbeatRepairStaleMultiplier` setting that timer-audit repair already uses; no extra stale-run knob exists
|
||||
- The base signal is the agent's `lastHeartbeatAt` / `recordHeartbeat(...)` freshness, not the scheduled timer interval
|
||||
- Default threshold is therefore **`2 × heartbeatTimeoutMs`** (default timeout `60s` → default reap threshold `120s`)
|
||||
|
||||
Layering with the existing recovery paths:
|
||||
- **`HeartbeatMonitor.reconcileOrphanedRunningAgents()`** still handles monitor-owned stale `running` agents and other tracked-session cleanup
|
||||
- **`HeartbeatTriggerScheduler.onTimerTick()`** now reaps a stale active run, logs `reason=tick-proceeded-after-reap`, and proceeds with the scheduled callback in the same tick
|
||||
- **`HeartbeatTriggerScheduler.auditTimerRegistrations()`** now reaps the stale active run first, then re-arms the missing timer in the same audit pass and logs `reason=timer-audit-rearmed`
|
||||
- Healthy active runs within threshold still keep the old `(active run)` skip behavior
|
||||
- Ephemeral/task-worker agents are never reaped by this path
|
||||
|
||||
Separation of responsibilities:
|
||||
- **HeartbeatMonitor recovery** handles **tracked stale sessions** (stuck in-memory run/session cleanup + pause/resume restart)
|
||||
- **HeartbeatTriggerScheduler audit** handles **untracked missing-timer registration drift** (re-arm scheduling)
|
||||
- **HeartbeatTriggerScheduler stale-run reaper** handles **orphaned persisted active runs** that would otherwise cause both tick and audit to skip forever on `(active run)`
|
||||
|
||||
## Dashboard Health Status
|
||||
|
||||
|
||||
@@ -37,6 +37,9 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
getActiveHeartbeatRun: vi.fn().mockResolvedValue(null),
|
||||
getBudgetStatus: vi.fn().mockResolvedValue(createBudgetStatus()),
|
||||
listAgents: vi.fn().mockResolvedValue([]),
|
||||
getRunDetail: vi.fn().mockResolvedValue(null),
|
||||
saveRun: vi.fn().mockResolvedValue(undefined),
|
||||
endHeartbeatRun: vi.fn().mockResolvedValue(undefined),
|
||||
on: vi.fn(),
|
||||
off: vi.fn(),
|
||||
updateAgent: vi.fn().mockImplementation(async (_id: string, updates: { metadata: Record<string, unknown> }) => ({
|
||||
@@ -232,6 +235,94 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
expect(scheduler.getRegisteredAgents()).not.toContain("agent-001");
|
||||
});
|
||||
|
||||
it("FN-4119 reaps a stale active run and re-arms the timer when audit finds a lost registration", async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date("2026-01-01T02:00:00.000Z"));
|
||||
const agent = {
|
||||
id: "agent-001",
|
||||
name: "Agent 001",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
lastHeartbeatAt: "2026-01-01T00:00:00.000Z",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 3_600_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
metadata: {},
|
||||
} as Agent;
|
||||
const activeRun = {
|
||||
id: "run-stale",
|
||||
agentId: "agent-001",
|
||||
startedAt: "2026-01-01T00:00:00.000Z",
|
||||
status: "active",
|
||||
} as any;
|
||||
vi.mocked(store.listAgents).mockResolvedValue([agent]);
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue(activeRun);
|
||||
vi.mocked(store.getRunDetail).mockResolvedValue(activeRun);
|
||||
|
||||
scheduler = new HeartbeatTriggerScheduler(store, callback);
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
expect(store.endHeartbeatRun).toHaveBeenCalledOnce();
|
||||
expect(store.endHeartbeatRun).toHaveBeenCalledWith("run-stale", "terminated");
|
||||
expect(scheduler.getRegisteredAgents()).toContain("agent-001");
|
||||
expect(heartbeatLog.warn).toHaveBeenCalledWith(expect.stringContaining("reason=orphaned-run-reaped agentId=agent-001 runId=run-stale"));
|
||||
expect(heartbeatLog.log).toHaveBeenCalledWith(expect.stringContaining("reason=timer-audit-rearmed agentId=agent-001 runId=run-stale"));
|
||||
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue(null);
|
||||
await vi.advanceTimersByTimeAsync(60_000);
|
||||
expect(store.endHeartbeatRun).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("FN-4119 leaves healthy active runs alone during audit", async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date("2026-01-01T00:00:15.000Z"));
|
||||
const agent = {
|
||||
id: "agent-001",
|
||||
name: "Agent 001",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
lastHeartbeatAt: "2026-01-01T00:00:10.000Z",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 3_600_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:10.000Z",
|
||||
metadata: {},
|
||||
} as Agent;
|
||||
vi.mocked(store.listAgents).mockResolvedValue([agent]);
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue({ id: "run-healthy", status: "active" } as any);
|
||||
|
||||
scheduler = new HeartbeatTriggerScheduler(store, callback);
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
expect(store.endHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(scheduler.getRegisteredAgents()).not.toContain("agent-001");
|
||||
expect(heartbeatLog.log).toHaveBeenCalledWith("Timer audit skipped re-arm for agent-001 (active run)");
|
||||
});
|
||||
|
||||
it("FN-4119 does not reap task-worker runs during audit", async () => {
|
||||
vi.useFakeTimers();
|
||||
const agent = {
|
||||
id: "executor-FN-999",
|
||||
name: "executor-FN-999",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 30_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
metadata: { agentKind: "task-worker" },
|
||||
} as Agent;
|
||||
vi.mocked(store.listAgents).mockResolvedValue([agent]);
|
||||
|
||||
scheduler = new HeartbeatTriggerScheduler(store, callback);
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(0);
|
||||
|
||||
expect(store.getActiveHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(store.endHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(scheduler.getRegisteredAgents()).not.toContain("executor-FN-999");
|
||||
});
|
||||
});
|
||||
|
||||
describe("registerAgent", () => {
|
||||
@@ -532,6 +623,100 @@ describe("HeartbeatTriggerScheduler", () => {
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("FN-4119 reaps a stale active run and proceeds with the timer tick", async () => {
|
||||
const staleAgent = {
|
||||
id: "agent-001",
|
||||
name: "Agent 001",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
lastHeartbeatAt: "2026-01-01T00:00:00.000Z",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 30_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
metadata: {},
|
||||
} as Agent;
|
||||
const activeRun = {
|
||||
id: "run-stale",
|
||||
agentId: "agent-001",
|
||||
startedAt: "2026-01-01T00:00:00.000Z",
|
||||
status: "active",
|
||||
} as any;
|
||||
vi.setSystemTime(new Date("2026-01-01T02:00:00.000Z"));
|
||||
vi.mocked(store.getAgent).mockResolvedValue(staleAgent);
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue(activeRun);
|
||||
vi.mocked(store.getRunDetail).mockResolvedValue(activeRun);
|
||||
|
||||
await (scheduler as any).onTimerTick("agent-001", 30_000);
|
||||
|
||||
expect(store.endHeartbeatRun).toHaveBeenCalledOnce();
|
||||
expect(store.endHeartbeatRun).toHaveBeenCalledWith("run-stale", "terminated");
|
||||
expect(callback).toHaveBeenCalledOnce();
|
||||
expect(callback).toHaveBeenCalledWith("agent-001", "timer", {
|
||||
wakeReason: "timer",
|
||||
triggerDetail: "scheduled",
|
||||
intervalMs: 30_000,
|
||||
});
|
||||
expect(heartbeatLog.log).toHaveBeenCalledWith(expect.stringContaining("reason=tick-proceeded-after-reap agentId=agent-001 runId=run-stale"));
|
||||
});
|
||||
|
||||
it("FN-4119 preserves the active-run skip when the run is still healthy", async () => {
|
||||
const healthyAgent = {
|
||||
id: "agent-001",
|
||||
name: "Agent 001",
|
||||
role: "executor",
|
||||
state: "active",
|
||||
lastHeartbeatAt: "2026-01-01T00:00:12.000Z",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 30_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:12.000Z",
|
||||
metadata: {},
|
||||
} as Agent;
|
||||
vi.setSystemTime(new Date("2026-01-01T00:00:15.000Z"));
|
||||
vi.mocked(store.getAgent).mockResolvedValue(healthyAgent);
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue({ id: "run-healthy", status: "active" } as any);
|
||||
|
||||
await (scheduler as any).onTimerTick("agent-001", 30_000);
|
||||
|
||||
expect(store.endHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
expect(heartbeatLog.log).toHaveBeenCalledWith("Timer tick skipped for agent-001 (active run)");
|
||||
});
|
||||
|
||||
it.each([
|
||||
{ name: "paused agent state", agentState: "paused" as const, settings: null, expectedLog: "Timer tick skipped for agent-001 (state=paused)" },
|
||||
{ name: "global pause", agentState: "active" as const, settings: { globalPause: true, enginePaused: false }, expectedLog: "Timer tick skipped for agent-001 (global pause active)" },
|
||||
{ name: "engine pause", agentState: "active" as const, settings: { globalPause: false, enginePaused: true }, expectedLog: "Timer tick skipped for agent-001 (engine paused)" },
|
||||
])("FN-4119 does not reap stale runs during $name", async ({ agentState, settings, expectedLog }) => {
|
||||
scheduler.stop();
|
||||
const taskStore = settings
|
||||
? ({ getSettings: vi.fn().mockResolvedValue(settings) } as unknown as TaskStore)
|
||||
: undefined;
|
||||
scheduler = new HeartbeatTriggerScheduler(store, callback, taskStore);
|
||||
scheduler.start();
|
||||
|
||||
const agent = {
|
||||
id: "agent-001",
|
||||
name: "Agent 001",
|
||||
role: "executor",
|
||||
state: agentState,
|
||||
lastHeartbeatAt: "2026-01-01T00:00:00.000Z",
|
||||
runtimeConfig: { enabled: true, heartbeatIntervalMs: 30_000, heartbeatTimeoutMs: 10_000 },
|
||||
createdAt: "2026-01-01T00:00:00.000Z",
|
||||
updatedAt: "2026-01-01T00:00:00.000Z",
|
||||
metadata: {},
|
||||
} as Agent;
|
||||
vi.setSystemTime(new Date("2026-01-01T02:00:00.000Z"));
|
||||
vi.mocked(store.getAgent).mockResolvedValue(agent);
|
||||
vi.mocked(store.getActiveHeartbeatRun).mockResolvedValue({ id: "run-stale", status: "active" } as any);
|
||||
|
||||
await (scheduler as any).onTimerTick("agent-001", 30_000);
|
||||
|
||||
expect(store.getActiveHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(store.endHeartbeatRun).not.toHaveBeenCalled();
|
||||
expect(callback).not.toHaveBeenCalled();
|
||||
expect(heartbeatLog.log).toHaveBeenCalledWith(expectedLog);
|
||||
});
|
||||
|
||||
it("skips timer dispatch when global pause is active", async () => {
|
||||
scheduler.stop();
|
||||
const taskStore = {
|
||||
|
||||
@@ -177,6 +177,30 @@ function formatRelativeTime(iso?: string | null): string {
|
||||
return `${formatDuration(elapsed)} ago`;
|
||||
}
|
||||
|
||||
function getHeartbeatAgeMs(agent: Agent, now: number = Date.now()): number {
|
||||
const lastTs = agent.lastHeartbeatAt ? Date.parse(agent.lastHeartbeatAt) : Number.NaN;
|
||||
return Number.isFinite(lastTs) ? Math.max(0, now - lastTs) : Number.NaN;
|
||||
}
|
||||
|
||||
async function terminatePersistedHeartbeatRun(
|
||||
store: AgentStore,
|
||||
agentId: string,
|
||||
runId: string,
|
||||
stderrExcerpt: string,
|
||||
): Promise<boolean> {
|
||||
const detail = await store.getRunDetail(agentId, runId);
|
||||
if (detail && detail.status !== "completed" && detail.status !== "failed" && detail.status !== "terminated") {
|
||||
await store.saveRun({
|
||||
...detail,
|
||||
endedAt: new Date().toISOString(),
|
||||
status: "terminated",
|
||||
stderrExcerpt,
|
||||
});
|
||||
}
|
||||
await store.endHeartbeatRun(runId, "terminated");
|
||||
return true;
|
||||
}
|
||||
|
||||
function isAutoClaimRelevantTasksEnabled(agent: Agent): boolean {
|
||||
const runtimeConfig = (agent.runtimeConfig ?? {}) as Record<string, unknown>;
|
||||
return runtimeConfig.autoClaimRelevantTasks !== false;
|
||||
@@ -750,24 +774,19 @@ export class HeartbeatMonitor {
|
||||
reason = "no active run";
|
||||
} else if (!this.trackedAgents.has(agent.id)) {
|
||||
const timeoutMs = this.resolveAgentConfig(agent.id).heartbeatTimeoutMs;
|
||||
const lastTs = agent.lastHeartbeatAt ? Date.parse(agent.lastHeartbeatAt) : NaN;
|
||||
const heartbeatAgeMs = Number.isFinite(lastTs) ? Math.max(0, now - lastTs) : Infinity;
|
||||
if (heartbeatAgeMs > timeoutMs * 3) {
|
||||
const heartbeatAgeMs = getHeartbeatAgeMs(agent, now);
|
||||
if (!Number.isFinite(heartbeatAgeMs) || heartbeatAgeMs > timeoutMs * 3) {
|
||||
try {
|
||||
const detail = await this.store.getRunDetail(agent.id, activeRun.id);
|
||||
if (detail && detail.status !== "completed" && detail.status !== "failed" && detail.status !== "terminated") {
|
||||
await this.store.saveRun({
|
||||
...detail,
|
||||
endedAt: new Date().toISOString(),
|
||||
status: "terminated",
|
||||
stderrExcerpt: `Reconciled stale run (no heartbeat for ${formatDuration(heartbeatAgeMs)}; threshold ${formatDuration(timeoutMs * 3)})`,
|
||||
});
|
||||
}
|
||||
await this.store.endHeartbeatRun(activeRun.id, "terminated");
|
||||
await terminatePersistedHeartbeatRun(
|
||||
this.store,
|
||||
agent.id,
|
||||
activeRun.id,
|
||||
`Reconciled stale run (no heartbeat for ${Number.isFinite(heartbeatAgeMs) ? formatDuration(heartbeatAgeMs) : "unknown"}; threshold ${formatDuration(timeoutMs * 3)})`,
|
||||
);
|
||||
} catch (runEndErr) {
|
||||
heartbeatLog.warn(`Failed to terminate stale run ${activeRun.id} for ${agent.id}: ${runEndErr instanceof Error ? runEndErr.message : String(runEndErr)}`);
|
||||
}
|
||||
reason = `stale heartbeat (${formatDuration(heartbeatAgeMs)} since lastHeartbeatAt)`;
|
||||
reason = `stale heartbeat (${Number.isFinite(heartbeatAgeMs) ? formatDuration(heartbeatAgeMs) : "unknown"} since lastHeartbeatAt)`;
|
||||
}
|
||||
}
|
||||
if (!reason) continue;
|
||||
@@ -2894,6 +2913,7 @@ export class HeartbeatTriggerScheduler {
|
||||
|
||||
private static readonly TIMER_AUDIT_INTERVAL_MS = 60_000;
|
||||
private static readonly DEFAULT_REPAIR_STALE_MULTIPLIER = 2;
|
||||
private static readonly DEFAULT_HEARTBEAT_TIMEOUT_MS = 60_000;
|
||||
|
||||
constructor(store: AgentStore, callback: TriggerCallback, taskStore?: TaskStore, options?: { isTaskExecuting?: (taskId: string) => boolean }) {
|
||||
this.store = store;
|
||||
@@ -3405,6 +3425,43 @@ export class HeartbeatTriggerScheduler {
|
||||
return Math.round(intervalMs * staleMultiplier);
|
||||
}
|
||||
|
||||
private getActiveRunStaleThresholdMs(agent: Agent, staleMultiplier: number): number {
|
||||
const runtimeConfig = (agent.runtimeConfig ?? {}) as { heartbeatTimeoutMs?: number };
|
||||
const rawTimeoutMs = runtimeConfig.heartbeatTimeoutMs;
|
||||
const timeoutMs = typeof rawTimeoutMs === "number" && Number.isFinite(rawTimeoutMs) && rawTimeoutMs > 0
|
||||
? Math.max(5000, Math.round(rawTimeoutMs))
|
||||
: HeartbeatTriggerScheduler.DEFAULT_HEARTBEAT_TIMEOUT_MS;
|
||||
return Math.round(timeoutMs * staleMultiplier);
|
||||
}
|
||||
|
||||
private async maybeReapStaleActiveRun(
|
||||
agent: Agent,
|
||||
activeRun: AgentHeartbeatRun,
|
||||
reason: "audit" | "timer",
|
||||
staleMultiplier: number,
|
||||
): Promise<{ reaped: boolean; elapsedMs: number; thresholdMs: number }> {
|
||||
if (!isHeartbeatManaged(agent)) {
|
||||
return { reaped: false, elapsedMs: Number.NaN, thresholdMs: Number.NaN };
|
||||
}
|
||||
|
||||
const thresholdMs = this.getActiveRunStaleThresholdMs(agent, staleMultiplier);
|
||||
const elapsedMs = getHeartbeatAgeMs(agent);
|
||||
if (!Number.isFinite(elapsedMs) || elapsedMs <= thresholdMs) {
|
||||
return { reaped: false, elapsedMs, thresholdMs };
|
||||
}
|
||||
|
||||
await terminatePersistedHeartbeatRun(
|
||||
this.store,
|
||||
agent.id,
|
||||
activeRun.id,
|
||||
`Reaped stale heartbeat run before next tick (no heartbeat for ${formatDuration(elapsedMs)}; threshold ${formatDuration(thresholdMs)})`,
|
||||
);
|
||||
heartbeatLog.warn(
|
||||
`Heartbeat stale-run reaped reason=orphaned-run-reaped agentId=${agent.id} runId=${activeRun.id} elapsedMs=${elapsedMs} thresholdMs=${thresholdMs} source=${reason}`,
|
||||
);
|
||||
return { reaped: true, elapsedMs, thresholdMs };
|
||||
}
|
||||
|
||||
private async markRepairMetadata(agent: Agent, staleAtRepair: boolean, staleRepairReason?: string): Promise<void> {
|
||||
const updater = (this.store as { updateAgent?: (agentId: string, updates: { metadata: Record<string, unknown> }) => Promise<unknown> }).updateAgent;
|
||||
if (typeof updater !== "function") {
|
||||
@@ -3447,9 +3504,23 @@ export class HeartbeatTriggerScheduler {
|
||||
if (this.timers.has(agent.id)) continue;
|
||||
|
||||
const activeRun = await this.store.getActiveHeartbeatRun(agent.id);
|
||||
const activeRunId = activeRun?.id ?? null;
|
||||
let reapedActiveRun = false;
|
||||
let activeRunElapsedMs = Number.NaN;
|
||||
let activeRunThresholdMs = Number.NaN;
|
||||
if (activeRun) {
|
||||
heartbeatLog.log(`Timer audit skipped re-arm for ${agent.id} (active run)`);
|
||||
continue;
|
||||
if (settings?.globalPause || settings?.enginePaused) {
|
||||
heartbeatLog.log(`Timer audit skipped re-arm for ${agent.id} (active run)`);
|
||||
continue;
|
||||
}
|
||||
const reapResult = await this.maybeReapStaleActiveRun(agent, activeRun, "audit", staleMultiplier);
|
||||
reapedActiveRun = reapResult.reaped;
|
||||
activeRunElapsedMs = reapResult.elapsedMs;
|
||||
activeRunThresholdMs = reapResult.thresholdMs;
|
||||
if (!reapedActiveRun) {
|
||||
heartbeatLog.log(`Timer audit skipped re-arm for ${agent.id} (active run)`);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
this.registerAgent(agent.id, this.getAgentTimerConfig(agent), {
|
||||
@@ -3457,8 +3528,7 @@ export class HeartbeatTriggerScheduler {
|
||||
});
|
||||
|
||||
const staleThresholdMs = this.getRepairStaleThresholdMs(agent, staleMultiplier);
|
||||
const lastHeartbeatMs = agent.lastHeartbeatAt ? Date.parse(agent.lastHeartbeatAt) : Number.NaN;
|
||||
const elapsedMs = Number.isFinite(lastHeartbeatMs) ? Date.now() - lastHeartbeatMs : Number.NaN;
|
||||
const elapsedMs = getHeartbeatAgeMs(agent);
|
||||
const staleAtRepair = Number.isFinite(elapsedMs) && elapsedMs > staleThresholdMs;
|
||||
const staleRepairReason = staleAtRepair
|
||||
? `No heartbeat for ${Math.round(elapsedMs / 1000)}s before timer audit repair (threshold ${Math.round(staleThresholdMs / 1000)}s)`
|
||||
@@ -3466,6 +3536,11 @@ export class HeartbeatTriggerScheduler {
|
||||
await this.markRepairMetadata(agent, staleAtRepair, staleRepairReason);
|
||||
|
||||
rearmedCount++;
|
||||
if (reapedActiveRun && activeRunId) {
|
||||
heartbeatLog.log(
|
||||
`Timer audit re-armed after stale-run reap reason=timer-audit-rearmed agentId=${agent.id} runId=${activeRunId} elapsedMs=${activeRunElapsedMs} thresholdMs=${activeRunThresholdMs}`,
|
||||
);
|
||||
}
|
||||
if (staleAtRepair) {
|
||||
heartbeatLog.warn(`Timer re-armed stale agent ${agent.id} (audit:${reason}): ${staleRepairReason ?? "heartbeat exceeded stale threshold before repair"}`);
|
||||
} else {
|
||||
@@ -3501,13 +3576,6 @@ export class HeartbeatTriggerScheduler {
|
||||
return;
|
||||
}
|
||||
|
||||
// Check for active runs
|
||||
const activeRun = await this.store.getActiveHeartbeatRun(agentId);
|
||||
if (activeRun) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (active run)`);
|
||||
return;
|
||||
}
|
||||
|
||||
// Guard: when parallel execution is disabled, skip if the agent's bound task is actively executing
|
||||
const timerRc = (agent.runtimeConfig ?? {}) as { allowParallelExecution?: boolean };
|
||||
if (timerRc.allowParallelExecution === false && agent.taskId && this.isTaskExecuting?.(agent.taskId)) {
|
||||
@@ -3517,16 +3585,28 @@ export class HeartbeatTriggerScheduler {
|
||||
|
||||
// Global/engine pause guard: scheduler should not dispatch timer callbacks
|
||||
// while globally paused (hard stop) or engine paused (soft stop for timers).
|
||||
if (this.taskStore) {
|
||||
const settings = await this.taskStore.getSettings();
|
||||
if (settings.globalPause) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (global pause active)`);
|
||||
return;
|
||||
}
|
||||
if (settings.enginePaused) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (engine paused)`);
|
||||
const settings = this.taskStore ? await this.taskStore.getSettings() : null;
|
||||
if (settings?.globalPause) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (global pause active)`);
|
||||
return;
|
||||
}
|
||||
if (settings?.enginePaused) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (engine paused)`);
|
||||
return;
|
||||
}
|
||||
|
||||
// Check for active runs
|
||||
const activeRun = await this.store.getActiveHeartbeatRun(agentId);
|
||||
if (activeRun) {
|
||||
const staleMultiplier = this.resolveRepairStaleMultiplier(settings);
|
||||
const reapResult = await this.maybeReapStaleActiveRun(agent, activeRun, "timer", staleMultiplier);
|
||||
if (!reapResult.reaped) {
|
||||
heartbeatLog.log(`Timer tick skipped for ${agentId} (active run)`);
|
||||
return;
|
||||
}
|
||||
heartbeatLog.log(
|
||||
`Heartbeat tick resumed after stale-run reap reason=tick-proceeded-after-reap agentId=${agentId} runId=${activeRun.id} elapsedMs=${reapResult.elapsedMs} thresholdMs=${reapResult.thresholdMs}`,
|
||||
);
|
||||
}
|
||||
|
||||
// Budget enforcement is handled in HeartbeatMonitor.executeHeartbeat() for timer sources.
|
||||
|
||||
Reference in New Issue
Block a user