diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 56eff323a..ceff0f677 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -485,6 +485,8 @@ export type { SettingsSyncPayload, SettingsSyncState, SettingsSyncResult, + SharedMeshStatePayload, + SnapshotBase, SystemMetrics, ProjectStatus, RegisteredProject, diff --git a/packages/core/src/types.ts b/packages/core/src/types.ts index 3c6038794..e2f55a608 100644 --- a/packages/core/src/types.ts +++ b/packages/core/src/types.ts @@ -2535,6 +2535,33 @@ export interface PeerInfo { } /** Request payload sent when a node initiates a peer sync. */ +export interface SnapshotBase { + version: number; + exportedAt: string; + checksum: string; +} + +export interface SharedMeshStatePayload { + taskMetadata?: SnapshotBase & { payload: { tasks: Task[] } }; + missionHierarchy?: SnapshotBase & { + payload: { + missions: import("./mission-types.js").Mission[]; + milestones: import("./mission-types.js").Milestone[]; + slices: import("./mission-types.js").Slice[]; + features: import("./mission-types.js").MissionFeature[]; + missionEvents: import("./mission-types.js").MissionEvent[]; + assertions: import("./mission-types.js").MissionContractAssertion[]; + featureAssertionLinks: import("./mission-types.js").FeatureAssertionLink[]; + }; + }; + agents?: SnapshotBase & { payload: { agents: Agent[]; blockedStates: { agentId: string; state: BlockedStateSnapshot }[] } }; + agentRuns?: SnapshotBase & { payload: { runs: AgentHeartbeatRun[] } }; + activityLog?: SnapshotBase & { payload: { entries: ActivityLogEntry[] } }; + runAudit?: SnapshotBase & { payload: { entries: RunAuditEvent[] } }; + projectSettings?: SnapshotBase & { payload: { global: GlobalSettings; projects?: Record } }; + authMaterial?: SnapshotBase & { payload: { providerAuth?: Record } }; +} + export interface PeerSyncRequest { /** Node ID of the sender. */ senderNodeId: string; @@ -2546,6 +2573,8 @@ export interface PeerSyncRequest { timestamp: string; /** Optional settings sync payload included in the request. */ settings?: SettingsSyncPayload; + /** Optional shared-state payload included in the request. */ + sharedState?: SharedMeshStatePayload; } /** Response payload returned after a peer sync exchange. */ @@ -2562,6 +2591,8 @@ export interface PeerSyncResponse { timestamp: string; /** Optional settings sync payload included in the response. */ settings?: SettingsSyncPayload; + /** Optional shared-state payload included in the response. */ + sharedState?: SharedMeshStatePayload; } /** A single provider's authentication credential for sync transport. */ diff --git a/packages/engine/src/__tests__/peer-exchange-service.test.ts b/packages/engine/src/__tests__/peer-exchange-service.test.ts index 8da4d5a86..854e3195b 100644 --- a/packages/engine/src/__tests__/peer-exchange-service.test.ts +++ b/packages/engine/src/__tests__/peer-exchange-service.test.ts @@ -49,6 +49,9 @@ describe("PeerExchangeService", () => { let mockReportMeshState: ReturnType; let mockGetSettingsForSync: ReturnType; let mockApplyRemoteSettings: ReturnType; + let mockGetProjectSettingsSnapshot: ReturnType; + let mockGetAuthMaterialSnapshot: ReturnType; + let mockApplyProjectSettingsSnapshot: ReturnType; beforeEach(() => { vi.useFakeTimers(); @@ -61,6 +64,9 @@ describe("PeerExchangeService", () => { mockReportMeshState = vi.fn(); mockGetSettingsForSync = vi.fn(); mockApplyRemoteSettings = vi.fn(); + mockGetProjectSettingsSnapshot = vi.fn(); + mockGetAuthMaterialSnapshot = vi.fn(); + mockApplyProjectSettingsSnapshot = vi.fn(); mockCentralCore = { listNodes: mockListNodes, @@ -69,6 +75,9 @@ describe("PeerExchangeService", () => { reportMeshState: mockReportMeshState, getSettingsForSync: mockGetSettingsForSync, applyRemoteSettings: mockApplyRemoteSettings, + getProjectSettingsSnapshot: mockGetProjectSettingsSnapshot, + getAuthMaterialSnapshot: mockGetAuthMaterialSnapshot, + applyProjectSettingsSnapshot: mockApplyProjectSettingsSnapshot, } as unknown as CentralCore; mockFetch = vi.fn(); @@ -386,6 +395,8 @@ describe("PeerExchangeService", () => { const service = new PeerExchangeService(mockCentralCore, { settingsSyncEnabled: true }); const payload = makeSettingsPayload({ checksum: "local-checksum-123" }); mockGetSettingsForSync.mockResolvedValue(payload); + mockGetProjectSettingsSnapshot.mockResolvedValue({ version: 1, exportedAt: payload.exportedAt, checksum: payload.checksum, payload: { global: {} } }); + mockGetAuthMaterialSnapshot.mockReturnValue({ version: 1, exportedAt: payload.exportedAt, checksum: "auth-checksum", payload: { providerAuth: {} } }); setupSuccessfulSync(); await service.syncWithNode(makeNode()); @@ -395,6 +406,7 @@ describe("PeerExchangeService", () => { const body = JSON.parse(call[1].body); expect(body.settings).toBeDefined(); expect(body.settings.checksum).toBe("local-checksum-123"); + expect(body.sharedState?.projectSettings?.checksum).toBe("local-checksum-123"); }); it("should include settings on first sync with a node (no throttle entry)", async () => { @@ -537,7 +549,7 @@ describe("PeerExchangeService", () => { const localPayload = makeSettingsPayload({ checksum: "local-checksum" }); const remotePayload = makeSettingsPayload({ checksum: "remote-checksum" }); mockGetSettingsForSync.mockResolvedValue(localPayload); - mockApplyRemoteSettings.mockResolvedValue({ + mockApplyProjectSettingsSnapshot.mockResolvedValue({ success: true, globalCount: 5, projectCount: 2, @@ -554,13 +566,16 @@ describe("PeerExchangeService", () => { knownPeers: [], newPeers: [], timestamp: "2026-04-01T12:00:00.000Z", - settings: remotePayload, + sharedState: { + projectSettings: remotePayload, + authMaterial: { ...remotePayload, payload: { providerAuth: {} } }, + }, }), }); const result = await service.syncWithNode(makeNode()); - expect(mockApplyRemoteSettings).toHaveBeenCalledWith(remotePayload); + expect(mockApplyProjectSettingsSnapshot).toHaveBeenCalledWith(remotePayload); expect(result.settingsApplied).toBe(true); expect(result.settingsVersion).toBe("remote-checksum"); }); diff --git a/packages/engine/src/peer-exchange-service.ts b/packages/engine/src/peer-exchange-service.ts index f10e7da7f..aa76dd568 100644 --- a/packages/engine/src/peer-exchange-service.ts +++ b/packages/engine/src/peer-exchange-service.ts @@ -1,4 +1,9 @@ -import type { CentralCore, GlobalSettings, SettingsSyncPayload } from "@fusion/core"; +import type { + CentralCore, + GlobalSettings, + SettingsSyncPayload, + SharedMeshStatePayload, +} from "@fusion/core"; import type { NodeConfig, PeerSyncRequest, PeerSyncResponse } from "@fusion/core"; import { peerExchangeLog } from "./logger.js"; @@ -57,6 +62,8 @@ export class PeerExchangeService { private lastSettingsSyncByNode = new Map(); /** Cached settings payload from the last successful getSettingsForSync call. */ private cachedSettingsPayload: SettingsSyncPayload | null = null; + /** Cached shared-state settings/auth payload built from canonical snapshots. */ + private cachedSharedStatePayload: SharedMeshStatePayload | null = null; /** Global settings provided via options. */ private globalSettings?: GlobalSettings; /** Provider auth credentials provided via options. */ @@ -87,6 +94,7 @@ export class PeerExchangeService { this.globalSettings = settings; // Invalidate cache to ensure fresh payload on next sync this.cachedSettingsPayload = null; + this.cachedSharedStatePayload = null; } /** @@ -309,6 +317,7 @@ export class PeerExchangeService { if (shouldIncludeSettings) { request.settings = this.cachedSettingsPayload; + request.sharedState = await this.getSharedStateSettingsBundle(); } } catch (err) { // Log error but continue with peer sync @@ -359,45 +368,46 @@ export class PeerExchangeService { const mergeResult = await this.centralCore.mergePeers(peerResponse.knownPeers); // ── Process remote settings if included in response ── - if (peerResponse.settings && this.settingsSyncEnabled) { - settingsVersion = peerResponse.settings.checksum; + if (this.settingsSyncEnabled && (peerResponse.sharedState || peerResponse.settings)) { + const remoteChecksum = + peerResponse.sharedState?.projectSettings?.checksum ?? + peerResponse.settings?.checksum; + settingsVersion = remoteChecksum; - // Check if we should apply remote settings - // Apply if remote checksum is different from our cached checksum const localChecksum = this.cachedSettingsPayload?.checksum ?? ""; - if (peerResponse.settings.checksum !== localChecksum) { + if (remoteChecksum && remoteChecksum !== localChecksum) { try { - const applyResult = await this.centralCore.applyRemoteSettings(peerResponse.settings); + const applyResult = await this.applyRemoteSharedState(peerResponse.sharedState, peerResponse.settings); if (applyResult.success) { settingsApplied = true; peerExchangeLog.log( - `Applied remote settings from ${node.name} (version: ${peerResponse.settings.checksum}, ` + - `global: ${applyResult.globalCount}, projects: ${applyResult.projectCount}, auth: ${applyResult.authCount})` + `Applied remote settings from ${node.name} (version: ${remoteChecksum}, ` + + `global: ${applyResult.globalCount}, projects: ${applyResult.projectCount}, auth: ${applyResult.authCount})`, ); // Invalidate cache to ensure fresh data on next sync this.cachedSettingsPayload = null; + this.cachedSharedStatePayload = null; } else { - peerExchangeLog.warn( - `Failed to apply remote settings from ${node.name}: ${applyResult.error}` - ); + peerExchangeLog.warn(`Failed to apply remote settings from ${node.name}: ${applyResult.error}`); } } catch (err) { const error = err instanceof Error ? err.message : String(err); peerExchangeLog.warn(`Settings sync error with ${node.name}: ${error}`); } - } else { + } else if (remoteChecksum) { peerExchangeLog.log( - `Remote settings from ${node.name} are up-to-date (version: ${peerResponse.settings.checksum})` + `Remote settings from ${node.name} are up-to-date (version: ${remoteChecksum})`, ); } - // Update throttle tracking - this.lastSettingsSyncByNode.set(node.id, { - version: peerResponse.settings.checksum, - timestamp: Date.now(), - }); + if (remoteChecksum) { + this.lastSettingsSyncByNode.set(node.id, { + version: remoteChecksum, + timestamp: Date.now(), + }); + } } peerExchangeLog.log( @@ -431,4 +441,53 @@ export class PeerExchangeService { return { nodeId: node.id, success: false, added: 0, updated: 0, error: message }; } } + + private async getSharedStateSettingsBundle(): Promise { + if (this.cachedSharedStatePayload) { + return this.cachedSharedStatePayload; + } + const globalSettings = this.globalSettings ?? {}; + const core = this.centralCore as CentralCore & { + getProjectSettingsSnapshot?: (settings: GlobalSettings) => Promise; + getAuthMaterialSnapshot?: ( + providerAuth?: Record, + ) => SharedMeshStatePayload["authMaterial"]; + }; + if (!core.getProjectSettingsSnapshot || !core.getAuthMaterialSnapshot) { + return undefined; + } + const projectSettings = await core.getProjectSettingsSnapshot(globalSettings); + const authMaterial = core.getAuthMaterialSnapshot(this.providerAuth); + this.cachedSharedStatePayload = { projectSettings, authMaterial }; + return this.cachedSharedStatePayload; + } + + private async applyRemoteSharedState( + sharedState: SharedMeshStatePayload | undefined, + fallbackSettings: SettingsSyncPayload | undefined, + ): Promise<{ success: boolean; globalCount: number; projectCount: number; authCount: number; error?: string }> { + const core = this.centralCore as CentralCore & { + applyProjectSettingsSnapshot?: (snapshot: NonNullable) => Promise<{ success: boolean; globalCount: number; projectCount: number; authCount: number; error?: string }>; + applyAuthMaterialSnapshot?: (snapshot: NonNullable) => { success?: boolean; authCount?: number; error?: string } | Record; + }; + + if (sharedState?.projectSettings && core.applyProjectSettingsSnapshot) { + const result = await core.applyProjectSettingsSnapshot(sharedState.projectSettings); + if (sharedState.authMaterial && core.applyAuthMaterialSnapshot) { + const authResult = core.applyAuthMaterialSnapshot(sharedState.authMaterial); + const authCount = + typeof (authResult as { authCount?: number }).authCount === "number" + ? (authResult as { authCount: number }).authCount + : Object.keys(sharedState.authMaterial.payload.providerAuth ?? {}).length; + return { ...result, authCount: Math.max(result.authCount, authCount) }; + } + return result; + } + + if (fallbackSettings) { + return this.centralCore.applyRemoteSettings(fallbackSettings); + } + + return { success: true, globalCount: 0, projectCount: 0, authCount: 0 }; + } }