feat(FN-1222): add distributed node discovery APIs

- Add discovery configuration and node metadata types in core, including exported interfaces for consumers
- Implement a NodeDiscovery service and wire it into CentralCore with lifecycle hooks and discovery helpers
- Add comprehensive unit coverage for NodeDiscovery behavior and CentralCore discovery integration
- Expose discovery REST endpoints in dashboard routes for config, node listing, and summary/health data
- Add frontend discovery API helpers and route-level tests, plus dependency updates in lockfiles
This commit is contained in:
gsxdsm
2026-04-08 14:58:34 -07:00
parent c1918c6819
commit af23d6ddb8
11 changed files with 1645 additions and 0 deletions

View File

@@ -3,6 +3,7 @@ import { mkdtempSync, rmSync, mkdirSync } from "node:fs";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { CentralCore } from "./central-core.js";
import { NodeDiscovery } from "./node-discovery.js";
import { NodeConnection, type ConnectionResult } from "./node-connection.js";
import * as systemMetrics from "./system-metrics.js";
import type {
@@ -11,6 +12,8 @@ import type {
CentralActivityLogEntry,
GlobalConcurrencyState,
SystemMetrics,
DiscoveryConfig,
DiscoveredNode,
} from "./types.js";
describe("CentralCore", () => {
@@ -874,6 +877,136 @@ describe("CentralCore", () => {
expect(emittedResult).toEqual(connectionResult);
});
it("should start and stop discovery lifecycle", async () => {
const startSpy = vi.spyOn(NodeDiscovery.prototype, "start").mockImplementation(() => {});
const stopSpy = vi.spyOn(NodeDiscovery.prototype, "stop").mockImplementation(() => {});
const config: DiscoveryConfig = {
broadcast: true,
listen: true,
serviceType: "_fusion._tcp",
port: 4040,
staleTimeoutMs: 300_000,
};
const discovery = await central.startDiscovery(config);
const local = (await central.listNodes()).find((node) => node.type === "local");
expect(discovery).toBeInstanceOf(NodeDiscovery);
expect(startSpy).toHaveBeenCalledWith(local?.id, local?.name);
expect(central.isDiscoveryActive()).toBe(true);
expect(central.getDiscoveryConfig()).toEqual(config);
central.stopDiscovery();
expect(stopSpy).toHaveBeenCalledTimes(1);
expect(central.isDiscoveryActive()).toBe(false);
expect(central.getDiscoveryConfig()).toBeNull();
});
it("should forward discovery events and track discovered nodes", async () => {
vi.spyOn(NodeDiscovery.prototype, "start").mockImplementation(() => {});
await central.startDiscovery({
broadcast: false,
listen: true,
serviceType: "_fusion._tcp",
port: 4040,
staleTimeoutMs: 300_000,
});
const discovery = (central as unknown as { nodeDiscovery: NodeDiscovery | null }).nodeDiscovery;
expect(discovery).toBeTruthy();
const discovered: DiscoveredNode = {
name: "mesh-peer",
host: "192.168.0.42",
port: 4040,
nodeType: "remote",
nodeId: "node_remote",
discoveredAt: "2026-04-01T12:00:00.000Z",
lastSeenAt: "2026-04-01T12:00:00.000Z",
};
let eventPayload: DiscoveredNode | undefined;
central.on("discovery:node:found", (node) => {
eventPayload = node;
});
discovery!.emit("node:discovered", discovered);
await Promise.resolve();
await Promise.resolve();
expect(eventPayload).toEqual(discovered);
expect(central.getDiscoveredNodes()).toEqual([discovered]);
const updated = {
...discovered,
lastSeenAt: "2026-04-01T12:01:00.000Z",
};
discovery!.emit("node:updated", updated);
await Promise.resolve();
await Promise.resolve();
expect(central.getDiscoveredNodes()).toEqual([updated]);
let lostName: string | undefined;
central.on("discovery:node:lost", (name) => {
lostName = name;
});
discovery!.emit("node:lost", discovered.name);
await Promise.resolve();
await Promise.resolve();
expect(lostName).toBe(discovered.name);
expect(central.getDiscoveredNodes()).toEqual([]);
});
it("should set registered nodes online/offline from discovery events", async () => {
vi.spyOn(NodeDiscovery.prototype, "start").mockImplementation(() => {});
const remote = await central.registerNode({
name: "remote-peer",
type: "remote",
url: "http://remote-peer:4040",
});
await central.startDiscovery({
broadcast: false,
listen: true,
serviceType: "_fusion._tcp",
port: 4040,
staleTimeoutMs: 300_000,
});
const discovery = (central as unknown as { nodeDiscovery: NodeDiscovery | null }).nodeDiscovery;
expect(discovery).toBeTruthy();
discovery!.emit("node:discovered", {
name: "remote-peer",
host: "192.168.0.22",
port: 4040,
nodeType: "remote",
nodeId: "node_remote_peer",
discoveredAt: "2026-04-01T12:00:00.000Z",
lastSeenAt: "2026-04-01T12:00:00.000Z",
} satisfies DiscoveredNode);
await Promise.resolve();
await Promise.resolve();
expect((await central.getNode(remote.id))?.status).toBe("online");
expect(central.getDiscoveredNodes()).toEqual([]);
discovery!.emit("node:lost", "remote-peer");
await Promise.resolve();
await Promise.resolve();
expect((await central.getNode(remote.id))?.status).toBe("offline");
});
it("should return empty discovered node list when discovery is inactive", () => {
expect(central.isDiscoveryActive()).toBe(false);
expect(central.getDiscoveredNodes()).toEqual([]);
expect(central.getDiscoveryConfig()).toBeNull();
});
it("should update node metrics and emit node:metrics:updated", async () => {
const local = (await central.listNodes()).find((node) => node.type === "local");
expect(local).toBeDefined();

View File

@@ -47,10 +47,13 @@ import type {
SystemMetrics,
NodeMeshState,
PeerNode,
DiscoveryConfig,
DiscoveredNode,
} from "./types.js";
import { CentralDatabase, toJson, toJsonNullable, fromJson } from "./central-db.js";
import { resolveGlobalDir } from "./global-settings.js";
import { NodeConnection } from "./node-connection.js";
import { NodeDiscovery } from "./node-discovery.js";
import { collectSystemMetrics } from "./system-metrics.js";
import type { ConnectionOptions, ConnectionResult } from "./node-connection.js";
@@ -85,6 +88,14 @@ export interface CentralCoreEvents {
"mesh:state:changed": [payload: { nodeId: string; state: NodeMeshState }];
/** Emitted after a remote node connection test completes */
"node:connection:test": [result: ConnectionResult];
/** Emitted when network discovery starts */
"discovery:started": [config: DiscoveryConfig];
/** Emitted when network discovery stops */
"discovery:stopped": [];
/** Emitted when a node is discovered via mDNS */
"discovery:node:found": [node: DiscoveredNode];
/** Emitted when a discovered node is lost */
"discovery:node:lost": [name: string];
/** Emitted when global concurrency state changes */
"concurrency:changed": [state: GlobalConcurrencyState];
}
@@ -95,6 +106,27 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
private db: CentralDatabase | null = null;
private readonly globalDir: string;
private initialized = false;
private nodeDiscovery: NodeDiscovery | null = null;
private discoveryConfig: DiscoveryConfig | null = null;
private readonly discoveredNodes = new Map<string, DiscoveredNode>();
private readonly onDiscoveryNodeDiscovered = (node: DiscoveredNode): void => {
void this.handleDiscoveryNodeDiscovered(node).catch((error) => {
console.warn("[central-core] Failed to process discovered node", error);
});
};
private readonly onDiscoveryNodeUpdated = (node: DiscoveredNode): void => {
void this.handleDiscoveryNodeUpdated(node).catch((error) => {
console.warn("[central-core] Failed to process discovery node update", error);
});
};
private readonly onDiscoveryNodeLost = (name: string): void => {
void this.handleDiscoveryNodeLost(name).catch((error) => {
console.warn("[central-core] Failed to process discovery node loss", error);
});
};
/**
* Create a CentralCore instance.
@@ -150,6 +182,10 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
* Closes database connections and releases resources.
*/
async close(): Promise<void> {
if (this.nodeDiscovery) {
this.stopDiscovery();
}
if (this.db) {
this.db.close();
this.db = null;
@@ -1031,6 +1067,84 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
return { result, node };
}
/**
* Start mDNS/DNS-SD node discovery for this process.
*/
async startDiscovery(config: DiscoveryConfig): Promise<NodeDiscovery> {
this.ensureInitialized();
if (this.nodeDiscovery) {
return this.nodeDiscovery;
}
const localNode = (await this.listNodes()).find((node) => node.type === "local");
if (!localNode) {
throw new Error("Local node not found");
}
this.discoveryConfig = {
...config,
};
this.discoveredNodes.clear();
const discovery = new NodeDiscovery(this.discoveryConfig);
this.nodeDiscovery = discovery;
discovery.on("node:discovered", this.onDiscoveryNodeDiscovered);
discovery.on("node:updated", this.onDiscoveryNodeUpdated);
discovery.on("node:lost", this.onDiscoveryNodeLost);
discovery.start(localNode.id, localNode.name);
this.emit("discovery:started", this.discoveryConfig);
return discovery;
}
/**
* Stop mDNS/DNS-SD node discovery.
*/
stopDiscovery(): void {
if (!this.nodeDiscovery) {
return;
}
const discovery = this.nodeDiscovery;
discovery.off("node:discovered", this.onDiscoveryNodeDiscovered);
discovery.off("node:updated", this.onDiscoveryNodeUpdated);
discovery.off("node:lost", this.onDiscoveryNodeLost);
discovery.stop();
this.nodeDiscovery = null;
this.discoveryConfig = null;
this.discoveredNodes.clear();
this.emit("discovery:stopped");
}
/**
* List currently discovered nodes.
*/
getDiscoveredNodes(): DiscoveredNode[] {
return Array.from(this.discoveredNodes.values());
}
/**
* Return whether discovery is currently active.
*/
isDiscoveryActive(): boolean {
return this.nodeDiscovery !== null;
}
/**
* Return active discovery config (if started).
*/
getDiscoveryConfig(): DiscoveryConfig | null {
if (!this.discoveryConfig) {
return null;
}
return { ...this.discoveryConfig };
}
/**
* Assign a project to a node.
*/
@@ -1597,6 +1711,41 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
// ── Private Helpers ─────────────────────────────────────────────────────
private async handleDiscoveryNodeDiscovered(node: DiscoveredNode): Promise<void> {
const existingNode = await this.getNodeByName(node.name);
if (!existingNode) {
this.discoveredNodes.set(node.name, node);
} else {
this.discoveredNodes.delete(node.name);
if (existingNode.status === "offline") {
await this.updateNode(existingNode.id, { status: "online" });
}
}
this.emit("discovery:node:found", node);
}
private async handleDiscoveryNodeUpdated(node: DiscoveredNode): Promise<void> {
if (!this.discoveredNodes.has(node.name)) {
return;
}
this.discoveredNodes.set(node.name, node);
}
private async handleDiscoveryNodeLost(name: string): Promise<void> {
this.discoveredNodes.delete(name);
const existingNode = await this.getNodeByName(name);
if (existingNode && existingNode.status !== "offline") {
await this.updateNode(existingNode.id, { status: "offline" });
}
this.emit("discovery:node:lost", name);
}
private ensureInitialized(): void {
if (!this.initialized || !this.db) {
throw new Error("CentralCore not initialized. Call init() first.");

View File

@@ -141,6 +141,7 @@ export { CentralCore } from "./central-core.js";
export type { CentralCoreEvents } from "./central-core.js";
export { CentralDatabase, createCentralDatabase } from "./central-db.js";
export { NodeConnection } from "./node-connection.js";
export { NodeDiscovery } from "./node-discovery.js";
export { collectSystemMetrics } from "./system-metrics.js";
export type {
ConnectionErrorType,
@@ -158,6 +159,9 @@ export type {
NodeConfig,
NodeMeshState,
NodeStatus,
NodeDiscoveryEvent,
DiscoveryConfig,
DiscoveredNode,
PeerNode,
ProjectHealth,
/** @deprecated Use RegisteredProject instead */

View File

@@ -0,0 +1,445 @@
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
import type { DiscoveryConfig, DiscoveredNode } from "./types.js";
interface MockBrowser {
on: ReturnType<typeof vi.fn>;
off: ReturnType<typeof vi.fn>;
stop: ReturnType<typeof vi.fn>;
emit: (event: string, ...args: unknown[]) => void;
}
const { BonjourMock, publishMock, findMock, destroyMock } = vi.hoisted(() => ({
BonjourMock: vi.fn(),
publishMock: vi.fn(),
findMock: vi.fn(),
destroyMock: vi.fn(),
}));
vi.mock("bonjour-service", () => ({
Bonjour: BonjourMock,
default: BonjourMock,
}));
import { NodeDiscovery } from "./node-discovery.js";
function createMockBrowser(): MockBrowser {
const listeners = new Map<string, Set<(...args: unknown[]) => void>>();
const browser: MockBrowser = {
on: vi.fn((event: string, callback: (...args: unknown[]) => void) => {
const callbacks = listeners.get(event) ?? new Set();
callbacks.add(callback);
listeners.set(event, callbacks);
return browser;
}),
off: vi.fn((event: string, callback: (...args: unknown[]) => void) => {
listeners.get(event)?.delete(callback);
return browser;
}),
stop: vi.fn(),
emit(event: string, ...args: unknown[]) {
for (const callback of listeners.get(event) ?? []) {
callback(...args);
}
},
};
return browser;
}
function defaultConfig(overrides: Partial<DiscoveryConfig> = {}): DiscoveryConfig {
return {
broadcast: false,
listen: false,
serviceType: "_fusion._tcp",
port: 4040,
staleTimeoutMs: 300_000,
...overrides,
};
}
function createService(overrides: Record<string, unknown> = {}): Record<string, unknown> {
return {
name: "peer-node",
port: 4040,
addresses: ["192.168.1.200"],
txt: {
nodeType: "remote",
nodeId: "node_remote_1",
},
...overrides,
};
}
describe("NodeDiscovery", () => {
let browser: MockBrowser;
let publishService: { stop: ReturnType<typeof vi.fn> };
beforeEach(() => {
vi.clearAllMocks();
vi.useRealTimers();
browser = createMockBrowser();
publishService = { stop: vi.fn() };
publishMock.mockReturnValue(publishService);
findMock.mockReturnValue(browser);
destroyMock.mockReturnValue(undefined);
BonjourMock.mockImplementation(() => ({
publish: publishMock,
find: findMock,
destroy: destroyMock,
}));
});
afterEach(() => {
vi.useRealTimers();
});
it("starts and stops broadcast mode", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true }));
discovery.start("node_local_1", "Local Node");
expect(publishMock).toHaveBeenCalledWith(
expect.objectContaining({
name: "Local Node",
type: "fusion",
protocol: "tcp",
port: 4040,
txt: expect.objectContaining({
nodeType: "local",
nodeId: "node_local_1",
version: expect.any(String),
}),
}),
);
discovery.stop();
expect(publishService.stop).toHaveBeenCalledTimes(1);
expect(destroyMock).toHaveBeenCalledTimes(1);
});
it("falls back to hostname when broadcast nodeName is empty", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true }));
discovery.start("node_local_1", " ");
expect(publishMock).toHaveBeenCalledWith(
expect.objectContaining({
name: expect.any(String),
}),
);
});
it("starts listen mode and emits node:discovered/node:lost", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
const lostHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.on("node:lost", lostHandler);
discovery.start("node_local_1", "Local");
expect(findMock).toHaveBeenCalledWith({
type: "fusion",
protocol: "tcp",
});
browser.emit("up", createService());
expect(discoveredHandler).toHaveBeenCalledTimes(1);
expect(discoveredHandler).toHaveBeenCalledWith(
expect.objectContaining({
name: "peer-node",
host: "192.168.1.200",
port: 4040,
nodeType: "remote",
nodeId: "node_remote_1",
discoveredAt: expect.any(String),
lastSeenAt: expect.any(String),
}),
);
browser.emit("down", createService());
expect(lostHandler).toHaveBeenCalledWith("peer-node");
});
it("ignores down events for unknown services", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const lostHandler = vi.fn();
discovery.on("node:lost", lostHandler);
discovery.start("node_local_1", "Local");
browser.emit("down", createService({ name: "missing-node" }));
expect(lostHandler).not.toHaveBeenCalled();
});
it("emits node:updated for already known services", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const updatedHandler = vi.fn();
discovery.on("node:updated", updatedHandler);
discovery.start("node_local_1", "Local");
browser.emit("up", createService());
browser.emit("up", createService({ addresses: ["192.168.1.201"] }));
expect(updatedHandler).toHaveBeenCalledTimes(1);
expect(updatedHandler).toHaveBeenCalledWith(
expect.objectContaining({ host: "192.168.1.201" }),
);
});
it("self-filters discovery events with matching nodeId", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.start("node_local_1", "Local");
browser.emit(
"up",
createService({
txt: {
nodeType: "local",
nodeId: "node_local_1",
},
}),
);
expect(discoveredHandler).not.toHaveBeenCalled();
expect(discovery.getDiscoveredNodes()).toEqual([]);
});
it("supports combined broadcast + listen mode", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true, listen: true }));
discovery.start("node_local_1", "Local Node");
expect(publishMock).toHaveBeenCalledTimes(1);
expect(findMock).toHaveBeenCalledTimes(1);
});
it("cleans up stale nodes after staleTimeoutMs", () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-04-08T12:00:00.000Z"));
const discovery = new NodeDiscovery(defaultConfig({ listen: true, staleTimeoutMs: 1_000 }));
const lostHandler = vi.fn();
discovery.on("node:lost", lostHandler);
discovery.start("node_local_1", "Local");
browser.emit("up", createService({ name: "stale-node" }));
expect(discovery.getDiscoveredNode("stale-node")).toBeDefined();
vi.advanceTimersByTime(61_000);
expect(discovery.getDiscoveredNode("stale-node")).toBeUndefined();
expect(lostHandler).toHaveBeenCalledWith("stale-node");
});
it("stop() is idempotent", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true, listen: true }));
discovery.start("node_local_1", "Local");
expect(() => discovery.stop()).not.toThrow();
expect(() => discovery.stop()).not.toThrow();
expect(publishService.stop).toHaveBeenCalledTimes(1);
expect(browser.stop).toHaveBeenCalledTimes(1);
expect(destroyMock).toHaveBeenCalledTimes(1);
});
it("startBroadcast() and startListening() are idempotent", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true, listen: true }));
discovery.startBroadcast("node_local_1", "Local");
discovery.startBroadcast("node_local_1", "Local");
discovery.startListening();
discovery.startListening();
expect(publishMock).toHaveBeenCalledTimes(1);
expect(findMock).toHaveBeenCalledTimes(1);
});
it("continues when broadcast publish throws and emits error", () => {
const error = new Error("multicast unavailable");
publishMock.mockImplementation(() => {
throw error;
});
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true }));
const errorHandler = vi.fn();
discovery.on("error", errorHandler);
expect(() => discovery.start("node_local_1", "Local")).not.toThrow();
expect(errorHandler).toHaveBeenCalledWith(error);
});
it("continues when listener startup throws and emits error", () => {
const error = new Error("listen failed");
findMock.mockImplementation(() => {
throw error;
});
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const errorHandler = vi.fn();
discovery.on("error", errorHandler);
expect(() => discovery.start("node_local_1", "Local")).not.toThrow();
expect(errorHandler).toHaveBeenCalledWith(error);
});
it("returns discovered nodes and specific discovered node by name", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
expect(discovery.getDiscoveredNodes()).toEqual([]);
expect(discovery.getDiscoveredNode("missing")).toBeUndefined();
discovery.start("node_local_1", "Local");
browser.emit("up", createService({ name: "peer-a" }));
browser.emit("up", createService({ name: "peer-b", addresses: ["192.168.1.201"] }));
const nodes = discovery.getDiscoveredNodes();
expect(nodes).toHaveLength(2);
expect(nodes.map((node) => node.name).sort()).toEqual(["peer-a", "peer-b"]);
expect(discovery.getDiscoveredNode("peer-b")).toEqual(
expect.objectContaining({ host: "192.168.1.201" }),
);
});
it("defaults nodeType to local when TXT nodeType is missing", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.start("node_local_1", "Local");
browser.emit(
"up",
createService({
txt: {
nodeId: "node_remote",
},
}),
);
expect(discoveredHandler).toHaveBeenCalledWith(
expect.objectContaining({ nodeType: "local" }),
);
});
it("uses fallback host resolution paths", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.start("node_local_1", "Local");
browser.emit(
"up",
createService({
addresses: ["fe80::1"],
referer: { address: "10.0.0.9" },
}),
);
expect(discoveredHandler).toHaveBeenCalledWith(
expect.objectContaining({ host: "10.0.0.9" }),
);
browser.emit(
"up",
createService({
name: "referer-only",
addresses: [],
referer: { address: "10.0.0.10" },
}),
);
expect(discovery.getDiscoveredNode("referer-only")).toEqual(
expect.objectContaining({ host: "10.0.0.10" }),
);
});
it("handles service type parsing for short format and uppercase", () => {
const shortTypeDiscovery = new NodeDiscovery(defaultConfig({ broadcast: true, serviceType: "fusion" }));
shortTypeDiscovery.start("node_local_1", "Local");
expect(publishMock).toHaveBeenCalledWith(
expect.objectContaining({ type: "fusion", protocol: "tcp" }),
);
const uppercaseDiscovery = new NodeDiscovery(
defaultConfig({
listen: true,
serviceType: "_FUSION._TCP",
}),
);
uppercaseDiscovery.start("node_local_1", "Local");
expect(findMock).toHaveBeenCalledWith({ type: "fusion", protocol: "tcp" });
});
it("stringifies numeric and boolean TXT values", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.start("node_local_1", "Local");
browser.emit(
"up",
createService({
txt: {
nodeType: true,
nodeId: 1234,
},
}),
);
expect(discoveredHandler).toHaveBeenCalledWith(
expect.objectContaining({ nodeId: "1234", nodeType: "local" }),
);
});
it("does not emit discovery when host resolution fails", () => {
const discovery = new NodeDiscovery(defaultConfig({ listen: true }));
const discoveredHandler = vi.fn();
discovery.on("node:discovered", discoveredHandler);
discovery.start("node_local_1", "Local");
browser.emit(
"up",
createService({
name: "no-host",
addresses: [],
referer: undefined,
}),
);
expect(discoveredHandler).not.toHaveBeenCalled();
expect(discovery.getDiscoveredNode("no-host")).toBeUndefined();
});
it("emits discovery start/stop lifecycle events", () => {
const discovery = new NodeDiscovery(defaultConfig({ broadcast: true }));
const startedHandler = vi.fn();
const stoppedHandler = vi.fn();
discovery.on("discovery:started", startedHandler);
discovery.on("discovery:stopped", stoppedHandler);
discovery.start("node_local_1", "Local");
discovery.stop();
expect(startedHandler).toHaveBeenCalledTimes(1);
expect(stoppedHandler).toHaveBeenCalledTimes(1);
});
});

