From 32ccd9c41072b83422a51cbe5ef01b0b4de57d9f Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Tue, 7 Apr 2026 19:49:52 +0200 Subject: [PATCH] feat: outbound S2S read state relay Queue read_state_update events when users ack DM messages on channels with a federatedId. Translates local message IDs to federation coordinates using sourceInstance/sourceMessageId. --- packages/server/src/utils/federationOutbox.ts | 70 +++++++++++++++++++ packages/server/src/ws/events.ts | 5 +- 2 files changed, 74 insertions(+), 1 deletion(-) diff --git a/packages/server/src/utils/federationOutbox.ts b/packages/server/src/utils/federationOutbox.ts index 26dd9e0c..35076479 100644 --- a/packages/server/src/utils/federationOutbox.ts +++ b/packages/server/src/utils/federationOutbox.ts @@ -472,6 +472,76 @@ export async function sendCallRelay( } } +/** + * Queue a read_state_update relay event for cross-instance read state sync. + * Translates a local channel ack into federation coordinates using the + * message's sourceInstance/sourceMessageId mapping. + */ +export function queueReadStateRelay( + channelId: string, + messageId: string, + userId: string, +): void { + if (!isFederationRelayEnabled()) return; + + const db = getDb(); + const ourOrigin = getOurOrigin(); + + // Channel must have a federatedId for cross-instance sync + const channel = db.select({ federatedId: schema.dmChannels.federatedId, ownerId: schema.dmChannels.ownerId }) + .from(schema.dmChannels) + .where(eq(schema.dmChannels.id, channelId)) + .get(); + if (!channel?.federatedId) return; + + // Resolve the user's federated identity + const user = db.select({ homeUserId: schema.users.homeUserId, homeInstance: schema.users.homeInstance }) + .from(schema.users) + .where(eq(schema.users.id, userId)) + .get(); + if (!user) return; + + const homeUserId = user.homeUserId || userId; + const homeInstance = user.homeInstance || ourOrigin; + + // Determine the acked message's federation coordinates + const msg = db.select({ sourceInstance: schema.dmMessages.sourceInstance, sourceMessageId: schema.dmMessages.sourceMessageId }) + .from(schema.dmMessages) + .where(eq(schema.dmMessages.id, messageId)) + .get(); + + let messageRef: { sourceInstance: string; sourceMessageId: string }; + if (msg?.sourceInstance && msg?.sourceMessageId) { + // Message was relayed here — use its original coordinates + messageRef = { sourceInstance: msg.sourceInstance, sourceMessageId: msg.sourceMessageId }; + } else { + // Message originated on this instance + messageRef = { sourceInstance: ourOrigin, sourceMessageId: messageId }; + } + + const payload: FederationRelayEvent = { + eventType: 'read_state_update', + dmChannelId: channelId, + messageId: `read_state:${userId}:${Date.now()}`, + federatedId: channel.federatedId, + encryptionVersion: 0, + timestamp: Date.now(), + readState: { + user: { homeUserId, homeInstance }, + messageRef, + }, + }; + + const targetOrigins = getGroupDmTargetOrigins(channelId); + queueOutboxEvent( + `read_state:${channel.federatedId}:${userId}`, + channelId, + 'read_state_update', + JSON.stringify(payload), + targetOrigins, + ); +} + /** * Send typing indicator events directly to remote peers (bypasses outbox). * Fire-and-forget — typing is ephemeral, lost packets are acceptable. diff --git a/packages/server/src/ws/events.ts b/packages/server/src/ws/events.ts index 47fea4e8..0874ebbb 100644 --- a/packages/server/src/ws/events.ts +++ b/packages/server/src/ws/events.ts @@ -11,7 +11,7 @@ import { ACTIVITY_LIMITS } from '@backspace/shared/src/activities.js'; import { sanitizeUser } from '../utils/sanitize.js'; import { deleteAttachmentFiles } from '../utils/fileCleanup.js'; import { resolveEmbeds, reResolveEmbeds, embedRowToEmbed } from '../utils/embedResolver.js'; -import { appendMutationLog, queueOutboxEvent, queueDmRelay, getGroupDmTargetOrigins, sendCallRelay, computeFederatedId, sendTypingRelay } from '../utils/federationOutbox.js'; +import { appendMutationLog, queueOutboxEvent, queueDmRelay, getGroupDmTargetOrigins, sendCallRelay, computeFederatedId, sendTypingRelay, queueReadStateRelay } from '../utils/federationOutbox.js'; import { getOurOrigin } from '../utils/federationAuth.js'; import { generateFederatedCallToken } from '../routes/livekit.js'; import { config } from '../config.js'; @@ -1320,6 +1320,9 @@ function handleChannelAck(event: Record, userId: string, isFede channelId, messageId, }); + + // Relay read state to federated peers for cross-instance sync + queueReadStateRelay(channelId, messageId, userId); } function handleMarkUnread(event: Record, userId: string, isFederated: boolean): void {