From a23e02339e79bd8a687326872a8bf7c95189bce9 Mon Sep 17 00:00:00 2001 From: Jannis Braun <151788261+TheZwiss@users.noreply.github.com> Date: Wed, 22 Apr 2026 00:48:49 +0200 Subject: [PATCH] feat(federation): capture dm_close/reopen/read_state/profile/file_rejected in mutation log MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Four event types previously bypassed appendMutationLog, making them unrecoverable via /api/federation/sync after peer inactivity: - queueDmCloseRelay (dm_close, dm_reopen) - queueReadStateRelay (read_state_update) - handleSizeRejection in federationWorker (file_rejected) - profile PATCH route (profile_update) — two call sites, one appendMutationLog per profile change (not per target origin) The /api/federation/sync response builder is extended to serialize these event types in the next task. --- packages/server/src/routes/users.ts | 16 ++- .../federationOutbox.mutationlog.test.ts | 135 ++++++++++++++++++ packages/server/src/utils/federationOutbox.ts | 12 ++ packages/server/src/utils/federationWorker.ts | 14 +- 4 files changed, 175 insertions(+), 2 deletions(-) create mode 100644 packages/server/src/utils/federationOutbox.mutationlog.test.ts diff --git a/packages/server/src/routes/users.ts b/packages/server/src/routes/users.ts index d49ba692..9d82e6a1 100644 --- a/packages/server/src/routes/users.ts +++ b/packages/server/src/routes/users.ts @@ -8,7 +8,7 @@ import { AVATAR_COLORS } from '@backspace/shared'; import { sanitizeUser } from '../utils/sanitize.js'; import { deleteUploadFile, deleteAttachmentByFilename } from '../utils/fileCleanup.js'; import { tombstoneUser, collectDeletionBroadcastTargets, collectProfileBroadcastTargetIds } from '../utils/userDeletion.js'; -import { queueOutboxEvent, isFederationRelayEnabled } from '../utils/federationOutbox.js'; +import { queueOutboxEvent, isFederationRelayEnabled, appendMutationLog } from '../utils/federationOutbox.js'; import { generateSnowflake } from '../utils/snowflake.js'; import { resizeProfileImage } from '../utils/thumbnail.js'; import { config } from '../config.js'; @@ -374,6 +374,13 @@ export async function userRoutes(app: FastifyInstance): Promise { bio: preUpdateUser.bio, }; + appendMutationLog( + preUpdateUser.id, + preUpdateUser.id, + 'profile_update', + JSON.stringify({ profileUpdate: profilePayload }), + 'profile', + ); for (const targetOrigin of newOrigins) { queueOutboxEvent( preUpdateUser.id, @@ -529,6 +536,13 @@ export async function userRoutes(app: FastifyInstance): Promise { bio: updatedUser!.bio, }; + appendMutationLog( + updatedUser!.id, + updatedUser!.id, + 'profile_update', + JSON.stringify({ profileUpdate: profilePayload }), + 'profile', + ); queueOutboxEvent( updatedUser!.id, // entityId — user's ID (coalesces rapid edits) updatedUser!.id, // contextId — user-scoped diff --git a/packages/server/src/utils/federationOutbox.mutationlog.test.ts b/packages/server/src/utils/federationOutbox.mutationlog.test.ts new file mode 100644 index 00000000..19c5242a --- /dev/null +++ b/packages/server/src/utils/federationOutbox.mutationlog.test.ts @@ -0,0 +1,135 @@ +import { describe, it, expect, beforeEach, vi } from 'vitest'; +import Database from 'better-sqlite3'; +import { drizzle } from 'drizzle-orm/better-sqlite3'; +import fs from 'node:fs'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import * as schema from '../db/schema.js'; +import { eq } from 'drizzle-orm'; + +const __dirname = path.dirname(fileURLToPath(import.meta.url)); +type TestDb = ReturnType>; +let sqlite: Database.Database; +let testDb: TestDb; + +vi.mock('../db/index.js', () => ({ + getDb: () => testDb, + schema, +})); + +vi.mock('../utils/federationAuth.js', () => ({ + getOurOrigin: () => 'https://test.example', + buildFederationHeaders: () => ({}), + generateHmacSecret: () => 'secret', +})); + +let _snowflakeCounter = 1; +vi.mock('../utils/snowflake.js', () => ({ + generateSnowflake: () => String(_snowflakeCounter++), + setWorkerId: vi.fn(), +})); + +function applyMigrations(db: Database.Database): void { + const migrationsDir = path.resolve(__dirname, '../../drizzle'); + const files = fs.readdirSync(migrationsDir).filter(f => f.endsWith('.sql')).sort(); + for (const f of files) { + const sql = fs.readFileSync(path.join(migrationsDir, f), 'utf8'); + const statements = sql.split(/-->\s*statement-breakpoint/); + for (const stmt of statements) { + const clean = stmt.trim(); + if (clean) db.exec(clean); + } + } +} + +function seedSettings(): void { + testDb.insert(schema.instanceSettings).values({ + id: 1, + federationRelayEnabled: 1, + updatedAt: Date.now(), + }).run(); +} + +function seedChannel(id: string, federatedId: string | null): void { + testDb.insert(schema.dmChannels).values({ + id, federatedId, ownerId: null, createdAt: Date.now(), + }).run(); +} + +function seedUser(id: string, username: string): void { + testDb.insert(schema.users).values({ + id, username, displayName: username, passwordHash: 'x', + createdAt: Date.now(), + }).run(); +} + +function seedDmMember(channelId: string, userId: string): void { + testDb.insert(schema.dmMembers).values({ + dmChannelId: channelId, userId, closed: 0, + }).run(); +} + +describe('queueDmCloseRelay — mutation log capture', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + seedSettings(); + }); + + it('appends a mutation log row for dm_close', async () => { + const { queueDmCloseRelay } = await import('./federationOutbox.js'); + seedUser('u-1', 'alice'); + seedChannel('ch-1', 'fed-1'); + seedDmMember('ch-1', 'u-1'); + + queueDmCloseRelay('ch-1', 'u-1', 'dm_close'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'dm_close')).all(); + expect(rows.length).toBe(1); + expect(rows[0]?.contextId).toBe('ch-1'); + expect(rows[0]?.contextType).toBe('dm'); + }); + + it('appends a mutation log row for dm_reopen', async () => { + const { queueDmCloseRelay } = await import('./federationOutbox.js'); + seedUser('u-1', 'alice'); + seedChannel('ch-2', 'fed-2'); + seedDmMember('ch-2', 'u-1'); + + queueDmCloseRelay('ch-2', 'u-1', 'dm_reopen'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'dm_reopen')).all(); + expect(rows.length).toBe(1); + }); +}); + +describe('queueReadStateRelay — mutation log capture', () => { + beforeEach(() => { + sqlite = new Database(':memory:'); + testDb = drizzle(sqlite, { schema }); + applyMigrations(sqlite); + seedSettings(); + }); + + it('appends a mutation log row for read_state_update', async () => { + const { queueReadStateRelay } = await import('./federationOutbox.js'); + seedUser('u-2', 'bob'); + seedChannel('ch-3', 'fed-3'); + seedDmMember('ch-3', 'u-2'); + testDb.insert(schema.dmMessages).values({ + id: 'm-1', dmChannelId: 'ch-3', userId: 'u-2', content: 'hi', + type: 'user', createdAt: Date.now(), + }).run(); + + queueReadStateRelay('ch-3', 'm-1', 'u-2'); + + const rows = testDb.select().from(schema.federationMutationLog) + .where(eq(schema.federationMutationLog.mutationType, 'read_state_update')).all(); + expect(rows.length).toBe(1); + expect(rows[0]?.contextId).toBe('ch-3'); + expect(rows[0]?.contextType).toBe('dm'); + }); +}); diff --git a/packages/server/src/utils/federationOutbox.ts b/packages/server/src/utils/federationOutbox.ts index 6e86414a..d73ee0df 100644 --- a/packages/server/src/utils/federationOutbox.ts +++ b/packages/server/src/utils/federationOutbox.ts @@ -700,6 +700,12 @@ export function queueReadStateRelay( }; const targetOrigins = getGroupDmTargetOrigins(channelId); + appendMutationLog( + `read_state:${channel.federatedId}:${userId}`, + channelId, + 'read_state_update', + JSON.stringify({ user: { homeUserId, homeInstance }, messageRef }), + ); queueOutboxEvent( `read_state:${channel.federatedId}:${userId}`, channelId, @@ -755,6 +761,12 @@ export function queueDmCloseRelay( }; const targetOrigins = getGroupDmTargetOrigins(dmChannelId); + appendMutationLog( + `${eventType}:${channel.federatedId}:${userId}`, + dmChannelId, + eventType, + JSON.stringify({ homeUserId, homeInstance }), + ); queueOutboxEvent( `${eventType}:${channel.federatedId}:${userId}`, dmChannelId, diff --git a/packages/server/src/utils/federationWorker.ts b/packages/server/src/utils/federationWorker.ts index d2beafde..e9c218e0 100644 --- a/packages/server/src/utils/federationWorker.ts +++ b/packages/server/src/utils/federationWorker.ts @@ -2,7 +2,7 @@ import { getDb } from '../db/index.js'; import * as schema from '../db/schema.js'; import { eq, and, lte, asc, inArray, sql } from 'drizzle-orm'; import { config } from '../config.js'; -import { isFederationRelayEnabled, queueOutboxEvent } from './federationOutbox.js'; +import { isFederationRelayEnabled, queueOutboxEvent, appendMutationLog } from './federationOutbox.js'; import { runFederationJanitor } from './storageJanitor.js'; import { buildFederationHeaders, getOurOrigin, generateHmacSecret, ROTATION_GRACE_PERIOD_MS } from './federationAuth.js'; import { evaluateAuthFailure, AUTH_FAILURE_THRESHOLD } from './federationAuthFailure.js'; @@ -703,6 +703,18 @@ function handleSizeRejection( affectedUserIds, }; + appendMutationLog( + localMsg.sourceMessageId, + localMsg.dmChannelId, + 'file_rejected', + JSON.stringify({ + attachmentId: att?.id ?? entry.sourceUrl, + sourceFilename, + rejectionReason: 'size_limit_exceeded', + rejectionLimit: maxUploadSize, + affectedUserIds, + }), + ); queueOutboxEvent( localMsg.sourceMessageId, localMsg.dmChannelId,