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:
Fusion
2026-04-29 13:49:05 -07:00
committed by gsxdsm
parent b91533ce43
commit 17a072c924
25 changed files with 819 additions and 205 deletions

View File

@@ -0,0 +1,5 @@
---
"@runfusion/fusion": patch
---
Fix agent heartbeat scheduling so disabled agents stay disabled and active timers are not reset by unrelated agent updates.

View File

@@ -0,0 +1,17 @@
---
"@runfusion/fusion": minor
"runfusion.ai": minor
"@fusion/core": minor
"@fusion/dashboard": minor
"@fusion/desktop": minor
"@fusion/engine": minor
"@fusion/mobile": minor
"@fusion/pi-claude-cli": minor
"@fusion/plugin-sdk": minor
---
Add `requirePrApproval` setting (related to [#21](https://github.com/Runfusion/Fusion/issues/21)).
When `mergeStrategy: "pull-request"`, GitHub's `required: true` flag for status checks only flows from branch protection — a Pro feature on private repos. On free private repos, `isPrMergeReady` reports every fresh PR as immediately mergeable, so `autoMerge: true` causes Fusion to auto-squash-merge the moment the PR opens with no chance for a human to review it.
The new `requirePrApproval` setting (project-level, default `false`) makes Fusion hold the merge until at least one approving GitHub review is present (`reviewDecision === "APPROVED"`), independent of GitHub's server-side enforcement. Surfaces in the dashboard's Merge settings panel under the Pull Request strategy. Lets you use Fusion's PR mode as "open the PR, wait for me to approve and merge" on any tier.

View File

@@ -70,6 +70,8 @@ Cloudflare **Quick Tunnel** startup gates (`quickTunnel = true`):
No Cloudflare account, tunnel token, named tunnel, or pre-created ingress URL is required. No Cloudflare account, tunnel token, named tunnel, or pre-created ingress URL is required.
Dashboard note: in Settings → Remote Access, selecting Cloudflare now performs a proactive `cloudflared` CLI detection check and shows a one-click **Install cloudflared** action (with manual command fallback) if the binary is missing.
Runtime command used by engine: Runtime command used by engine:
- `cloudflared tunnel --url http://localhost:<dashboardPort>` - `cloudflared tunnel --url http://localhost:<dashboardPort>`

View File

@@ -268,6 +268,8 @@ The canonical persisted shape is a nested `remoteAccess` object.
Use **[Remote Access runbook](./remote-access.md)** for setup prerequisites (Tailscale/Cloudflare), tokenized login-link security caveats, and operational troubleshooting. Keep this section as a schema reference. Use **[Remote Access runbook](./remote-access.md)** for setup prerequisites (Tailscale/Cloudflare), tokenized login-link security caveats, and operational troubleshooting. Keep this section as a schema reference.
When `remoteAccess.activeProvider` is `cloudflare`, the Settings UI fetches `/api/remote/status` and surfaces `cloudflaredAvailable` to show installed/missing state plus a one-click `POST /api/remote/install-cloudflared` action.
| Setting | Type | Default | Description | | Setting | Type | Default | Description |
|---|---|---:|---| |---|---|---:|---|
| `remoteAccess.enabled` | `boolean` | `false` | Master toggle for remote access orchestration. | | `remoteAccess.enabled` | `boolean` | `false` | Master toggle for remote access orchestration. |

View File

@@ -26,11 +26,12 @@ interface MockTask {
column: string; column: string;
} }
function makeStore(task: MockTask) { function makeStore(task: MockTask, settings: Record<string, unknown> = {}) {
const emitter = new EventEmitter(); const emitter = new EventEmitter();
const updates: Array<{ id: string; patch: Record<string, unknown> }> = []; const updates: Array<{ id: string; patch: Record<string, unknown> }> = [];
return Object.assign(emitter, { return Object.assign(emitter, {
getTask: vi.fn().mockResolvedValue(task), getTask: vi.fn().mockResolvedValue(task),
getSettings: vi.fn().mockResolvedValue({ requirePrApproval: false, ...settings }),
updateTask: vi.fn(async (id: string, patch: Record<string, unknown>) => { updateTask: vi.fn(async (id: string, patch: Record<string, unknown>) => {
updates.push({ id, patch }); updates.push({ id, patch });
}), }),
@@ -187,4 +188,150 @@ describe("processPullRequestMergeTask", () => {
expect(github.createPr).not.toHaveBeenCalled(); expect(github.createPr).not.toHaveBeenCalled();
}); });
describe("requirePrApproval", () => {
function makeReadyMergeStatus(reviewDecision: string | null) {
const prInfo = {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "open" as const,
headBranch: "fusion/fn-9100",
baseBranch: "main",
};
// Simulate the "free private repo" case: GitHub reports no required
// checks and no blocking review state, so isPrMergeReady returns
// mergeReady: true. Without the gate this would auto-merge.
return {
prInfo,
reviewDecision,
checks: [],
mergeReady: true,
blockingReasons: [],
};
}
it("holds the merge when requirePrApproval is true and reviewDecision is not APPROVED", async () => {
const task: MockTask = {
id: "FN-9100",
title: "test",
description: "desc",
column: "in-review",
prInfo: {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "open",
headBranch: "fusion/fn-9100",
baseBranch: "main",
},
};
const store = makeStore(task, { requirePrApproval: true });
const github = {
findPrForBranch: vi.fn(),
createPr: vi.fn(),
getPrMergeStatus: vi.fn(async () => makeReadyMergeStatus(null)),
mergePr: vi.fn(),
};
const result = await processPullRequestMergeTask(
store as never,
"/repo",
task.id,
github as never,
() => undefined,
);
expect(result).toBe("waiting");
expect(github.mergePr).not.toHaveBeenCalled();
const lastUpdate = (store as { _updates: Array<{ patch: Record<string, unknown> }> })._updates.at(-1);
expect(lastUpdate?.patch).toEqual({ status: "awaiting-pr-checks" });
});
it("merges when requirePrApproval is true and reviewDecision is APPROVED", async () => {
const task: MockTask = {
id: "FN-9101",
title: "test",
description: "desc",
column: "in-review",
prInfo: {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "open",
headBranch: "fusion/fn-9101",
baseBranch: "main",
},
};
const store = makeStore(task, { requirePrApproval: true });
const merged = {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "merged" as const,
headBranch: "fusion/fn-9101",
baseBranch: "main",
};
const github = {
findPrForBranch: vi.fn(),
createPr: vi.fn(),
getPrMergeStatus: vi.fn(async () => makeReadyMergeStatus("APPROVED")),
mergePr: vi.fn(async () => merged),
};
const result = await processPullRequestMergeTask(
store as never,
"/repo",
task.id,
github as never,
() => undefined,
);
expect(result).toBe("merged");
expect(github.mergePr).toHaveBeenCalledWith({ number: 100, method: "squash" });
});
it("preserves existing behavior when requirePrApproval is false", async () => {
const task: MockTask = {
id: "FN-9102",
title: "test",
description: "desc",
column: "in-review",
prInfo: {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "open",
headBranch: "fusion/fn-9102",
baseBranch: "main",
},
};
const store = makeStore(task, { requirePrApproval: false });
const merged = {
number: 100,
url: "https://github.com/x/y/pull/100",
status: "merged" as const,
headBranch: "fusion/fn-9102",
baseBranch: "main",
};
const github = {
findPrForBranch: vi.fn(),
createPr: vi.fn(),
// reviewDecision: null but mergeReady: true — without the gate,
// this should still merge (the buggy default that #21's reviewer
// flagged as too aggressive on free private repos).
getPrMergeStatus: vi.fn(async () => makeReadyMergeStatus(null)),
mergePr: vi.fn(async () => merged),
};
const result = await processPullRequestMergeTask(
store as never,
"/repo",
task.id,
github as never,
() => undefined,
);
expect(result).toBe("merged");
expect(github.mergePr).toHaveBeenCalled();
});
});
}); });

View File

@@ -230,6 +230,18 @@ export async function processPullRequestMergeTask(
return "merged"; return "merged";
} }
// Optional approval gate. GitHub's `required: true` flag for checks only
// flows from branch protection (Pro feature on private repos), so on free
// private repos every fresh PR is "merge ready" and would auto-squash
// immediately. `requirePrApproval` lets users keep PR mode as "open the
// PR, wait for me to approve and merge it" by holding the merge until
// reviewDecision === "APPROVED".
const settings = await store.getSettings();
if (settings.requirePrApproval && mergeStatus.reviewDecision !== "APPROVED") {
await store.updateTask(task.id, { status: "awaiting-pr-checks" });
return "waiting";
}
if (!mergeStatus.mergeReady) { if (!mergeStatus.mergeReady) {
if (mergeStatus.prInfo.status === "open") { if (mergeStatus.prInfo.status === "open") {
await store.updateTask(task.id, { status: "awaiting-pr-checks" }); await store.updateTask(task.id, { status: "awaiting-pr-checks" });

View File

@@ -12,7 +12,6 @@
*/ */
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
import { AgentStore } from "../agent-store.js"; import { AgentStore } from "../agent-store.js";
import { Database } from "../db.js";
import { TaskStore } from "../store.js"; import { TaskStore } from "../store.js";
import { rm } from "node:fs/promises"; import { rm } from "node:fs/promises";
import { join } from "node:path"; import { join } from "node:path";
@@ -104,10 +103,7 @@ describe("AgentStore", () => {
} }
}); });
it("normalizes legacy durable agents to heartbeat enabled once", async () => { it("preserves disabled heartbeat config for durable agents across restart", async () => {
// Migration test: opens a raw Database on disk to seed a meta key,
// then re-opens the AgentStore to assert migration ran. Needs both
// the store and the raw DB to be disk-backed.
store.close(); store.close();
store = new AgentStore({ rootDir }); store = new AgentStore({ rootDir });
await store.init(); await store.init();
@@ -124,20 +120,12 @@ describe("AgentStore", () => {
}, },
}); });
const db = new Database(rootDir);
db.init();
db.prepare(`
INSERT INTO __meta (key, value)
VALUES ('agentHeartbeatDefaultVersion', '0')
ON CONFLICT(key) DO UPDATE SET value = '0'
`).run();
store.close(); store.close();
store = new AgentStore({ rootDir }); store = new AgentStore({ rootDir });
await store.init(); await store.init();
const migrated = await store.getAgent(agent.id); const persisted = await store.getAgent(agent.id);
expect((migrated?.runtimeConfig as Record<string, unknown> | undefined)?.enabled).toBe(true); expect((persisted?.runtimeConfig as Record<string, unknown> | undefined)?.enabled).toBe(false);
}); });
}); });

