feat(FN-3960): add mesh replay outage contracts and runtime hooks
- Add core outage schema and mesh replay contract types, with CentralCore/CentralDB support for queue and snapshot persistence - Expose runtime replay hooks in engine health monitoring and peer exchange paths to drive outage recovery flows - Expand protocol and architecture docs for shared mesh replay behavior and finalized contract expectations - Add regression tests across core and engine for queue sequencing, snapshot contracts, and replay integration Fusion-Task-Id: FN-3960
This commit is contained in:
@@ -117,6 +117,23 @@ describe("NodeHealthMonitor", () => {
|
||||
expect(monitor.getNodeHealth("node-remote")).toBe("online");
|
||||
});
|
||||
|
||||
it("invokes recovery hook once per non-online to online transition", async () => {
|
||||
const onNodeRecovered = vi.fn();
|
||||
monitor = new NodeHealthMonitor(mockCentralCore, { checkIntervalMs: 1_000, onNodeRecovered });
|
||||
checkNodeHealthMock
|
||||
.mockResolvedValueOnce("offline")
|
||||
.mockResolvedValueOnce("online")
|
||||
.mockResolvedValueOnce("online");
|
||||
|
||||
await monitor.start();
|
||||
await monitor.checkAllNodes();
|
||||
await monitor.checkAllNodes();
|
||||
await monitor.checkAllNodes();
|
||||
|
||||
expect(onNodeRecovered).toHaveBeenCalledTimes(1);
|
||||
expect(onNodeRecovered).toHaveBeenCalledWith("node-remote", "offline");
|
||||
});
|
||||
|
||||
it("is a no-op when no remote nodes are registered", async () => {
|
||||
listNodesMock.mockResolvedValue([
|
||||
createNode({ id: "node-local-only", name: "Local Only", type: "local", status: "online" }),
|
||||
|
||||
@@ -52,6 +52,12 @@ describe("PeerExchangeService", () => {
|
||||
let mockGetProjectSettingsSnapshot: ReturnType<typeof vi.fn>;
|
||||
let mockGetAuthMaterialSnapshot: ReturnType<typeof vi.fn>;
|
||||
let mockApplyProjectSettingsSnapshot: ReturnType<typeof vi.fn>;
|
||||
let mockEnqueueMeshWrite: ReturnType<typeof vi.fn>;
|
||||
let mockListPendingMeshWrites: ReturnType<typeof vi.fn>;
|
||||
let mockMarkMeshWriteReplayStarted: ReturnType<typeof vi.fn>;
|
||||
let mockMarkMeshWriteApplied: ReturnType<typeof vi.fn>;
|
||||
let mockMarkMeshWriteFailed: ReturnType<typeof vi.fn>;
|
||||
let mockGetNode: ReturnType<typeof vi.fn>;
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers();
|
||||
@@ -67,6 +73,12 @@ describe("PeerExchangeService", () => {
|
||||
mockGetProjectSettingsSnapshot = vi.fn();
|
||||
mockGetAuthMaterialSnapshot = vi.fn();
|
||||
mockApplyProjectSettingsSnapshot = vi.fn();
|
||||
mockEnqueueMeshWrite = vi.fn();
|
||||
mockListPendingMeshWrites = vi.fn();
|
||||
mockMarkMeshWriteReplayStarted = vi.fn();
|
||||
mockMarkMeshWriteApplied = vi.fn();
|
||||
mockMarkMeshWriteFailed = vi.fn();
|
||||
mockGetNode = vi.fn();
|
||||
|
||||
mockCentralCore = {
|
||||
listNodes: mockListNodes,
|
||||
@@ -78,10 +90,23 @@ describe("PeerExchangeService", () => {
|
||||
getProjectSettingsSnapshot: mockGetProjectSettingsSnapshot,
|
||||
getAuthMaterialSnapshot: mockGetAuthMaterialSnapshot,
|
||||
applyProjectSettingsSnapshot: mockApplyProjectSettingsSnapshot,
|
||||
enqueueMeshWrite: mockEnqueueMeshWrite,
|
||||
listPendingMeshWrites: mockListPendingMeshWrites,
|
||||
markMeshWriteReplayStarted: mockMarkMeshWriteReplayStarted,
|
||||
markMeshWriteApplied: mockMarkMeshWriteApplied,
|
||||
markMeshWriteFailed: mockMarkMeshWriteFailed,
|
||||
getNode: mockGetNode,
|
||||
} as unknown as CentralCore;
|
||||
|
||||
mockFetch = vi.fn();
|
||||
globalThis.fetch = mockFetch;
|
||||
|
||||
mockListPendingMeshWrites.mockResolvedValue([]);
|
||||
mockGetNode.mockResolvedValue(makeNode());
|
||||
mockEnqueueMeshWrite.mockResolvedValue({ id: "mq-1" });
|
||||
mockMarkMeshWriteReplayStarted.mockResolvedValue(undefined);
|
||||
mockMarkMeshWriteApplied.mockResolvedValue(undefined);
|
||||
mockMarkMeshWriteFailed.mockResolvedValue(undefined);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
@@ -333,6 +358,53 @@ describe("PeerExchangeService", () => {
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
expect(result.error).toContain("HTTP 401");
|
||||
expect(mockEnqueueMeshWrite).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("queues retryable failures and returns queued write id", async () => {
|
||||
const node = makeNode();
|
||||
mockListNodes.mockResolvedValue([makeNode({ id: "node_local", type: "local", status: "online" })]);
|
||||
mockGetAllKnownPeerInfo.mockResolvedValue([]);
|
||||
mockReportMeshState.mockResolvedValue({});
|
||||
mockFetch.mockResolvedValue({ ok: false, status: 503, statusText: "Service Unavailable" });
|
||||
mockEnqueueMeshWrite.mockResolvedValue({ id: "mq-queued-1" });
|
||||
|
||||
const service = new PeerExchangeService(mockCentralCore);
|
||||
const result = await service.syncWithNode(node);
|
||||
|
||||
expect(result.success).toBe(false);
|
||||
expect(result.queuedWriteId).toBe("mq-queued-1");
|
||||
expect(mockEnqueueMeshWrite).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it("reports replay summary after successful sync", async () => {
|
||||
const node = makeNode();
|
||||
setupSuccessfulSync(node);
|
||||
mockListPendingMeshWrites.mockResolvedValue([
|
||||
{
|
||||
id: "mq-1",
|
||||
originNodeId: "node_local",
|
||||
targetNodeId: node.id,
|
||||
projectId: null,
|
||||
scope: "mesh.sync",
|
||||
entityType: "shared-state-sync",
|
||||
entityId: node.id,
|
||||
operation: "sync",
|
||||
payload: { request: { senderNodeId: "node_local", knownPeers: [], senderNodeUrl: "", timestamp: "2026-04-01T12:00:00.000Z" } },
|
||||
intentVersion: "1.0",
|
||||
status: "pending",
|
||||
attemptCount: 0,
|
||||
createdAt: "2026-04-01T12:00:00.000Z",
|
||||
updatedAt: "2026-04-01T12:00:00.000Z",
|
||||
},
|
||||
]);
|
||||
|
||||
const service = new PeerExchangeService(mockCentralCore);
|
||||
const result = await service.syncWithNode(node);
|
||||
|
||||
expect(result.replaySummary).toEqual({ replayed: 1, applied: 1, failed: 0, queuedWriteIds: ["mq-1"] });
|
||||
expect(mockMarkMeshWriteReplayStarted).toHaveBeenCalledWith("mq-1");
|
||||
expect(mockMarkMeshWriteApplied).toHaveBeenCalledWith("mq-1", {});
|
||||
});
|
||||
});
|
||||
|
||||
|
||||
@@ -3,6 +3,7 @@ import { nodeHealthMonitorLog } from "./logger.js";
|
||||
|
||||
export interface NodeHealthMonitorOptions {
|
||||
checkIntervalMs?: number;
|
||||
onNodeRecovered?: (nodeId: string, previousStatus: NodeStatus) => Promise<void> | void;
|
||||
}
|
||||
|
||||
export interface NodeHealthCheckSummary {
|
||||
@@ -19,12 +20,14 @@ export class NodeHealthMonitor {
|
||||
private running = false;
|
||||
private lastKnownStatus = new Map<string, NodeStatus>();
|
||||
private activeCheck: Promise<NodeHealthCheckSummary> | null = null;
|
||||
private readonly onNodeRecovered?: (nodeId: string, previousStatus: NodeStatus) => Promise<void> | void;
|
||||
|
||||
constructor(
|
||||
private readonly centralCore: CentralCore,
|
||||
options: NodeHealthMonitorOptions = {}
|
||||
) {
|
||||
this.checkIntervalMs = options.checkIntervalMs ?? 60_000;
|
||||
this.onNodeRecovered = options.onNodeRecovered;
|
||||
}
|
||||
|
||||
async start(): Promise<void> {
|
||||
@@ -134,6 +137,9 @@ export class NodeHealthMonitor {
|
||||
nodeHealthMonitorLog.log(
|
||||
`Remote node ${node.name} (${node.id}) recovered: ${previousStatus} → online`
|
||||
);
|
||||
if (this.onNodeRecovered) {
|
||||
await this.onNodeRecovered(node.id, previousStatus);
|
||||
}
|
||||
}
|
||||
}
|
||||
} catch (error) {
|
||||
|
||||
@@ -1,9 +1,11 @@
|
||||
import type {
|
||||
CentralCore,
|
||||
GlobalSettings,
|
||||
MeshWriteReplaySummary,
|
||||
SettingsSyncPayload,
|
||||
SharedMeshStatePayload,
|
||||
} from "@fusion/core";
|
||||
import { isRetryableMeshWriteFailure } from "@fusion/core";
|
||||
import type { NodeConfig, PeerSyncRequest, PeerSyncResponse } from "@fusion/core";
|
||||
import { peerExchangeLog } from "./logger.js";
|
||||
|
||||
@@ -42,6 +44,8 @@ export interface SyncResult {
|
||||
settingsApplied?: boolean;
|
||||
/** The settings version (checksum) observed on the remote node. */
|
||||
settingsVersion?: string;
|
||||
queuedWriteId?: string;
|
||||
replaySummary?: MeshWriteReplaySummary;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -357,12 +361,15 @@ export class PeerExchangeService {
|
||||
clearTimeout(timeoutId);
|
||||
|
||||
if (!response.ok) {
|
||||
const errorMessage = `HTTP ${response.status}: ${response.statusText}`;
|
||||
const queuedWriteId = await this.enqueueRetryableSyncWrite(node, request, response.status, errorMessage);
|
||||
return {
|
||||
nodeId: node.id,
|
||||
success: false,
|
||||
added: 0,
|
||||
updated: 0,
|
||||
error: `HTTP ${response.status}: ${response.statusText}`,
|
||||
error: errorMessage,
|
||||
queuedWriteId,
|
||||
};
|
||||
}
|
||||
|
||||
@@ -420,11 +427,14 @@ export class PeerExchangeService {
|
||||
`${peerResponse.newPeers.length} new to sender`
|
||||
);
|
||||
|
||||
const replaySummary = await this.replayPendingWritesForNode(node.id);
|
||||
|
||||
const result: SyncResult = {
|
||||
nodeId: node.id,
|
||||
success: true,
|
||||
added: mergeResult.added.length,
|
||||
updated: mergeResult.updated.length,
|
||||
replaySummary,
|
||||
};
|
||||
|
||||
// Only include settings fields when settingsSyncEnabled is true
|
||||
@@ -440,13 +450,88 @@ export class PeerExchangeService {
|
||||
}
|
||||
} 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 };
|
||||
const normalizedMessage = message.includes("abort") ? "Timeout (10s)" : message;
|
||||
const queuedWriteId = await this.enqueueRetryableSyncWrite(node, undefined, undefined, normalizedMessage);
|
||||
return { nodeId: node.id, success: false, added: 0, updated: 0, error: normalizedMessage, queuedWriteId };
|
||||
}
|
||||
}
|
||||
|
||||
async replayPendingWritesForNode(targetNodeId: string): Promise<MeshWriteReplaySummary> {
|
||||
const node = await this.centralCore.getNode(targetNodeId);
|
||||
if (!node?.url) {
|
||||
return { replayed: 0, applied: 0, failed: 0, queuedWriteIds: [] };
|
||||
}
|
||||
|
||||
const pending = await this.centralCore.listPendingMeshWrites({ targetNodeId, status: "pending" });
|
||||
let applied = 0;
|
||||
let failed = 0;
|
||||
|
||||
for (const entry of pending) {
|
||||
await this.centralCore.markMeshWriteReplayStarted(entry.id);
|
||||
try {
|
||||
const headers: Record<string, string> = { "Content-Type": "application/json" };
|
||||
if (node.apiKey) headers.Authorization = `Bearer ${node.apiKey}`;
|
||||
const response = await fetch(`${node.url}/api/mesh/sync`, {
|
||||
method: "POST",
|
||||
headers,
|
||||
body: JSON.stringify(entry.payload?.request ?? {}),
|
||||
});
|
||||
if (!response.ok) {
|
||||
await this.centralCore.markMeshWriteFailed(entry.id, { lastError: `HTTP ${response.status}: ${response.statusText}` });
|
||||
failed += 1;
|
||||
continue;
|
||||
}
|
||||
await this.centralCore.markMeshWriteApplied(entry.id, {});
|
||||
applied += 1;
|
||||
} catch (error) {
|
||||
await this.centralCore.markMeshWriteFailed(entry.id, {
|
||||
lastError: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
failed += 1;
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
replayed: pending.length,
|
||||
applied,
|
||||
failed,
|
||||
queuedWriteIds: pending.map((entry) => entry.id),
|
||||
};
|
||||
}
|
||||
|
||||
private async enqueueRetryableSyncWrite(
|
||||
node: NodeConfig,
|
||||
request: PeerSyncRequest | undefined,
|
||||
statusCode: number | undefined,
|
||||
message: string,
|
||||
): Promise<string | undefined> {
|
||||
if (!isRetryableMeshWriteFailure(statusCode, message)) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
const localNode = (await this.centralCore.listNodes()).find((n) => n.type === "local");
|
||||
if (!localNode) {
|
||||
return undefined;
|
||||
}
|
||||
|
||||
const entry = await this.centralCore.enqueueMeshWrite({
|
||||
originNodeId: localNode.id,
|
||||
targetNodeId: node.id,
|
||||
projectId: null,
|
||||
scope: "mesh.sync",
|
||||
entityType: "shared-state-sync",
|
||||
entityId: node.id,
|
||||
operation: "sync",
|
||||
payload: {
|
||||
request: request ?? {},
|
||||
error: message,
|
||||
},
|
||||
intentVersion: "1.0",
|
||||
});
|
||||
|
||||
return entry.id;
|
||||
}
|
||||
|
||||
private async getSharedStateSettingsBundle(): Promise<SharedMeshStatePayload | undefined> {
|
||||
if (this.cachedSharedStatePayload) {
|
||||
return this.cachedSharedStatePayload;
|
||||
|
||||
Reference in New Issue
Block a user