diff --git a/packages/server/src/routes/dm.ts b/packages/server/src/routes/dm.ts index 23416951..a7e603e1 100644 --- a/packages/server/src/routes/dm.ts +++ b/packages/server/src/routes/dm.ts @@ -21,6 +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'; /** * Batch-fetch reactions for a set of DM message IDs. @@ -970,6 +971,12 @@ export async function dmRoutes(app: FastifyInstance): Promise { // Broadcast to all DM members (including those who closed the channel) broadcastDmMessage(id, message); + // Federation: log mutation and queue for relay + appendMutationLog(messageId, id, 'create'); + queueOutboxEvent(messageId, id, 'create', JSON.stringify({ + message: { ...buildRelayPayload(message, message.user), attachments: [] }, + })); + // Resolve embeds asynchronously after responding setImmediate(() => { resolveEmbeds(messageId, content?.trim() || null, id, true, null).catch(() => {}); @@ -1031,6 +1038,12 @@ export async function dmRoutes(app: FastifyInstance): Promise { }); } + // Federation: log mutation and queue for relay + appendMutationLog(id, msg.dmChannelId, 'update'); + queueOutboxEvent(id, msg.dmChannelId, 'update', JSON.stringify({ + message: buildRelayPayload(updated, updated.user), + })); + // Resolve new embeds asynchronously (old ones already deleted above) setImmediate(() => { resolveEmbeds(id, content.trim(), msg.dmChannelId, true, null).catch(() => {}); @@ -1093,6 +1106,10 @@ export async function dmRoutes(app: FastifyInstance): Promise { }); } + // Federation: log mutation and queue for relay + appendMutationLog(id, msg.dmChannelId, 'delete'); + queueOutboxEvent(id, msg.dmChannelId, 'delete', JSON.stringify({ deleted: true })); + return reply.code(200).send({ success: true }); }); } diff --git a/packages/server/src/utils/federationOutbox.ts b/packages/server/src/utils/federationOutbox.ts index f4c60f9e..c030c8cd 100644 --- a/packages/server/src/utils/federationOutbox.ts +++ b/packages/server/src/utils/federationOutbox.ts @@ -217,8 +217,8 @@ export function buildRelayPayload( message: { id: string; content: string | null; - replyToId: string | null; - editedAt: number | null; + replyToId?: string | null; + editedAt?: number | null; createdAt: number; }, user: { @@ -232,8 +232,8 @@ export function buildRelayPayload( homeUserId: user.homeUserId || user.id, homeInstance: user.homeInstance || '', content: message.content, - replyToId: message.replyToId, - editedAt: message.editedAt, + replyToId: message.replyToId ?? null, + editedAt: message.editedAt ?? null, createdAt: message.createdAt, }; } diff --git a/packages/server/src/ws/events.ts b/packages/server/src/ws/events.ts index b62b3597..b98e2e0b 100644 --- a/packages/server/src/ws/events.ts +++ b/packages/server/src/ws/events.ts @@ -11,6 +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'; /** * Re-evaluate SPEAK permission for all participants in voice channels @@ -885,6 +886,12 @@ function handleDmMessageCreate(event: Record, userId: string): // Broadcast to all DM members (including those who closed the channel) broadcastDmMessage(dmChannelId, dmMessage); + // Federation: log mutation and queue for relay + appendMutationLog(messageId, dmChannelId, 'create'); + queueOutboxEvent(messageId, dmChannelId, 'create', JSON.stringify({ + message: { ...buildRelayPayload(dmMessage, dmMessage.user), attachments: [] }, + })); + // Resolve embeds asynchronously setImmediate(() => { resolveEmbeds(messageId, hasContent ? content!.trim() : null, dmChannelId, true, null).catch(() => {}); @@ -983,6 +990,12 @@ function handleDmMessageEdit(event: Record, userId: string): vo }); } + // Federation: log mutation and queue for relay + appendMutationLog(messageId, msg.dmChannelId, 'update'); + queueOutboxEvent(messageId, msg.dmChannelId, 'update', JSON.stringify({ + message: buildRelayPayload(updated, updated.user), + })); + // Resolve new embeds asynchronously (old ones already deleted above) setImmediate(() => { resolveEmbeds(messageId, content.trim(), msg.dmChannelId, true, null).catch(() => {}); @@ -1041,6 +1054,10 @@ function handleDmMessageDelete(event: Record, userId: string): dmChannelId: msg.dmChannelId, }); } + + // Federation: log mutation and queue for relay + appendMutationLog(messageId, msg.dmChannelId, 'delete'); + queueOutboxEvent(messageId, msg.dmChannelId, 'delete', JSON.stringify({ deleted: true })); } // ─── Reaction Handlers ───────────────────────────────────────────────────── @@ -1114,6 +1131,22 @@ function handleReactionAdd(event: Record, userId: string): void messageId, reaction: { id: reactionId, messageId, userId, emoji, createdAt: now, user: userObj }, }); + + // Federation: log reaction mutation and queue for relay + appendMutationLog(messageId, dmMsg.dmChannelId, 'reaction_add', JSON.stringify({ + userId, + homeUserId: reactionUser?.homeUserId || userId, + emoji, + createdAt: now, + })); + queueOutboxEvent(reactionId, dmMsg.dmChannelId, 'reaction_add', JSON.stringify({ + reaction: { + userId, + homeUserId: reactionUser?.homeUserId || userId, + emoji, + createdAt: now, + }, + })); } catch (err) { // Unique constraint violation (already reacted) } @@ -1171,6 +1204,26 @@ function handleReactionRemove(event: Record, userId: string): v userId, emoji, }); + + // Federation: log reaction removal and queue for relay + const removingUser = db.select().from(schema.users).where(eq(schema.users.id, userId)).get(); + appendMutationLog(messageId, dmMsg.dmChannelId, 'reaction_remove', JSON.stringify({ + userId, + homeUserId: removingUser?.homeUserId || userId, + emoji, + })); + queueOutboxEvent( + `${messageId}:${userId}:${emoji}`, + dmMsg.dmChannelId, + 'reaction_remove', + JSON.stringify({ + reaction: { + userId, + homeUserId: removingUser?.homeUserId || userId, + emoji, + }, + }), + ); } }