View File

@@ -213,7 +213,6 @@ export class AgentStore extends EventEmitter {
const _ = this.db; const _ = this.db;
await mkdir(this.agentsDir, { recursive: true }); await mkdir(this.agentsDir, { recursive: true });
await this.importLegacyFileDataOnce(); await this.importLegacyFileDataOnce();
await this.normalizeHeartbeatDefaultsOnce();
} }
/** /**
@@ -345,60 +344,6 @@ export class AgentStore extends EventEmitter {
this.db.bumpLastModified(); this.db.bumpLastModified();
} }
/**
* One-time normalization for durable agents created before the heartbeat
* toggle was exposed in the UI. Those agents could persist
* `runtimeConfig.enabled = false` even though users had no supported way to
* manage that flag, which caused timers to stay disabled after restart.
*
* We normalize only once per project. After this migration lands, explicit
* user choices are preserved because the version gate prevents reruns.
*/
private async normalizeHeartbeatDefaultsOnce(): Promise<void> {
const migrationKey = "agentHeartbeatDefaultVersion";
const migrationVersion = "1";
const row = this.db.prepare("SELECT value FROM __meta WHERE key = ?").get(migrationKey) as
| { value: string }
| undefined;
if (row?.value === migrationVersion) {
return;
}
const agents = await this.listAgents({ includeEphemeral: true });
let changed = 0;
for (const agent of agents) {
if (isEphemeralAgent(agent)) {
continue;
}
const nextRuntimeConfig = {
...(resolveCreationRuntimeConfig(agent.runtimeConfig, agent.metadata) ?? {}),
enabled: true,
};
const currentRuntimeConfig = agent.runtimeConfig ?? undefined;
if (JSON.stringify(nextRuntimeConfig) === JSON.stringify(currentRuntimeConfig)) {
continue;
}
await this.writeAgent({
...agent,
runtimeConfig: nextRuntimeConfig,
});
changed++;
}
this.db.prepare(`
INSERT INTO __meta (key, value)
VALUES (?, ?)
ON CONFLICT(key) DO UPDATE SET value = excluded.value
`).run(migrationKey, migrationVersion);
if (changed > 0) {
this.db.bumpLastModified();
}
}
/** /**
* Create a new agent with "idle" state. * Create a new agent with "idle" state.
* *

View File

@@ -80,6 +80,7 @@ export const DEFAULT_PROJECT_SETTINGS = {
overlapIgnorePaths: [], overlapIgnorePaths: [],
autoMerge: true, autoMerge: true,
mergeStrategy: "direct", mergeStrategy: "direct",
requirePrApproval: false,
pushAfterMerge: false, pushAfterMerge: false,
pushRemote: "origin", pushRemote: "origin",
unavailableNodePolicy: "block", unavailableNodePolicy: "block",

View File

@@ -1471,6 +1471,12 @@ export interface ProjectSettings {
* before merging through GitHub * before merging through GitHub
* Default: "direct" for backward compatibility. */ * Default: "direct" for backward compatibility. */
mergeStrategy?: MergeStrategy; mergeStrategy?: MergeStrategy;
/** When true, only auto-merge a pull request after it has at least one approving
* review (`reviewDecision === "APPROVED"`). Independent of GitHub's branch-protection
* `required` flag, so this works on free private repos where required reviewers can't
* be enforced server-side. Only applies when `mergeStrategy === "pull-request"`.
* Default: false. */
requirePrApproval?: boolean;
/** When true, automatically push to the configured remote after a successful direct merge. /** When true, automatically push to the configured remote after a successful direct merge.
* The push process includes pulling the latest from the remote (rebase) first. * The push process includes pulling the latest from the remote (rebase) first.
* If conflicts arise during the pull, they are resolved using the AI conflict resolution pipeline. * If conflicts arise during the pull, they are resolved using the AI conflict resolution pipeline.

View File

@@ -462,6 +462,7 @@ export interface RemoteStatus {
url: string | null; url: string | null;
lastError: string | null; lastError: string | null;
lastErrorCode?: string | null; lastErrorCode?: string | null;
cloudflaredAvailable?: boolean | null;
restore?: { restore?: {
outcome: "applied" | "skipped" | "failed"; outcome: "applied" | "skipped" | "failed";
reason: string; reason: string;
@@ -489,6 +490,12 @@ export function fetchRemoteStatus(projectId?: string): Promise<RemoteStatus> {
return api<RemoteStatus>(withProjectId("/remote/status", projectId)); return api<RemoteStatus>(withProjectId("/remote/status", projectId));
} }
export function installCloudflared(projectId?: string): Promise<{ success: boolean; command: string; error?: string }> {
return api(withProjectId("/remote/install-cloudflared", projectId), {
method: "POST",
});
}
export function activateRemoteProvider(provider: "tailscale" | "cloudflare", projectId?: string): Promise<{ activeProvider: "tailscale" | "cloudflare" }> { export function activateRemoteProvider(provider: "tailscale" | "cloudflare", projectId?: string): Promise<{ activeProvider: "tailscale" | "cloudflare" }> {
return api<{ activeProvider: "tailscale" | "cloudflare" }>(withProjectId("/remote/provider/activate", projectId), { return api<{ activeProvider: "tailscale" | "cloudflare" }>(withProjectId("/remote/provider/activate", projectId), {
method: "POST", method: "POST",

View File

@@ -615,6 +615,44 @@
gap: var(--space-sm); gap: var(--space-sm);
} }
.remote-cli-detection {
display: flex;
align-items: flex-start;
gap: var(--space-sm);
margin: 0 var(--space-xl) var(--space-md);
padding: var(--space-sm) var(--space-md);
border-radius: var(--radius-sm);
}
.remote-cli-detection--available {
background: color-mix(in srgb, var(--color-success) 10%, transparent);
color: var(--color-success);
}
.remote-cli-detection--missing {
background: color-mix(in srgb, var(--color-warning) 10%, transparent);
color: var(--color-warning);
}
.remote-cli-detection-content {
display: flex;
flex-direction: column;
gap: var(--space-xs);
color: var(--text);
}
.remote-cli-detection .btn {
align-self: flex-start;
}
.remote-cli-install-error {
color: var(--color-error);
}
.remote-cli-manual {
color: var(--text-dim);
}
.remote-advanced-details { .remote-advanced-details {
margin: var(--space-lg) var(--space-xl) 0; margin: var(--space-lg) var(--space-xl) 0;
border: var(--btn-border-width) solid var(--border); border: var(--btn-border-width) solid var(--border);
@@ -791,6 +829,10 @@
.remote-tunnel-actions .btn { .remote-tunnel-actions .btn {
min-height: 36px; min-height: 36px;
} }
.remote-cli-detection {
margin: 0 var(--space-lg) var(--space-md);
}
} }
/* === Notifications Settings === */ /* === Notifications Settings === */

View File

