feat: restart the development engine on source changes (#3329)
## Summary Add an opt-in source-development loop that restarts the dashboard and engine when runtime TypeScript or JSON changes. Use `pnpm dev:watch`; `pnpm dev:hmr` now combines Vite UI HMR with the same supervised API/engine restart path. The watcher filters tests, fixtures, generated declarations, build output, and task state. It coalesces bursts with a two-second maximum wait, waits for the child to acknowledge its IPC listener, and rebuilds runtime dist artifacts before a source-triggered respawn. ## Safety model - Close scheduler, triage, heartbeat, mission, routine, self-healing, and merge admission before checking for active work. - Let already-running agents reach a safe boundary; do not mutate durable pause settings. - Enter the existing graceful exit-code-86 shutdown and supervised respawn path. - Retry failed liveness reads and declined restart requests instead of dropping the pending change. - Keep ordinary `pnpm dev` behavior unchanged; inherited watch state does not break nested non-dashboard development commands. A development restart intentionally replaces the dashboard process, so transient dashboard connections and project dev-server children reconnect or restart with it. Agent work is the protected boundary. ## Validation - `pnpm lint` - `pnpm test:gate` (753 tests passed across engine, core, PostgreSQL gate, and CI-shape suites) - Focused CLI watcher/restart/supervision suites: 40 tests passed - Focused engine drain/manager suites: 52 tests passed - `pnpm --filter @runfusion/fusion typecheck` - `pnpm --filter @fusion/engine typecheck` - `pnpm verify:fast` (13 steps passed, including CLI build and real health boot smoke) - Manual unsupported-command probe confirms explicit `--watch` fails clearly outside the dashboard command ## Post-Deploy Monitoring & Validation - Watch for `[fusion:dev] source changed`, `source restart deferred`, `active work drained`, and `restart requested` logs during the first watched development session. - Healthy behavior is one exit-86 respawn per edit batch, no interrupted active agents, refreshed dist artifacts, and a healthy dashboard after respawn. - Investigate repeated restart loops, watcher attachment warnings, declined restart retries, or liveness-read failures. - Immediate mitigation is to use ordinary `pnpm dev` without `--watch`; no production runtime behavior or durable setting needs rollback. - Validation owner: Fusion maintainers during the first source edit after merge. --- [](https://github.com/EveryInc/compound-engineering-plugin) <!-- This is an auto-generated comment: release notes by coderabbit.ai --> ## Summary by CodeRabbit * **New Features** * Added `pnpm dev:watch` to automatically restart development runtime processes when source files change. * Development restarts now wait for active work to finish, preventing new work from starting during the transition. * Enhanced `pnpm dev:hmr` with graceful runtime source restarts while keeping the dashboard available. * Rapid source changes are grouped to avoid unnecessary restarts. * **Documentation** * Updated development setup and contribution guides with the new watch workflow. <!-- end of auto-generated comment: release notes by coderabbit.ai -->
This commit is contained in:
7
.changeset/dev-engine-watch.md
Normal file
7
.changeset/dev-engine-watch.md
Normal file
@@ -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.
|
||||||
@@ -665,6 +665,7 @@ pnpm local --no-engine # Start local dashboard/API only
|
|||||||
pnpm build # Build default workspace packages (excludes desktop/mobile)
|
pnpm build # Build default workspace packages (excludes desktop/mobile)
|
||||||
pnpm build:all # Build all packages (including desktop/mobile)
|
pnpm build:all # Build all packages (including desktop/mobile)
|
||||||
pnpm dev dashboard # Run dashboard + AI engine
|
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 dev:ui # Dashboard only (no AI engine)
|
||||||
pnpm lint # Lint all packages
|
pnpm lint # Lint all packages
|
||||||
pnpm typecheck # Type-check all packages
|
pnpm typecheck # Type-check all packages
|
||||||
|
|||||||
@@ -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 --no-engine # fast local dashboard/API-only startup
|
||||||
pnpm local --prebuild <none|client|full> # local dashboard/API + AI engine startup with an explicit prebuild level
|
pnpm local --prebuild <none|client|full> # 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 # 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
|
FUSION_DEV_PREBUILD=full pnpm dev dashboard # production-like full workspace prebuild
|
||||||
pnpm dev:ui # dashboard dev server only
|
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 lint # lint all packages
|
||||||
pnpm test # merge-gate suite + changed-only affected tests (bounded; never full-suite)
|
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
|
pnpm test:gate # the merge gate: curated engine-core suite + CI-shape test
|
||||||
|
|||||||
@@ -34,6 +34,7 @@
|
|||||||
"smoke:boot": "node scripts/boot-smoke.mjs",
|
"smoke:boot": "node scripts/boot-smoke.mjs",
|
||||||
"local": "node scripts/start-local.mjs",
|
"local": "node scripts/start-local.mjs",
|
||||||
"dev": "node scripts/dev-with-memory.mjs",
|
"dev": "node scripts/dev-with-memory.mjs",
|
||||||
|
"dev:watch": "node scripts/dev-with-memory.mjs --watch",
|
||||||
"start": "node scripts/dev-with-memory.mjs",
|
"start": "node scripts/dev-with-memory.mjs",
|
||||||
"dev:ui": "pnpm --filter @fusion/dashboard dev",
|
"dev:ui": "pnpm --filter @fusion/dashboard dev",
|
||||||
"dev:hmr": "node scripts/dev-hmr.mjs",
|
"dev:hmr": "node scripts/dev-hmr.mjs",
|
||||||
|
|||||||
@@ -1,12 +1,21 @@
|
|||||||
import { describe, expect, it } from "vitest";
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||||
import {
|
import {
|
||||||
buildForwardedDevArgs,
|
buildForwardedDevArgs,
|
||||||
buildDevNodeArgs,
|
buildDevNodeArgs,
|
||||||
|
createDevWatchRestartCoordinator,
|
||||||
getPrebuildCommand,
|
getPrebuildCommand,
|
||||||
normalizePrebuildMode,
|
normalizePrebuildMode,
|
||||||
parseDevWrapperArgs,
|
parseDevWrapperArgs,
|
||||||
resolvePrebuildMode,
|
resolvePrebuildMode,
|
||||||
} from "../../../../scripts/dev-with-memory-lib.mjs";
|
} from "../../../../scripts/dev-with-memory-lib.mjs";
|
||||||
|
import {
|
||||||
|
createDevSourceWatcher,
|
||||||
|
isRestartableSourceFile,
|
||||||
|
} from "../../../../scripts/lib/dev-source-watch.mjs";
|
||||||
|
|
||||||
|
afterEach(() => {
|
||||||
|
vi.useRealTimers();
|
||||||
|
});
|
||||||
|
|
||||||
describe("buildDevNodeArgs", () => {
|
describe("buildDevNodeArgs", () => {
|
||||||
it("enables source-condition resolution before loading the tsx runtime", () => {
|
it("enables source-condition resolution before loading the tsx runtime", () => {
|
||||||
@@ -44,6 +53,26 @@ describe("dev-with-memory prebuild options", () => {
|
|||||||
inspectFlags: ["--inspect=9230"],
|
inspectFlags: ["--inspect=9230"],
|
||||||
args: ["dashboard", "--port", "4050"],
|
args: ["dashboard", "--port", "4050"],
|
||||||
requestedPrebuild: "none",
|
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"));
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
149
packages/cli/src/commands/__tests__/dev-source-restart.test.ts
Normal file
149
packages/cli/src/commands/__tests__/dev-source-restart.test.ts
Normal file
@@ -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();
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -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 { 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 { DASHBOARD_STARTUP_STATUS, runTuiStartupPrelude } from "./dashboard-startup-chain.js";
|
||||||
import { phaseTime } from "../startup-phase.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
|
// Re-export for backward compatibility with tests
|
||||||
export { promptForPort };
|
export { promptForPort };
|
||||||
@@ -1295,6 +1299,23 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
|||||||
getRecent: (limit?: number) => logSink.getRecentEntries(limit),
|
getRecent: (limit?: number) => logSink.getRecentEntries(limit),
|
||||||
subscribe: (listener: (entry: import("./dashboard-tui/log-ring-buffer.js").LogEntry) => void) => logSink.subscribeEntries(listener),
|
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:
|
* FNXC:DashboardShutdown 2026-06-27-10:32:
|
||||||
@@ -2440,6 +2461,7 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
|||||||
}, 300);
|
}, 300);
|
||||||
return true;
|
return true;
|
||||||
};
|
};
|
||||||
|
bindDevSourceRestart(centralCoreForEngine, () => engineManager.beginDrain());
|
||||||
registerHandler(process, "SIGINT", () => void shutdown("SIGINT"));
|
registerHandler(process, "SIGINT", () => void shutdown("SIGINT"));
|
||||||
registerHandler(process, "SIGTERM", () => void shutdown("SIGTERM"));
|
registerHandler(process, "SIGTERM", () => void shutdown("SIGTERM"));
|
||||||
|
|
||||||
@@ -2775,6 +2797,12 @@ export async function runDashboard(port: number, opts: { paused?: boolean; dev?:
|
|||||||
}, 300);
|
}, 300);
|
||||||
return true;
|
return true;
|
||||||
};
|
};
|
||||||
|
if (centralCoreForMesh) {
|
||||||
|
bindDevSourceRestart(centralCoreForMesh, () => {
|
||||||
|
triggerScheduler?.stop();
|
||||||
|
heartbeatMonitorImpl?.stop();
|
||||||
|
});
|
||||||
|
}
|
||||||
registerHandler(process, "SIGINT", () => void devShutdown("SIGINT"));
|
registerHandler(process, "SIGINT", () => void devShutdown("SIGINT"));
|
||||||
registerHandler(process, "SIGTERM", () => void devShutdown("SIGTERM"));
|
registerHandler(process, "SIGTERM", () => void devShutdown("SIGTERM"));
|
||||||
|
|
||||||
|
|||||||
121
packages/cli/src/commands/dev-source-restart.ts
Normal file
121
packages/cli/src/commands/dev-source-restart.ts
Normal file
@@ -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<EventEmitter, "on" | "off">;
|
||||||
|
notifyArmed?: () => void;
|
||||||
|
beginDrain: () => Promise<void> | void;
|
||||||
|
getLiveRunningAgentCounts: () => Promise<{ currentlyActive: number }>;
|
||||||
|
requestRestart: (reason: string) => boolean;
|
||||||
|
logger?: Pick<Console, "log" | "warn">;
|
||||||
|
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<typeof setTimeout> | 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);
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -2138,9 +2138,14 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
|
|||||||
incremented by real work and the durable counter it maintained was fiction.
|
incremented by real work and the durable counter it maintained was fiction.
|
||||||
|
|
||||||
`getLiveRunningAgentCounts` SURVIVES and is now the only global readout. It is
|
`getLiveRunningAgentCounts` SURVIVES and is now the only global readout. It is
|
||||||
TELEMETRY, not a limiter: it derives live per-project agent counts from the
|
TELEMETRY, not a capacity limiter: it derives live per-project agent counts
|
||||||
registered side-effect-safe source so the dashboard can show "N running (all
|
from the registered side-effect-safe source so the dashboard can show "N
|
||||||
projects)". Nothing gates on it.
|
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<RunningAgentCounts> {
|
async getLiveRunningAgentCounts(options?: { source?: RunningAgentCountSource }): Promise<RunningAgentCounts> {
|
||||||
this.ensureInitialized();
|
this.ensureInitialized();
|
||||||
|
|||||||
@@ -18,6 +18,7 @@ vi.mock("../project-engine.js", () => {
|
|||||||
ProjectEngine: vi.fn().mockImplementation(function (config: any) {
|
ProjectEngine: vi.fn().mockImplementation(function (config: any) {
|
||||||
return {
|
return {
|
||||||
start: vi.fn().mockResolvedValue(undefined),
|
start: vi.fn().mockResolvedValue(undefined),
|
||||||
|
beginDrain: vi.fn(),
|
||||||
stop: vi.fn().mockResolvedValue(undefined),
|
stop: vi.fn().mockResolvedValue(undefined),
|
||||||
getTaskStore: vi.fn().mockReturnValue({ projectId: config.projectId }),
|
getTaskStore: vi.fn().mockReturnValue({ projectId: config.projectId }),
|
||||||
getHeartbeatMonitor: vi.fn().mockReturnValue(undefined),
|
getHeartbeatMonitor: vi.fn().mockReturnValue(undefined),
|
||||||
@@ -295,6 +296,37 @@ describe("ProjectEngineManager", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
describe("stopAll", () => {
|
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 () => {
|
it("stops all engines and clears state", async () => {
|
||||||
const manager = new ProjectEngineManager(centralCore);
|
const manager = new ProjectEngineManager(centralCore);
|
||||||
await manager.startAll();
|
await manager.startAll();
|
||||||
|
|||||||
@@ -15,6 +15,55 @@ describe("FN-5403 reliability interactions: engine stop aborts execution", () =>
|
|||||||
beforeEach(() => vi.useFakeTimers());
|
beforeEach(() => vi.useFakeTimers());
|
||||||
afterEach(() => vi.useRealTimers());
|
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 () => {
|
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;
|
const runtime = new InProcessRuntime({ projectId: "p", workingDirectory: "/tmp", isolationMode: "in-process" } as any, {} as any) as any;
|
||||||
let aborted = false;
|
let aborted = false;
|
||||||
|
|||||||
@@ -256,8 +256,8 @@ export class ProjectEngineManager {
|
|||||||
runtimeLog.log(`Engine startup complete: ${started} started, ${failed} failed`);
|
runtimeLog.log(`Engine startup complete: ${started} started, ${failed} failed`);
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Gracefully stop all engines and reconciliation. */
|
/** Close admission on every engine and stop reconciliation without tearing engines down. */
|
||||||
async stopAll(): Promise<void> {
|
beginDrain(): void {
|
||||||
this.stopped = true;
|
this.stopped = true;
|
||||||
this.reconciliationStopped = true;
|
this.reconciliationStopped = true;
|
||||||
|
|
||||||
@@ -267,6 +267,24 @@ export class ProjectEngineManager {
|
|||||||
this.reconciliationInterval = null;
|
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<void> {
|
||||||
|
this.beginDrain();
|
||||||
|
|
||||||
/*
|
/*
|
||||||
FNXC:PostgresResourceLifecycle 2026-07-14-18:42:
|
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.
|
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.
|
||||||
|
|||||||
@@ -1459,6 +1459,12 @@ export class ProjectEngine {
|
|||||||
runtimeLog.log(`ProjectEngine stopped for ${this.config.projectId}`);
|
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 ──
|
// ── Public accessors ──
|
||||||
|
|
||||||
/** Get the underlying InProcessRuntime. */
|
/** Get the underlying InProcessRuntime. */
|
||||||
|
|||||||
@@ -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<readonly [string, () => 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.
|
* Stop the runtime with graceful shutdown.
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -11,8 +11,8 @@
|
|||||||
* (serves app/ with HMR; proxies /api and WS to the API)
|
* (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.
|
* 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
|
* Edits to packages/dashboard/app/** hot-reload in Vite. Runtime source edits
|
||||||
* code) still require restarting this script.
|
* gracefully restart the API/engine child while Vite stays available.
|
||||||
*
|
*
|
||||||
* Env:
|
* Env:
|
||||||
* FUSION_API_PORT API port (default 4050). Vite's proxy reads the same
|
* FUSION_API_PORT API port (default 4050). Vite's proxy reads the same
|
||||||
@@ -94,7 +94,7 @@ launch(
|
|||||||
"api",
|
"api",
|
||||||
"32",
|
"32",
|
||||||
"pnpm",
|
"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 } },
|
{ cwd: repoRoot, env: { ...globalThis.process.env, FUSION_API_PORT: API_PORT } },
|
||||||
);
|
);
|
||||||
|
|
||||||
|
|||||||
@@ -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"]);
|
const VALID_PREBUILD_MODES = new Set(["auto", "none", "client", "full"]);
|
||||||
|
|
||||||
export function normalizePrebuildMode(value) {
|
export function normalizePrebuildMode(value) {
|
||||||
@@ -50,6 +124,8 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) {
|
|||||||
const inspectFlags = [];
|
const inspectFlags = [];
|
||||||
const args = [];
|
const args = [];
|
||||||
let requestedPrebuild = env.FUSION_DEV_PREBUILD ?? "auto";
|
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) {
|
for (let i = 0; i < rawArgs.length; i += 1) {
|
||||||
const arg = rawArgs[i];
|
const arg = rawArgs[i];
|
||||||
@@ -78,6 +154,12 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (arg === "--watch") {
|
||||||
|
watchSource = true;
|
||||||
|
watchSourceFromFlag = true;
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
args.push(arg);
|
args.push(arg);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -85,6 +167,8 @@ export function parseDevWrapperArgs(rawArgs, env = process.env) {
|
|||||||
inspectFlags,
|
inspectFlags,
|
||||||
args,
|
args,
|
||||||
requestedPrebuild: normalizePrebuildMode(requestedPrebuild),
|
requestedPrebuild: normalizePrebuildMode(requestedPrebuild),
|
||||||
|
watchSource,
|
||||||
|
watchSourceFromFlag,
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -11,10 +11,12 @@
|
|||||||
import {
|
import {
|
||||||
buildForwardedDevArgs,
|
buildForwardedDevArgs,
|
||||||
buildDevNodeArgs,
|
buildDevNodeArgs,
|
||||||
|
createDevWatchRestartCoordinator,
|
||||||
getPrebuildCommand,
|
getPrebuildCommand,
|
||||||
parseDevWrapperArgs,
|
parseDevWrapperArgs,
|
||||||
resolvePrebuildMode,
|
resolvePrebuildMode,
|
||||||
} from "./dev-with-memory-lib.mjs";
|
} 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
|
// Set increased heap size (8GB) to prevent OOM during initial build/start
|
||||||
const MEMORY_MB = process.env.FUSION_DEV_MEMORY_MB || "8192";
|
const MEMORY_MB = process.env.FUSION_DEV_MEMORY_MB || "8192";
|
||||||
@@ -29,7 +31,8 @@ try {
|
|||||||
console.error(error instanceof Error ? error.message : String(error));
|
console.error(error instanceof Error ? error.message : String(error));
|
||||||
process.exit(1);
|
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 +
|
// NODE_OPTIONS is shared with every spawned node process (build + run +
|
||||||
// agents). Heap size belongs here. Inspector flags do NOT — see comment above.
|
// 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
|
// 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.
|
// the dashboard via `pnpm dev dashboard` and only if no --host was passed.
|
||||||
const forwardedArgs = buildForwardedDevArgs(args);
|
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 prebuildMode = resolvePrebuildMode(requestedPrebuild, forwardedArgs);
|
||||||
const prebuildCommand = getPrebuildCommand(prebuildMode);
|
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).
|
propagates unchanged (no crash-restart loop here — `--supervise` owns that).
|
||||||
*/
|
*/
|
||||||
const RESTART_EXIT_CODE = 86;
|
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) {
|
function runApp(extraArgs) {
|
||||||
const tsx = spawn(process.execPath, buildDevNodeArgs({
|
const tsx = spawn(process.execPath, buildDevNodeArgs({
|
||||||
@@ -80,22 +102,46 @@ function runApp(extraArgs) {
|
|||||||
entry: ENTRY,
|
entry: ENTRY,
|
||||||
args: extraArgs,
|
args: extraArgs,
|
||||||
}), {
|
}), {
|
||||||
stdio: "inherit",
|
stdio: watchSource ? ["inherit", "inherit", "inherit", "ipc"] : "inherit",
|
||||||
// FNXC:SystemPanel 2026-07-25-10:05: stamp the supervisor pid alongside the
|
// 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
|
// flag so the child can tell a real supervising parent from an inherited
|
||||||
// copy of the variable (see hasLiveSupervisingParent in commands/dashboard.ts).
|
// 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) => {
|
tsx.on("close", (c) => {
|
||||||
|
const sourceRestart = watchRestart.detach(tsx);
|
||||||
|
if (appChild === tsx) appChild = undefined;
|
||||||
if (c === RESTART_EXIT_CODE) {
|
if (c === RESTART_EXIT_CODE) {
|
||||||
console.log("[fusion:dev] restart requested — restarting…");
|
console.log("[fusion:dev] restart requested — restarting…");
|
||||||
runApp(extraArgs);
|
if (sourceRestart && prebuildCommand) {
|
||||||
|
runPrebuild(() => runApp(extraArgs));
|
||||||
|
} else {
|
||||||
|
runApp(extraArgs);
|
||||||
|
}
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
process.exit(c ?? 1);
|
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() {
|
async function warnIfSourceVersionBehind() {
|
||||||
if (process.env.FUSION_SKIP_STARTUP_UPDATE_PREFLIGHT === "1") {
|
if (process.env.FUSION_SKIP_STARTUP_UPDATE_PREFLIGHT === "1") {
|
||||||
return;
|
return;
|
||||||
@@ -177,10 +223,5 @@ await warnIfDistStale();
|
|||||||
if (!prebuildCommand) {
|
if (!prebuildCommand) {
|
||||||
runApp(forwardedArgs);
|
runApp(forwardedArgs);
|
||||||
} else {
|
} else {
|
||||||
console.log(`[fusion] Running ${prebuildCommand.label} (${prebuildMode}) before source startup...`);
|
runPrebuild(() => runApp(forwardedArgs));
|
||||||
const build = spawn(prebuildCommand.command, prebuildCommand.args, { stdio: "inherit", shell: true });
|
|
||||||
build.on("close", (code) => {
|
|
||||||
if (code !== 0) process.exit(code ?? 1);
|
|
||||||
runApp(forwardedArgs);
|
|
||||||
});
|
|
||||||
}
|
}
|
||||||
|
|||||||
107
scripts/lib/dev-source-watch.mjs
Normal file
107
scripts/lib/dev-source-watch.mjs
Normal file
@@ -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)}`);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user