feat(FN-1224): add peer gossip protocol for mesh network synchronization

- Add peer exchange types (MeshSyncState, PeerExchangeMessage, PeerGossipConfig) to core types
- Add peer merge and sync methods to CentralCore (mergeProject, syncWithPeer, getMeshState)
- Add mesh sync API endpoints (GET /api/mesh/state, POST /api/mesh/sync)
- Implement PeerExchangeService background gossip engine for periodic peer synchronization
- Add comprehensive tests for CentralCore mesh methods and PeerExchangeService
- Update memory docs with peer gossip protocol documentation
This commit is contained in:
gsxdsm
2026-04-09 13:22:29 -07:00
parent fc23e522a7
commit 008f664183
11 changed files with 1473 additions and 1 deletions

View File

@@ -11,6 +11,7 @@
- `HeartbeatMonitor.executeHeartbeat()` uses the Paperclip wake→check→work→exit model. The lazy `import("./pi.js")` pattern keeps pi SDK out of the module graph when only monitoring (not execution) is needed.
- Agent tool factories (`createTaskCreateTool`, `createTaskLogTool`) live in `agent-tools.ts` and are shared between `TaskExecutor` and `HeartbeatMonitor` to avoid duplication.
- Dashboard SSE clients (planning/subtask/mission interview) now use a shared keep-alive pattern: start a 25s `setInterval` in stream `onOpen` that `POST`s `/api/ai-sessions/:id/ping`, and always stop it on stream `close`, `complete`, and fatal errors.
- **Peer Gossip Protocol (FN-1224)**: Nodes exchange peer information via `POST /api/mesh/sync` endpoint. `PeerExchangeService` runs periodic sync cycles (default 60s interval) with all online remote nodes. `CentralCore.mergePeers()` handles peer data merging — new peers are registered via `registerGossipPeer()`, stale peers are updated with fresher data, and the local node is never overwritten. The service uses single-flight pattern to prevent overlapping syncs and refreshes local metrics before each sync.
## Conventions

View File

