From 2a741a0dc777cfc6940d5a765e04d14d9b8ad37d Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Fri, 27 Mar 2026 00:46:16 +0100 Subject: [PATCH] refactor(federation): update worker and janitor for generalized outbox columns Replace dmChannelId/messageId column references with contextId/entityId/contextType in federationWorker outbox delivery, spread all payload fields (membership, ownership, group, friendship), add friend-context initial sync pass, and fix storageJanitor DM purge queries. --- packages/server/src/utils/federationWorker.ts | 70 ++++++++++++++++--- packages/server/src/utils/storageJanitor.ts | 4 +- 2 files changed, 62 insertions(+), 12 deletions(-) diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index c15933fc..5587b216 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -108,8 +108,9 @@ async function processOutboxTick(): Promise { .select({ outboxId: schema.federationOutbox.id, peerId: schema.federationOutbox.peerId, - dmChannelId: schema.federationOutbox.dmChannelId, - messageId: schema.federationOutbox.messageId, + contextId: schema.federationOutbox.contextId, + entityId: schema.federationOutbox.entityId, + contextType: schema.federationOutbox.contextType, eventType: schema.federationOutbox.eventType, payload: schema.federationOutbox.payload, encryptionVersion: schema.federationOutbox.encryptionVersion, @@ -161,17 +162,24 @@ async function processOutboxTick(): Promise { // Build relay events from outbox entries const events: FederationRelayEvent[] = peerEntries.map((entry) => { const parsed = JSON.parse(entry.payload) as Partial; - return { + const isDm = entry.contextType === 'dm' || !entry.contextType; + const evt: FederationRelayEvent = { eventType: entry.eventType as FederationRelayEvent['eventType'], - dmChannelId: entry.dmChannelId, - messageId: entry.messageId, + contextType: (entry.contextType ?? 'dm') as 'dm' | 'friend', + messageId: entry.entityId ?? '', encryptionVersion: (entry.encryptionVersion ?? 0) as 0, timestamp: entry.createdAt, - ...(parsed.participants ? { participants: parsed.participants } : {}), - ...(parsed.message ? { message: parsed.message } : {}), - ...(parsed.reactions ? { reactions: parsed.reactions } : {}), - ...(parsed.reaction ? { reaction: parsed.reaction } : {}), }; + if (isDm && entry.contextId) evt.dmChannelId = entry.contextId; + if (parsed.participants) evt.participants = parsed.participants; + if (parsed.message) evt.message = parsed.message; + if (parsed.reactions) evt.reactions = parsed.reactions; + if (parsed.reaction) evt.reaction = parsed.reaction; + if (parsed.membership) evt.membership = parsed.membership; + if (parsed.ownership) evt.ownership = parsed.ownership; + if (parsed.group) evt.group = parsed.group; + if (parsed.friendship) evt.friendship = parsed.friendship; + return evt; }); const request: FederationRelayRequest = { @@ -205,7 +213,7 @@ async function processOutboxTick(): Promise { // Map accepted messageIds to outbox IDs const acceptedSet = new Set(result.accepted); const acceptedOutboxIds = peerEntries - .filter((e) => acceptedSet.has(e.messageId)) + .filter((e) => acceptedSet.has(e.entityId)) .map((e) => e.outboxId); if (acceptedOutboxIds.length > 0) { @@ -679,6 +687,48 @@ async function runInitialSyncForNewPeers(): Promise { 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, peer.hmacSecret, 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; + + const friendRelayBody = JSON.stringify({ + version: 1, + sourceInstance: peer.origin, + events: friendData.events, + }); + const friendRelayHeaders = buildFederationHeaders(friendRelayBody, peer.hmacSecret, peer.origin); + + await fetch(`${ourOrigin}/api/federation/relay`, { + method: 'POST', + headers: friendRelayHeaders, + body: friendRelayBody, + signal: AbortSignal.timeout(30_000), + }); + + 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() }) diff --git a/packages/server/src/utils/storageJanitor.ts b/packages/server/src/utils/storageJanitor.ts index 50efa356..aadb83cf 100644 --- a/packages/server/src/utils/storageJanitor.ts +++ b/packages/server/src/utils/storageJanitor.ts @@ -515,12 +515,12 @@ export function cleanupSoftDeletedDmChannels(): number { // Delete federation outbox entries tx.delete(schema.federationOutbox) - .where(eq(schema.federationOutbox.dmChannelId, channel.id)) + .where(eq(schema.federationOutbox.contextId, channel.id)) .run(); // Delete mutation log entries tx.delete(schema.federationMutationLog) - .where(eq(schema.federationMutationLog.dmChannelId, channel.id)) + .where(eq(schema.federationMutationLog.contextId, channel.id)) .run(); // Delete the channel itself