@@ -1,5 +1,5 @@
import { useState, useEffect, useCallback, useRef, lazy, Suspense, type MouseEvent } from "react"; import { useState, useEffect, useCallback, useRef, lazy, Suspense, type MouseEvent } from "react";
import { Globe, Folder, RefreshCw, Star, HelpCircle, Loader2 } from "lucide-react"; import { Globe, Folder, RefreshCw, Star, HelpCircle, Loader2, CheckCircle, AlertTriangle } from "lucide-react";
import { import {
THINKING_LEVELS, THINKING_LEVELS,
getErrorMessage, getErrorMessage,
@@ -10,7 +10,7 @@ import {
resolveTitleSummarizerSettingsModel, resolveTitleSummarizerSettingsModel,
} from "@fusion/core"; } from "@fusion/core";
import type { Settings, GlobalSettings, ThemeMode, ColorTheme, ModelPreset, NtfyNotificationEvent, AgentPromptsConfig, ThinkingLevel } from "@fusion/core"; import type { Settings, GlobalSettings, ThemeMode, ColorTheme, ModelPreset, NtfyNotificationEvent, AgentPromptsConfig, ThinkingLevel } from "@fusion/core";
import { fetchSettings, fetchSettingsByScope, updateSettings, updateGlobalSettings, fetchAuthStatus, loginProvider, logoutProvider, saveApiKey, clearApiKey, fetchModels, testNotification, fetchBackups, createBackup, exportSettings, importSettings, fetchMemoryFile, fetchMemoryFiles, saveMemoryFile, compactMemory, fetchGlobalConcurrency, updateGlobalConcurrency, installQmd, testMemoryRetrieval, triggerMemoryDreams, fetchGitRemotesDetailed, fetchDashboardHealth, checkForUpdates, fetchRemoteSettings, updateRemoteSettings, fetchRemoteStatus, startRemoteTunnel, stopRemoteTunnel, regenerateRemotePersistentToken, generateShortLivedRemoteToken, fetchRemoteQr, fetchRemoteUrl } from "../api"; import { fetchSettings, fetchSettingsByScope, updateSettings, updateGlobalSettings, fetchAuthStatus, loginProvider, logoutProvider, saveApiKey, clearApiKey, fetchModels, testNotification, fetchBackups, createBackup, exportSettings, importSettings, fetchMemoryFile, fetchMemoryFiles, saveMemoryFile, compactMemory, fetchGlobalConcurrency, updateGlobalConcurrency, installQmd, testMemoryRetrieval, triggerMemoryDreams, fetchGitRemotesDetailed, fetchDashboardHealth, checkForUpdates, fetchRemoteSettings, updateRemoteSettings, fetchRemoteStatus, installCloudflared, startRemoteTunnel, stopRemoteTunnel, regenerateRemotePersistentToken, generateShortLivedRemoteToken, fetchRemoteQr, fetchRemoteUrl } from "../api";
import type { AuthProvider, ModelInfo, BackupListResponse, SettingsExportData, MemoryFileInfo, MemoryRetrievalTestResult, GitRemoteDetailed, RemoteSettings, RemoteStatus, UpdateCheckResponse } from "../api"; import type { AuthProvider, ModelInfo, BackupListResponse, SettingsExportData, MemoryFileInfo, MemoryRetrievalTestResult, GitRemoteDetailed, RemoteSettings, RemoteStatus, UpdateCheckResponse } from "../api";
import { useMemoryBackendStatus } from "../hooks/useMemoryBackendStatus"; import { useMemoryBackendStatus } from "../hooks/useMemoryBackendStatus";
import { useOverlayDismiss } from "../hooks/useOverlayDismiss"; import { useOverlayDismiss } from "../hooks/useOverlayDismiss";
@@ -424,6 +424,8 @@ export function SettingsModal({
// Remote access state // Remote access state
const [remoteStatus, setRemoteStatus] = useState<RemoteStatus | null>(null); const [remoteStatus, setRemoteStatus] = useState<RemoteStatus | null>(null);
const [remoteBusyAction, setRemoteBusyAction] = useState<string | null>(null); const [remoteBusyAction, setRemoteBusyAction] = useState<string | null>(null);
const [cloudflaredInstalling, setCloudflaredInstalling] = useState(false);
const [cloudflaredInstallError, setCloudflaredInstallError] = useState<string | null>(null);
const [remoteAuthLinkTokenType, setRemoteAuthLinkTokenType] = useState<"persistent" | "short-lived">("persistent"); const [remoteAuthLinkTokenType, setRemoteAuthLinkTokenType] = useState<"persistent" | "short-lived">("persistent");
const [remoteUrlPreview, setRemoteUrlPreview] = useState<{ url: string; expiresAt: string | null; tokenType: "persistent" | "short-lived" } | null>(null); const [remoteUrlPreview, setRemoteUrlPreview] = useState<{ url: string; expiresAt: string | null; tokenType: "persistent" | "short-lived" } | null>(null);
const [remoteQrSvg, setRemoteQrSvg] = useState<string | null>(null); const [remoteQrSvg, setRemoteQrSvg] = useState<string | null>(null);
@@ -1538,6 +1540,35 @@ export function SettingsModal({
} }
}, [addToast, loadRemoteData]); }, [addToast, loadRemoteData]);
const cloudflaredManualInstallCommand = useCallback(() => {
if (typeof navigator !== "undefined" && navigator.userAgent.includes("Windows")) {
return "winget install Cloudflare.cloudflared";
}
if (typeof navigator !== "undefined" && /(Mac|iPhone|iPad|iPod)/i.test(navigator.platform)) {
return "brew install cloudflared";
}
return "curl -L --output /usr/local/bin/cloudflared https://github.com/cloudflare/cloudflared/releases/latest/download/cloudflared-linux-amd64 && chmod +x /usr/local/bin/cloudflared";
}, []);
const handleInstallCloudflared = useCallback(async () => {
setCloudflaredInstalling(true);
setCloudflaredInstallError(null);
try {
const result = await installCloudflared(projectId);
if (!result.success) {
setCloudflaredInstallError(result.error ?? "Installation failed");
return;
}
const status = await fetchRemoteStatus(projectId);
setRemoteStatus(status);
addToast("cloudflared installed successfully", "success");
} catch (err) {
setCloudflaredInstallError(err instanceof Error ? err.message : "Installation failed");
} finally {
setCloudflaredInstalling(false);
}
}, [addToast, projectId]);
/** Render a scope indicator banner for the current section with theme-aware Lucide icons */ /** Render a scope indicator banner for the current section with theme-aware Lucide icons */
const renderScopeBanner = () => { const renderScopeBanner = () => {
if (activeSectionScope === "global") { if (activeSectionScope === "global") {
@@ -2972,6 +3003,27 @@ export function SettingsModal({
</small> </small>
</details> </details>
</div> </div>
{form.mergeStrategy === "pull-request" && (
<div className="form-group">
<label htmlFor="requirePrApproval" className="checkbox-label">
<input
id="requirePrApproval"
type="checkbox"
checked={form.requirePrApproval ?? false}
onChange={(e) =>
setForm((f) => ({ ...f, requirePrApproval: e.target.checked }))
}
/>
Wait for an approving review before merging the PR
</label>
<details className="settings-option-details">
<summary>More details</summary>
<small>
When enabled, Fusion holds the PR in In Review until at least one approving GitHub review has been submitted. Useful on free private repos where GitHub&apos;s required-reviewer enforcement isn&apos;t available — without this, a fresh PR with no required checks is treated as immediately mergeable.
</small>
</details>
</div>
)}
<div className="form-group"> <div className="form-group">
<label htmlFor="includeTaskIdInCommit" className="checkbox-label"> <label htmlFor="includeTaskIdInCommit" className="checkbox-label">
<input <input
@@ -4050,6 +4102,32 @@ export function SettingsModal({
{!activeProvider && <small>Select a provider above to configure remote access.</small>} {!activeProvider && <small>Select a provider above to configure remote access.</small>}
</div> </div>
{activeProvider === "cloudflare" && remoteStatus?.cloudflaredAvailable === true && (
<div className="remote-cli-detection remote-cli-detection--available" role="status">
<CheckCircle aria-hidden="true" />
<span>cloudflared is installed</span>
</div>
)}
{activeProvider === "cloudflare" && remoteStatus?.cloudflaredAvailable === false && (
<div className="remote-cli-detection remote-cli-detection--missing" role="status">
<AlertTriangle aria-hidden="true" />
<div className="remote-cli-detection-content">
<span>cloudflared is not installed</span>
<button
type="button"
className="btn btn-sm"
disabled={cloudflaredInstalling || remoteBusyAction !== null}
onClick={() => void handleInstallCloudflared()}
>
{cloudflaredInstalling ? "Installing…" : "Install cloudflared"}
</button>
{cloudflaredInstallError && <small className="remote-cli-install-error">{cloudflaredInstallError}</small>}
<small className="remote-cli-manual">Manual install: <code>{cloudflaredManualInstallCommand()}</code></small>
</div>
</div>
)}
{activeProvider && ( {activeProvider && (
<div className="form-group remote-provider-settings"> <div className="form-group remote-provider-settings">
{activeProvider === "tailscale" ? ( {activeProvider === "tailscale" ? (
@@ -4108,32 +4186,37 @@ export function SettingsModal({
{remoteBusyAction === "stop" ? "Stopping…" : "Stop Tunnel"} {remoteBusyAction === "stop" ? "Stopping…" : "Stop Tunnel"}
</button> </button>
) : ( ) : (
<button type="button" className="btn btn-primary" disabled={!activeProvider || remoteBusyAction !== null} onClick={() => void runRemoteAction("start", async () => { <>
const formState = form as Record<string, unknown>; <button type="button" className="btn btn-primary" disabled={!activeProvider || remoteBusyAction !== null} onClick={() => void runRemoteAction("start", async () => {
const savePayload: Partial<RemoteSettings> = { const formState = form as Record<string, unknown>;
remoteActiveProvider: activeProvider, const savePayload: Partial<RemoteSettings> = {
remoteTailscaleEnabled: activeProvider === "tailscale", remoteActiveProvider: activeProvider,
remoteTailscaleHostname: String(formState.remoteTailscaleHostname ?? ""), remoteTailscaleEnabled: activeProvider === "tailscale",
// Server overrides this with req.socket.localPort remoteTailscaleHostname: String(formState.remoteTailscaleHostname ?? ""),
// when starting the tunnel; the value sent here is // Server overrides this with req.socket.localPort
// only a fallback if that override doesn't fire. // when starting the tunnel; the value sent here is
remoteTailscaleTargetPort: Number(formState.remoteTailscaleTargetPort ?? 4040), // only a fallback if that override doesn't fire.
remoteTailscaleAcceptRoutes: Boolean(formState.remoteTailscaleAcceptRoutes), remoteTailscaleTargetPort: Number(formState.remoteTailscaleTargetPort ?? 4040),
remoteCloudflareEnabled: activeProvider === "cloudflare", remoteTailscaleAcceptRoutes: Boolean(formState.remoteTailscaleAcceptRoutes),
remoteCloudflareQuickTunnel: Boolean(formState.remoteCloudflareQuickTunnel ?? true), remoteCloudflareEnabled: activeProvider === "cloudflare",
remoteCloudflareTunnelName: String(formState.remoteCloudflareTunnelName ?? ""), remoteCloudflareQuickTunnel: Boolean(formState.remoteCloudflareQuickTunnel ?? true),
remoteCloudflareTunnelToken: (formState.remoteCloudflareTunnelToken as string | null) || null, remoteCloudflareTunnelName: String(formState.remoteCloudflareTunnelName ?? ""),
remoteCloudflareIngressUrl: String(formState.remoteCloudflareIngressUrl ?? ""), remoteCloudflareTunnelToken: (formState.remoteCloudflareTunnelToken as string | null) || null,
remoteShortLivedEnabled: Boolean(formState.remoteShortLivedEnabled), remoteCloudflareIngressUrl: String(formState.remoteCloudflareIngressUrl ?? ""),
remoteShortLivedTtlMs: Number(formState.remoteShortLivedTtlMs ?? 900000), remoteShortLivedEnabled: Boolean(formState.remoteShortLivedEnabled),
remoteRememberLastRunning: Boolean(formState.remoteRememberLastRunning), remoteShortLivedTtlMs: Number(formState.remoteShortLivedTtlMs ?? 900000),
}; remoteRememberLastRunning: Boolean(formState.remoteRememberLastRunning),
await updateRemoteSettings(savePayload, projectId); };
await startRemoteTunnel(projectId); await updateRemoteSettings(savePayload, projectId);
addToast("Remote tunnel started", "success"); await startRemoteTunnel(projectId);
})}> addToast("Remote tunnel started", "success");
{remoteBusyAction === "start" ? "Starting…" : "Start Tunnel"} })}>
</button> {remoteBusyAction === "start" ? "Starting…" : "Start Tunnel"}
</button>
{activeProvider === "cloudflare" && remoteStatus?.cloudflaredAvailable === false ? (
<small className="field-error">cloudflared must be installed to start the tunnel</small>
) : null}
</>
)} )}
</div> </div>

View File

@@ -34,6 +34,7 @@ const mockCheckForUpdates = vi.fn();
const mockFetchRemoteSettings = vi.fn(); const mockFetchRemoteSettings = vi.fn();
const mockUpdateRemoteSettings = vi.fn(); const mockUpdateRemoteSettings = vi.fn();
const mockFetchRemoteStatus = vi.fn(); const mockFetchRemoteStatus = vi.fn();
const mockInstallCloudflared = vi.fn();
const mockStartRemoteTunnel = vi.fn(); const mockStartRemoteTunnel = vi.fn();
const mockStopRemoteTunnel = vi.fn(); const mockStopRemoteTunnel = vi.fn();
const mockRegenerateRemotePersistentToken = vi.fn(); const mockRegenerateRemotePersistentToken = vi.fn();
@@ -73,6 +74,7 @@ vi.mock("../../api", () => ({
fetchRemoteSettings: (...args: unknown[]) => mockFetchRemoteSettings(...args), fetchRemoteSettings: (...args: unknown[]) => mockFetchRemoteSettings(...args),
updateRemoteSettings: (...args: unknown[]) => mockUpdateRemoteSettings(...args), updateRemoteSettings: (...args: unknown[]) => mockUpdateRemoteSettings(...args),
fetchRemoteStatus: (...args: unknown[]) => mockFetchRemoteStatus(...args), fetchRemoteStatus: (...args: unknown[]) => mockFetchRemoteStatus(...args),
installCloudflared: (...args: unknown[]) => mockInstallCloudflared(...args),
startRemoteTunnel: (...args: unknown[]) => mockStartRemoteTunnel(...args), startRemoteTunnel: (...args: unknown[]) => mockStartRemoteTunnel(...args),
stopRemoteTunnel: (...args: unknown[]) => mockStopRemoteTunnel(...args), stopRemoteTunnel: (...args: unknown[]) => mockStopRemoteTunnel(...args),
regenerateRemotePersistentToken: (...args: unknown[]) => mockRegenerateRemotePersistentToken(...args), regenerateRemotePersistentToken: (...args: unknown[]) => mockRegenerateRemotePersistentToken(...args),
@@ -282,6 +284,7 @@ describe("SettingsModal", () => {
}, },
}); });
mockFetchRemoteStatus.mockResolvedValue({ provider: null, state: "stopped", url: null, lastError: null }); mockFetchRemoteStatus.mockResolvedValue({ provider: null, state: "stopped", url: null, lastError: null });
mockInstallCloudflared.mockResolvedValue({ success: true, command: "brew install cloudflared" });
mockStartRemoteTunnel.mockResolvedValue({ state: "starting", provider: "tailscale" }); mockStartRemoteTunnel.mockResolvedValue({ state: "starting", provider: "tailscale" });
mockStopRemoteTunnel.mockResolvedValue({ state: "stopped", provider: null }); mockStopRemoteTunnel.mockResolvedValue({ state: "stopped", provider: null });
mockRegenerateRemotePersistentToken.mockResolvedValue({ token: "token", maskedToken: "****" }); mockRegenerateRemotePersistentToken.mockResolvedValue({ token: "token", maskedToken: "****" });
@@ -1695,6 +1698,69 @@ describe("SettingsModal", () => {
}); });
}); });
it("shows cloudflared available indicator when Cloudflare is selected and cloudflared is installed", async () => {
mockFetchRemoteStatus.mockResolvedValue({ provider: "cloudflare", state: "stopped", url: null, lastError: null, cloudflaredAvailable: true });
renderModal();
await waitForSettingsModalReady();
await openRemoteSection();
await userEvent.click(screen.getByLabelText("Cloudflare"));
expect(await screen.findByText("cloudflared is installed")).toBeInTheDocument();
expect(screen.queryByRole("button", { name: "Install cloudflared" })).not.toBeInTheDocument();
});
it("shows install button when Cloudflare is selected and cloudflared is not available", async () => {
mockFetchRemoteStatus.mockResolvedValue({ provider: "cloudflare", state: "stopped", url: null, lastError: null, cloudflaredAvailable: false });
renderModal();
await waitForSettingsModalReady();
await openRemoteSection();
await userEvent.click(screen.getByLabelText("Cloudflare"));
expect(await screen.findByText("cloudflared is not installed")).toBeInTheDocument();
expect(screen.getByRole("button", { name: "Install cloudflared" })).toBeInTheDocument();
expect(screen.getByText("cloudflared must be installed to start the tunnel")).toBeInTheDocument();
});
it("install button triggers install and refreshes status", async () => {
const addToast = vi.fn();
mockFetchRemoteStatus
.mockResolvedValueOnce({ provider: "cloudflare", state: "stopped", url: null, lastError: null, cloudflaredAvailable: false })
.mockResolvedValueOnce({ provider: "cloudflare", state: "stopped", url: null, lastError: null, cloudflaredAvailable: true });
mockInstallCloudflared.mockResolvedValueOnce({ success: true, command: "brew install cloudflared" });
renderModal({ addToast });
await waitForSettingsModalReady();
await openRemoteSection();
await userEvent.click(screen.getByLabelText("Cloudflare"));
const installButton = await screen.findByRole("button", { name: "Install cloudflared" });
await userEvent.click(installButton);
await waitFor(() => {
expect(mockInstallCloudflared).toHaveBeenCalledWith(undefined);
});
await waitFor(() => {
expect(addToast).toHaveBeenCalledWith("cloudflared installed successfully", "success");
});
expect(await screen.findByText("cloudflared is installed")).toBeInTheDocument();
});
it("install button shows error on failure", async () => {
mockFetchRemoteStatus.mockResolvedValue({ provider: "cloudflare", state: "stopped", url: null, lastError: null, cloudflaredAvailable: false });
mockInstallCloudflared.mockResolvedValueOnce({ success: false, command: "brew install cloudflared", error: "Command failed" });
renderModal();
await waitForSettingsModalReady();
await openRemoteSection();
await userEvent.click(screen.getByLabelText("Cloudflare"));
await userEvent.click(await screen.findByRole("button", { name: "Install cloudflared" }));
expect(await screen.findByText("Command failed")).toBeInTheDocument();
});
it("shows lifecycle state changes for start and stop actions, including error state", async () => { it("shows lifecycle state changes for start and stop actions, including error state", async () => {
mockFetchRemoteStatus mockFetchRemoteStatus
.mockResolvedValueOnce({ provider: null, state: "stopped", url: null, lastError: null }) .mockResolvedValueOnce({ provider: null, state: "stopped", url: null, lastError: null })

View File

@@ -132,10 +132,9 @@ describe("getAgentHealthStatus", () => {
}); });
}); });
// Heartbeat scheduling is driven by agent.state on the server; there is no // Heartbeat disabled is a real durable-agent state in the UI. Task workers
// separate "disabled" UI concept anymore. Non-task-worker agents with a // still follow execution-state health because their runtimeConfig.enabled
// legacy `runtimeConfig.enabled === false` on disk are rendered by state // flag only opts them out of scheduler timers.
// just like any other agent.
describe("task worker health classification", () => { describe("task worker health classification", () => {
it('returns "Running" for metadata-marked task workers with disabled heartbeat', () => { it('returns "Running" for metadata-marked task workers with disabled heartbeat', () => {
@@ -173,7 +172,7 @@ describe("getAgentHealthStatus", () => {
expect(status.color).toBe("var(--state-active-text)"); expect(status.color).toBe("var(--state-active-text)");
}); });
it('ignores legacy runtimeConfig.enabled=false on non-task-worker agents', () => { it('returns "Heartbeat Disabled" for non-task-worker agents with heartbeat disabled', () => {
const agent = makeAgent({ const agent = makeAgent({
name: "Reviewer", name: "Reviewer",
role: "reviewer", role: "reviewer",
@@ -181,8 +180,21 @@ describe("getAgentHealthStatus", () => {
runtimeConfig: { enabled: false }, runtimeConfig: { enabled: false },
}); });
const status = getAgentHealthStatus(agent); const status = getAgentHealthStatus(agent);
// No persisted heartbeat, no lastHeartbeatAt → Starting... not Disabled. expect(status.label).toBe("Heartbeat Disabled");
expect(status.label).toBe("Starting..."); expect(status.stateDerived).toBe(false);
expect(status.color).toBe("var(--state-paused-text)");
});
it('returns "Heartbeat Disabled" even when a disabled durable agent has a recent heartbeat', () => {
const agent = makeAgent({
name: "Reviewer",
role: "reviewer",
state: "active",
lastHeartbeatAt: new Date(FIXED_NOW - 1_000).toISOString(),
runtimeConfig: { enabled: false, heartbeatIntervalMs: 60_000 },
});
const status = getAgentHealthStatus(agent);
expect(status.label).toBe("Heartbeat Disabled");
}); });
}); });
@@ -436,9 +448,9 @@ describe("getAgentHealthStatus", () => {
expect(getAgentHealthStatus(agent).label).toBe("Healthy"); expect(getAgentHealthStatus(agent).label).toBe("Healthy");
}); });
it("ignores runtimeConfig.enabled and uses interval-based staleness", () => { it("uses interval-based staleness for enabled durable agents", () => {
// 6 minute interval → 24 minute threshold. 25 minutes elapsed is stale regardless of any // 6 minute interval → 24 minute threshold. 25 minutes elapsed is stale
// legacy enabled flag or per-run timeout. // regardless of any per-run timeout.
const agent = makeAgent({ const agent = makeAgent({
state: "active", state: "active",
lastHeartbeatAt: new Date(FIXED_NOW - 25 * 60 * 1000).toISOString(), // 25 minutes ago lastHeartbeatAt: new Date(FIXED_NOW - 25 * 60 * 1000).toISOString(), // 25 minutes ago
@@ -454,8 +466,7 @@ describe("getAgentHealthStatus", () => {
{ agent: makeAgent({ state: "paused" }), expectedIconType: "Pause" }, { agent: makeAgent({ state: "paused" }), expectedIconType: "Pause" },
{ agent: makeAgent({ state: "running" }), expectedIconType: "Activity" }, { agent: makeAgent({ state: "running" }), expectedIconType: "Activity" },
{ agent: makeAgent({ state: "idle" }), expectedIconType: "Bot" }, { agent: makeAgent({ state: "idle" }), expectedIconType: "Bot" },
// state=active + no lastHeartbeatAt → "Starting..." → Bot icon { agent: makeAgent({ state: "active", runtimeConfig: { enabled: false } }), expectedIconType: "Pause" },
{ agent: makeAgent({ state: "active", runtimeConfig: { enabled: false } }), expectedIconType: "Bot" },
{ {
agent: makeAgent({ agent: makeAgent({
name: "executor-FN-1661", name: "executor-FN-1661",

View File

@@ -3,11 +3,10 @@ import { Bot, Heart, Activity, Pause, Square } from "lucide-react";
import type { Agent } from "../api"; import type { Agent } from "../api";
import { resolveHeartbeatIntervalMs } from "./heartbeatIntervals"; import { resolveHeartbeatIntervalMs } from "./heartbeatIntervals";
// Heartbeat scheduling is driven by `agent.state` on the server — active and // Heartbeat scheduling depends on both state and `runtimeConfig.enabled`.
// running tick, everything else does not. There is no separate "heartbeat // Durable agents with heartbeat disabled should render distinctly from healthy
// enabled" flag surfaced in the UI, so this file derives freshness straight // or merely-starting agents, while task-worker agents still follow their
// from state + lastHeartbeatAt and ignores any legacy `runtimeConfig.enabled` // execution lifecycle regardless of the scheduler toggle.
// value that may still be persisted on older agent records.
/** /**
* Grace multiplier applied to an agent's configured interval before flagging * Grace multiplier applied to an agent's configured interval before flagging
@@ -86,6 +85,7 @@ function isTaskWorkerAgent(agent: AgentHealthInput): boolean {
* - "Error" — agent.state === "error" (uses lastError if available) * - "Error" — agent.state === "error" (uses lastError if available)
* - "Paused" — agent.state === "paused" (uses pauseReason if available) * - "Paused" — agent.state === "paused" (uses pauseReason if available)
* - "Running" — agent.state === "running", or a detected task worker in "active" * - "Running" — agent.state === "running", or a detected task worker in "active"
* - "Heartbeat Disabled" — durable agent with `runtimeConfig.enabled === false`
* - "Starting..." — state === "active" && no lastHeartbeatAt * - "Starting..." — state === "active" && no lastHeartbeatAt
* - "Idle" — state !== "active" && no lastHeartbeatAt * - "Idle" — state !== "active" && no lastHeartbeatAt
* - "Healthy" — heartbeat is fresh within 2× the configured interval * - "Healthy" — heartbeat is fresh within 2× the configured interval
@@ -97,6 +97,7 @@ function isTaskWorkerAgent(agent: AgentHealthInput): boolean {
export function getAgentHealthStatus(agent: AgentHealthInput): AgentHealthStatus { export function getAgentHealthStatus(agent: AgentHealthInput): AgentHealthStatus {
const { state, lastHeartbeatAt, lastError, pauseReason, runtimeConfig } = agent; const { state, lastHeartbeatAt, lastError, pauseReason, runtimeConfig } = agent;
const isTaskWorker = isTaskWorkerAgent(agent); const isTaskWorker = isTaskWorkerAgent(agent);
const isHeartbeatEnabled = isTaskWorker || runtimeConfig?.enabled !== false;
// Terminal states - these always take precedence // Terminal states - these always take precedence
if (state === "terminated") { if (state === "terminated") {
@@ -136,6 +137,15 @@ export function getAgentHealthStatus(agent: AgentHealthInput): AgentHealthStatus
}; };
} }
if (!isHeartbeatEnabled) {
return {
label: "Heartbeat Disabled",
icon: <Pause size={14} />,
color: "var(--state-paused-text)",
stateDerived: false,
};
}
// No heartbeat data yet // No heartbeat data yet
if (!lastHeartbeatAt) { if (!lastHeartbeatAt) {
return { return {

View File

@@ -1,8 +1,20 @@
// @vitest-environment node // @vitest-environment node
import { describe, expect, it, vi } from "vitest"; import { beforeEach, describe, expect, it, vi } from "vitest";
import express from "express"; import express from "express";
import type { TaskStore } from "@fusion/core"; import type { TaskStore } from "@fusion/core";
const { mockExecFile } = vi.hoisted(() => ({
mockExecFile: vi.fn(),
}));
vi.mock("node:child_process", async (importOriginal) => {
const actual = await importOriginal<typeof import("node:child_process")>();
return {
...actual,
execFile: mockExecFile,
};
});
import { createApiRoutes } from "../routes.js"; import { createApiRoutes } from "../routes.js";
import { request as performRequest } from "../test-request.js"; import { request as performRequest } from "../test-request.js";
@@ -84,6 +96,16 @@ async function REQUEST(app: express.Express, method: string, path: string, body?
); );
} }
beforeEach(() => {
mockExecFile.mockReset();
mockExecFile.mockImplementation((command: string, _args: string[], optionsOrCallback: unknown, maybeCallback?: (error: Error | null, stdout?: string, stderr?: string) => void) => {
const callback = typeof optionsOrCallback === "function"
? optionsOrCallback as (error: Error | null, stdout?: string, stderr?: string) => void
: maybeCallback;
callback?.(null, command === "where" || command === "which" ? "/usr/local/bin/cloudflared" : "", "");
});
});
describe("remote access provider/lifecycle contracts", () => { describe("remote access provider/lifecycle contracts", () => {
it("switches active provider and rejects invalid provider values", async () => { it("switches active provider and rejects invalid provider values", async () => {
const updateSettings = vi.fn().mockResolvedValue(undefined); const updateSettings = vi.fn().mockResolvedValue(undefined);
@@ -180,4 +202,87 @@ describe("remote access provider/lifecycle contracts", () => {
details: { code: "REMOTE_TUNNEL_PREREQUISITE_MISSING" }, details: { code: "REMOTE_TUNNEL_PREREQUISITE_MISSING" },
}); });
}); });
it("includes cloudflaredAvailable in remote status for cloudflare provider", async () => {
const { app } = createApp();
const status = await REQUEST(app, "GET", "/api/remote/status");
expect(status.status).toBe(200);
expect(status.body).toEqual(expect.objectContaining({
provider: "cloudflare",
cloudflaredAvailable: true,
}));
});
it("returns cloudflaredAvailable false when cloudflared check fails", async () => {
mockExecFile.mockImplementation((command: string, _args: string[], optionsOrCallback: unknown, maybeCallback?: (error: Error | null) => void) => {
const callback = typeof optionsOrCallback === "function"
? optionsOrCallback as (error: Error | null) => void
: maybeCallback;
if (command === "which" || command === "where") {
callback?.(new Error("missing"));
return;
}
callback?.(null);
});
const { app } = createApp();
const status = await REQUEST(app, "GET", "/api/remote/status");
expect(status.status).toBe(200);
expect(status.body).toEqual(expect.objectContaining({ cloudflaredAvailable: false }));
});
it("returns cloudflaredAvailable null for non-cloudflare provider", async () => {
const store = createMockStore({
getSettings: vi.fn().mockResolvedValue({
remoteAccess: buildRemoteAccessSettings({ activeProvider: "tailscale" }),
}),
});
const { app } = createApp({ store });
const status = await REQUEST(app, "GET", "/api/remote/status");
expect(status.status).toBe(200);
expect(status.body).toEqual(expect.objectContaining({
provider: "tailscale",
cloudflaredAvailable: null,
}));
});
it("installs cloudflared via endpoint and returns install command metadata", async () => {
const { app } = createApp();
const result = await REQUEST(app, "POST", "/api/remote/install-cloudflared", {});
expect(result.status).toBe(200);
expect(result.body).toEqual(expect.objectContaining({
success: true,
command: expect.any(String),
}));
});
it("returns install failure details when cloudflared installation command fails", async () => {
mockExecFile.mockImplementation((command: string, _args: string[], optionsOrCallback: unknown, maybeCallback?: (error: Error | null, stdout?: string, stderr?: string) => void) => {
const callback = typeof optionsOrCallback === "function"
? optionsOrCallback as (error: Error | null, stdout?: string, stderr?: string) => void
: maybeCallback;
if (command === "sh" || command === "cmd") {
callback?.(new Error("Command failed"), "", "Command failed");
return;
}
callback?.(null, "/usr/local/bin/cloudflared", "");
});
const { app } = createApp();
const result = await REQUEST(app, "POST", "/api/remote/install-cloudflared", {});
expect(result.status).toBe(200);
expect(result.body).toEqual(expect.objectContaining({
success: false,
command: expect.any(String),
error: expect.stringContaining("Command failed"),
}));
});
}); });

View File

@@ -474,6 +474,28 @@ describe("Agent runs routes (with HeartbeatMonitor)", () => {
}, },
}); });
}); });
it("resuming to active does not auto-trigger heartbeat when disabled", async () => {
mockGetAgent.mockResolvedValue({
id: "agent-001",
state: "paused",
runtimeConfig: { enabled: false },
});
mockUpdateAgentState.mockResolvedValue({ id: "agent-001", state: "active" });
const response = await request(
app,
"POST",
"/api/agents/agent-001/state",
JSON.stringify({ state: "active" }),
{ "content-type": "application/json" },
);
expect(response.status).toBe(200);
expect(response.body).toEqual({ id: "agent-001", state: "active" });
await Promise.resolve();
expect(mockExecuteHeartbeat).not.toHaveBeenCalled();
});
}); });
describe("POST /api/agents/:id/runs", () => { describe("POST /api/agents/:id/runs", () => {

View File

@@ -380,8 +380,8 @@ describe("createServer health and headless mode", () => {
const dashStatusBody = dashStatus.body as Record<string, unknown>; const dashStatusBody = dashStatus.body as Record<string, unknown>;
const headlessStatusBody = headlessStatus.body as Record<string, unknown>; const headlessStatusBody = headlessStatus.body as Record<string, unknown>;
expect(Object.keys(dashStatusBody).sort()).toEqual(["lastError", "lastErrorCode", "provider", "restore", "state", "url"]); expect(Object.keys(dashStatusBody).sort()).toEqual(["cloudflaredAvailable", "lastError", "lastErrorCode", "provider", "restore", "state", "url"]);
expect(Object.keys(headlessStatusBody).sort()).toEqual(["lastError", "lastErrorCode", "provider", "restore", "state", "url"]); expect(Object.keys(headlessStatusBody).sort()).toEqual(["cloudflaredAvailable", "lastError", "lastErrorCode", "provider", "restore", "state", "url"]);
expect(headlessRoot.status).toBe(404); expect(headlessRoot.status).toBe(404);
}); });
}); });

View File

@@ -432,7 +432,8 @@ export function registerAgentRuntimeRoutes(ctx: ApiRoutesContext, deps: AgentRun
} }
} }
if (nextState === "active" && projectHeartbeatMonitor) { const isHeartbeatEnabled = currentAgent.runtimeConfig?.enabled !== false;
if (nextState === "active" && isHeartbeatEnabled && projectHeartbeatMonitor) {
await projectHeartbeatMonitor.executeHeartbeat({ await projectHeartbeatMonitor.executeHeartbeat({
agentId, agentId,
source: "on_demand", source: "on_demand",

View File

@@ -42,6 +42,8 @@ import {
} from "@fusion/core"; } from "@fusion/core";
import { createFnAgent as engineCreateFnAgent } from "@fusion/engine"; import { createFnAgent as engineCreateFnAgent } from "@fusion/engine";
import QRCode from "qrcode"; import QRCode from "qrcode";
import { execFile } from "node:child_process";
import { promisify } from "node:util";
import { ApiError, badRequest } from "../api-error.js"; import { ApiError, badRequest } from "../api-error.js";
import { generateRemoteToken, issueRemoteAuthToken, maskRemoteToken } from "../remote-auth.js"; import { generateRemoteToken, issueRemoteAuthToken, maskRemoteToken } from "../remote-auth.js";
import { invalidateAllGlobalSettingsCaches } from "../project-store-resolver.js"; import { invalidateAllGlobalSettingsCaches } from "../project-store-resolver.js";
@@ -57,6 +59,44 @@ interface SettingsMemoryRouteDeps {
export function registerSettingsMemoryRoutes(ctx: ApiRoutesContext, deps: SettingsMemoryRouteDeps): void { export function registerSettingsMemoryRoutes(ctx: ApiRoutesContext, deps: SettingsMemoryRouteDeps): void {
const { router, options, store, runtimeLogger, getProjectContext, rethrowAsApiError } = ctx; const { router, options, store, runtimeLogger, getProjectContext, rethrowAsApiError } = ctx;
const { githubToken, validateModelPresets, sanitizeOverlapIgnorePaths, discoverDashboardPiExtensions } = deps; const { githubToken, validateModelPresets, sanitizeOverlapIgnorePaths, discoverDashboardPiExtensions } = deps;
const execFileAsync = promisify(execFile);
async function isCloudflaredAvailable(): Promise<boolean> {
const command = process.platform === "win32" ? "where" : "which";
try {
await execFileAsync(command, ["cloudflared"]);
return true;
} catch {
return false;
}
}
function resolveCloudflaredInstallCommand(): string {
if (process.platform === "darwin") {
return "brew install cloudflared";
}
if (process.platform === "win32") {
return "winget install Cloudflare.cloudflared";
}
return "curl -L --output /usr/local/bin/cloudflared https://github.com/cloudflare/cloudflared/releases/latest/download/cloudflared-linux-amd64 && chmod +x /usr/local/bin/cloudflared";
}
async function installCloudflared(): Promise<{ success: boolean; command: string; error?: string }> {
const command = resolveCloudflaredInstallCommand();
const shell = process.platform === "win32" ? "cmd" : "sh";
const shellArgs = process.platform === "win32" ? ["/c", command] : ["-c", command];
try {
await execFileAsync(shell, shellArgs, { timeout: 120_000 });
return { success: true, command };
} catch (error) {
return {
success: false,
command,
error: error instanceof Error ? error.message : String(error),
};
}
}
function resolveRemoteBaseUrl( function resolveRemoteBaseUrl(
remoteAccess: NonNullable<Awaited<ReturnType<typeof store.getSettings>>["remoteAccess"]>, remoteAccess: NonNullable<Awaited<ReturnType<typeof store.getSettings>>["remoteAccess"]>,
@@ -422,12 +462,19 @@ export function registerSettingsMemoryRoutes(ctx: ApiRoutesContext, deps: Settin
const tunnelStatus = manager?.getStatus(); const tunnelStatus = manager?.getStatus();
const restore = engine?.getRemoteTunnelRestoreDiagnostics(); const restore = engine?.getRemoteTunnelRestoreDiagnostics();
const activeProvider = tunnelStatus?.provider ?? settings.remoteAccess?.activeProvider ?? null;
let cloudflaredAvailable: boolean | null = null;
if (activeProvider === "cloudflare") {
cloudflaredAvailable = await isCloudflaredAvailable();
}
res.json({ res.json({
provider: tunnelStatus?.provider ?? settings.remoteAccess?.activeProvider ?? null, provider: activeProvider,
state: tunnelStatus?.state ?? "stopped", state: tunnelStatus?.state ?? "stopped",
url: tunnelStatus?.url ?? null, url: tunnelStatus?.url ?? null,
lastError: tunnelStatus?.lastError?.message ?? null, lastError: tunnelStatus?.lastError?.message ?? null,
lastErrorCode: tunnelStatus?.lastError?.code ?? null, lastErrorCode: tunnelStatus?.lastError?.code ?? null,
cloudflaredAvailable,
restore: restore ?? { restore: restore ?? {
outcome: "skipped", outcome: "skipped",
reason: "not_attempted", reason: "not_attempted",
@@ -441,6 +488,16 @@ export function registerSettingsMemoryRoutes(ctx: ApiRoutesContext, deps: Settin
} }
}); });
router.post("/remote/install-cloudflared", async (_req, res) => {
try {
const result = await installCloudflared();
res.json(result);
} catch (err: unknown) {
if (err instanceof ApiError) throw err;
rethrowAsApiError(err, "Failed to install cloudflared");
}
});
router.post("/remote/provider/activate", async (req, res) => { router.post("/remote/provider/activate", async (req, res) => {
try { try {
const provider = req.body?.provider; const provider = req.body?.provider;

View File

@@ -4217,13 +4217,9 @@ describe("HeartbeatTriggerScheduler", () => {
expect(scheduler.getRegisteredAgents()).toContain("agent-001"); expect(scheduler.getRegisteredAgents()).toContain("agent-001");
}); });
it("registers regardless of the legacy enabled flag (state is the source of truth)", () => { it("does not register when heartbeat is explicitly disabled", () => {
// runtimeConfig.enabled is no longer honored by the scheduler — pause
// and resume happen through agent.state, and the agent:updated listener
// drives register/unregister. Callers that still pass `enabled: false`
// should not silently lose the timer.
scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 10000, enabled: false }); scheduler.registerAgent("agent-001", { heartbeatIntervalMs: 10000, enabled: false });
expect(scheduler.getRegisteredAgents()).toContain("agent-001"); expect(scheduler.getRegisteredAgents()).not.toContain("agent-001");
}); });
it("applies default 3600-second interval when intervalMs is undefined", async () => { it("applies default 3600-second interval when intervalMs is undefined", async () => {

View File

@@ -17,7 +17,7 @@
* - onTerminated: Called when an unresponsive agent is terminated * - onTerminated: Called when an unresponsive agent is terminated
*/ */
import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail, AgentRole, Agent, InboxTask, BlockedStateSnapshot, RunMutationContext, Settings } from "@fusion/core"; import type { AgentStore, AgentHeartbeatRun, HeartbeatInvocationSource, AgentHeartbeatConfig, AgentBudgetStatus, Message, MessageStore, TaskStore, TaskDetail, AgentRole, Agent, InboxTask, BlockedStateSnapshot, RunMutationContext, Settings, AgentConfigRevision } from "@fusion/core";
import { buildExecutionMemoryInstructions, isEphemeralAgent, hasAgentIdentity } from "@fusion/core"; import { buildExecutionMemoryInstructions, isEphemeralAgent, hasAgentIdentity } from "@fusion/core";
import type { ToolDefinition } from "@mariozechner/pi-coding-agent"; import type { ToolDefinition } from "@mariozechner/pi-coding-agent";
import { Type, type Static } from "@mariozechner/pi-ai"; import { Type, type Static } from "@mariozechner/pi-ai";
@@ -1750,11 +1750,14 @@ function isHeartbeatManaged(agent: Agent): boolean {
/** /**
* HeartbeatTriggerScheduler manages timer-based heartbeat triggers for agents. * HeartbeatTriggerScheduler manages timer-based heartbeat triggers for agents.
* *
* State is the source of truth: state ∈ {active, running} on a non-ephemeral * Timers are armed only for durable agents where all of the following hold:
* agent arms the timer; any other state or any ephemeral agent doesn't. The * - `runtimeConfig.enabled !== false`
* `runtimeConfig.enabled` flag is no longer consulted here — pause/resume * - `state ∈ {active, running, idle}`
* happens through `agent.state`, and the `agent:updated` listener arms or *
* clears the timer on transitions. * Any other state, or any ephemeral/task-worker agent, clears the timer.
* State changes and heartbeat config updates are observed via AgentStore
* lifecycle events, while callers can still explicitly register existing
* agents during startup bootstrap.
* *
* Other config knobs still apply: * Other config knobs still apply:
* - `heartbeatIntervalMs`: Timer interval (default 1h) * - `heartbeatIntervalMs`: Timer interval (default 1h)
@@ -1768,7 +1771,9 @@ export class HeartbeatTriggerScheduler {
private registrationEpochs: Map<string, number> = new Map(); private registrationEpochs: Map<string, number> = new Map();
private running = false; private running = false;
private assignedListener: ((agent: import("@fusion/core").Agent, taskId: string) => void) | null = null; private assignedListener: ((agent: import("@fusion/core").Agent, taskId: string) => void) | null = null;
private createdListener: ((agent: import("@fusion/core").Agent) => void) | null = null;
private updatedListener: ((agent: import("@fusion/core").Agent) => void) | null = null; private updatedListener: ((agent: import("@fusion/core").Agent) => void) | null = null;
private configRevisionListener: ((agentId: string, revision: AgentConfigRevision) => void) | null = null;
private deletedListener: ((agentId: string) => void) | null = null; private deletedListener: ((agentId: string) => void) | null = null;
constructor(store: AgentStore, callback: TriggerCallback, taskStore?: TaskStore) { constructor(store: AgentStore, callback: TriggerCallback, taskStore?: TaskStore) {
@@ -1779,7 +1784,7 @@ export class HeartbeatTriggerScheduler {
/** /**
* Start the scheduler. Enables assignment watching. * Start the scheduler. Enables assignment watching.
* Individual agents must be registered separately via registerAgent(). * Existing agents still need one startup bootstrap pass via registerAgent().
*/ */
start(): void { start(): void {
if (this.running) return; if (this.running) return;
@@ -1826,9 +1831,10 @@ export class HeartbeatTriggerScheduler {
* @param config - Per-agent heartbeat config * @param config - Per-agent heartbeat config
*/ */
registerAgent(agentId: string, config: AgentHeartbeatConfig): void { registerAgent(agentId: string, config: AgentHeartbeatConfig): void {
// State drives whether an agent ticks; this method no longer honors if (config.enabled === false) {
// `config.enabled` as a registration gate. Callers filter based on this.unregisterAgent(agentId);
// state + ephemeral classification before calling through. return;
}
// Apply default interval if not explicitly configured // Apply default interval if not explicitly configured
// This ensures agents with heartbeat monitoring enabled but no explicit interval // This ensures agents with heartbeat monitoring enabled but no explicit interval
@@ -2042,46 +2048,112 @@ export class HeartbeatTriggerScheduler {
} }
} }
private isTimerEligibleAgent(agent: Agent): boolean {
return isHeartbeatManaged(agent)
&& agent.runtimeConfig?.enabled !== false
&& isTickableState(agent.state);
}
private getAgentTimerConfig(agent: Agent): AgentHeartbeatConfig {
const rc = (agent.runtimeConfig ?? {}) as {
enabled?: boolean;
heartbeatIntervalMs?: number;
maxConcurrentRuns?: number;
};
return {
enabled: rc.enabled,
heartbeatIntervalMs: rc.heartbeatIntervalMs,
maxConcurrentRuns: rc.maxConcurrentRuns,
};
}
private syncTimerForAgent(agent: Agent, reason: string): void {
if (!this.isTimerEligibleAgent(agent)) {
this.unregisterAgent(agent.id);
return;
}
if (this.timers.has(agent.id)) {
// Already ticking — non-config updates should not reset the interval.
return;
}
this.registerAgent(agent.id, this.getAgentTimerConfig(agent));
heartbeatLog.log(`Timer armed for ${agent.id} (${reason})`);
}
private async syncTimerForAgentFromStore(agentId: string, reason: string): Promise<void> {
const agent = await this.store.getAgent(agentId);
if (!agent) {
this.unregisterAgent(agentId);
return;
}
if (!this.isTimerEligibleAgent(agent)) {
this.unregisterAgent(agentId);
return;
}
this.registerAgent(agent.id, this.getAgentTimerConfig(agent));
heartbeatLog.log(`Timer refreshed for ${agent.id} (${reason})`);
}
private didHeartbeatScheduleChange(revision: AgentConfigRevision): boolean {
const before = (revision.before.runtimeConfig ?? {}) as Record<string, unknown>;
const after = (revision.after.runtimeConfig ?? {}) as Record<string, unknown>;
const pickScheduleFields = (runtimeConfig: Record<string, unknown>) => ({
enabled: runtimeConfig.enabled,
heartbeatIntervalMs: runtimeConfig.heartbeatIntervalMs,
maxConcurrentRuns: runtimeConfig.maxConcurrentRuns,
});
return JSON.stringify(pickScheduleFields(before)) !== JSON.stringify(pickScheduleFields(after));
}
private watchAgentLifecycle(): void { private watchAgentLifecycle(): void {
if (this.updatedListener || this.deletedListener) return; if (this.createdListener || this.updatedListener || this.configRevisionListener || this.deletedListener) return;
this.createdListener = (agent) => {
this.syncTimerForAgent(agent, `created:${agent.state}`);
};
// State-driven registration: when an agent transitions into a tickable // State-driven registration: when an agent transitions into a tickable
// state (active/running) arm the timer; transitioning out clears it. // state arm the timer; transitioning out clears it. Existing timers are
// left alone here so unrelated agent updates do not reset the interval.
this.updatedListener = (agent) => { this.updatedListener = (agent) => {
if (!isHeartbeatManaged(agent) || !isTickableState(agent.state)) { this.syncTimerForAgent(agent, `state:${agent.state}`);
this.unregisterAgent(agent.id); };
this.configRevisionListener = (agentId, revision) => {
if (!this.didHeartbeatScheduleChange(revision)) {
return; return;
} }
if (this.timers.has(agent.id)) {
// Already ticking — re-registering would reset the interval mid-cycle void this.syncTimerForAgentFromStore(agentId, "runtime-config-updated");
// on every unrelated agent update.
return;
}
const rc = (agent.runtimeConfig ?? {}) as {
heartbeatIntervalMs?: number;
maxConcurrentRuns?: number;
};
this.registerAgent(agent.id, {
heartbeatIntervalMs: rc.heartbeatIntervalMs,
maxConcurrentRuns: rc.maxConcurrentRuns,
});
heartbeatLog.log(
`State-driven registration: ${agent.id} is ${agent.state} — timer armed`,
);
}; };
this.deletedListener = (agentId) => { this.deletedListener = (agentId) => {
this.unregisterAgent(agentId); this.unregisterAgent(agentId);
}; };
this.store.on("agent:created", this.createdListener);
this.store.on("agent:updated", this.updatedListener); this.store.on("agent:updated", this.updatedListener);
this.store.on("agent:configRevision", this.configRevisionListener);
this.store.on("agent:deleted", this.deletedListener); this.store.on("agent:deleted", this.deletedListener);
} }
private unwatchAgentLifecycle(): void { private unwatchAgentLifecycle(): void {
if (this.createdListener) {
this.store.off("agent:created", this.createdListener);
this.createdListener = null;
}
if (this.updatedListener) { if (this.updatedListener) {
this.store.off("agent:updated", this.updatedListener); this.store.off("agent:updated", this.updatedListener);
this.updatedListener = null; this.updatedListener = null;
} }
if (this.configRevisionListener) {
this.store.off("agent:configRevision", this.configRevisionListener);
this.configRevisionListener = null;
}
if (this.deletedListener) { if (this.deletedListener) {
this.store.off("agent:deleted", this.deletedListener); this.store.off("agent:deleted", this.deletedListener);
this.deletedListener = null; this.deletedListener = null;

View File

@@ -862,24 +862,46 @@ describe("InProcessRuntime", () => {
expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id); expect(scheduler!.getRegisteredAgents()).not.toContain(agent.id);
}); });
it("re-registers an existing agent when agent:updated event is emitted", async () => { it("does not reset an armed timer on unrelated agent updates", async () => {
// Create a new agent
const store = getAgentStore(runtime); 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({ const agent = await store.createAgent({
name: "test-agent-update", name: "test-agent-update",
role: "executor", role: "executor",
runtimeConfig: {
enabled: true,
heartbeatIntervalMs: 1000,
},
}); });
const scheduler = runtime.getTriggerScheduler(); const scheduler = runtime.getTriggerScheduler();
expect(scheduler!.getRegisteredAgents()).toContain(agent.id); expect(scheduler!.getRegisteredAgents()).toContain(agent.id);
// Update the agent await vi.advanceTimersByTimeAsync(400);
await store.updateAgent(agent.id, { await store.updateAgent(agent.id, {
name: "test-agent-update-renamed", name: "test-agent-update-renamed",
}); });
// Verify the agent is still registered (re-registration succeeded) await vi.advanceTimersByTimeAsync(599);
expect(scheduler!.getRegisteredAgents()).toContain(agent.id); 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 () => { 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); 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 () => { it("clears timers on pause and re-arms from resume without stale pre-pause firing", async () => {
const store = getAgentStore(runtime); const store = getAgentStore(runtime);
const monitor = runtime.getHeartbeatMonitor(); const monitor = runtime.getHeartbeatMonitor();

View File

@@ -98,8 +98,6 @@ export class InProcessRuntime
private triageProcessor?: TriageProcessor; private triageProcessor?: TriageProcessor;
private messageStore?: MessageStore; private messageStore?: MessageStore;
private concurrencyChangedListener?: (state: { globalMaxConcurrent: number }) => void; 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) */ /** Set of agent IDs with scheduled ephemeral cleanup (prevents duplicate deletion) */
private pendingEphemeralDeletions = new Set<string>(); private pendingEphemeralDeletions = new Set<string>();
/** Map of agent IDs to their cleanup timer IDs */ /** Map of agent IDs to their cleanup timer IDs */
@@ -506,9 +504,8 @@ export class InProcessRuntime
); );
this.triggerScheduler.start(); this.triggerScheduler.start();
// Dynamic registration follows per-agent heartbeat enablement and tickable state. // Startup bootstrap for already-persisted agents. Ongoing lifecycle
// Non-ephemeral agents are managed unless runtimeConfig.enabled is explicitly false. // updates are handled inside HeartbeatTriggerScheduler itself.
// Paused/error/terminated states are never timer-armed.
const isHeartbeatEnabledAgent = (agent: import("@fusion/core").Agent) => const isHeartbeatEnabledAgent = (agent: import("@fusion/core").Agent) =>
!isEphemeralAgent(agent) && agent.runtimeConfig?.enabled !== false; !isEphemeralAgent(agent) && agent.runtimeConfig?.enabled !== false;
const isTickableHeartbeatState = (state: import("@fusion/core").AgentState) => const isTickableHeartbeatState = (state: import("@fusion/core").AgentState) =>
@@ -516,34 +513,6 @@ export class InProcessRuntime
const isTimerManagedAgent = (agent: import("@fusion/core").Agent) => const isTimerManagedAgent = (agent: import("@fusion/core").Agent) =>
isHeartbeatEnabledAgent(agent) && isTickableHeartbeatState(agent.state); 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. // Listen for agent state transitions to clean up terminated ephemeral agents.
// This catches cases where ephemeral agents (task-workers, spawned children) are // This catches cases where ephemeral agents (task-workers, spawned children) are
// terminated by HeartbeatMonitor or other pathways outside of onComplete/onError callbacks. // terminated by HeartbeatMonitor or other pathways outside of onComplete/onError callbacks.
@@ -595,6 +564,7 @@ export class InProcessRuntime
if (!isTimerManagedAgent(agent)) continue; if (!isTimerManagedAgent(agent)) continue;
const rc = agent.runtimeConfig; const rc = agent.runtimeConfig;
this.triggerScheduler.registerAgent(agent.id, { this.triggerScheduler.registerAgent(agent.id, {
enabled: rc?.enabled as boolean | undefined,
heartbeatIntervalMs: rc?.heartbeatIntervalMs as number | undefined, heartbeatIntervalMs: rc?.heartbeatIntervalMs as number | undefined,
maxConcurrentRuns: rc?.maxConcurrentRuns 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) // 3. Remove agent event listeners (before stopping trigger scheduler)
// Guard on this.agentStore being defined - it may not exist if AgentStore init failed // 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) { if (this.ephemeralTerminationListener && this.agentStore) {
this.agentStore.off("agent:stateChanged", this.ephemeralTerminationListener); this.agentStore.off("agent:stateChanged", this.ephemeralTerminationListener);
this.ephemeralTerminationListener = undefined; this.ephemeralTerminationListener = undefined;