feat(FN-2523): add safe restart restore lifecycle diagnostics

- Add ProjectEngine restore lifecycle core to perform safe restarts and surface detailed restore state transitions
- Expose restore diagnostics through remote-access status types and settings/memory route context, including legacy API mapping updates
- Add comprehensive regression coverage for restore lifecycle behavior in engine and dashboard headless remote-access tests
- Document the restore lifecycle contract in architecture/settings docs and include a patch changeset for @runfusion/fusion
This commit is contained in:
Fusion
2026-04-26 05:18:12 -07:00
committed by gsxdsm
parent a4c219f89b
commit e5ed696922
12 changed files with 956 additions and 19 deletions

View File

@@ -15,6 +15,7 @@ const mocks = vi.hoisted(() => ({
runtimeStart: vi.fn(async () => undefined),
runtimeStop: vi.fn(async () => undefined),
aiMergeTask: vi.fn(),
execFile: vi.fn(),
currentStore: null as Record<string, unknown> | null,
}));
@@ -47,6 +48,14 @@ vi.mock("../merger.js", () => ({
aiMergeTask: mocks.aiMergeTask,
}));
vi.mock("node:child_process", async (importOriginal) => {
const actual = await importOriginal<typeof import("node:child_process")>();
return {
...actual,
execFile: mocks.execFile,
};
});
vi.mock("../pr-monitor.js", () => ({
PrMonitor: vi.fn().mockImplementation(() => ({
onNewComments: vi.fn(),
@@ -86,15 +95,22 @@ type SettingsHandlerPayload = {
};
function createMockStore(initialSettings: Record<string, unknown>) {
let settings = { ...initialSettings };
let settings = structuredClone(initialSettings);
const settingsHandlers = new Set<(payload: SettingsHandlerPayload) => void | Promise<void>>();
const store = {
getSettings: vi.fn(async () => ({ ...settings })),
getSettings: vi.fn(async () => structuredClone(settings)),
listTasks: vi.fn(async () => []),
getTask: vi.fn(async (taskId: string) => ({ id: taskId, column: "in-review", mergeRetries: 0, status: null })),
updateTask: vi.fn(async () => undefined),
moveTask: vi.fn(async () => undefined),
updateSettings: vi.fn(async (patch: Record<string, unknown>) => {
settings = {
...settings,
...patch,
};
return structuredClone(settings);
}),
logEntry: vi.fn(async () => undefined),
addTaskComment: vi.fn(async () => undefined),
getActiveMergingTask: vi.fn(() => null),
@@ -116,15 +132,52 @@ function createMockStore(initialSettings: Record<string, unknown>) {
next: Record<string, unknown>,
previous: Record<string, unknown>,
) => {
settings = { ...next };
settings = structuredClone(next);
for (const handler of settingsHandlers) {
await handler({ settings: { ...next }, previous: { ...previous } });
await handler({ settings: structuredClone(next), previous: structuredClone(previous) });
}
};
return { store, emitSettingsUpdated };
const getCurrentSettings = () => structuredClone(settings);
return { store, emitSettingsUpdated, getCurrentSettings };
}
const baseRemoteAccess = {
enabled: true,
activeProvider: "cloudflare" as const,
providers: {
tailscale: {
enabled: true,
hostname: "tail.example.ts.net",
targetPort: 4040,
acceptRoutes: false,
},
cloudflare: {
enabled: true,
tunnelName: "demo",
tunnelToken: "cf-secret-token",
ingressUrl: "https://remote.example.com",
},
},
tokenStrategy: {
persistent: {
enabled: true,
token: "frt_persistent",
},
shortLived: {
enabled: true,
ttlMs: 120_000,
maxTtlMs: 86_400_000,
},
},
lifecycle: {
rememberLastRunning: true,
wasRunningOnShutdown: false,
lastRunningProvider: null,
},
};
const baseSettings: Record<string, unknown> = {
autoMerge: false,
globalPause: false,
@@ -139,6 +192,7 @@ const baseSettings: Record<string, unknown> = {
insightExtractionEnabled: false,
insightExtractionSchedule: "0 3 * * *",
insightExtractionMinIntervalMs: 0,
remoteAccess: baseRemoteAccess,
};
function createEngine() {
@@ -155,6 +209,26 @@ function createEngine() {
);
}
beforeEach(() => {
mocks.execFile.mockImplementation((
_file: string,
_args: string[],
_options: unknown,
callback?: (error: Error | null, result: { stdout: string; stderr: string }) => void,
) => {
if (typeof _options === "function") {
(_options as (error: Error | null, result: { stdout: string; stderr: string }) => void)(null, {
stdout: "/usr/bin/mock\n",
stderr: "",
});
return {} as never;
}
callback?.(null, { stdout: "/usr/bin/mock\n", stderr: "" });
return {} as never;
});
});
describe("ProjectEngine auto-summarize wiring", () => {
beforeEach(() => {
vi.clearAllMocks();
@@ -284,6 +358,228 @@ describe("ProjectEngine remote tunnel manager wiring", () => {
});
});
describe("ProjectEngine remote lifecycle restore policy", () => {
beforeEach(() => {
vi.clearAllMocks();
const mockStore = createMockStore(baseSettings);
mocks.currentStore = mockStore.store;
});
it("attempts restore on startup when rememberLastRunning and prior-running markers are set", async () => {
const restoreSettings = {
...baseSettings,
remoteAccess: {
...baseRemoteAccess,
lifecycle: {
...baseRemoteAccess.lifecycle,
rememberLastRunning: true,
wasRunningOnShutdown: true,
lastRunningProvider: "cloudflare" as const,
},
},
};
const mockStore = createMockStore(restoreSettings);
mocks.currentStore = mockStore.store;
const startSpy = vi.spyOn(TunnelProcessManager.prototype, "start").mockResolvedValue(undefined);
const engine = createEngine();
await engine.start();
expect(startSpy).toHaveBeenCalledTimes(1);
expect(startSpy.mock.calls[0]?.[0]).toBe("cloudflare");
expect(engine.getRemoteTunnelRestoreDiagnostics()).toMatchObject({
outcome: "applied",
reason: "restore_started",
provider: "cloudflare",
});
await engine.stop();
startSpy.mockRestore();
});
it("skips restore when rememberLastRunning is disabled", async () => {
const restoreSettings = {
...baseSettings,
remoteAccess: {
...baseRemoteAccess,
lifecycle: {
...baseRemoteAccess.lifecycle,
rememberLastRunning: false,
wasRunningOnShutdown: true,
lastRunningProvider: "cloudflare" as const,
},
},
};
const mockStore = createMockStore(restoreSettings);
mocks.currentStore = mockStore.store;
const startSpy = vi.spyOn(TunnelProcessManager.prototype, "start").mockResolvedValue(undefined);
const engine = createEngine();
await engine.start();
expect(startSpy).not.toHaveBeenCalled();
expect(engine.getRemoteTunnelRestoreDiagnostics()).toMatchObject({
outcome: "skipped",
reason: "remember_last_running_disabled",
provider: null,
});
expect(mockStore.store.updateSettings).toHaveBeenCalledWith(expect.objectContaining({
remoteAccess: expect.objectContaining({
lifecycle: expect.objectContaining({
wasRunningOnShutdown: false,
lastRunningProvider: null,
}),
}),
}));
await engine.stop();
startSpy.mockRestore();
});
it("skips restore with explicit reason and clears stale marker when prerequisites are missing", async () => {
const restoreSettings = {
...baseSettings,
remoteAccess: {
...baseRemoteAccess,
providers: {
...baseRemoteAccess.providers,
cloudflare: {
...baseRemoteAccess.providers.cloudflare,
tunnelToken: null,
},
},
lifecycle: {
...baseRemoteAccess.lifecycle,
rememberLastRunning: true,
wasRunningOnShutdown: true,
lastRunningProvider: "cloudflare" as const,
},
},
};
const mockStore = createMockStore(restoreSettings);
mocks.currentStore = mockStore.store;
const startSpy = vi.spyOn(TunnelProcessManager.prototype, "start").mockResolvedValue(undefined);
const engine = createEngine();
await engine.start();
expect(startSpy).not.toHaveBeenCalled();
expect(engine.getRemoteTunnelRestoreDiagnostics()).toMatchObject({
outcome: "skipped",
reason: "provider_not_configured",
provider: "cloudflare",
});
expect(mockStore.store.updateSettings).toHaveBeenCalledWith(expect.objectContaining({
remoteAccess: expect.objectContaining({
lifecycle: expect.objectContaining({
wasRunningOnShutdown: false,
lastRunningProvider: null,
}),
}),
}));
await engine.stop();
startSpy.mockRestore();
});
it("reconciles stale persisted running marker to avoid restore loops", async () => {
const restoreSettings = {
...baseSettings,
remoteAccess: {
...baseRemoteAccess,
providers: {
...baseRemoteAccess.providers,
cloudflare: {
...baseRemoteAccess.providers.cloudflare,
tunnelToken: null,
},
},
lifecycle: {
...baseRemoteAccess.lifecycle,
rememberLastRunning: true,
wasRunningOnShutdown: true,
lastRunningProvider: "cloudflare" as const,
},
},
};
const mockStore = createMockStore(restoreSettings);
mocks.currentStore = mockStore.store;
const startSpy = vi.spyOn(TunnelProcessManager.prototype, "start").mockResolvedValue(undefined);
const firstEngine = createEngine();
await firstEngine.start();
await firstEngine.stop();
const secondEngine = createEngine();
await secondEngine.start();
await secondEngine.stop();
expect(startSpy).not.toHaveBeenCalled();
expect(secondEngine.getRemoteTunnelRestoreDiagnostics()).toMatchObject({
outcome: "skipped",
reason: "no_prior_running_marker",
});
startSpy.mockRestore();
});
it("does not auto-start on settings updates and manual stop clears future restore intent", async () => {
const mockStore = createMockStore(baseSettings);
mocks.currentStore = mockStore.store;
const startSpy = vi.spyOn(TunnelProcessManager.prototype, "start").mockResolvedValue(undefined);
const stopSpy = vi.spyOn(TunnelProcessManager.prototype, "stop").mockResolvedValue(undefined);
const statusSpy = vi.spyOn(TunnelProcessManager.prototype, "getStatus").mockReturnValue({
provider: null,
state: "stopped",
pid: null,
startedAt: null,
stoppedAt: null,
url: null,
lastError: null,
});
const engine = createEngine();
await engine.start();
await mockStore.emitSettingsUpdated(
{
...baseSettings,
remoteAccess: {
...baseRemoteAccess,
activeProvider: "tailscale" as const,
},
},
baseSettings,
);
const startsBeforeManualAction = startSpy.mock.calls.length;
await engine.startRemoteTunnel();
await engine.stopRemoteTunnel();
expect(startSpy.mock.calls.length).toBe(startsBeforeManualAction + 1);
expect(stopSpy).toHaveBeenCalled();
expect(mockStore.store.updateSettings).toHaveBeenCalledWith(expect.objectContaining({
remoteAccess: expect.objectContaining({
lifecycle: expect.objectContaining({
wasRunningOnShutdown: false,
lastRunningProvider: null,
}),
}),
}));
await engine.stop();
startSpy.mockRestore();
stopSpy.mockRestore();
statusSpy.mockRestore();
});
});
describe("ProjectEngine shutdown merge handling", () => {
beforeEach(() => {
vi.clearAllMocks();

View File

@@ -102,6 +102,9 @@ export {
type TunnelProviderAdapter,
type TunnelProviderConfig,
type TunnelReadinessEvent,
type TunnelRestoreDiagnostics,
type TunnelRestoreOutcome,
type TunnelRestoreReasonCode,
type TunnelStatusListener,
type TunnelStatusSnapshot,
} from "./remote-access/index.js";

View File

@@ -9,6 +9,8 @@ import type {
ScheduledTask,
AutomationRunResult,
} from "@fusion/core";
import { execFile } from "node:child_process";
import { promisify } from "node:util";
import { InProcessRuntime } from "./runtimes/in-process-runtime.js";
import type { ProjectRuntimeConfig } from "./project-runtime.js";
import { PrMonitor } from "./pr-monitor.js";
@@ -21,6 +23,13 @@ import { PRIORITY_MERGE } from "./concurrency.js";
import { runtimeLog } from "./logger.js";
import type { HeartbeatTriggerScheduler } from "./agent-heartbeat.js";
import { TunnelProcessManager } from "./remote-access/tunnel-process-manager.js";
import type {
TunnelProvider,
TunnelProviderConfig,
TunnelRestoreDiagnostics,
TunnelRestoreReasonCode,
TunnelStatusSnapshot,
} from "./remote-access/types.js";
/**
* Callback for processing pull-request merge strategy.
@@ -32,6 +41,15 @@ export type ProcessPullRequestMergeFn = (
taskId: string,
) => Promise<"merged" | "waiting" | "skipped">;
const execFileAsync = promisify(execFile);
interface RemoteLifecycleEvaluation {
provider: TunnelProvider;
config?: TunnelProviderConfig;
reason?: TunnelRestoreReasonCode;
message?: string;
}
export interface ProjectEngineOptions {
/** Project identifier for notification deep links */
projectId?: string;
@@ -92,6 +110,12 @@ export class ProjectEngine {
private cronRunner?: CronRunner;
private automationStore?: AutomationStoreType;
private remoteTunnelManager?: TunnelProcessManager;
private remoteTunnelRestoreDiagnostics: TunnelRestoreDiagnostics = {
outcome: "skipped",
reason: "not_attempted",
at: new Date().toISOString(),
provider: null,
};
// ── Auto-merge state ──
private mergeQueue: string[] = [];
@@ -143,6 +167,13 @@ export class ProjectEngine {
const cwd = this.config.workingDirectory;
this.remoteTunnelManager = new TunnelProcessManager();
try {
await this.restoreRemoteTunnelIfNeeded(store);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
this.setRestoreDiagnostics("failed", "restore_start_failed", null, message);
runtimeLog.warn(`Remote tunnel restore evaluation failed (continuing startup): ${message}`);
}
// 2. Initialize PrMonitor + PrCommentHandler
this.prMonitor = new PrMonitor();
@@ -288,6 +319,22 @@ export class ProjectEngine {
const tunnelManager = this.remoteTunnelManager;
this.remoteTunnelManager = undefined;
if (tunnelManager) {
let shutdownStore: TaskStore | null = null;
try {
shutdownStore = this.runtime.getTaskStore();
} catch {
shutdownStore = null;
}
if (shutdownStore) {
try {
await this.persistShutdownRemoteLifecycle(shutdownStore, tunnelManager.getStatus());
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
runtimeLog.warn(`Failed to persist remote lifecycle shutdown markers: ${message}`);
}
}
try {
await tunnelManager.stop();
} catch (error) {
@@ -359,6 +406,80 @@ export class ProjectEngine {
return this.remoteTunnelManager;
}
getRemoteTunnelRestoreDiagnostics(): TunnelRestoreDiagnostics {
return { ...this.remoteTunnelRestoreDiagnostics };
}
async startRemoteTunnel(): Promise<TunnelStatusSnapshot> {
const manager = this.remoteTunnelManager;
if (!manager) {
throw new Error("remote_tunnel_unavailable:remote tunnel manager is not initialized");
}
const store = this.runtime.getTaskStore();
const settings = await store.getSettings();
const remoteAccess = settings.remoteAccess;
if (!remoteAccess?.enabled) {
throw new Error("invalid_config:remote access is disabled");
}
const provider = remoteAccess.activeProvider;
if (!provider) {
throw new Error("invalid_config:no active remote provider configured");
}
const lifecycle = await this.evaluateRemoteLifecycle(settings, provider);
if (!lifecycle.config) {
throw new Error(`${lifecycle.reason ?? "invalid_config"}:${lifecycle.message ?? "remote provider prerequisites are not met"}`);
}
const current = manager.getStatus();
if (current.state === "running" && current.provider === provider) {
await this.writeRemoteLifecycleState(store, remoteAccess, {
...remoteAccess.lifecycle,
wasRunningOnShutdown: true,
lastRunningProvider: provider,
});
return manager.getStatus();
}
if (current.state === "running" && current.provider && current.provider !== provider) {
await manager.switchProvider(provider, lifecycle.config);
} else {
await manager.start(provider, lifecycle.config);
}
await this.writeRemoteLifecycleState(store, remoteAccess, {
...remoteAccess.lifecycle,
wasRunningOnShutdown: true,
lastRunningProvider: provider,
});
return manager.getStatus();
}
async stopRemoteTunnel(): Promise<TunnelStatusSnapshot> {
const manager = this.remoteTunnelManager;
if (!manager) {
throw new Error("remote_tunnel_unavailable:remote tunnel manager is not initialized");
}
await manager.stop();
const store = this.runtime.getTaskStore();
const settings = await store.getSettings();
const remoteAccess = settings.remoteAccess;
if (remoteAccess) {
await this.writeRemoteLifecycleState(store, remoteAccess, {
...remoteAccess.lifecycle,
wasRunningOnShutdown: false,
lastRunningProvider: null,
});
}
return manager.getStatus();
}
/** Get the RoutineRunner (if initialized). */
getRoutineRunner(): RoutineRunner | undefined {
return this.runtime.getRoutineRunner();
@@ -401,6 +522,206 @@ export class ProjectEngine {
});
}
private setRestoreDiagnostics(
outcome: TunnelRestoreDiagnostics["outcome"],
reason: TunnelRestoreReasonCode,
provider: TunnelProvider | null,
message?: string,
): void {
this.remoteTunnelRestoreDiagnostics = {
outcome,
reason,
provider,
message,
at: new Date().toISOString(),
};
}
private async restoreRemoteTunnelIfNeeded(store: TaskStore): Promise<void> {
const manager = this.remoteTunnelManager;
if (!manager) {
return;
}
const settings = await store.getSettings();
const remoteAccess = settings.remoteAccess;
if (!remoteAccess?.enabled) {
this.setRestoreDiagnostics("skipped", "remote_access_disabled", null);
return;
}
const lifecycle = remoteAccess.lifecycle;
if (!lifecycle.rememberLastRunning) {
this.setRestoreDiagnostics("skipped", "remember_last_running_disabled", null);
if (lifecycle.wasRunningOnShutdown || lifecycle.lastRunningProvider) {
await this.writeRemoteLifecycleState(store, remoteAccess, {
...lifecycle,
wasRunningOnShutdown: false,
lastRunningProvider: null,
});
}
return;
}
if (!lifecycle.wasRunningOnShutdown) {
this.setRestoreDiagnostics("skipped", "no_prior_running_marker", null);
return;
}
const provider = lifecycle.lastRunningProvider ?? remoteAccess.activeProvider;
if (!provider) {
this.setRestoreDiagnostics("skipped", "provider_missing", null);
await this.writeRemoteLifecycleState(store, remoteAccess, {
...lifecycle,
wasRunningOnShutdown: false,
lastRunningProvider: null,
});
return;
}
const evaluation = await this.evaluateRemoteLifecycle(settings, provider);
if (!evaluation.config) {
this.setRestoreDiagnostics("skipped", evaluation.reason ?? "provider_not_configured", provider, evaluation.message);
await this.writeRemoteLifecycleState(store, remoteAccess, {
...lifecycle,
wasRunningOnShutdown: false,
lastRunningProvider: null,
});
return;
}
try {
await manager.start(provider, evaluation.config);
this.setRestoreDiagnostics("applied", "restore_started", provider);
await this.writeRemoteLifecycleState(store, remoteAccess, {
...lifecycle,
wasRunningOnShutdown: true,
lastRunningProvider: provider,
}, provider);
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
this.setRestoreDiagnostics("failed", "restore_start_failed", provider, message);
runtimeLog.warn(`Remote tunnel restore failed for ${provider}: ${message}`);
await this.writeRemoteLifecycleState(store, remoteAccess, {
...lifecycle,
wasRunningOnShutdown: false,
lastRunningProvider: null,
});
}
}
private async persistShutdownRemoteLifecycle(
store: TaskStore,
status: TunnelStatusSnapshot,
): Promise<void> {
const settings = await store.getSettings();
const remoteAccess = settings.remoteAccess;
if (!remoteAccess) {
return;
}
const shouldRememberRunning =
(status.state === "running" || status.state === "starting" || status.state === "stopping") &&
status.provider !== null;
await this.writeRemoteLifecycleState(store, remoteAccess, {
...remoteAccess.lifecycle,
wasRunningOnShutdown: shouldRememberRunning,
lastRunningProvider: shouldRememberRunning ? status.provider : null,
}, shouldRememberRunning ? status.provider : remoteAccess.activeProvider);
}
private async writeRemoteLifecycleState(
store: TaskStore,
remoteAccess: NonNullable<Settings["remoteAccess"]>,
lifecycle: NonNullable<Settings["remoteAccess"]>["lifecycle"],
activeProviderOverride?: TunnelProvider | null,
): Promise<void> {
await store.updateSettings({
remoteAccess: {
...remoteAccess,
activeProvider: activeProviderOverride === undefined ? remoteAccess.activeProvider : activeProviderOverride,
lifecycle,
},
});
}
private async evaluateRemoteLifecycle(
settings: Settings,
provider: TunnelProvider,
): Promise<RemoteLifecycleEvaluation> {
const remoteAccess = settings.remoteAccess;
if (!remoteAccess?.enabled) {
return { provider, reason: "remote_access_disabled", message: "Remote access is disabled" };
}
if (provider === "tailscale") {
const tailscale = remoteAccess.providers.tailscale;
if (!tailscale.enabled) {
return { provider, reason: "provider_not_enabled", message: "Tailscale provider is disabled" };
}
if (!tailscale.hostname?.trim() || !Number.isFinite(tailscale.targetPort) || tailscale.targetPort <= 0) {
return { provider, reason: "provider_not_configured", message: "Tailscale hostname and target port must be configured" };
}
const executable = await this.checkExecutableAvailable("tailscale");
if (!executable.available) {
return { provider, reason: "runtime_prerequisite_missing", message: executable.message };
}
return {
provider,
config: {
provider: "tailscale",
executablePath: "tailscale",
args: ["funnel", String(Math.floor(tailscale.targetPort))],
},
};
}
const cloudflare = remoteAccess.providers.cloudflare;
if (!cloudflare.enabled) {
return { provider, reason: "provider_not_enabled", message: "Cloudflare provider is disabled" };
}
if (!cloudflare.tunnelName?.trim() || !cloudflare.ingressUrl?.trim()) {
return { provider, reason: "provider_not_configured", message: "Cloudflare tunnel name and ingress URL must be configured" };
}
if (!cloudflare.tunnelToken?.trim()) {
return { provider, reason: "provider_not_configured", message: "Cloudflare tunnel token is required" };
}
const executable = await this.checkExecutableAvailable("cloudflared");
if (!executable.available) {
return { provider, reason: "runtime_prerequisite_missing", message: executable.message };
}
return {
provider,
config: {
provider: "cloudflare",
executablePath: "cloudflared",
args: ["tunnel", "--no-autoupdate", "run", cloudflare.tunnelName.trim()],
tokenEnvVar: "TUNNEL_TOKEN",
env: {
TUNNEL_TOKEN: cloudflare.tunnelToken,
},
},
};
}
private async checkExecutableAvailable(command: string): Promise<{ available: boolean; message?: string }> {
const checker = process.platform === "win32" ? "where" : "which";
try {
await execFileAsync(checker, [command]);
return { available: true };
} catch {
return {
available: false,
message: `${command} is not available on PATH`,
};
}
}
// ── Merge eligibility helpers (richer logic from dashboard.ts) ──
/**

View File

@@ -22,6 +22,9 @@ export type {
TunnelProviderAdapter,
TunnelProviderConfig,
TunnelReadinessEvent,
TunnelRestoreDiagnostics,
TunnelRestoreOutcome,
TunnelRestoreReasonCode,
TunnelStatusListener,
TunnelStatusSnapshot,
} from "./types.js";

View File

@@ -37,6 +37,28 @@ export interface TunnelStatusSnapshot {
lastError: TunnelError | null;
}
export type TunnelRestoreOutcome = "applied" | "skipped" | "failed";
export type TunnelRestoreReasonCode =
| "not_attempted"
| "remote_access_disabled"
| "remember_last_running_disabled"
| "no_prior_running_marker"
| "provider_missing"
| "provider_not_enabled"
| "provider_not_configured"
| "runtime_prerequisite_missing"
| "restore_start_failed"
| "restore_started";
export interface TunnelRestoreDiagnostics {
outcome: TunnelRestoreOutcome;
reason: TunnelRestoreReasonCode;
at: string;
provider: TunnelProvider | null;
message?: string;
}
export type TunnelLogLevel = "info" | "warn" | "error";
export interface TunnelLogEntry {