diff --git a/.changeset/dev-engine-watch.md b/.changeset/dev-engine-watch.md new file mode 100644 index 0000000000..79dea7b880 --- /dev/null +++ b/.changeset/dev-engine-watch.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": minor +--- + +summary: Add safe development source watching with automatic engine restarts. +category: feature +dev: Use `pnpm dev:watch`; restarts close new admission, drain active agents, rebuild dist, and respawn the supervised child. diff --git a/README.md b/README.md index 1ec174c4ef..3f3bb3874e 100644 --- a/README.md +++ b/README.md @@ -665,6 +665,7 @@ pnpm local --no-engine # Start local dashboard/API only pnpm build # Build default workspace packages (excludes desktop/mobile) pnpm build:all # Build all packages (including desktop/mobile) pnpm dev dashboard # Run dashboard + AI engine +pnpm dev:watch # Dashboard + AI engine; restart on source edits after agents drain pnpm dev:ui # Dashboard only (no AI engine) pnpm lint # Lint all packages pnpm typecheck # Type-check all packages diff --git a/docs/contributing.md b/docs/contributing.md index 2fa192ff44..372dd6327e 100644 --- a/docs/contributing.md +++ b/docs/contributing.md @@ -72,9 +72,10 @@ pnpm local # fast local dashboard/API + AI engine startup on a safe pnpm local --no-engine # fast local dashboard/API-only startup pnpm local --prebuild # local dashboard/API + AI engine startup with an explicit prebuild level pnpm dev # source-mode CLI; dashboard gets a client-only prebuild, other commands skip it +pnpm dev:watch # dashboard + engine; gracefully restart runtime source after active agents drain FUSION_DEV_PREBUILD=full pnpm dev dashboard # production-like full workspace prebuild pnpm dev:ui # dashboard dev server only -pnpm dev:hmr # dashboard API + Vite HMR UI, with no startup prebuild +pnpm dev:hmr # dashboard API + Vite UI HMR + graceful runtime source restarts pnpm lint # lint all packages pnpm test # merge-gate suite + changed-only affected tests (bounded; never full-suite) pnpm test:gate # the merge gate: curated engine-core suite + CI-shape test diff --git a/package.json b/package.json index a042a2f1ad..0344181e5f 100644 --- a/package.json +++ b/package.json @@ -34,6 +34,7 @@ "smoke:boot": "node scripts/boot-smoke.mjs", "local": "node scripts/start-local.mjs", "dev": "node scripts/dev-with-memory.mjs", + "dev:watch": "node scripts/dev-with-memory.mjs --watch", "start": "node scripts/dev-with-memory.mjs", "dev:ui": "pnpm --filter @fusion/dashboard dev", "dev:hmr": "node scripts/dev-hmr.mjs", diff --git a/packages/cli/src/__tests__/dev-with-memory-lib.test.ts b/packages/cli/src/__tests__/dev-with-memory-lib.test.ts index 25e1d65db1..418a3d8ef6 100644 --- a/packages/cli/src/__tests__/dev-with-memory-lib.test.ts +++ b/packages/cli/src/__tests__/dev-with-memory-lib.test.ts @@ -1,12 +1,21 @@ -import { describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { buildForwardedDevArgs, buildDevNodeArgs, + createDevWatchRestartCoordinator, getPrebuildCommand, normalizePrebuildMode, parseDevWrapperArgs, resolvePrebuildMode, } from "../../../../scripts/dev-with-memory-lib.mjs"; +import { + createDevSourceWatcher, + isRestartableSourceFile, +} from "../../../../scripts/lib/dev-source-watch.mjs"; + +afterEach(() => { + vi.useRealTimers(); +}); describe("buildDevNodeArgs", () => { it("enables source-condition resolution before loading the tsx runtime", () => { @@ -44,6 +53,26 @@ describe("dev-with-memory prebuild options", () => { inspectFlags: ["--inspect=9230"], args: ["dashboard", "--port", "4050"], requestedPrebuild: "none", + watchSource: false, + watchSourceFromFlag: false, + }); + }); + + it("keeps the wrapper-only watch flag out of CLI arguments", () => { + expect(parseDevWrapperArgs(["--watch", "dashboard", "--port", "4050"], {})).toEqual({ + inspectFlags: [], + args: ["dashboard", "--port", "4050"], + requestedPrebuild: "auto", + watchSource: true, + watchSourceFromFlag: true, + }); + }); + + it("allows source watching to be enabled through the development environment", () => { + expect(parseDevWrapperArgs(["dashboard"], { FUSION_DEV_WATCH: "1" })).toMatchObject({ + args: ["dashboard"], + watchSource: true, + watchSourceFromFlag: false, }); }); @@ -120,3 +149,175 @@ describe("dev-with-memory prebuild options", () => { }); }); }); + +describe("development source restart watcher", () => { + it("holds source changes until the child acknowledges its IPC listener", () => { + const send = vi.fn((_message: unknown, callback?: (error?: Error) => void) => callback?.()); + const child = { connected: true, send }; + const coordinator = createDevWatchRestartCoordinator({ log: vi.fn(), warn: vi.fn() }); + coordinator.attach(child); + + coordinator.request(["packages/core/src/store.ts"]); + expect(send).not.toHaveBeenCalled(); + + coordinator.onMessage({ type: "fusion:dev-source-restart-armed" }); + expect(send).toHaveBeenCalledWith( + { type: "fusion:dev-source-changed" }, + expect.any(Function), + ); + expect(coordinator.detach(child)).toBe(true); + }); + + it("re-arms after the watched child is replaced", () => { + const first = { connected: true, send: vi.fn((_message: unknown, callback?: (error?: Error) => void) => callback?.()) }; + const second = { connected: true, send: vi.fn((_message: unknown, callback?: (error?: Error) => void) => callback?.()) }; + const coordinator = createDevWatchRestartCoordinator({ log: vi.fn(), warn: vi.fn() }); + + coordinator.attach(first); + coordinator.onMessage({ type: "fusion:dev-source-restart-armed" }); + coordinator.request(["packages/core/src/first.ts"]); + expect(coordinator.detach(first)).toBe(true); + + coordinator.attach(second); + coordinator.request(["packages/core/src/second.ts"]); + expect(second.send).not.toHaveBeenCalled(); + coordinator.onMessage({ type: "fusion:dev-source-restart-armed" }); + expect(second.send).toHaveBeenCalledTimes(1); + }); + + it("retains changes while the child is disconnected and sends them later", () => { + const log = vi.fn(); + const warn = vi.fn(); + const send = vi.fn((_message: unknown, callback?: (error?: Error) => void) => callback?.()); + const child = { connected: false, send }; + const coordinator = createDevWatchRestartCoordinator({ log, warn }); + coordinator.attach(child); + coordinator.onMessage({ type: "fusion:dev-source-restart-armed" }); + + coordinator.request(["packages/core/src/first.ts"]); + child.connected = true; + coordinator.request(["packages/core/src/second.ts"]); + + expect(warn).toHaveBeenCalledWith(expect.stringContaining("not connected")); + expect(send).toHaveBeenCalledWith( + { type: "fusion:dev-source-changed" }, + expect.any(Function), + ); + expect(log).toHaveBeenLastCalledWith( + expect.stringContaining("packages/core/src/first.ts, packages/core/src/second.ts"), + ); + }); + + it("retains changes after a send callback error and retries later", () => { + const log = vi.fn(); + const warn = vi.fn(); + const send = vi.fn() + .mockImplementationOnce((_message: unknown, callback?: (error?: Error) => void) => callback?.(new Error("send failed"))) + .mockImplementationOnce((_message: unknown, callback?: (error?: Error) => void) => callback?.()); + const child = { connected: true, send }; + const coordinator = createDevWatchRestartCoordinator({ log, warn }); + coordinator.attach(child); + coordinator.onMessage({ type: "fusion:dev-source-restart-armed" }); + + coordinator.request(["packages/core/src/first.ts"]); + coordinator.request(["packages/core/src/second.ts"]); + + expect(warn).toHaveBeenCalledWith(expect.stringContaining("send failed")); + expect(send).toHaveBeenCalledTimes(2); + expect(log).toHaveBeenLastCalledWith( + expect.stringContaining("packages/core/src/first.ts, packages/core/src/second.ts"), + ); + }); + + it("restarts for runtime sources but ignores tests and generated declarations", () => { + expect(isRestartableSourceFile("store.ts")).toBe(true); + expect(isRestartableSourceFile("routes/system.tsx")).toBe(true); + expect(isRestartableSourceFile("__tests__/store.test.ts")).toBe(false); + expect(isRestartableSourceFile("store.spec.ts")).toBe(false); + expect(isRestartableSourceFile("generated/runtime.d.ts")).toBe(false); + expect(isRestartableSourceFile("README.md")).toBe(false); + }); + + it("coalesces source events into one restart and reports the changed paths", () => { + vi.useFakeTimers(); + const listeners: Array<(eventType: string, filename: string | Buffer | null) => void> = []; + const close = vi.fn(); + const onRestart = vi.fn(); + + const watcher = createDevSourceWatcher({ + rootDir: "/repo", + watchPaths: ["packages/core/src", "packages/engine/src"], + debounceMs: 250, + watch: (_path, _options, listener) => { + listeners.push(listener); + return { close, on: vi.fn() }; + }, + onRestart, + }); + + listeners[0]?.("change", "store.ts"); + listeners[1]?.("rename", "triage.ts"); + listeners[1]?.("change", "__tests__/triage.test.ts"); + + vi.advanceTimersByTime(249); + expect(onRestart).not.toHaveBeenCalled(); + vi.advanceTimersByTime(1); + expect(onRestart).toHaveBeenCalledTimes(1); + expect(onRestart).toHaveBeenCalledWith([ + "packages/core/src/store.ts", + "packages/engine/src/triage.ts", + ]); + + watcher.close(); + expect(close).toHaveBeenCalledTimes(2); + }); + + it("does not starve a restart during a sustained source event stream", () => { + vi.useFakeTimers(); + let listener: ((eventType: string, filename: string | Buffer | null) => void) | undefined; + const onRestart = vi.fn(); + const watcher = createDevSourceWatcher({ + rootDir: "/repo", + watchPaths: ["packages/core/src"], + debounceMs: 350, + maxWaitMs: 1_000, + watch: (_path, _options, nextListener) => { + listener = nextListener; + return { close: vi.fn(), on: vi.fn() } as never; + }, + onRestart, + }); + + for (let elapsed = 0; elapsed < 1_000; elapsed += 200) { + listener?.("change", `file-${elapsed}.ts`); + vi.advanceTimersByTime(200); + } + vi.advanceTimersByTime(1); + + expect(onRestart).toHaveBeenCalledTimes(1); + watcher.close(); + }); + + it("continues closing source watchers when one close throws", () => { + const logger = { warn: vi.fn() }; + const closes = [ + vi.fn(() => { + throw new Error("close failed"); + }), + vi.fn(), + ]; + let index = 0; + const watcher = createDevSourceWatcher({ + rootDir: "/repo", + watchPaths: ["packages/core/src", "packages/engine/src"], + watch: () => ({ close: closes[index++]!, on: vi.fn() }) as never, + onRestart: vi.fn(), + logger, + }); + + expect(() => watcher.close()).not.toThrow(); + expect(closes[0]).toHaveBeenCalledOnce(); + expect(closes[1]).toHaveBeenCalledOnce(); + expect(logger.warn).toHaveBeenCalledWith(expect.stringContaining("close failed")); + }); +}); diff --git a/packages/cli/src/commands/__tests__/dev-source-restart.test.ts b/packages/cli/src/commands/__tests__/dev-source-restart.test.ts new file mode 100644 index 0000000000..f9441426e1 --- /dev/null +++ b/packages/cli/src/commands/__tests__/dev-source-restart.test.ts @@ -0,0 +1,149 @@ +import { EventEmitter } from "node:events"; +import { afterEach, describe, expect, it, vi } from "vitest"; + +import { + DEV_SOURCE_CHANGE_MESSAGE, + registerDevSourceRestart, +} from "../dev-source-restart.js"; + +afterEach(() => { + vi.useRealTimers(); +}); + +function createHarness(currentlyActive: number) { + const processEvents = new EventEmitter(); + const getLiveRunningAgentCounts = vi.fn().mockResolvedValue({ currentlyActive, projectsActive: {} }); + const beginDrain = vi.fn(); + const requestRestart = vi.fn().mockReturnValue(true); + const log = vi.fn(); + const warn = vi.fn(); + const notifyArmed = vi.fn(); + const dispose = registerDevSourceRestart({ + enabled: true, + processEvents, + beginDrain, + notifyArmed, + getLiveRunningAgentCounts, + requestRestart, + logger: { log, warn }, + recheckIntervalMs: 1_000, + }); + return { processEvents, beginDrain, notifyArmed, getLiveRunningAgentCounts, requestRestart, log, warn, dispose }; +} + +describe("registerDevSourceRestart", () => { + it("acknowledges that the child restart listener is armed", () => { + const harness = createHarness(0); + expect(harness.notifyArmed).toHaveBeenCalledOnce(); + harness.dispose(); + }); + + it("requests the canonical restart immediately when no agents are active", async () => { + vi.useFakeTimers(); + const harness = createHarness(0); + + harness.processEvents.emit("message", { + type: DEV_SOURCE_CHANGE_MESSAGE, + }); + await vi.advanceTimersByTimeAsync(0); + + expect(harness.beginDrain).toHaveBeenCalledTimes(1); + expect(harness.beginDrain.mock.invocationCallOrder[0]).toBeLessThan( + harness.getLiveRunningAgentCounts.mock.invocationCallOrder[0]!, + ); + expect(harness.getLiveRunningAgentCounts).toHaveBeenCalledTimes(1); + expect(harness.requestRestart).toHaveBeenCalledWith("dev-source-change"); + harness.dispose(); + }); + + it("coalesces edits and waits for active agents to reach a safe boundary", async () => { + vi.useFakeTimers(); + const harness = createHarness(2); + + harness.processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + harness.processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + await vi.advanceTimersByTimeAsync(0); + + expect(harness.getLiveRunningAgentCounts).toHaveBeenCalledTimes(1); + expect(harness.requestRestart).not.toHaveBeenCalled(); + expect(harness.log).toHaveBeenCalledWith(expect.stringContaining("2 active agents")); + + harness.getLiveRunningAgentCounts.mockResolvedValue({ currentlyActive: 0, projectsActive: {} }); + await vi.advanceTimersByTimeAsync(1_000); + + expect(harness.requestRestart).toHaveBeenCalledTimes(1); + expect(harness.requestRestart).toHaveBeenCalledWith("dev-source-change"); + harness.dispose(); + }); + + it("recovers from a liveness read error and retries after the drain is closed", async () => { + vi.useFakeTimers(); + const harness = createHarness(0); + harness.getLiveRunningAgentCounts.mockRejectedValueOnce(new Error("database unavailable")); + + harness.processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + await vi.advanceTimersByTimeAsync(0); + + expect(harness.warn).toHaveBeenCalledWith(expect.stringContaining("database unavailable")); + expect(harness.requestRestart).not.toHaveBeenCalled(); + + harness.getLiveRunningAgentCounts.mockResolvedValueOnce({ currentlyActive: 0, projectsActive: {} }); + await vi.advanceTimersByTimeAsync(1_000); + expect(harness.getLiveRunningAgentCounts).toHaveBeenCalledTimes(2); + expect(harness.requestRestart).toHaveBeenCalledWith("dev-source-change"); + harness.dispose(); + }); + + it("contains a synchronous drain failure in the IPC handler", async () => { + const harness = createHarness(0); + harness.beginDrain.mockImplementation(() => { + throw new Error("drain failed"); + }); + + expect(() => { + harness.processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + }).not.toThrow(); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + + expect(harness.warn).toHaveBeenCalledWith(expect.stringContaining("drain failed")); + expect(harness.getLiveRunningAgentCounts).not.toHaveBeenCalled(); + harness.dispose(); + }); + + it("keeps the restart pending when the host initially declines it", async () => { + vi.useFakeTimers(); + const harness = createHarness(0); + harness.requestRestart.mockReturnValueOnce(false).mockReturnValueOnce(true); + + harness.processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + await vi.advanceTimersByTimeAsync(0); + expect(harness.requestRestart).toHaveBeenCalledTimes(1); + + await vi.advanceTimersByTimeAsync(1_000); + expect(harness.requestRestart).toHaveBeenCalledTimes(2); + harness.dispose(); + }); + + it("does not bind when the caller safety gate is disabled", () => { + const processEvents = new EventEmitter(); + const requestRestart = vi.fn(); + const beginDrain = vi.fn(); + const getLiveRunningAgentCounts = vi.fn(); + + const dispose = registerDevSourceRestart({ + enabled: false, + processEvents, + beginDrain, + getLiveRunningAgentCounts, + requestRestart, + }); + processEvents.emit("message", { type: DEV_SOURCE_CHANGE_MESSAGE }); + + expect(getLiveRunningAgentCounts).not.toHaveBeenCalled(); + expect(beginDrain).not.toHaveBeenCalled(); + expect(requestRestart).not.toHaveBeenCalled(); + dispose(); + }); +}); diff --git a/packages/cli/src/commands/dashboard.ts b/packages/cli/src/commands/dashboard.ts index 70f400cfcf..3a808f483c 100644 --- a/packages/cli/src/commands/dashboard.ts +++ b/packages/cli/src/commands/dashboard.ts @@ -145,6 +145,10 @@ import { handleOpencodeGoApiKeySaved, syncStartupModels } from "./startup-model- import { DashboardTUI, DashboardLogSink, isTTYAvailable, type SystemInfo, type GitStatus, type GitCommit, type GitCommitDetail, type GitBranch, type GitWorktree, type FileEntry, type FileReadResult, type TaskStep as TUITaskStep, type TaskLogEntry as TUITaskLogEntry, type TaskDetailData, type TaskEvent } from "./dashboard-tui/index.js"; import { DASHBOARD_STARTUP_STATUS, runTuiStartupPrelude } from "./dashboard-startup-chain.js"; import { phaseTime } from "../startup-phase.js"; +import { + DEV_SOURCE_RESTART_ARMED_MESSAGE, + registerDevSourceRestart, +} from "./dev-source-restart.js"; // Re-export for backward compatibility with tests export { promptForPort }; @@ -1295,6 +1299,23 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?: getRecent: (limit?: number) => logSink.getRecentEntries(limit), subscribe: (listener: (entry: import("./dashboard-tui/log-ring-buffer.js").LogEntry) => void) => logSink.subscribeEntries(listener), }; + const bindDevSourceRestart = (centralCore: CentralCore, beginDrain: () => void) => { + disposeCallbacks.push(registerDevSourceRestart({ + enabled: process.env.FUSION_DEV_WATCH === "1" + && systemControlForServer.supervised + && Boolean(systemControlForServer.sourceWorkspaceRoot), + beginDrain, + notifyArmed: () => { + process.send?.({ type: DEV_SOURCE_RESTART_ARMED_MESSAGE }); + }, + getLiveRunningAgentCounts: () => centralCore.getLiveRunningAgentCounts(), + requestRestart: (reason) => requestSelfRestart?.(reason) ?? false, + logger: { + log: (message) => logSink.log(message, "dashboard"), + warn: (message) => logSink.warn(message, "dashboard"), + }, + })); + }; /* * FNXC:DashboardShutdown 2026-06-27-10:32: @@ -2440,6 +2461,7 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?: }, 300); return true; }; + bindDevSourceRestart(centralCoreForEngine, () => engineManager.beginDrain()); registerHandler(process, "SIGINT", () => void shutdown("SIGINT")); registerHandler(process, "SIGTERM", () => void shutdown("SIGTERM")); @@ -2775,6 +2797,12 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?: }, 300); return true; }; + if (centralCoreForMesh) { + bindDevSourceRestart(centralCoreForMesh, () => { + triggerScheduler?.stop(); + heartbeatMonitorImpl?.stop(); + }); + } registerHandler(process, "SIGINT", () => void devShutdown("SIGINT")); registerHandler(process, "SIGTERM", () => void devShutdown("SIGTERM")); diff --git a/packages/cli/src/commands/dev-source-restart.ts b/packages/cli/src/commands/dev-source-restart.ts new file mode 100644 index 0000000000..5609ea7a47 --- /dev/null +++ b/packages/cli/src/commands/dev-source-restart.ts @@ -0,0 +1,121 @@ +import type { EventEmitter } from "node:events"; + +export const DEV_SOURCE_CHANGE_MESSAGE = "fusion:dev-source-changed"; +export const DEV_SOURCE_RESTART_ARMED_MESSAGE = "fusion:dev-source-restart-armed"; + +interface DevSourceChangeMessage { + type: typeof DEV_SOURCE_CHANGE_MESSAGE; +} + +interface DevSourceRestartOptions { + enabled: boolean; + processEvents?: Pick; + notifyArmed?: () => void; + beginDrain: () => Promise | void; + getLiveRunningAgentCounts: () => Promise<{ currentlyActive: number }>; + requestRestart: (reason: string) => boolean; + logger?: Pick; + recheckIntervalMs?: number; +} + +function isDevSourceChangeMessage(message: unknown): message is DevSourceChangeMessage { + return Boolean(message) + && typeof message === "object" + && (message as { type?: unknown }).type === DEV_SOURCE_CHANGE_MESSAGE; +} + +/** + * FNXC:DevEngineWatch 2026-08-04-01:25: + * A source edit is an intentional restart, not a crash or operating-system + * signal. Wait for live top-level agents to drain, then enter the dashboard's + * canonical exit-86 path so its existing graceful teardown and supervisor own + * the lifecycle. Liveness read failures defer rather than aborting active work. + */ +export function registerDevSourceRestart({ + enabled, + processEvents = process, + notifyArmed, + beginDrain, + getLiveRunningAgentCounts, + requestRestart, + logger = console, + recheckIntervalMs = 15_000, +}: DevSourceRestartOptions): () => void { + if (!enabled) return () => undefined; + + let disposed = false; + let pending = false; + let checking = false; + let restartAccepted = false; + let deferredForActiveWork = false; + let drainStarted = false; + let timer: ReturnType | undefined; + + const schedule = (delayMs: number) => { + if (disposed || timer) return; + timer = setTimeout(() => { + timer = undefined; + void checkRestartSafety(); + }, delayMs); + timer.unref?.(); + }; + + const checkRestartSafety = async () => { + if (disposed || restartAccepted || !pending || checking) return; + checking = true; + try { + const { currentlyActive } = await getLiveRunningAgentCounts(); + if (currentlyActive > 0) { + if (!deferredForActiveWork) { + const noun = currentlyActive === 1 ? "agent" : "agents"; + logger.log(`[fusion:dev] source restart deferred until ${currentlyActive} active ${noun} reach a safe boundary`); + deferredForActiveWork = true; + } + schedule(recheckIntervalMs); + return; + } + + if (deferredForActiveWork) { + logger.log("[fusion:dev] active work drained — applying pending source restart"); + deferredForActiveWork = false; + } + restartAccepted = requestRestart("dev-source-change"); + if (restartAccepted) { + pending = false; + } else { + logger.warn("[fusion:dev] source restart request was declined by the host lifecycle"); + schedule(recheckIntervalMs); + } + } catch (error) { + logger.warn(`[fusion:dev] could not verify restart safety: ${error instanceof Error ? error.message : String(error)}`); + schedule(recheckIntervalMs); + } finally { + checking = false; + } + }; + + const onMessage = (message: unknown) => { + if (restartAccepted || !isDevSourceChangeMessage(message)) return; + pending = true; + if (drainStarted) return; + drainStarted = true; + // FNXC:DevSourceRestart 2026-08-04-09:19: + // Enter the promise chain before draining so synchronous failures are logged instead of escaping IPC dispatch. + void Promise.resolve() + .then(beginDrain) + .then(() => schedule(0)) + .catch((error) => { + logger.warn(`[fusion:dev] could not begin source restart drain: ${error instanceof Error ? error.message : String(error)}`); + }); + }; + + processEvents.on("message", onMessage); + notifyArmed?.(); + return () => { + disposed = true; + pending = false; + if (timer) clearTimeout(timer); + timer = undefined; + processEvents.off("message", onMessage); + }; +} diff --git a/packages/core/src/central/central-core.ts b/packages/core/src/central/central-core.ts index 38547837aa..85e5d0bcee 100644 --- a/packages/core/src/central/central-core.ts +++ b/packages/core/src/central/central-core.ts @@ -2138,9 +2138,14 @@ export class CentralCore extends EventEmitter { incremented by real work and the durable counter it maintained was fiction. `getLiveRunningAgentCounts` SURVIVES and is now the only global readout. It is - TELEMETRY, not a limiter: it derives live per-project agent counts from the - registered side-effect-safe source so the dashboard can show "N running (all - projects)". Nothing gates on it. + TELEMETRY, not a capacity limiter: it derives live per-project agent counts + from the registered side-effect-safe source so the dashboard can show "N + running (all projects)". Task admission does not gate on it. + + FNXC:DevEngineWatch 2026-08-04-01:25: + The opt-in source-development restart also reads this conservative telemetry + to wait for active work to drain. That is shutdown safety, not capacity + arbitration; it never reserves slots or changes task admission. */ async getLiveRunningAgentCounts(options?: { source?: RunningAgentCountSource }): Promise { this.ensureInitialized(); diff --git a/packages/engine/src/__tests__/project-engine-manager.test.ts b/packages/engine/src/__tests__/project-engine-manager.test.ts index cdb0a2f059..359d6c9331 100644 --- a/packages/engine/src/__tests__/project-engine-manager.test.ts +++ b/packages/engine/src/__tests__/project-engine-manager.test.ts @@ -18,6 +18,7 @@ vi.mock("../project-engine.js", () => { ProjectEngine: vi.fn().mockImplementation(function (config: any) { return { start: vi.fn().mockResolvedValue(undefined), + beginDrain: vi.fn(), stop: vi.fn().mockResolvedValue(undefined), getTaskStore: vi.fn().mockReturnValue({ projectId: config.projectId }), getHeartbeatMonitor: vi.fn().mockReturnValue(undefined), @@ -295,6 +296,37 @@ describe("ProjectEngineManager", () => { }); describe("stopAll", () => { + it("closes admission on every engine before asynchronous shutdown work", async () => { + const manager = new ProjectEngineManager(centralCore); + await manager.startAll(); + const engineA = manager.getEngine("proj_aaa")!; + engineA.beginDrain = vi.fn(); + + manager.beginDrain(); + + expect(engineA.beginDrain).toHaveBeenCalledOnce(); + await expect(manager.ensureEngine("proj_aaa")).rejects.toThrow( + "ProjectEngineManager is stopped", + ); + expect(centralCore.markLocalNodeOffline).not.toHaveBeenCalled(); + }); + + it("continues draining other engines when one engine throws", async () => { + const manager = new ProjectEngineManager(centralCore); + await manager.startAll(); + const engineA = manager.getEngine("proj_aaa")!; + const engineB = manager.getEngine("proj_bbb")!; + engineA.beginDrain = vi.fn(() => { + throw new Error("drain failed"); + }); + engineB.beginDrain = vi.fn(); + + expect(() => manager.beginDrain()).not.toThrow(); + + expect(engineA.beginDrain).toHaveBeenCalledOnce(); + expect(engineB.beginDrain).toHaveBeenCalledOnce(); + }); + it("stops all engines and clears state", async () => { const manager = new ProjectEngineManager(centralCore); await manager.startAll(); diff --git a/packages/engine/src/__tests__/reliability-interactions/engine-stop-aborts-execution.test.ts b/packages/engine/src/__tests__/reliability-interactions/engine-stop-aborts-execution.test.ts index e7427a0c90..89468a0700 100644 --- a/packages/engine/src/__tests__/reliability-interactions/engine-stop-aborts-execution.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/engine-stop-aborts-execution.test.ts @@ -15,6 +15,55 @@ describe("FN-5403 reliability interactions: engine stop aborts execution", () => beforeEach(() => vi.useFakeTimers()); afterEach(() => vi.useRealTimers()); + it("closes every runtime admission source without aborting active execution", () => { + const runtime = new InProcessRuntime({ projectId: "p", workingDirectory: "/tmp", isolationMode: "in-process" } as any, {} as any) as any; + runtime.status = "active"; + runtime.workflowContinuationTimer = setInterval(() => undefined, 1_000); + runtime.selfHealingManager = { stop: vi.fn() }; + runtime.routineScheduler = { stop: vi.fn() }; + runtime.triggerScheduler = { stop: vi.fn() }; + runtime.stuckTaskDetector = { stop: vi.fn() }; + runtime.heartbeatMonitor = { stop: vi.fn() }; + runtime.triageProcessor = { stop: vi.fn() }; + runtime.scheduler = { stop: vi.fn() }; + runtime.missionAutopilot = { stop: vi.fn() }; + runtime.missionExecutionLoop = { stop: vi.fn() }; + runtime.executor = makeExecutor(); + + runtime.beginDrain(); + + expect(runtime.getStatus()).toBe("paused"); + expect(runtime.workflowContinuationTimer).toBeUndefined(); + for (const source of [ + runtime.selfHealingManager, + runtime.routineScheduler, + runtime.triggerScheduler, + runtime.stuckTaskDetector, + runtime.heartbeatMonitor, + runtime.triageProcessor, + runtime.scheduler, + runtime.missionAutopilot, + runtime.missionExecutionLoop, + ]) { + expect(source.stop).toHaveBeenCalledOnce(); + } + expect(runtime.executor.abortAllInFlight).not.toHaveBeenCalled(); + }); + + it("continues closing admission sources when one stop throws", () => { + const runtime = new InProcessRuntime({ projectId: "p", workingDirectory: "/tmp", isolationMode: "in-process" } as any, {} as any) as any; + runtime.status = "active"; + runtime.selfHealingManager = { stop: vi.fn(() => { + throw new Error("stop failed"); + }) }; + runtime.routineScheduler = { stop: vi.fn() }; + + expect(() => runtime.beginDrain()).not.toThrow(); + + expect(runtime.selfHealingManager.stop).toHaveBeenCalledOnce(); + expect(runtime.routineScheduler.stop).toHaveBeenCalledOnce(); + }); + it("FN-5403: engine stop aborts executor AI sessions before drain completes", async () => { const runtime = new InProcessRuntime({ projectId: "p", workingDirectory: "/tmp", isolationMode: "in-process" } as any, {} as any) as any; let aborted = false; diff --git a/packages/engine/src/project-engine-manager.ts b/packages/engine/src/project-engine-manager.ts index 9dacff4f96..769d13030b 100644 --- a/packages/engine/src/project-engine-manager.ts +++ b/packages/engine/src/project-engine-manager.ts @@ -256,8 +256,8 @@ export class ProjectEngineManager { runtimeLog.log(`Engine startup complete: ${started} started, ${failed} failed`); } - /** Gracefully stop all engines and reconciliation. */ - async stopAll(): Promise { + /** Close admission on every engine and stop reconciliation without tearing engines down. */ + beginDrain(): void { this.stopped = true; this.reconciliationStopped = true; @@ -267,6 +267,24 @@ export class ProjectEngineManager { this.reconciliationInterval = null; } + for (const engine of this.engines.values()) { + try { + engine.beginDrain?.(); + } catch (error) { + runtimeLog.warn( + `Engine drain error: ${error instanceof Error ? error.message : String(error)}`, + ); + } + } + for (const starting of this.starting.values()) { + void starting.then((engine) => engine.beginDrain?.()).catch(() => undefined); + } + } + + /** Gracefully stop all engines and reconciliation. */ + async stopAll(): Promise { + this.beginDrain(); + /* FNXC:PostgresResourceLifecycle 2026-07-14-18:42: Project runtimes own the PostgreSQL pools CentralCore may have adopted. Persist mesh-offline state before stopping any engine so runtime backend shutdown cannot race the final central write against a closed pool. diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index 5275d83f11..b3ca6d60a1 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -1459,6 +1459,12 @@ export class ProjectEngine { runtimeLog.log(`ProjectEngine stopped for ${this.config.projectId}`); } + /** Stop new lifecycle admission while allowing active work to finish. */ + beginDrain(): void { + this.shuttingDown = true; + this.runtime.beginDrain(); + } + // ── Public accessors ── /** Get the underlying InProcessRuntime. */ diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index c50046ea96..c4b4fff5fe 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -1934,6 +1934,42 @@ export class InProcessRuntime } } + /** + * Close every process-local admission source without aborting work that is + * already running. Development source reload uses this boundary before it + * waits for the live-agent count to reach zero. + */ + beginDrain(): void { + if (this.status !== "active") return; + + this.setStatus("paused"); + if (this.workflowContinuationTimer) { + clearInterval(this.workflowContinuationTimer); + this.workflowContinuationTimer = undefined; + } + const admissionStops: Array void]> = [ + ["self-healing manager", () => this.selfHealingManager?.stop()], + ["routine scheduler", () => this.routineScheduler?.stop()], + ["trigger scheduler", () => this.triggerScheduler?.stop()], + ["stuck task detector", () => this.stuckTaskDetector?.stop()], + ["heartbeat monitor", () => this.heartbeatMonitor?.stop()], + ["triage processor", () => this.triageProcessor?.stop()], + ["scheduler", () => this.scheduler?.stop()], + ["mission autopilot", () => this.missionAutopilot?.stop()], + ["mission execution loop", () => this.missionExecutionLoop?.stop()], + ]; + for (const [label, stop] of admissionStops) { + try { + stop(); + } catch (error) { + runtimeLog.warn( + `Failed to stop ${label} while draining ${this.config.projectId}: ${error instanceof Error ? error.message : String(error)}`, + ); + } + } + runtimeLog.log(`InProcessRuntime draining for ${this.config.projectId}`); + } + /** * Stop the runtime with graceful shutdown. * diff --git a/scripts/dev-hmr.mjs b/scripts/dev-hmr.mjs index 7db9d2b2de..91d66692d8 100644 --- a/scripts/dev-hmr.mjs +++ b/scripts/dev-hmr.mjs @@ -11,8 +11,8 @@ * (serves app/ with HMR; proxies /api and WS to the API) * * Open the URL Vite prints (e.g. http://localhost:5173), NOT the API URL. - * Edits to packages/dashboard/app/** hot-reload. Edits to src/** (server - * code) still require restarting this script. + * Edits to packages/dashboard/app/** hot-reload in Vite. Runtime source edits + * gracefully restart the API/engine child while Vite stays available. * * Env: * FUSION_API_PORT API port (default 4050). Vite's proxy reads the same @@ -94,7 +94,7 @@ launch( "api", "32", "pnpm", - ["dev", "--prebuild=none", "dashboard", "--no-auth", "--port", API_PORT, "--host", "127.0.0.1"], + ["dev", "--watch", "--prebuild=none", "dashboard", "--no-auth", "--port", API_PORT, "--host", "127.0.0.1"], { cwd: repoRoot, env: { ...globalThis.process.env, FUSION_API_PORT: API_PORT } }, ); diff --git a/scripts/dev-with-memory-lib.mjs b/scripts/dev-with-memory-lib.mjs index 54cb0bf321..51b531faa2 100644 --- a/scripts/dev-with-memory-lib.mjs +++ b/scripts/dev-with-memory-lib.mjs @@ -17,6 +17,80 @@ export function buildDevNodeArgs({ ]; } +export function createDevWatchRestartCoordinator({ log = console.log, warn = console.warn } = {}) { + let child; + let armed = false; + let queued = false; + let pendingPaths = []; + + const requeuePaths = (changedPaths) => { + pendingPaths = [...new Set([...pendingPaths, ...changedPaths])]; + }; + + const sendRestart = (changedPaths) => { + if (!child?.connected) { + requeuePaths(changedPaths); + warn("[fusion:dev] source restart deferred; the engine child is not connected"); + return; + } + const preview = changedPaths.slice(0, 3).join(", "); + const remainder = Math.max(0, changedPaths.length - 3); + log(`[fusion:dev] source changed (${preview}${remainder > 0 ? ` +${remainder} more` : ""}) — restart queued…`); + queued = true; + try { + child.send({ type: "fusion:dev-source-changed" }, (error) => { + if (!error) return; + queued = false; + requeuePaths(changedPaths); + warn(`[fusion:dev] source restart message failed: ${error.message}`); + }); + } catch (error) { + queued = false; + requeuePaths(changedPaths); + warn(`[fusion:dev] source restart message failed: ${error instanceof Error ? error.message : String(error)}`); + } + }; + + return { + attach(nextChild) { + child = nextChild; + armed = false; + queued = false; + }, + request(changedPaths) { + if (queued) return; + if (!armed) { + pendingPaths = [...new Set([...pendingPaths, ...changedPaths])]; + log("[fusion:dev] source changed while the engine child is starting — restart will queue when watch is armed"); + return; + } + if (!child?.connected) { + requeuePaths(changedPaths); + warn("[fusion:dev] source restart deferred; the engine child is not connected"); + return; + } + const paths = [...new Set([...pendingPaths, ...changedPaths])]; + pendingPaths = []; + sendRestart(paths); + }, + onMessage(message) { + if (!message || typeof message !== "object" || message.type !== "fusion:dev-source-restart-armed") return; + armed = true; + if (pendingPaths.length === 0) return; + const paths = pendingPaths; + pendingPaths = []; + sendRestart(paths); + }, + detach(nextChild) { + if (child !== nextChild) return false; + const sourceRestart = queued; + child = undefined; + armed = false; + return sourceRestart; + }, + }; +} + const VALID_PREBUILD_MODES = new Set(["auto", "none", "client", "full"]); export function normalizePrebuildMode(value) { @@ -50,6 +124,8 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) { const inspectFlags = []; const args = []; let requestedPrebuild = env.FUSION_DEV_PREBUILD ?? "auto"; + let watchSource = env.FUSION_DEV_WATCH === "1"; + let watchSourceFromFlag = false; for (let i = 0; i < rawArgs.length; i += 1) { const arg = rawArgs[i]; @@ -78,6 +154,12 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) { continue; } + if (arg === "--watch") { + watchSource = true; + watchSourceFromFlag = true; + continue; + } + args.push(arg); } @@ -85,6 +167,8 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) { inspectFlags, args, requestedPrebuild: normalizePrebuildMode(requestedPrebuild), + watchSource, + watchSourceFromFlag, }; } diff --git a/scripts/dev-with-memory.mjs b/scripts/dev-with-memory.mjs index 6ed6831bcd..3c62ff84ec 100644 --- a/scripts/dev-with-memory.mjs +++ b/scripts/dev-with-memory.mjs @@ -11,10 +11,12 @@ import { buildForwardedDevArgs, buildDevNodeArgs, + createDevWatchRestartCoordinator, getPrebuildCommand, parseDevWrapperArgs, resolvePrebuildMode, } from "./dev-with-memory-lib.mjs"; +import { createDevSourceWatcher } from "./lib/dev-source-watch.mjs"; // Set increased heap size (8GB) to prevent OOM during initial build/start const MEMORY_MB = process.env.FUSION_DEV_MEMORY_MB || "8192"; @@ -29,7 +31,8 @@ try { console.error(error instanceof Error ? error.message : String(error)); process.exit(1); } -const { inspectFlags, args, requestedPrebuild } = parsedArgs; +const { inspectFlags, args, requestedPrebuild, watchSourceFromFlag } = parsedArgs; +let { watchSource } = parsedArgs; // NODE_OPTIONS is shared with every spawned node process (build + run + // agents). Heap size belongs here. Inspector flags do NOT — see comment above. @@ -41,6 +44,13 @@ process.env.NODE_OPTIONS = nodeOptions; // builds default to 127.0.0.1; this override only applies when starting // the dashboard via `pnpm dev dashboard` and only if no --host was passed. const forwardedArgs = buildForwardedDevArgs(args); +if (watchSource && forwardedArgs[0] !== "dashboard") { + if (watchSourceFromFlag) { + console.error("[fusion:dev] --watch is supported for the dashboard engine process only"); + process.exit(1); + } + watchSource = false; +} const prebuildMode = resolvePrebuildMode(requestedPrebuild, forwardedArgs); const prebuildCommand = getPrebuildCommand(prebuildMode); @@ -71,6 +81,18 @@ is what makes the dashboard advertise restart support. Any other exit code propagates unchanged (no crash-restart loop here — `--supervise` owns that). */ const RESTART_EXIT_CODE = 86; +let appChild; +let sourceWatcher; +const watchRestart = createDevWatchRestartCoordinator(); + +function ensureSourceWatcher() { + if (!watchSource || sourceWatcher) return; + sourceWatcher = createDevSourceWatcher({ + rootDir: process.cwd(), + onRestart: (paths) => watchRestart.request(paths), + }); + console.log(`[fusion:dev] source watch active (${sourceWatcher.watchedPaths.join(", ")})`); +} function runApp(extraArgs) { const tsx = spawn(process.execPath, buildDevNodeArgs({ @@ -80,22 +102,46 @@ function runApp(extraArgs) { entry: ENTRY, args: extraArgs, }), { - stdio: "inherit", + stdio: watchSource ? ["inherit", "inherit", "inherit", "ipc"] : "inherit", // FNXC:SystemPanel 2026-07-25-10:05: stamp the supervisor pid alongside the // flag so the child can tell a real supervising parent from an inherited // copy of the variable (see hasLiveSupervisingParent in commands/dashboard.ts). - env: { ...process.env, FUSION_RESTART_SUPERVISED: "1", FUSION_SUPERVISOR_PID: String(process.pid) }, + env: { + ...process.env, + FUSION_RESTART_SUPERVISED: "1", + FUSION_SUPERVISOR_PID: String(process.pid), + ...(watchSource ? { FUSION_DEV_WATCH: "1" } : {}), + }, }); + appChild = tsx; + watchRestart.attach(tsx); + tsx.on("message", (message) => watchRestart.onMessage(message)); + ensureSourceWatcher(); tsx.on("close", (c) => { + const sourceRestart = watchRestart.detach(tsx); + if (appChild === tsx) appChild = undefined; if (c === RESTART_EXIT_CODE) { console.log("[fusion:dev] restart requested — restarting…"); - runApp(extraArgs); + if (sourceRestart && prebuildCommand) { + runPrebuild(() => runApp(extraArgs)); + } else { + runApp(extraArgs); + } return; } process.exit(c ?? 1); }); } +function runPrebuild(onSuccess) { + console.log(`[fusion] Running ${prebuildCommand.label} (${prebuildMode}) before source startup...`); + const build = spawn(prebuildCommand.command, prebuildCommand.args, { stdio: "inherit", shell: true }); + build.on("close", (code) => { + if (code !== 0) process.exit(code ?? 1); + onSuccess(); + }); +} + async function warnIfSourceVersionBehind() { if (process.env.FUSION_SKIP_STARTUP_UPDATE_PREFLIGHT === "1") { return; @@ -177,10 +223,5 @@ await warnIfDistStale(); if (!prebuildCommand) { runApp(forwardedArgs); } else { - console.log(`[fusion] Running ${prebuildCommand.label} (${prebuildMode}) before source startup...`); - const build = spawn(prebuildCommand.command, prebuildCommand.args, { stdio: "inherit", shell: true }); - build.on("close", (code) => { - if (code !== 0) process.exit(code ?? 1); - runApp(forwardedArgs); - }); + runPrebuild(() => runApp(forwardedArgs)); } diff --git a/scripts/lib/dev-source-watch.mjs b/scripts/lib/dev-source-watch.mjs new file mode 100644 index 0000000000..098482c9f6 --- /dev/null +++ b/scripts/lib/dev-source-watch.mjs @@ -0,0 +1,107 @@ +import { watch as fsWatch } from "node:fs"; +import { join } from "node:path"; + +export const DEFAULT_DEV_SOURCE_WATCH_PATHS = [ + "packages/core/src", + "packages/engine/src", + "packages/dashboard/src", + "packages/cli/src", +]; + +const RUNTIME_SOURCE_EXTENSION = /\.(?:[cm]?[jt]sx?|json)$/i; +const NON_RUNTIME_SEGMENTS = new Set([ + "__fixtures__", + "__generated__", + "__tests__", + "fixtures", + "generated", + "test", + "tests", +]); + +export function isRestartableSourceFile(filename) { + if (filename === null || filename === undefined) return false; + const normalized = String(filename).replaceAll("\\", "/"); + if (!normalized || !RUNTIME_SOURCE_EXTENSION.test(normalized)) return false; + if (/\.(?:test|spec)\.[cm]?[jt]sx?$/i.test(normalized)) return false; + if (/\.d\.[cm]?ts$/i.test(normalized)) return false; + return !normalized.split("/").some((segment) => NON_RUNTIME_SEGMENTS.has(segment)); +} + +/** + * FNXC:DevEngineWatch 2026-08-04-01:25: + * Watch only source roots that execute inside the long-lived dashboard/engine + * process. Tests, fixtures, generated declarations, build output, task state, + * and worktrees must not create restart storms. The supervisor owns process + * replacement; this helper only coalesces source events and reports paths. + */ +export function createDevSourceWatcher({ + rootDir, + onRestart, + watchPaths = DEFAULT_DEV_SOURCE_WATCH_PATHS, + debounceMs = 350, + maxWaitMs = 2_000, + watch = fsWatch, + logger = console, +}) { + const watchers = []; + const watchedPaths = []; + const changedPaths = new Set(); + let debounceTimer; + let firstChangeAt; + let closed = false; + + const scheduleRestart = (watchPath, filename) => { + if (closed || !isRestartableSourceFile(filename)) return; + const normalized = String(filename).replaceAll("\\", "/"); + changedPaths.add(`${watchPath}/${normalized}`); + firstChangeAt ??= Date.now(); + if (debounceTimer) clearTimeout(debounceTimer); + const elapsedMs = Date.now() - firstChangeAt; + const delayMs = Math.min(debounceMs, Math.max(0, maxWaitMs - elapsedMs)); + debounceTimer = setTimeout(() => { + debounceTimer = undefined; + firstChangeAt = undefined; + const paths = [...changedPaths]; + changedPaths.clear(); + Promise.resolve(onRestart(paths)).catch((error) => { + logger.warn(`[fusion:dev] source restart request failed: ${error instanceof Error ? error.message : String(error)}`); + }); + }, delayMs); + debounceTimer.unref?.(); + }; + + for (const watchPath of watchPaths) { + const absolutePath = join(rootDir, watchPath); + try { + const watcher = watch(absolutePath, { recursive: true }, (_eventType, filename) => { + scheduleRestart(watchPath, filename); + }); + watcher.on?.("error", (error) => { + logger.warn(`[fusion:dev] source watcher error for ${watchPath}: ${error instanceof Error ? error.message : String(error)}`); + }); + watchers.push(watcher); + watchedPaths.push(watchPath); + } catch (error) { + logger.warn(`[fusion:dev] could not watch ${watchPath}: ${error instanceof Error ? error.message : String(error)}`); + } + } + + return { + watchedPaths, + close() { + closed = true; + if (debounceTimer) clearTimeout(debounceTimer); + debounceTimer = undefined; + firstChangeAt = undefined; + changedPaths.clear(); + for (const watcher of watchers) { + try { + watcher.close(); + } catch (error) { + logger.warn(`[fusion:dev] source watcher close failed: ${error instanceof Error ? error.message : String(error)}`); + } + } + }, + }; +}