feat(server): aggregate Path-1 call fan-out failures and surface to originator

This commit is contained in:
Jannis Braun
2026-04-23 23:12:55 +02:00
parent 13241345de
commit 3cb80d110d
+112 -25
View File
@@ -1512,8 +1512,14 @@ async function handleDmCallAccept(event: Record<string, unknown>, 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<string, unknown>, 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<string, unknown>, 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<void> {
/**
* 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<CallFanoutFailure[]> {
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<string, string | null>();
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<void> {
async function sendFederatedCallEnd(
dmChannelId: string,
endedByUserId: string,
): Promise<CallFanoutFailure[]> {
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<string, string | null>();
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 ──────────────────────────────────────────────