View File

@@ -0,0 +1,341 @@
import { EventEmitter } from "node:events";
import os from "node:os";
import { Bonjour, type Browser, type Service } from "bonjour-service";
import type {
DiscoveredNode,
DiscoveryConfig,
NodeDiscoveryEvent,
} from "./types.js";
const DEFAULT_DISCOVERY_CONFIG: DiscoveryConfig = {
broadcast: true,
listen: true,
serviceType: "_fusion._tcp",
port: 4040,
staleTimeoutMs: 300_000,
};
const STALE_CLEANUP_INTERVAL_MS = 60_000;
const FUSION_VERSION = "0.1.0";
interface NodeDiscoveryEvents {
"node:discovered": [node: DiscoveredNode];
"node:updated": [node: DiscoveredNode];
"node:lost": [name: string];
"discovery:started": [];
"discovery:stopped": [];
error: [error: Error];
}
interface ParsedServiceType {
type: string;
protocol: "tcp" | "udp";
}
/**
* mDNS/DNS-SD discovery service for local-network Fusion nodes.
*/
export class NodeDiscovery extends EventEmitter<NodeDiscoveryEvents> {
private readonly config: DiscoveryConfig;
private bonjour: Bonjour | null = null;
private broadcastService: Service | null = null;
private browser: Browser | null = null;
private staleCleanupInterval: NodeJS.Timeout | null = null;
private readonly discoveredNodes = new Map<string, DiscoveredNode>();
private localNodeId: string | null = null;
private started = false;
constructor(config: DiscoveryConfig) {
super();
this.config = {
...DEFAULT_DISCOVERY_CONFIG,
...config,
staleTimeoutMs: Math.max(1_000, config.staleTimeoutMs),
};
}
start(nodeId: string, nodeName: string): void {
this.localNodeId = nodeId;
if (this.config.broadcast) {
this.startBroadcast(nodeId, nodeName);
}
if (this.config.listen) {
this.startListening();
this.startStaleCleanup();
}
this.started = true;
this.emit("discovery:started");
}
stop(): void {
if (this.config.listen) {
this.stopListening();
}
if (this.config.broadcast) {
this.stopBroadcast();
}
if (this.staleCleanupInterval) {
clearInterval(this.staleCleanupInterval);
this.staleCleanupInterval = null;
}
for (const name of this.discoveredNodes.keys()) {
this.emit("node:lost", name);
}
this.discoveredNodes.clear();
if (this.bonjour) {
this.bonjour.destroy();
this.bonjour = null;
}
this.localNodeId = null;
if (this.started) {
this.started = false;
this.emit("discovery:stopped");
}
}
startBroadcast(nodeId: string, nodeName: string): void {
this.localNodeId = nodeId;
if (this.broadcastService) {
return;
}
const bonjour = this.getBonjour();
const serviceType = this.parseServiceType(this.config.serviceType);
try {
this.broadcastService = bonjour.publish({
name: nodeName.trim() || os.hostname(),
type: serviceType.type,
protocol: serviceType.protocol,
port: this.config.port,
txt: {
nodeType: "local",
nodeId,
version: FUSION_VERSION,
},
});
} catch (error) {
this.warn("Failed to start mDNS broadcast", error);
this.emit("error", this.asError(error));
}
}
stopBroadcast(): void {
if (!this.broadcastService) {
return;
}
try {
this.broadcastService.stop?.();
} catch (error) {
this.warn("Failed to stop mDNS broadcast", error);
this.emit("error", this.asError(error));
} finally {
this.broadcastService = null;
}
}
startListening(): void {
if (this.browser) {
return;
}
const bonjour = this.getBonjour();
const serviceType = this.parseServiceType(this.config.serviceType);
try {
this.browser = bonjour.find({
type: serviceType.type,
protocol: serviceType.protocol,
});
this.browser.on("up", this.onServiceUp);
this.browser.on("down", this.onServiceDown);
} catch (error) {
this.warn("Failed to start mDNS listening", error);
this.emit("error", this.asError(error));
}
}
stopListening(): void {
if (!this.browser) {
return;
}
const browser = this.browser;
try {
browser.off("up", this.onServiceUp);
browser.off("down", this.onServiceDown);
browser.stop();
} catch (error) {
this.warn("Failed to stop mDNS listening", error);
this.emit("error", this.asError(error));
} finally {
this.browser = null;
}
}
getDiscoveredNodes(): DiscoveredNode[] {
return Array.from(this.discoveredNodes.values());
}
getDiscoveredNode(name: string): DiscoveredNode | undefined {
return this.discoveredNodes.get(name);
}
private startStaleCleanup(): void {
if (this.staleCleanupInterval) {
return;
}
this.staleCleanupInterval = setInterval(() => {
const now = Date.now();
const lostNodes: string[] = [];
for (const [name, node] of this.discoveredNodes.entries()) {
const lastSeenMs = Date.parse(node.lastSeenAt);
if (now - lastSeenMs > this.config.staleTimeoutMs) {
this.discoveredNodes.delete(name);
lostNodes.push(name);
}
}
for (const name of lostNodes) {
this.emit("node:lost", name);
}
}, STALE_CLEANUP_INTERVAL_MS);
}
private onServiceUp = (service: Service): void => {
const nodeId = this.getServiceText(service, "nodeId");
if (nodeId && this.localNodeId && nodeId === this.localNodeId) {
return;
}
const host = this.resolveServiceHost(service);
if (!host) {
return;
}
const existing = this.discoveredNodes.get(service.name);
const now = new Date().toISOString();
const nextNode: DiscoveredNode = {
name: service.name,
host,
port: service.port,
nodeType: this.getServiceText(service, "nodeType") === "remote" ? "remote" : "local",
nodeId,
discoveredAt: existing?.discoveredAt ?? now,
lastSeenAt: now,
};
this.discoveredNodes.set(service.name, nextNode);
this.emit(existing ? "node:updated" : "node:discovered", nextNode);
};
private onServiceDown = (service: Service): void => {
if (!this.discoveredNodes.has(service.name)) {
return;
}
this.discoveredNodes.delete(service.name);
this.emit("node:lost", service.name);
};
private getBonjour(): Bonjour {
if (!this.bonjour) {
this.bonjour = new Bonjour();
}
return this.bonjour;
}
private parseServiceType(serviceType: string): ParsedServiceType {
const normalized = serviceType.trim();
const mdnsPattern = /^_([^._]+)\._(tcp|udp)$/i;
const match = normalized.match(mdnsPattern);
if (match) {
return {
type: match[1].toLowerCase(),
protocol: match[2].toLowerCase() as "tcp" | "udp",
};
}
return {
type: normalized.replace(/^_/, "").split(".")[0].toLowerCase() || "fusion",
protocol: "tcp",
};
}
private getServiceText(service: Service, key: string): string | undefined {
const txt = service.txt;
if (!txt || typeof txt !== "object" || !(key in txt)) {
return undefined;
}
const value = (txt as Record<string, unknown>)[key];
if (typeof value === "string") {
return value;
}
if (typeof value === "number" || typeof value === "boolean") {
return String(value);
}
return undefined;
}
private resolveServiceHost(service: Service): string | undefined {
const addresses = service.addresses ?? [];
const ipv4 = addresses.find((address) => this.isIpv4(address));
if (ipv4) {
return ipv4;
}
const nonLinkLocal = addresses.find((address) => !address.startsWith("fe80:"));
if (nonLinkLocal) {
return nonLinkLocal;
}
const refererAddress = service.referer?.address;
if (refererAddress) {
return refererAddress;
}
return undefined;
}
private isIpv4(value: string): boolean {
return /^\d{1,3}(\.\d{1,3}){3}$/.test(value);
}
private warn(message: string, error?: unknown): void {
if (error) {
console.warn(`[node-discovery] ${message}`, error);
return;
}
console.warn(`[node-discovery] ${message}`);
}
private asError(error: unknown): Error {
if (error instanceof Error) {
return error;
}
return new Error(String(error));
}
}
export type { NodeDiscoveryEvent };

