diff --git a/packages/server/src/ws/events.ts b/packages/server/src/ws/events.ts index edcafe04..4dc5415c 100644 --- a/packages/server/src/ws/events.ts +++ b/packages/server/src/ws/events.ts @@ -1512,8 +1512,14 @@ async function handleDmCallAccept(event: Record, userId: string userId, action: 'join', }); - sendFederatedCallAccept(dmChannelId, userId) - .catch(err => console.error('[federation] sendFederatedCallAccept error:', err)); + const acceptFanoutFailures = await sendFederatedCallAccept(dmChannelId, userId); + emitFanoutUndeliverable( + userId, + dmChannelId, + fedIdRow?.federatedId ?? null, + 'accept', + acceptFanoutFailures, + ); return; } } @@ -1598,8 +1604,16 @@ async function handleDmCallReject(event: Record, userId: string connectionManager.clearVoiceWs(meta.callerId); connectionManager.destroyRoom(dmChannelId); connectionManager.sendToDmMembers(dmChannelId, { type: 'dm_call_rejected', dmChannelId }); - sendFederatedCallEnd(dmChannelId, userId) - .catch(err => console.error('[federation] sendFederatedCallEnd error:', err)); + const fedIdRejectRow = getDb().select({ federatedId: schema.dmChannels.federatedId }) + .from(schema.dmChannels).where(eq(schema.dmChannels.id, dmChannelId)).get(); + const rejectFanoutFailures = await sendFederatedCallEnd(dmChannelId, userId); + emitFanoutUndeliverable( + userId, + dmChannelId, + fedIdRejectRow?.federatedId ?? null, + 'reject', + rejectFanoutFailures, + ); return; } } @@ -1684,10 +1698,18 @@ async function handleDmCallEnd(event: Record, userId: string): connectionManager.clearVoiceUserStatus(participantId); connectionManager.clearVoiceWs(participantId); } + const fedIdEndRow = getDb().select({ federatedId: schema.dmChannels.federatedId }) + .from(schema.dmChannels).where(eq(schema.dmChannels.id, dmChannelId)).get(); connectionManager.destroyRoom(dmChannelId); connectionManager.sendToDmMembers(dmChannelId, { type: 'dm_call_ended', dmChannelId }); - sendFederatedCallEnd(dmChannelId, userId) - .catch(err => console.error('[federation] sendFederatedCallEnd error:', err)); + const endFanoutFailures = await sendFederatedCallEnd(dmChannelId, userId); + emitFanoutUndeliverable( + userId, + dmChannelId, + fedIdEndRow?.federatedId ?? null, + 'end', + endFanoutFailures, + ); return; } } @@ -2026,13 +2048,44 @@ function emitUndeliverableAndMaybeDestroy(args: { }); } -async function sendFederatedCallAccept(dmChannelId: string, acceptorUserId: string): Promise { +/** + * Emit a non-terminal dm_call_undeliverable to the local user whose Path-1 action + * (accept / reject / end) had one or more fan-out failures reach peers. + * No-op when there were no failures or no federatedId to reference. + */ +function emitFanoutUndeliverable( + userId: string, + dmChannelId: string | null, + federatedId: string | null | undefined, + phase: 'accept' | 'reject' | 'end', + fanoutFailures: CallFanoutFailure[], +): void { + if (fanoutFailures.length === 0 || !federatedId) return; + const failures: DmCallUndeliverableFailure[] = fanoutFailures.map(f => ({ + reason: f.reason, + peerOrigin: f.origin, + peerLabel: f.peerLabel, + })); + connectionManager.sendToUser(userId, { + type: 'dm_call_undeliverable', + dmChannelId, + federatedCallId: federatedId, + terminal: false, + phase, + failures, + }); +} + +async function sendFederatedCallAccept( + dmChannelId: string, + acceptorUserId: string, +): Promise { const db = getDb(); const channel = db.select({ federatedId: schema.dmChannels.federatedId }) .from(schema.dmChannels) .where(eq(schema.dmChannels.id, dmChannelId)) .get(); - if (!channel?.federatedId) return; + if (!channel?.federatedId) return []; const members = db.select({ homeInstance: schema.users.homeInstance }) .from(schema.dmMembers) @@ -2048,7 +2101,7 @@ async function sendFederatedCallAccept(dmChannelId: string, acceptorUserId: stri if (normalized !== ourOrigin) targets.add(normalized); } } - if (targets.size === 0) return; + if (targets.size === 0) return []; const user = db.select({ homeUserId: schema.users.homeUserId }) .from(schema.users) @@ -2067,22 +2120,41 @@ async function sendFederatedCallAccept(dmChannelId: string, acceptorUserId: stri }, }; - await Promise.all( - Array.from(targets).map(origin => - sendCallRelay(origin, [event]).catch(err => - console.error(`[federation] Failed to send dm_call_accept to ${origin}:`, err) - ) - ) + const labelByOrigin = new Map(); + for (const r of db.select({ origin: schema.federationPeers.origin, instanceName: schema.federationPeers.instanceName }) + .from(schema.federationPeers) + .all()) { + labelByOrigin.set(r.origin, r.instanceName ?? null); + } + + const results = await Promise.all( + Array.from(targets).map(async origin => ({ origin, result: await sendCallRelay(origin, [event]) })), ); + + const failures: CallFanoutFailure[] = []; + for (const { origin, result } of results) { + if (!result.ok) { + console.error(`[federation] dm_call_accept fanout to ${origin} failed (${result.reason}): ${result.error}`); + failures.push({ + origin, + peerLabel: labelByOrigin.get(origin) ?? undefined, + reason: mapCallReasonToEventReason(result.reason), + }); + } + } + return failures; } -async function sendFederatedCallEnd(dmChannelId: string, endedByUserId: string): Promise { +async function sendFederatedCallEnd( + dmChannelId: string, + endedByUserId: string, +): Promise { const db = getDb(); const channel = db.select({ federatedId: schema.dmChannels.federatedId }) .from(schema.dmChannels) .where(eq(schema.dmChannels.id, dmChannelId)) .get(); - if (!channel?.federatedId) return; + if (!channel?.federatedId) return []; const members = db.select({ homeInstance: schema.users.homeInstance }) .from(schema.dmMembers) @@ -2098,9 +2170,8 @@ async function sendFederatedCallEnd(dmChannelId: string, endedByUserId: string): if (normalized !== ourOrigin) targets.add(normalized); } } - if (targets.size === 0) return; + if (targets.size === 0) return []; - // Resolve actual homeUserId from DB const endUser = db.select({ homeUserId: schema.users.homeUserId }) .from(schema.users) .where(eq(schema.users.id, endedByUserId)) @@ -2118,13 +2189,29 @@ async function sendFederatedCallEnd(dmChannelId: string, endedByUserId: string): }, }; - await Promise.all( - Array.from(targets).map(origin => - sendCallRelay(origin, [event]).catch(err => - console.error(`[federation] Failed to send dm_call_end to ${origin}:`, err) - ) - ) + const labelByOrigin = new Map(); + for (const r of db.select({ origin: schema.federationPeers.origin, instanceName: schema.federationPeers.instanceName }) + .from(schema.federationPeers) + .all()) { + labelByOrigin.set(r.origin, r.instanceName ?? null); + } + + const results = await Promise.all( + Array.from(targets).map(async origin => ({ origin, result: await sendCallRelay(origin, [event]) })), ); + + const failures: CallFanoutFailure[] = []; + for (const { origin, result } of results) { + if (!result.ok) { + console.error(`[federation] dm_call_end fanout to ${origin} failed (${result.reason}): ${result.error}`); + failures.push({ + origin, + peerLabel: labelByOrigin.get(origin) ?? undefined, + reason: mapCallReasonToEventReason(result.reason), + }); + } + } + return failures; } // ─── Voice Moderation Handlers ──────────────────────────────────────────────