@@ -1188,6 +1188,241 @@ describe("CentralCore", () => {
expect(state.metrics).toEqual(metrics);
expect(state.knownPeers).toEqual([]);
});
describe("peer exchange methods", () => {
it("should register a gossip peer and preserve its nodeId", async () => {
const peerInfo = {
nodeId: "node_remote_gossip",
nodeName: "Gossip Peer",
nodeUrl: "https://gossip.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 3,
};
const registered = await central.registerGossipPeer(peerInfo);
expect(registered.id).toBe("node_remote_gossip");
expect(registered.name).toBe("Gossip Peer");
expect(registered.type).toBe("remote");
expect(registered.url).toBe("https://gossip.example.com");
expect(registered.status).toBe("online");
expect(registered.maxConcurrent).toBe(3);
// Verify it can be retrieved by the preserved ID
const fetched = await central.getNode("node_remote_gossip");
expect(fetched?.id).toBe("node_remote_gossip");
});
it("should handle duplicate peer names by appending suffix", async () => {
// First, register a local node with the same name
await central.registerNode({ name: "Same Name", type: "local" });
const peerInfo = {
nodeId: "node_same_1",
nodeName: "Same Name",
nodeUrl: "https://same1.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
};
const registered = await central.registerGossipPeer(peerInfo);
// Should have suffix added to avoid collision
expect(registered.name).toBe("Same Name-2");
});
it("should merge peers - add new peers", async () => {
const peerInfo = {
nodeId: "node_new_peer",
nodeName: "New Peer",
nodeUrl: "https://new-peer.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
};
const result = await central.mergePeers([peerInfo]);
expect(result.added).toContain("node_new_peer");
expect(result.updated).toEqual([]);
expect(await central.getNode("node_new_peer")).toBeDefined();
});
it("should merge peers - update stale peers", async () => {
// First, register a peer
const peerInfo = {
nodeId: "node_stale_peer",
nodeName: "Stale Peer",
nodeUrl: "https://stale-peer.example.com",
status: "offline" as const,
metrics: null,
lastSeen: "2026-04-01T11:00:00.000Z",
maxConcurrent: 2,
};
await central.registerGossipPeer(peerInfo);
// Now merge with fresher data
const fresherPeer = {
...peerInfo,
status: "online" as const,
lastSeen: "2026-04-01T12:30:00.000Z",
};
const result = await central.mergePeers([fresherPeer]);
expect(result.added).toEqual([]);
expect(result.updated).toContain("node_stale_peer");
const updated = await central.getNode("node_stale_peer");
expect(updated?.status).toBe("online");
});
it("should merge peers - skip fresher local data", async () => {
// First, register a peer
const peerInfo = {
nodeId: "node_fresher_local",
nodeName: "Fresher Local",
nodeUrl: "https://fresher-local.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T11:00:00.000Z",
maxConcurrent: 2,
};
await central.registerGossipPeer(peerInfo);
// Manually update to be fresher
await central.updateNode("node_fresher_local", {
status: "offline",
});
// Now merge with older data - should not update
const olderPeer = {
...peerInfo,
status: "online" as const,
lastSeen: "2026-04-01T10:00:00.000Z",
};
const result = await central.mergePeers([olderPeer]);
expect(result.updated).toEqual([]);
const updated = await central.getNode("node_fresher_local");
expect(updated?.status).toBe("offline");
});
it("should merge peers - never overwrite local node", async () => {
const local = (await central.listNodes()).find((node) => node.type === "local");
expect(local).toBeDefined();
// Create a fake peer info with the local node's ID
const fakePeerInfo = {
nodeId: local!.id,
nodeName: "Fake Local",
nodeUrl: "https://fake-local.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 10,
};
const result = await central.mergePeers([fakePeerInfo]);
// Should not add or update
expect(result.added).toEqual([]);
expect(result.updated).toEqual([]);
// Local node should be unchanged
const unchanged = await central.getNode(local!.id);
expect(unchanged?.maxConcurrent).toBe(4); // Default local node maxConcurrent
});
it("should merge peers - emit events correctly", async () => {
let gossipEvent: { nodeId: string; peer: unknown } | undefined;
let stateChangedEvent: { nodeId: string } | undefined;
central.on("gossip:peer:registered", (payload) => {
gossipEvent = payload;
});
central.on("mesh:state:changed", (payload) => {
stateChangedEvent = payload;
});
const peerInfo = {
nodeId: "node_event_peer",
nodeName: "Event Peer",
nodeUrl: "https://event-peer.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
};
await central.mergePeers([peerInfo]);
expect(gossipEvent?.nodeId).toBe("node_event_peer");
expect(stateChangedEvent?.nodeId).toBeDefined();
});
it("should merge peers - empty input returns empty result", async () => {
const result = await central.mergePeers([]);
expect(result.added).toEqual([]);
expect(result.updated).toEqual([]);
});
it("should get local peer info", async () => {
const peerInfo = await central.getLocalPeerInfo();
expect(peerInfo.nodeId).toBeDefined();
expect(peerInfo.nodeName).toBe("local");
expect(peerInfo.nodeUrl).toBe("");
expect(peerInfo.status).toBe("online");
expect(peerInfo.lastSeen).toBe("2026-04-01T12:00:00.000Z");
expect(peerInfo.maxConcurrent).toBe(4);
});
it("should get all known peer info", async () => {
// Register some peers
await central.registerGossipPeer({
nodeId: "node_all_peer_1",
nodeName: "All Peer 1",
nodeUrl: "https://all-peer-1.example.com",
status: "online" as const,
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
});
await central.registerGossipPeer({
nodeId: "node_all_peer_2",
nodeName: "All Peer 2",
nodeUrl: "https://all-peer-2.example.com",
status: "offline" as const,
metrics: null,
lastSeen: "2026-04-01T11:00:00.000Z",
maxConcurrent: 3,
});
const allPeers = await central.getAllKnownPeerInfo();
// Should include local node plus 2 registered peers
expect(allPeers.length).toBeGreaterThanOrEqual(3);
expect(allPeers.map((p) => p.nodeId)).toContain("node_all_peer_1");
expect(allPeers.map((p) => p.nodeId)).toContain("node_all_peer_2");
});
it("should get all known peer info - empty list", async () => {
// Don't register any peers, just check the local node
const allPeers = await central.getAllKnownPeerInfo();
// Should at least include the local node
expect(allPeers.length).toBeGreaterThanOrEqual(1);
expect(allPeers.some((p) => p.nodeName === "local")).toBe(true);
});
});
});
describe("project health", () => {

View File

@@ -46,6 +46,7 @@ import type {
NodeStatus,
SystemMetrics,
NodeMeshState,
PeerInfo,
PeerNode,
DiscoveryConfig,
DiscoveredNode,
@@ -78,7 +79,7 @@ export interface CentralCoreEvents {
"node:updated": [node: NodeConfig];
/** Emitted when node health status changes */
"node:health:changed": [node: NodeConfig];
/** Emitted when node metrics are updated */
/** Emitted when node metrics is updated */
"node:metrics:updated": [payload: { nodeId: string; metrics: SystemMetrics }];
/** Emitted when a mesh peer is added for a node */
"mesh:peer:added": [payload: { nodeId: string; peer: PeerNode }];
@@ -86,6 +87,8 @@ export interface CentralCoreEvents {
"mesh:peer:removed": [payload: { nodeId: string; peerNodeId: string }];
/** Emitted when a node mesh snapshot changes */
"mesh:state:changed": [payload: { nodeId: string; state: NodeMeshState }];
/** Emitted when a new node is discovered via gossip peer exchange */
"gossip:peer:registered": [payload: { nodeId: string; peer: PeerInfo }];
/** Emitted after a remote node connection test completes */
"node:connection:test": [result: ConnectionResult];
/** Emitted when network discovery starts */
@@ -563,6 +566,68 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
return node;
}
/**
* Register a remote peer node from gossip exchange.
*
* This method is used during peer merge to register nodes discovered via
* the gossip protocol. It preserves the remote node's ID (rather than
* generating a new one) so that cross-node lookups work correctly.
*
* @param peer — Peer info from the gossip exchange
* @returns The registered node
*/
async registerGossipPeer(peer: PeerInfo): Promise<NodeConfig> {
this.ensureInitialized();
const now = new Date().toISOString();
// Handle name uniqueness by appending suffix if needed
let name = peer.nodeName;
let suffix = 1;
while (true) {
const existing = await this.getNodeByName(name);
if (!existing) break;
suffix++;
name = `${peer.nodeName}-${suffix}`;
}
// Determine URL - use provided URL or empty string for local-style
const normalizedUrl = peer.nodeUrl || undefined;
const node: NodeConfig = {
id: peer.nodeId,
name,
type: "remote",
url: normalizedUrl,
status: peer.status,
capabilities: peer.capabilities,
systemMetrics: peer.metrics ?? undefined,
maxConcurrent: peer.maxConcurrent,
createdAt: now,
updatedAt: now,
};
this.db!.prepare(
`INSERT INTO nodes (id, name, type, url, status, capabilities, systemMetrics, maxConcurrent, createdAt, updatedAt)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`
).run(
node.id,
node.name,
node.type,
node.url ?? null,
node.status,
toJsonNullable(node.capabilities),
toJsonNullable(node.systemMetrics),
node.maxConcurrent,
node.createdAt,
node.updatedAt
);
this.db!.bumpLastModified();
this.emit("node:registered", node);
return node;
}
/**
* Unregister a runtime node.
*
@@ -1001,6 +1066,116 @@ export class CentralCore extends EventEmitter<CentralCoreEvents> {
return this.getMeshState(localNode.id);
}
/**
* Merge incoming peer information from a gossip exchange.
*
* This method processes a list of peers received from another node during
* the gossip protocol. It adds new peers, updates stale entries, and
* emits appropriate events for mesh state changes.
*
* @param incomingPeers — List of peer info from gossip exchange
* @returns Object with lists of added and updated node IDs
*/
async mergePeers(incomingPeers: PeerInfo[]): Promise<{ added: string[]; updated: string[] }> {
this.ensureInitialized();
const added: string[] = [];
const updated: string[] = [];
for (const peer of incomingPeers) {
const existing = await this.getNode(peer.nodeId);
if (!existing) {
// New peer - register it
const newNode = await this.registerGossipPeer(peer);
added.push(newNode.id);
this.emit("gossip:peer:registered", { nodeId: newNode.id, peer });
} else if (existing.type === "local") {
// Never overwrite the local node from incoming peer data
continue;
} else {
// Existing remote node - check if incoming data is fresher
const incomingLastSeen = new Date(peer.lastSeen);
const localUpdatedAt = new Date(existing.updatedAt);
if (incomingLastSeen > localUpdatedAt) {
// Incoming data is fresher - update the node
await this.updateNode(existing.id, {
status: peer.status,
url: peer.nodeUrl || undefined,
capabilities: peer.capabilities,
maxConcurrent: peer.maxConcurrent,
});
// Update metrics if provided
if (peer.metrics) {
await this.updateNodeMetrics(existing.id, peer.metrics);
}
updated.push(existing.id);
}
}
}
// Emit mesh state changed if any modifications were made
if (added.length > 0 || updated.length > 0) {
const localNode = await this.getLocalNode();
if (localNode) {
const state = await this.getMeshState(localNode.id);
this.emit("mesh:state:changed", { nodeId: localNode.id, state });
}
}
return { added, updated };
}
/**
* Get a PeerInfo snapshot of the local node for gossip transmission.
*
* @returns PeerInfo for the local node with current metrics
*/
async getLocalPeerInfo(): Promise<PeerInfo> {
this.ensureInitialized();
const localNode = await this.getLocalNode();
if (!localNode) {
throw new Error("Local node not found");
}
return {
nodeId: localNode.id,
nodeName: localNode.name,
nodeUrl: localNode.url || "",
status: localNode.status,
metrics: localNode.systemMetrics ?? null,
lastSeen: new Date().toISOString(),
capabilities: localNode.capabilities,
maxConcurrent: localNode.maxConcurrent,
};
}
/**
* Get PeerInfo snapshots for all known nodes.
*
* @returns Array of PeerInfo for all nodes in the registry
*/
async getAllKnownPeerInfo(): Promise<PeerInfo[]> {
this.ensureInitialized();
const nodes = await this.listNodes();
return nodes.map((node) => ({
nodeId: node.id,
nodeName: node.name,
nodeUrl: node.url || "",
status: node.status,
metrics: node.systemMetrics ?? null,
lastSeen: node.updatedAt,
capabilities: node.capabilities,
maxConcurrent: node.maxConcurrent,
}));
}
/**
* Test connectivity to a remote Fusion node without registering it.
*/

View File

@@ -164,7 +164,10 @@ export type {
NodeDiscoveryEvent,
DiscoveryConfig,
DiscoveredNode,
PeerInfo,
PeerNode,
PeerSyncRequest,
PeerSyncResponse,
ProjectHealth,
/** @deprecated Use RegisteredProject instead */
ProjectInfo,

View File

@@ -1549,6 +1549,52 @@ export interface MeshDiscovery {
discoveryVersion: number;
}
/** Lightweight snapshot of a known node suitable for gossip transmission. */
export interface PeerInfo {
/** Unique node identifier. */
nodeId: string;
/** Display name of the node. */
nodeName: string;
/** Base URL of the node (empty string for local nodes). */
nodeUrl: string;
/** Current node status. */
status: NodeStatus;
/** Latest system metrics snapshot, if available. */
metrics: SystemMetrics | null;
/** ISO timestamp of when this info was last updated. */
lastSeen: string;
/** Optional capabilities available on this node. */
capabilities?: AgentCapability[];
/** Maximum concurrent tasks/runtimes this node can host. */
maxConcurrent: number;
}
/** Request payload sent when a node initiates a peer sync. */
export interface PeerSyncRequest {
/** Node ID of the sender. */
senderNodeId: string;
/** Base URL of the sender node. */
senderNodeUrl: string;
/** List of peers known by the sender. */
knownPeers: PeerInfo[];
/** ISO timestamp of when this sync request was generated. */
timestamp: string;
}
/** Response payload returned after a peer sync exchange. */
export interface PeerSyncResponse {
/** Node ID of the responding node (local node). */
senderNodeId: string;
/** Base URL of the responding node. */
senderNodeUrl: string;
/** Full list of peers known by the responding node. */
knownPeers: PeerInfo[];
/** Peers in the local list that the sender didn't know about. */
newPeers: PeerInfo[];
/** ISO timestamp of when this response was generated. */
timestamp: string;
}
/** A runtime node that can host project execution (local machine or remote host) */
export interface NodeConfig {
/** Unique node ID (e.g., "node_abc123") */

View File

@@ -0,0 +1,379 @@
import { describe, it, expect, vi, beforeEach } from "vitest";
import { EventEmitter } from "node:events";
import type { Task } from "@fusion/core";
import { request } from "../test-request.js";
import { createServer } from "../server.js";
// Request helper type for the test-request module
type TestRequestFn = (
app: (req: import("node:http").IncomingMessage, res: import("node:http").ServerResponse) => void,
method: string,
path: string,
body?: string,
headers?: Record<string, string>
) => Promise<{ status: number; body: unknown; headers: Record<string, string | string[] | undefined> }>;
const mockInit = vi.fn().mockResolvedValue(undefined);
const mockClose = vi.fn().mockResolvedValue(undefined);
const mockMergePeers = vi.fn().mockResolvedValue({ added: [], updated: [] });
const mockGetAllKnownPeerInfo = vi.fn().mockResolvedValue([]);
const mockGetLocalPeerInfo = vi.fn();
const mockGetNode = vi.fn();
const mockUpdateNode = vi.fn();
const mockGetLocalNode = vi.fn();
vi.mock("@fusion/core", async () => {
const actual = await vi.importActual<typeof import("@fusion/core")>("@fusion/core");
return {
...actual,
CentralCore: vi.fn().mockImplementation(() => ({
init: mockInit,
close: mockClose,
mergePeers: mockMergePeers,
getAllKnownPeerInfo: mockGetAllKnownPeerInfo,
getLocalPeerInfo: mockGetLocalPeerInfo,
getNode: mockGetNode,
updateNode: mockUpdateNode,
getLocalNode: mockGetLocalNode,
})),
};
});
class MockStore extends EventEmitter {
getRootDir(): string {
return "/tmp/fn-1224";
}
getDatabase() {
return {
exec: vi.fn(),
prepare: vi.fn().mockReturnValue({ run: vi.fn().mockReturnValue({ changes: 0 }), get: vi.fn(), all: vi.fn().mockReturnValue([]) }),
};
}
getMissionStore() {
return {
listMissions: vi.fn().mockResolvedValue([]),
createMission: vi.fn(),
getMission: vi.fn(),
updateMission: vi.fn(),
deleteMission: vi.fn(),
listTemplates: vi.fn().mockResolvedValue([]),
createTemplate: vi.fn(),
getTemplate: vi.fn(),
updateTemplate: vi.fn(),
deleteTemplate: vi.fn(),
instantiateMission: vi.fn(),
};
}
async listTasks(): Promise<Task[]> {
return [];
}
}
function makePeerInfo(overrides: Partial<Record<string, unknown>> = {}) {
return {
nodeId: "node_peer_1",
nodeName: "Peer Node 1",
nodeUrl: "https://peer-1.example.com",
status: "online",
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
...overrides,
};
}
function makeNodeConfig(overrides: Partial<Record<string, unknown>> = {}) {
return {
id: "node_remote_1",
name: "Remote Node",
type: "remote",
url: "https://remote.example.com",
apiKey: undefined,
status: "online",
maxConcurrent: 2,
createdAt: "2026-04-01T10:00:00.000Z",
updatedAt: "2026-04-01T12:00:00.000Z",
...overrides,
};
}
describe("POST /api/mesh/sync", () => {
let app: ReturnType<typeof createServer>;
beforeEach(async () => {
vi.clearAllMocks();
// Reset mock implementations
mockInit.mockResolvedValue(undefined);
mockClose.mockResolvedValue(undefined);
mockMergePeers.mockResolvedValue({ added: [], updated: [] });
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockGetLocalPeerInfo.mockResolvedValue({
nodeId: "node_local",
nodeName: "local",
nodeUrl: "",
status: "online",
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 4,
});
mockGetNode.mockResolvedValue(undefined);
mockUpdateNode.mockResolvedValue({ id: "node_remote", status: "online" });
mockGetLocalNode.mockResolvedValue({
id: "node_local",
name: "local",
type: "local",
status: "online",
maxConcurrent: 4,
createdAt: "2026-04-01T10:00:00.000Z",
updatedAt: "2026-04-01T12:00:00.000Z",
});
const store = new MockStore();
app = createServer(store);
});
it("should merge peers and return sync response", async () => {
const peers = [makePeerInfo({ nodeId: "node_new" })];
const allKnownPeers = [
makePeerInfo({ nodeId: "node_local", nodeName: "local" }),
makePeerInfo({ nodeId: "node_new" }),
makePeerInfo({ nodeId: "node_existing" }),
];
mockMergePeers.mockResolvedValue({ added: ["node_new"], updated: [] });
mockGetAllKnownPeerInfo.mockResolvedValue(allKnownPeers);
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: peers,
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
expect(mockMergePeers).toHaveBeenCalledWith(peers);
expect(response.body).toMatchObject({
senderNodeId: "node_local",
knownPeers: allKnownPeers,
timestamp: expect.any(String),
});
expect((response.body as any).newPeers).toHaveLength(2); // node_local and node_existing (node_new was in knownPeers)
});
it("should reject missing senderNodeId with 400", async () => {
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({ knownPeers: [] }),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(400);
expect(response.body).toMatchObject({ error: "senderNodeId is required" });
});
it("should reject non-array knownPeers with 400", async () => {
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
knownPeers: "not-an-array",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(400);
expect(response.body).toMatchObject({ error: "knownPeers must be an array" });
});
it("should reject malformed peer entries with 400", async () => {
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
knownPeers: [
{ nodeId: "valid" }, // missing nodeName and status
],
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(400);
expect((response.body as any).error).toContain("Each knownPeers entry must have");
});
it("should update sender node status to online", async () => {
mockGetNode.mockResolvedValue(makeNodeConfig({ id: "node_remote", status: "offline" }));
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
expect(mockUpdateNode).toHaveBeenCalledWith("node_remote", { status: "online" });
});
it("should silently skip update if sender node not found", async () => {
mockGetNode.mockResolvedValue(undefined);
mockUpdateNode.mockRejectedValue(new Error("Node not found"));
// Should not throw, just silently skip
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_unknown",
senderNodeUrl: "https://unknown.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
});
it("should validate API key when sender has one configured", async () => {
mockGetNode.mockResolvedValue(makeNodeConfig({ apiKey: "secret-key" }));
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json", "Authorization": "Bearer wrong-key" }
);
expect(response.status).toBe(401);
expect(response.body).toMatchObject({ error: "Unauthorized" });
});
it("should accept request with correct API key", async () => {
mockGetNode.mockResolvedValue(makeNodeConfig({ apiKey: "correct-key" }));
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json", "Authorization": "Bearer correct-key" }
);
expect(response.status).toBe(200);
expect(mockMergePeers).toHaveBeenCalled();
});
it("should allow request without auth when sender has no API key", async () => {
mockGetNode.mockResolvedValue(makeNodeConfig({ apiKey: undefined }));
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
expect(mockMergePeers).toHaveBeenCalled();
});
it("should handle empty knownPeers array", async () => {
const localPeer = makePeerInfo({ nodeId: "node_local", nodeName: "local" });
mockGetAllKnownPeerInfo.mockResolvedValue([localPeer]);
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
expect(mockMergePeers).toHaveBeenCalledWith([]);
expect((response.body as any).newPeers).toHaveLength(1); // All local peers are "new" to sender
expect((response.body as any).newPeers[0].nodeId).toBe("node_local");
});
it("should compute newPeers correctly - sender knows some peers", async () => {
const allKnownPeers = [
makePeerInfo({ nodeId: "node_local", nodeName: "local" }),
makePeerInfo({ nodeId: "node_a" }),
makePeerInfo({ nodeId: "node_b" }),
makePeerInfo({ nodeId: "node_c" }),
];
mockGetAllKnownPeerInfo.mockResolvedValue(allKnownPeers);
// Sender knows node_a and node_b, but not node_c or node_local
const response = await request(
app,
"POST",
"/api/mesh/sync",
JSON.stringify({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [
makePeerInfo({ nodeId: "node_a" }),
makePeerInfo({ nodeId: "node_b" }),
],
timestamp: "2026-04-01T12:00:00.000Z",
}),
{ "Content-Type": "application/json" }
);
expect(response.status).toBe(200);
expect((response.body as any).newPeers).toHaveLength(2);
const newPeerIds = (response.body as any).newPeers.map((p: { nodeId: string }) => p.nodeId);
expect(newPeerIds).toContain("node_local");
expect(newPeerIds).toContain("node_c");
expect(newPeerIds).not.toContain("node_a");
expect(newPeerIds).not.toContain("node_b");
});
});

View File

@@ -10671,6 +10671,90 @@ Output ONLY the prompt text (no markdown, no explanations).`;
}
});
/**
* POST /api/mesh/sync
* Exchange peer information with another node for gossip protocol.
*
* Request body: PeerSyncRequest
* Response body: PeerSyncResponse
*/
router.post("/mesh/sync", async (req, res) => {
try {
const { CentralCore } = await import("@fusion/core");
const central = new CentralCore();
await central.init();
// Validate required fields
const senderNodeId = req.body?.senderNodeId;
if (!senderNodeId) {
throw badRequest("senderNodeId is required");
}
const knownPeers = req.body?.knownPeers;
if (!Array.isArray(knownPeers)) {
throw badRequest("knownPeers must be an array");
}
// Optional: validate knownPeers entries have required fields
for (const peer of knownPeers) {
if (!peer?.nodeId || !peer?.nodeName || typeof peer?.status !== "string") {
throw badRequest("Each knownPeers entry must have nodeId, nodeName, and status");
}
}
// Get sender node from registry to validate auth
const senderNode = await central.getNode(senderNodeId);
// Auth validation: if sender is registered with an apiKey, validate it
if (senderNode?.apiKey) {
const authHeader = req.headers.authorization;
const token = authHeader?.startsWith("Bearer ") ? authHeader.slice(7) : undefined;
if (!token || token !== senderNode.apiKey) {
await central.close();
res.status(401).json({ error: "Unauthorized" });
return;
}
}
// Merge incoming peer data
await central.mergePeers(knownPeers);
// Update sender node status to online (it sent us a request, so it's alive)
try {
await central.updateNode(senderNodeId, { status: "online" });
} catch {
// Silently skip if sender node not found in local registry
}
// Get all known peers
const allKnownPeers = await central.getAllKnownPeerInfo();
// Calculate newPeers - peers the sender doesn't know about
const senderKnownIds = new Set(knownPeers.map((p: { nodeId: string }) => p.nodeId));
const newPeers = allKnownPeers.filter((peer) => !senderKnownIds.has(peer.nodeId));
// Get local node info
const localPeer = await central.getLocalPeerInfo();
await central.close();
// Return sync response
res.json({
senderNodeId: localPeer.nodeId,
senderNodeUrl: localPeer.nodeUrl,
knownPeers: allKnownPeers,
newPeers,
timestamp: new Date().toISOString(),
});
} catch (err: any) {
if (err instanceof ApiError) {
throw err;
}
rethrowAsApiError(err);
}
});
// ── Node Discovery Routes (mDNS / DNS-SD) ────────────────────────────────
/**

View File

@@ -32,6 +32,7 @@ export { TokenCapDetector, type TokenCapCheckResult } from "./token-cap-detector
export { SelfHealingManager, type SelfHealingOptions } from "./self-healing.js";
export { ProjectManager } from "./project-manager.js";
export { NodeHealthMonitor } from "./node-health-monitor.js";
export { PeerExchangeService, type PeerExchangeServiceOptions, type SyncResult } from "./peer-exchange-service.js";
export { RemoteNodeClient } from "./runtimes/remote-node-client.js";
export { RemoteNodeRuntime, type RemoteNodeRuntimeConfig } from "./runtimes/remote-node-runtime.js";
export { StepSessionExecutor } from "./step-session-executor.js";

View File

@@ -88,3 +88,6 @@ export const remoteNodeLog = createLogger("remote-node");
/** Logger for periodic node health monitor subsystem. */
export const nodeHealthMonitorLog = createLogger("node-health-monitor");
/** Logger for the peer exchange (gossip) subsystem. */
export const peerExchangeLog = createLogger("peer-exchange");

View File

@@ -0,0 +1,278 @@
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
import type { CentralCore, NodeConfig, PeerInfo } from "@fusion/core";
import { PeerExchangeService } from "./peer-exchange-service.js";
function makeNode(overrides: Partial<NodeConfig> = {}): NodeConfig {
return {
id: "node_remote",
name: "Remote Node",
type: "remote",
url: "https://remote.example.com",
apiKey: undefined,
status: "online",
maxConcurrent: 2,
createdAt: "2026-04-01T10:00:00.000Z",
updatedAt: "2026-04-01T12:00:00.000Z",
...overrides,
};
}
function makePeerInfo(overrides: Partial<PeerInfo> = {}): PeerInfo {
return {
nodeId: "node_peer",
nodeName: "Peer Node",
nodeUrl: "https://peer.example.com",
status: "online",
metrics: null,
lastSeen: "2026-04-01T12:00:00.000Z",
maxConcurrent: 2,
...overrides,
};
}
describe("PeerExchangeService", () => {
let mockCentralCore: CentralCore;
let mockFetch: ReturnType<typeof vi.fn>;
let mockListNodes: ReturnType<typeof vi.fn>;
let mockGetAllKnownPeerInfo: ReturnType<typeof vi.fn>;
let mockMergePeers: ReturnType<typeof vi.fn>;
let mockReportMeshState: ReturnType<typeof vi.fn>;
beforeEach(() => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-04-01T12:00:00.000Z"));
// Create individual mocks
mockListNodes = vi.fn();
mockGetAllKnownPeerInfo = vi.fn();
mockMergePeers = vi.fn();
mockReportMeshState = vi.fn();
mockCentralCore = {
listNodes: mockListNodes,
getAllKnownPeerInfo: mockGetAllKnownPeerInfo,
mergePeers: mockMergePeers,
reportMeshState: mockReportMeshState,
} as unknown as CentralCore;
mockFetch = vi.fn();
globalThis.fetch = mockFetch;
});
afterEach(() => {
vi.useRealTimers();
vi.restoreAllMocks();
});
describe("constructor", () => {
it("should create service instance", () => {
const service = new PeerExchangeService(mockCentralCore);
expect(service).toBeDefined();
});
it("should accept custom sync interval", () => {
const service = new PeerExchangeService(mockCentralCore, { syncIntervalMs: 30_000 });
expect(service).toBeDefined();
});
});
describe("syncWithNode()", () => {
it("should send correct request body with auth header when apiKey is set", async () => {
const node = makeNode({ apiKey: "secret-key" });
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([
makePeerInfo({ nodeId: "node_local", nodeName: "local" }),
makePeerInfo({ nodeId: "node_remote" }),
]);
mockMergePeers.mockResolvedValue({ added: [], updated: [] });
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: true,
status: 200,
json: async () => ({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
newPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
});
const service = new PeerExchangeService(mockCentralCore);
const result = await service.syncWithNode(node);
expect(result.success).toBe(true);
expect(mockFetch).toHaveBeenCalledWith(
"https://remote.example.com/api/mesh/sync",
expect.objectContaining({
method: "POST",
headers: {
"Content-Type": "application/json",
"Authorization": "Bearer secret-key",
},
body: expect.stringContaining('"senderNodeId":"node_local"'),
})
);
});
it("should send request without auth header when no apiKey", async () => {
const node = makeNode({ apiKey: undefined });
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockMergePeers.mockResolvedValue({ added: [], updated: [] });
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: true,
status: 200,
json: async () => ({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
newPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
});
const service = new PeerExchangeService(mockCentralCore);
await service.syncWithNode(node);
expect(mockFetch).toHaveBeenCalledWith(
expect.any(String),
expect.objectContaining({
headers: expect.not.objectContaining({ "Authorization": expect.anything() }),
})
);
});
it("should merge response.knownPeers (not just newPeers)", async () => {
const node = makeNode();
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
const allPeersFromResponse = [
makePeerInfo({ nodeId: "node_local", nodeName: "local", status: "online" }),
makePeerInfo({ nodeId: "node_peer_a", status: "online" }),
makePeerInfo({ nodeId: "node_peer_b", status: "offline" }),
];
mockGetAllKnownPeerInfo.mockResolvedValue([
makePeerInfo({ nodeId: "node_local", nodeName: "local" }),
]);
mockMergePeers.mockResolvedValue({ added: ["node_peer_a"], updated: ["node_peer_b"] });
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: true,
status: 200,
json: async () => ({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: allPeersFromResponse,
newPeers: [makePeerInfo({ nodeId: "node_peer_c" })],
timestamp: "2026-04-01T12:00:00.000Z",
}),
});
const service = new PeerExchangeService(mockCentralCore);
const result = await service.syncWithNode(node);
expect(result.success).toBe(true);
// Verify merge was called with all knownPeers, not just newPeers
expect(mockMergePeers).toHaveBeenCalledWith(allPeersFromResponse);
});
it("should refresh local metrics before sending request", async () => {
const node = makeNode();
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockMergePeers.mockResolvedValue({ added: [], updated: [] });
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: true,
status: 200,
json: async () => ({
senderNodeId: "node_remote",
senderNodeUrl: "https://remote.example.com",
knownPeers: [],
newPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
});
const service = new PeerExchangeService(mockCentralCore);
await service.syncWithNode(node);
expect(mockReportMeshState).toHaveBeenCalled();
});
it("should handle network error gracefully", async () => {
const node = makeNode();
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockReportMeshState.mockResolvedValue({});
mockFetch.mockRejectedValue(new Error("Network error"));
const service = new PeerExchangeService(mockCentralCore);
const result = await service.syncWithNode(node);
expect(result.success).toBe(false);
expect(result.error).toContain("Network error");
});
it("should handle non-2xx response", async () => {
const node = makeNode();
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: false,
status: 401,
statusText: "Unauthorized",
});
const service = new PeerExchangeService(mockCentralCore);
const result = await service.syncWithNode(node);
expect(result.success).toBe(false);
expect(result.error).toContain("HTTP 401");
});
});
describe("triggerSync()", () => {
it("should trigger sync when called", async () => {
mockListNodes.mockResolvedValue([
makeNode({ id: "node_local", type: "local", status: "online" }),
makeNode({ id: "node_1", name: "Remote 1", status: "online", url: "https://remote1.example.com" }),
]);
mockGetAllKnownPeerInfo.mockResolvedValue([]);
mockMergePeers.mockResolvedValue({ added: [], updated: [] });
mockReportMeshState.mockResolvedValue({});
mockFetch.mockResolvedValue({
ok: true,
status: 200,
json: async () => ({
senderNodeId: "node_1",
senderNodeUrl: "https://remote1.example.com",
knownPeers: [],
newPeers: [],
timestamp: "2026-04-01T12:00:00.000Z",
}),
});
const service = new PeerExchangeService(mockCentralCore, { syncIntervalMs: 60_000 });
const results = await service.triggerSync();
expect(mockFetch).toHaveBeenCalled();
});
});
});

View File

@@ -0,0 +1,267 @@
import type { CentralCore } from "@fusion/core";
import type { NodeConfig, PeerInfo, PeerSyncRequest, PeerSyncResponse } from "@fusion/core";
import { peerExchangeLog } from "./logger.js";
export interface PeerExchangeServiceOptions {
/** Interval between peer sync cycles in milliseconds. Default: 60000 (1 minute) */
syncIntervalMs?: number;
}
/**
* Result of syncing with a single node.
*/
export interface SyncResult {
/** Node ID that was synced with */
nodeId: string;
/** Whether the sync was successful */
success: boolean;
/** Number of new peers discovered */
added: number;
/** Number of peers updated */
updated: number;
/** Error message if sync failed */
error?: string;
}
/**
* Background service that implements the peer gossip protocol.
*
* Periodically exchanges peer information with connected remote nodes
* to keep the mesh state up-to-date across all nodes.
*/
export class PeerExchangeService {
private centralCore: CentralCore;
private syncIntervalMs: number;
private interval: ReturnType<typeof setInterval> | null = null;
private activeSync: Promise<void> | null = null;
private stopped = false;
/**
* Create a PeerExchangeService.
*
* @param centralCore - CentralCore instance for node registry access
* @param options - Configuration options
*/
constructor(centralCore: CentralCore, options: PeerExchangeServiceOptions = {}) {
this.centralCore = centralCore;
this.syncIntervalMs = options.syncIntervalMs ?? 60_000; // 1 minute default
}
/**
* Start the peer exchange service.
* Begins periodic gossip with all online remote nodes.
*/
start(): void {
if (this.stopped) {
peerExchangeLog.warn("Cannot start - service has been stopped");
return;
}
// Get initial peer count for logging (async call)
this.centralCore.listNodes().then((nodes) => {
const onlineRemoteCount = nodes.filter(
(n) => n.type === "remote" && n.status === "online" && n.url
).length;
peerExchangeLog.log(`Starting peer exchange service (sync interval: ${this.syncIntervalMs}ms, ${onlineRemoteCount} online remote peers)`);
}).catch((err) => {
peerExchangeLog.warn(`Failed to get initial peer count: ${err}`);
});
// Start periodic sync
this.interval = setInterval(() => {
void this.syncWithAllPeers();
}, this.syncIntervalMs);
}
/**
* Stop the peer exchange service.
* Clears the sync interval and prevents further syncs.
*/
stop(): void {
if (this.interval) {
clearInterval(this.interval);
this.interval = null;
}
this.stopped = true;
peerExchangeLog.log("Stopped peer exchange service");
}
/**
* Trigger an immediate sync with all peers, bypassing the interval.
*
* If a sync is already in progress, returns the in-progress sync.
*
* @returns Promise that resolves when the sync completes
*/
async triggerSync(): Promise<SyncResult[]> {
return this.syncWithAllPeers();
}
/**
* Sync with all online remote nodes.
*
* Uses single-flight pattern to prevent overlapping syncs.
* If a sync is already in progress, returns that sync's promise.
*/
async syncWithAllPeers(): Promise<SyncResult[]> {
// Single-flight: if a sync is already running, return that
if (this.activeSync) {
peerExchangeLog.log("Sync already in progress, skipping");
await this.activeSync;
return [];
}
this.activeSync = this.runSyncWithAllPeers();
try {
await this.activeSync;
} finally {
this.activeSync = null;
}
return [];
}
private async runSyncWithAllPeers(): Promise<void> {
try {
// Get all online remote nodes with URLs
const nodes = await this.centralCore.listNodes();
const onlineRemoteNodes = nodes.filter(
(node) => node.type === "remote" && node.status === "online" && node.url
);
if (onlineRemoteNodes.length === 0) {
peerExchangeLog.log("No online remote nodes to sync with");
return;
}
peerExchangeLog.log(`Starting sync with ${onlineRemoteNodes.length} peers`);
// Sync with each node sequentially (not in parallel to avoid thundering herd)
let totalAdded = 0;
let totalUpdated = 0;
const errors: string[] = [];
for (const node of onlineRemoteNodes) {
try {
const result = await this.syncWithNode(node);
totalAdded += result.added;
totalUpdated += result.updated;
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
errors.push(`${node.name}: ${message}`);
peerExchangeLog.warn(`Sync with ${node.name} failed: ${message}`);
}
}
// Log summary
if (errors.length > 0) {
peerExchangeLog.log(
`Sync complete: ${onlineRemoteNodes.length - errors.length} succeeded, ${errors.length} failed. ` +
`${totalAdded} new peers discovered, ${totalUpdated} updated. Errors: ${errors.join("; ")}`
);
} else {
peerExchangeLog.log(
`Sync complete: ${onlineRemoteNodes.length} peers synced. ` +
`${totalAdded} new peers discovered, ${totalUpdated} updated.`
);
}
} catch (error) {
peerExchangeLog.error("Unexpected error in sync loop:", error);
}
}
/**
* Sync with a single remote node.
*
* Sends our known peers and merges the response.
*
* @param node - Remote node configuration
* @returns Sync result with counts and any errors
*/
async syncWithNode(node: NodeConfig): Promise<SyncResult> {
try {
// Build the sync request
// Refresh local metrics first to ensure freshness
await this.centralCore.reportMeshState();
// Get local node info
const nodes = await this.centralCore.listNodes();
const localNode = nodes.find((n) => n.type === "local");
if (!localNode) {
return { nodeId: node.id, success: false, added: 0, updated: 0, error: "Local node not found" };
}
// Get all known peers for the request
const allKnownPeers = await this.centralCore.getAllKnownPeerInfo();
const request: PeerSyncRequest = {
senderNodeId: localNode.id,
senderNodeUrl: localNode.url || "",
knownPeers: allKnownPeers,
timestamp: new Date().toISOString(),
};
// Build headers
const headers: Record<string, string> = {
"Content-Type": "application/json",
};
if (node.apiKey) {
headers["Authorization"] = `Bearer ${node.apiKey}`;
}
// Send the sync request with 10-second timeout
const controller = new AbortController();
const timeoutId = setTimeout(() => controller.abort(), 10_000);
try {
const response = await fetch(`${node.url}/api/mesh/sync`, {
method: "POST",
headers,
body: JSON.stringify(request),
signal: controller.signal,
});
clearTimeout(timeoutId);
if (!response.ok) {
return {
nodeId: node.id,
success: false,
added: 0,
updated: 0,
error: `HTTP ${response.status}: ${response.statusText}`,
};
}
const peerResponse: PeerSyncResponse = await response.json();
// Merge ALL known peers from the response (not just newPeers)
// This ensures we get updates for existing peers too
const mergeResult = await this.centralCore.mergePeers(peerResponse.knownPeers);
peerExchangeLog.log(
`Synced with ${node.name}: ${mergeResult.added.length} new, ${mergeResult.updated.length} updated, ` +
`${peerResponse.newPeers.length} new to sender`
);
return {
nodeId: node.id,
success: true,
added: mergeResult.added.length,
updated: mergeResult.updated.length,
};
} catch (fetchError) {
clearTimeout(timeoutId);
throw fetchError;
}
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
if (message.includes("abort")) {
return { nodeId: node.id, success: false, added: 0, updated: 0, error: "Timeout (10s)" };
}
return { nodeId: node.id, success: false, added: 0, updated: 0, error: message };
}
}
}