refactor(federation): replace runInitialSyncForNewPeers with startupBootstrapSync
The per-peer sync body is now syncPeerMutationLog (in the new peer-activation module), invoked via onPeerActivated. The startup path scans for status='active' AND lastSyncedAt=0 and calls the unified handler for each — same trigger condition as before, unified code path with runtime transitions.
This commit is contained in:
@@ -164,5 +164,16 @@ export async function syncPeerMutationLog(
|
|||||||
* Invoked from startFederationWorkers.
|
* Invoked from startFederationWorkers.
|
||||||
*/
|
*/
|
||||||
export async function startupBootstrapSync(): Promise<void> {
|
export async function startupBootstrapSync(): Promise<void> {
|
||||||
// Stub — implemented in Task 5.
|
if (!isFederationRelayEnabled()) return;
|
||||||
|
|
||||||
|
const db = getDb();
|
||||||
|
const peers = db.select().from(schema.federationPeers)
|
||||||
|
.where(and(
|
||||||
|
eq(schema.federationPeers.status, 'active'),
|
||||||
|
eq(schema.federationPeers.lastSyncedAt, 0),
|
||||||
|
)).all();
|
||||||
|
|
||||||
|
for (const peer of peers) {
|
||||||
|
await onPeerActivated(peer.id, 'startup_bootstrap');
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import { connectionManager } from '../ws/handler.js';
|
|||||||
import { generateThumbnail } from './thumbnail.js';
|
import { generateThumbnail } from './thumbnail.js';
|
||||||
import { processRelayEvents } from '../routes/federation.js';
|
import { processRelayEvents } from '../routes/federation.js';
|
||||||
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
|
import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent } from '@backspace/shared';
|
||||||
|
import { startupBootstrapSync } from './federationPeerActivation.js';
|
||||||
import fs from 'node:fs';
|
import fs from 'node:fs';
|
||||||
import path from 'node:path';
|
import path from 'node:path';
|
||||||
import crypto from 'node:crypto';
|
import crypto from 'node:crypto';
|
||||||
@@ -1076,124 +1077,15 @@ function scheduleJanitorTick(): void {
|
|||||||
|
|
||||||
// ─── Lifecycle ──────────────────────────────────────────────────────────────
|
// ─── Lifecycle ──────────────────────────────────────────────────────────────
|
||||||
|
|
||||||
/**
|
|
||||||
* Trigger checkpoint sync for peers that have never been synced (lastSyncedAt === 0).
|
|
||||||
* This catches historical messages that existed before the relay was enabled.
|
|
||||||
*/
|
|
||||||
async function runInitialSyncForNewPeers(): Promise<void> {
|
|
||||||
if (!isFederationRelayEnabled()) return;
|
|
||||||
|
|
||||||
const db = getDb();
|
|
||||||
const unsyncedPeers = db
|
|
||||||
.select()
|
|
||||||
.from(schema.federationPeers)
|
|
||||||
.where(and(
|
|
||||||
eq(schema.federationPeers.status, 'active'),
|
|
||||||
eq(schema.federationPeers.lastSyncedAt, 0),
|
|
||||||
))
|
|
||||||
.all();
|
|
||||||
|
|
||||||
if (unsyncedPeers.length === 0) return;
|
|
||||||
|
|
||||||
const ourOrigin = getOurOrigin();
|
|
||||||
|
|
||||||
for (const peer of unsyncedPeers) {
|
|
||||||
const signingSecret = (peer.pendingHmacSecret && peer.secretRotationAt)
|
|
||||||
? peer.pendingHmacSecret
|
|
||||||
: peer.hmacSecret;
|
|
||||||
try {
|
|
||||||
console.log(`[federation-worker] Running initial sync with ${peer.origin}...`);
|
|
||||||
let sinceTimestamp = 0;
|
|
||||||
let totalEvents = 0;
|
|
||||||
|
|
||||||
// Paginate through all events from the peer
|
|
||||||
while (true) {
|
|
||||||
const body = JSON.stringify({ sinceTimestamp, limit: 100 });
|
|
||||||
const headers = buildFederationHeaders(body, signingSecret, ourOrigin);
|
|
||||||
|
|
||||||
const response = await fetch(`${peer.origin}/api/federation/sync`, {
|
|
||||||
method: 'POST',
|
|
||||||
headers,
|
|
||||||
body,
|
|
||||||
signal: AbortSignal.timeout(30_000),
|
|
||||||
});
|
|
||||||
|
|
||||||
if (!response.ok) {
|
|
||||||
console.error(`[federation-worker] Sync with ${peer.origin} failed: ${response.status}`);
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
const data = await response.json() as { events: FederationRelayEvent[]; hasMore: boolean; checkpoint: number };
|
|
||||||
|
|
||||||
if (data.events.length === 0) break;
|
|
||||||
|
|
||||||
// Process events directly — no HTTP round-trip (FED-005)
|
|
||||||
await processRelayEvents(data.events, peer.origin, peer.origin, db);
|
|
||||||
|
|
||||||
totalEvents += data.events.length;
|
|
||||||
sinceTimestamp = data.checkpoint;
|
|
||||||
|
|
||||||
if (!data.hasMore) break;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Second pass: sync friend events
|
|
||||||
let friendSinceTimestamp = 0;
|
|
||||||
while (true) {
|
|
||||||
const friendBody = JSON.stringify({ sinceTimestamp: friendSinceTimestamp, contextType: 'friend', limit: 100 });
|
|
||||||
const friendHeaders = buildFederationHeaders(friendBody, signingSecret, ourOrigin);
|
|
||||||
|
|
||||||
const friendResponse = await fetch(`${peer.origin}/api/federation/sync`, {
|
|
||||||
method: 'POST',
|
|
||||||
headers: friendHeaders,
|
|
||||||
body: friendBody,
|
|
||||||
signal: AbortSignal.timeout(30_000),
|
|
||||||
});
|
|
||||||
|
|
||||||
if (!friendResponse.ok) {
|
|
||||||
console.error(`[federation-worker] Friend sync with ${peer.origin} failed: ${friendResponse.status}`);
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
const friendData = await friendResponse.json() as { events: FederationRelayEvent[]; hasMore: boolean; checkpoint: number };
|
|
||||||
|
|
||||||
if (friendData.events.length === 0) break;
|
|
||||||
|
|
||||||
// Process events directly — no HTTP round-trip (FED-005)
|
|
||||||
await processRelayEvents(friendData.events, peer.origin, peer.origin, db);
|
|
||||||
|
|
||||||
totalEvents += friendData.events.length;
|
|
||||||
friendSinceTimestamp = friendData.checkpoint;
|
|
||||||
|
|
||||||
if (!friendData.hasMore) break;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Update lastSyncedAt so this doesn't run again
|
|
||||||
db.update(schema.federationPeers)
|
|
||||||
.set({ lastSyncedAt: Date.now() })
|
|
||||||
.where(eq(schema.federationPeers.id, peer.id))
|
|
||||||
.run();
|
|
||||||
|
|
||||||
if (totalEvents > 0) {
|
|
||||||
console.log(`[federation-worker] Initial sync with ${peer.origin}: ${totalEvents} events synced`);
|
|
||||||
} else {
|
|
||||||
console.log(`[federation-worker] Initial sync with ${peer.origin}: no events to sync`);
|
|
||||||
}
|
|
||||||
} catch (err) {
|
|
||||||
console.error(`[federation-worker] Initial sync with ${peer.origin} failed:`, err);
|
|
||||||
// Don't update lastSyncedAt — will retry next startup
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
export function startFederationWorkers(): void {
|
export function startFederationWorkers(): void {
|
||||||
console.log('[federation-worker] Federation workers started');
|
console.log('[federation-worker] Federation workers started');
|
||||||
scheduleOutboxTick();
|
scheduleOutboxTick();
|
||||||
scheduleFileQueueTick();
|
scheduleFileQueueTick();
|
||||||
scheduleHealthCheckTick();
|
scheduleHealthCheckTick();
|
||||||
scheduleJanitorTick();
|
scheduleJanitorTick();
|
||||||
// Run initial sync for newly peered instances (async, non-blocking)
|
// Bootstrap sync for freshly-peered rows (async, non-blocking)
|
||||||
runInitialSyncForNewPeers().catch((err) => {
|
startupBootstrapSync().catch((err) => {
|
||||||
console.error('[federation-worker] Initial sync error:', err);
|
console.error('[federation-worker] Startup bootstrap sync error:', err);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user