diff --git a/packages/server/src/routes/dm.ts b/packages/server/src/routes/dm.ts index a7e603e1..ba63d7b8 100644 --- a/packages/server/src/routes/dm.ts +++ b/packages/server/src/routes/dm.ts @@ -21,7 +21,7 @@ import { import { sanitizeUser } from '../utils/sanitize.js'; import { deleteUploadFile, deleteAttachmentFiles } from '../utils/fileCleanup.js'; import { fetchDmEmbedsForMessages, resolveEmbeds, reResolveEmbeds, embedRowToEmbed } from '../utils/embedResolver.js'; -import { appendMutationLog, queueOutboxEvent, buildRelayPayload } from '../utils/federationOutbox.js'; +import { appendMutationLog, queueOutboxEvent, buildRelayPayload, getDmParticipants } from '../utils/federationOutbox.js'; /** * Batch-fetch reactions for a set of DM message IDs. @@ -973,8 +973,10 @@ export async function dmRoutes(app: FastifyInstance): Promise { // Federation: log mutation and queue for relay appendMutationLog(messageId, id, 'create'); + const participants = getDmParticipants(id); queueOutboxEvent(messageId, id, 'create', JSON.stringify({ message: { ...buildRelayPayload(message, message.user), attachments: [] }, + participants, })); // Resolve embeds asynchronously after responding @@ -1042,6 +1044,7 @@ export async function dmRoutes(app: FastifyInstance): Promise { appendMutationLog(id, msg.dmChannelId, 'update'); queueOutboxEvent(id, msg.dmChannelId, 'update', JSON.stringify({ message: buildRelayPayload(updated, updated.user), + participants: getDmParticipants(msg.dmChannelId), })); // Resolve new embeds asynchronously (old ones already deleted above) diff --git a/packages/server/src/routes/federation.ts b/packages/server/src/routes/federation.ts index 7053c7ef..a6af62cf 100644 --- a/packages/server/src/routes/federation.ts +++ b/packages/server/src/routes/federation.ts @@ -9,7 +9,7 @@ import { config } from '../config.js'; import { connectionManager } from '../ws/handler.js'; import { sanitizeUser } from '../utils/sanitize.js'; import { deleteAttachmentFiles } from '../utils/fileCleanup.js'; -import { canonicalDmPairId } from '../utils/federationOutbox.js'; +import { canonicalDmPairId, getDmParticipants } from '../utils/federationOutbox.js'; import { broadcastDmMessage } from './dm.js'; import type { FederationRelayRequest, FederationRelayResponse, FederationRelayEvent, FederationRelayAttachment, FederationSyncRequest, FederationSyncResponse, DmMessageWithUser } from '@backspace/shared'; @@ -748,6 +748,7 @@ export async function federationRoutes(app: FastifyInstance): Promise { messageId: message.id, encryptionVersion: 0, timestamp: mutation.mutated_at, + participants: getDmParticipants(mutation.dm_channel_id), message: { userId: message.userId, homeUserId, @@ -795,16 +796,25 @@ function resolveLocalUser( homeUserId: string, db: ReturnType, ): typeof schema.users.$inferSelect | undefined { - return db + const candidates = db .select() .from(schema.users) .where( - or( - eq(schema.users.homeUserId, homeUserId), - and(eq(schema.users.id, homeUserId), isNull(schema.users.homeInstance)), + and( + or( + eq(schema.users.homeUserId, homeUserId), + and(eq(schema.users.id, homeUserId), isNull(schema.users.homeInstance)), + ), + eq(schema.users.isDeleted, 0), ), ) - .get(); + .all(); + + // Prefer non-deleted active users; if multiple, prefer the one with homeUserId set + // (replicated user) over a local user match + if (candidates.length === 0) return undefined; + if (candidates.length === 1) return candidates[0]; + return candidates.find(u => u.homeUserId === homeUserId) ?? candidates[0]; } /** @@ -933,10 +943,8 @@ function processCreateEvent( return; } - // Resolve the message author to a local user - const authorUser = resolveLocalUser(event.message.homeUserId, db); - if (!authorUser) { - rejected.push({ messageId: event.messageId, reason: 'user_not_found' }); + if (!event.participants || event.participants.length < 2) { + rejected.push({ messageId: event.messageId, reason: 'missing_participants' }); return; } @@ -957,126 +965,46 @@ function processCreateEvent( return; } - // Resolve the DM recipient. Federated DMs are 1-on-1: one side is the author, - // the other is a local user on this instance. We match via canonical_pair_id. - const authorHomeUserId = event.message.homeUserId; + // Resolve ALL participants to local users + const resolvedParticipants: Array<{ + localUser: typeof schema.users.$inferSelect; + homeUserId: string; + }> = []; - // First, search existing DM channels where the author is already a member - const authorMemberships = db - .select({ dmChannelId: schema.dmMembers.dmChannelId }) - .from(schema.dmMembers) - .where(eq(schema.dmMembers.userId, authorUser.id)) - .all(); - - let localDmChannelId: string | null = null; - - // Check each of the author's DM channels to find the matching one - for (const membership of authorMemberships) { - const channelMembers = db - .select() - .from(schema.dmMembers) - .where(eq(schema.dmMembers.dmChannelId, membership.dmChannelId)) - .all(); - - // For 1-on-1 DMs, there should be exactly 2 members - if (channelMembers.length === 2) { - const otherMember = channelMembers.find(m => m.userId !== authorUser.id); - if (otherMember) { - const otherUser = db - .select() - .from(schema.users) - .where(eq(schema.users.id, otherMember.userId)) - .get(); - - if (otherUser) { - const otherHomeUserId = otherUser.homeUserId || otherUser.id; - const pairId = canonicalDmPairId(authorHomeUserId, otherHomeUserId); - const channel = db - .select() - .from(schema.dmChannels) - .where(eq(schema.dmChannels.id, membership.dmChannelId)) - .get(); - - if (channel?.canonicalPairId === pairId) { - localDmChannelId = membership.dmChannelId; - break; - } - } - } + for (const p of event.participants) { + const localUser = resolveLocalUser(p.homeUserId, db); + if (localUser) { + resolvedParticipants.push({ localUser, homeUserId: p.homeUserId }); } } - // If no existing channel found, search the author's friends for the recipient. - // On cold start (first federated DM), we use the friends list as a hint to - // find the local user and create the DM channel. - if (!localDmChannelId) { - const friendRows = db - .select() - .from(schema.friends) - .where( - or( - eq(schema.friends.userId, authorUser.id), - eq(schema.friends.friendId, authorUser.id), - ), - ) - .all(); - - const friendIds = friendRows.map(f => - f.userId === authorUser.id ? f.friendId : f.userId, - ); - - for (const friendId of friendIds) { - const friendUser = db - .select() - .from(schema.users) - .where(eq(schema.users.id, friendId)) - .get(); - - if (friendUser) { - const friendHomeUserId = friendUser.homeUserId || friendUser.id; - const pairId = canonicalDmPairId(authorHomeUserId, friendHomeUserId); - - // Check if a channel already exists with this pair ID - const existingChannel = db - .select() - .from(schema.dmChannels) - .where(eq(schema.dmChannels.canonicalPairId, pairId)) - .get(); - - if (existingChannel) { - localDmChannelId = existingChannel.id; - break; - } - } - } - - // If still no channel with a matching canonical pair ID, create one. - // The recipient must be a local (non-federated) user who is friends with the author. - if (!localDmChannelId && friendIds.length > 0) { - for (const friendId of friendIds) { - const friendUser = db - .select() - .from(schema.users) - .where(eq(schema.users.id, friendId)) - .get(); - - if (friendUser && !friendUser.homeInstance) { - // This is a local user — they're a candidate recipient - const friendHomeUserId = friendUser.homeUserId || friendUser.id; - const pairId = canonicalDmPairId(authorHomeUserId, friendHomeUserId); - - localDmChannelId = findOrCreateDmChannel(pairId, authorUser.id, friendId, db); - break; - } - } - } - } - - if (!localDmChannelId) { - rejected.push({ messageId: event.messageId, reason: 'recipient_not_found' }); + if (resolvedParticipants.length < 2) { + rejected.push({ messageId: event.messageId, reason: 'participant_not_found' }); return; } + // Find the author among the resolved participants + const authorEntry = resolvedParticipants.find( + p => p.homeUserId === event.message!.homeUserId, + ); + if (!authorEntry) { + rejected.push({ messageId: event.messageId, reason: 'author_not_found' }); + return; + } + const authorUser = authorEntry.localUser; + + // Compute canonical pair ID from participants' home user IDs and find/create channel + const pairId = canonicalDmPairId( + resolvedParticipants[0]!.homeUserId, + resolvedParticipants[1]!.homeUserId, + ); + const localDmChannelId = findOrCreateDmChannel( + pairId, + resolvedParticipants[0]!.localUser.id, + resolvedParticipants[1]!.localUser.id, + db, + ); + // Insert the message const localMessageId = generateSnowflake(); db.insert(schema.dmMessages) diff --git a/packages/server/src/utils/federationOutbox.ts b/packages/server/src/utils/federationOutbox.ts index 96716b9d..81ab6a46 100644 --- a/packages/server/src/utils/federationOutbox.ts +++ b/packages/server/src/utils/federationOutbox.ts @@ -3,7 +3,7 @@ import * as schema from '../db/schema.js'; import { eq, and } from 'drizzle-orm'; import { generateSnowflake } from './snowflake.js'; import crypto from 'node:crypto'; -import type { FederationRelayEvent } from '@backspace/shared'; +import type { FederationRelayEvent, FederationRelayParticipant } from '@backspace/shared'; import { config } from '../config.js'; // ─── Settings Cache ────────────────────────────────────────────────────────── @@ -210,6 +210,33 @@ export function canonicalDmPairId(homeUserIdA: string, homeUserIdB: string): str return crypto.createHash('sha256').update(sorted.join(':')).digest('hex').slice(0, 32); } +/** + * Look up all members of a DM channel and return their federated identities. + * Used to include participants in relay events so the receiving instance can + * resolve both parties without relying on the friends list. + */ +export function getDmParticipants(dmChannelId: string): FederationRelayParticipant[] { + const db = getDb(); + const members = db + .select({ + userId: schema.dmMembers.userId, + homeUserId: schema.users.homeUserId, + homeInstance: schema.users.homeInstance, + id: schema.users.id, + }) + .from(schema.dmMembers) + .innerJoin(schema.users, eq(schema.dmMembers.userId, schema.users.id)) + .where(eq(schema.dmMembers.dmChannelId, dmChannelId)) + .all(); + + const domainOrigin = config.domain ? `https://${config.domain}` : ''; + + return members.map(m => ({ + homeUserId: m.homeUserId || m.id, + homeInstance: m.homeInstance || domainOrigin, + })); +} + /** * Build the relay payload object for a DM message. * The caller may augment the returned object with attachments before serialization. diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index 2edd1034..df756678 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -172,6 +172,7 @@ async function processOutboxTick(): Promise { messageId: entry.messageId, 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 } : {}), diff --git a/packages/server/src/ws/events.ts b/packages/server/src/ws/events.ts index b98e2e0b..e2493a44 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, buildRelayPayload } from '../utils/federationOutbox.js'; +import { appendMutationLog, queueOutboxEvent, buildRelayPayload, getDmParticipants } from '../utils/federationOutbox.js'; /** * Re-evaluate SPEAK permission for all participants in voice channels @@ -888,8 +888,10 @@ function handleDmMessageCreate(event: Record, userId: string): // Federation: log mutation and queue for relay appendMutationLog(messageId, dmChannelId, 'create'); + const participants = getDmParticipants(dmChannelId); queueOutboxEvent(messageId, dmChannelId, 'create', JSON.stringify({ message: { ...buildRelayPayload(dmMessage, dmMessage.user), attachments: [] }, + participants, })); // Resolve embeds asynchronously @@ -994,6 +996,7 @@ function handleDmMessageEdit(event: Record, userId: string): vo appendMutationLog(messageId, msg.dmChannelId, 'update'); queueOutboxEvent(messageId, msg.dmChannelId, 'update', JSON.stringify({ message: buildRelayPayload(updated, updated.user), + participants: getDmParticipants(msg.dmChannelId), })); // Resolve new embeds asynchronously (old ones already deleted above) diff --git a/packages/shared/src/types.ts b/packages/shared/src/types.ts index 970d7c6e..95ca25a5 100644 --- a/packages/shared/src/types.ts +++ b/packages/shared/src/types.ts @@ -717,12 +717,18 @@ export interface AdminResetPasswordResponse { // ─── Federation Relay Types ────────────────────────────────────────────────── +export interface FederationRelayParticipant { + homeUserId: string; + homeInstance: string; +} + export interface FederationRelayEvent { eventType: 'create' | 'update' | 'delete' | 'reaction_add' | 'reaction_remove'; dmChannelId: string; messageId: string; encryptionVersion: 0; timestamp: number; + participants?: FederationRelayParticipant[]; message?: { userId: string; homeUserId: string;