feat(federation): extend catch-up sync to include group DMs and membership mutations
This commit is contained in:
@@ -529,6 +529,7 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
|
|
||||||
const sinceTimestamp = body.sinceTimestamp;
|
const sinceTimestamp = body.sinceTimestamp;
|
||||||
const dmChannelIdFilter = body.dmChannelId && typeof body.dmChannelId === 'string' ? body.dmChannelId : null;
|
const dmChannelIdFilter = body.dmChannelId && typeof body.dmChannelId === 'string' ? body.dmChannelId : null;
|
||||||
|
const federatedIdFilter = body.federatedId && typeof body.federatedId === 'string' ? body.federatedId : null;
|
||||||
|
|
||||||
// Clamp limit: min 1, max 500, default 100
|
// Clamp limit: min 1, max 500, default 100
|
||||||
let limit = typeof body.limit === 'number' ? body.limit : 100;
|
let limit = typeof body.limit === 'number' ? body.limit : 100;
|
||||||
@@ -545,6 +546,17 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
|
|
||||||
const sharedChannelIds = sharedChannelRows.map(r => r.dm_channel_id);
|
const sharedChannelIds = sharedChannelRows.map(r => r.dm_channel_id);
|
||||||
|
|
||||||
|
// If filtering by federatedId, resolve to local channel ID
|
||||||
|
let effectiveChannelFilter = dmChannelIdFilter;
|
||||||
|
if (federatedIdFilter && !effectiveChannelFilter) {
|
||||||
|
const fedChannel = rawDb.prepare(`
|
||||||
|
SELECT id FROM dm_channels WHERE federated_id = ?
|
||||||
|
`).get(federatedIdFilter) as { id: string } | undefined;
|
||||||
|
if (fedChannel) {
|
||||||
|
effectiveChannelFilter = fedChannel.id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (sharedChannelIds.length === 0) {
|
if (sharedChannelIds.length === 0) {
|
||||||
const syncResponse: FederationSyncResponse = {
|
const syncResponse: FederationSyncResponse = {
|
||||||
events: [],
|
events: [],
|
||||||
@@ -567,9 +579,9 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
payload: string | null;
|
payload: string | null;
|
||||||
}>;
|
}>;
|
||||||
|
|
||||||
if (dmChannelIdFilter) {
|
if (effectiveChannelFilter) {
|
||||||
// Validate that the requested channel is actually shared with this peer
|
// Validate that the requested channel is actually shared with this peer
|
||||||
if (!sharedChannelIds.includes(dmChannelIdFilter)) {
|
if (!sharedChannelIds.includes(effectiveChannelFilter)) {
|
||||||
const syncResponse: FederationSyncResponse = {
|
const syncResponse: FederationSyncResponse = {
|
||||||
events: [],
|
events: [],
|
||||||
hasMore: false,
|
hasMore: false,
|
||||||
@@ -581,12 +593,13 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
mutationRows = rawDb.prepare(`
|
mutationRows = rawDb.prepare(`
|
||||||
SELECT ml.id, ml.dm_message_id, ml.dm_channel_id, ml.mutation_type, ml.mutated_at, ml.payload
|
SELECT ml.id, ml.dm_message_id, ml.dm_channel_id, ml.mutation_type, ml.mutated_at, ml.payload
|
||||||
FROM federation_mutation_log ml
|
FROM federation_mutation_log ml
|
||||||
JOIN dm_messages dm ON ml.dm_message_id = dm.id
|
LEFT JOIN dm_messages dm ON ml.dm_message_id = dm.id
|
||||||
WHERE ml.dm_channel_id = ?
|
WHERE ml.dm_channel_id = ?
|
||||||
AND ml.mutated_at > ?
|
AND ml.mutated_at > ?
|
||||||
|
AND (dm.id IS NOT NULL OR ml.mutation_type IN ('delete', 'member_add', 'member_remove', 'ownership_transfer'))
|
||||||
ORDER BY ml.mutated_at ASC
|
ORDER BY ml.mutated_at ASC
|
||||||
LIMIT ?
|
LIMIT ?
|
||||||
`).all(dmChannelIdFilter, sinceTimestamp, limit) as typeof mutationRows;
|
`).all(effectiveChannelFilter, sinceTimestamp, limit) as typeof mutationRows;
|
||||||
|
|
||||||
// For delete mutations, the dm_messages row won't exist — handle separately
|
// For delete mutations, the dm_messages row won't exist — handle separately
|
||||||
const deleteMutations = rawDb.prepare(`
|
const deleteMutations = rawDb.prepare(`
|
||||||
@@ -598,7 +611,7 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
AND ml.dm_message_id NOT IN (SELECT dm.id FROM dm_messages dm WHERE dm.id = ml.dm_message_id)
|
AND ml.dm_message_id NOT IN (SELECT dm.id FROM dm_messages dm WHERE dm.id = ml.dm_message_id)
|
||||||
ORDER BY ml.mutated_at ASC
|
ORDER BY ml.mutated_at ASC
|
||||||
LIMIT ?
|
LIMIT ?
|
||||||
`).all(dmChannelIdFilter, sinceTimestamp, limit) as typeof mutationRows;
|
`).all(effectiveChannelFilter, sinceTimestamp, limit) as typeof mutationRows;
|
||||||
|
|
||||||
// Merge, deduplicate, sort, and re-limit
|
// Merge, deduplicate, sort, and re-limit
|
||||||
const seen = new Set(mutationRows.map(r => r.id));
|
const seen = new Set(mutationRows.map(r => r.id));
|
||||||
@@ -619,9 +632,10 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
mutationRows = rawDb.prepare(`
|
mutationRows = rawDb.prepare(`
|
||||||
SELECT ml.id, ml.dm_message_id, ml.dm_channel_id, ml.mutation_type, ml.mutated_at, ml.payload
|
SELECT ml.id, ml.dm_message_id, ml.dm_channel_id, ml.mutation_type, ml.mutated_at, ml.payload
|
||||||
FROM federation_mutation_log ml
|
FROM federation_mutation_log ml
|
||||||
JOIN dm_messages dm ON ml.dm_message_id = dm.id
|
LEFT JOIN dm_messages dm ON ml.dm_message_id = dm.id
|
||||||
WHERE ml.dm_channel_id IN (${placeholders})
|
WHERE ml.dm_channel_id IN (${placeholders})
|
||||||
AND ml.mutated_at > ?
|
AND ml.mutated_at > ?
|
||||||
|
AND (dm.id IS NOT NULL OR ml.mutation_type IN ('delete', 'member_add', 'member_remove', 'ownership_transfer'))
|
||||||
ORDER BY ml.mutated_at ASC
|
ORDER BY ml.mutated_at ASC
|
||||||
LIMIT ?
|
LIMIT ?
|
||||||
`).all(...sharedChannelIds, sinceTimestamp, limit) as typeof mutationRows;
|
`).all(...sharedChannelIds, sinceTimestamp, limit) as typeof mutationRows;
|
||||||
@@ -656,7 +670,21 @@ export async function federationRoutes(app: FastifyInstance): Promise<void> {
|
|||||||
const events: FederationRelayEvent[] = [];
|
const events: FederationRelayEvent[] = [];
|
||||||
|
|
||||||
for (const mutation of mutationRows) {
|
for (const mutation of mutationRows) {
|
||||||
const mutationType = mutation.mutation_type as 'create' | 'update' | 'delete' | 'reaction_add' | 'reaction_remove';
|
const mutationType = mutation.mutation_type as 'create' | 'update' | 'delete' | 'reaction_add' | 'reaction_remove' | 'member_add' | 'member_remove' | 'ownership_transfer';
|
||||||
|
|
||||||
|
if (['member_add', 'member_remove', 'ownership_transfer'].includes(mutationType)) {
|
||||||
|
// Membership mutations store the full event in the payload
|
||||||
|
const payload = mutation.payload ? JSON.parse(mutation.payload) : {};
|
||||||
|
events.push({
|
||||||
|
eventType: mutationType as FederationRelayEvent['eventType'],
|
||||||
|
dmChannelId: mutation.dm_channel_id,
|
||||||
|
messageId: mutation.dm_message_id,
|
||||||
|
encryptionVersion: 0,
|
||||||
|
timestamp: mutation.mutated_at,
|
||||||
|
...payload,
|
||||||
|
});
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
if (mutationType === 'delete') {
|
if (mutationType === 'delete') {
|
||||||
// For deletes, we don't need the message content — just the ID and channel
|
// For deletes, we don't need the message content — just the ID and channel
|
||||||
|
|||||||
Reference in New Issue
Block a user