View File

@@ -1397,6 +1397,48 @@ export type ProjectStatus = "active" | "paused" | "errored" | "initializing";
/** Node connectivity/health status in the central registry */
export type NodeStatus = "online" | "offline" | "connecting" | "error";
/** A node discovered on the local network via mDNS/DNS-SD */
export interface DiscoveredNode {
/** Node name from the mDNS service instance name */
name: string;
/** Host address (IP address) */
host: string;
/** Port the Fusion dashboard is running on */
port: number;
/** Node type from TXT record */
nodeType: "local" | "remote";
/** Node ID from TXT record (if the node has registered itself) */
nodeId?: string;
/** When this node was first discovered */
discoveredAt: string;
/** When this node was last seen (updated on each mDNS response) */
lastSeenAt: string;
}
/** Configuration for network node discovery */
export interface DiscoveryConfig {
/** Whether to broadcast this node's presence on the network */
broadcast: boolean;
/** Whether to listen for other nodes on the network */
listen: boolean;
/** mDNS service type name (default: "_fusion._tcp") */
serviceType: string;
/** Port to advertise (defaults to the dashboard port) */
port: number;
/**
* How long (ms) to remember a discovered node after last seeing it.
* Default: 300000 (5 minutes).
*/
staleTimeoutMs: number;
}
export type NodeDiscoveryEvent =
| { type: "node:discovered"; node: DiscoveredNode }
| { type: "node:updated"; node: DiscoveredNode }
| { type: "node:lost"; name: string }
| { type: "discovery:started" }
| { type: "discovery:stopped" };
/** Host-level resource and uptime metrics reported by a node. */
export interface SystemMetrics {
/** CPU utilization percentage (0-100). */