Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions docs/calls-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,13 @@ Call `calls/capabilities` before joining. A 1.x client requires protocol major 1

## Command and event flow

1. Join with `calls/rooms/join`, fetch `calls/rooms/show`, then subscribe to `callsRoom` using `{ "roomId": "..." }`.
1. Join with `calls/rooms/join`, retain the returned provisional participant, fetch `calls/rooms/show`, then subscribe to `callsRoom` using `{ "roomId": "..." }`. Public snapshots and user Calls presence exclude participants whose media connection is not ready.
2. Treat WebSocket events as the normal state path. Each event body contains a monotonically increasing `sequence`, authoritative `roomRevision`, and `occurredAt` timestamp. The subscribed channel identifies the room.
3. Ignore an event whose sequence and room revision were already applied. On a forward gap, or when a lower sequence arrives with a newer room revision after Redis-state loss, fetch a snapshot and `calls/media/reconcile` before applying later events.
4. Create a media session and retain its short-lived media credential. The credential is bound to instance, user, application, room, participant, connection generation, capability, expiry, and nonce.
5. Publish, subscribe, or renegotiate over HTTPS. Apply an answer directly; when an offer or `requiresImmediateRenegotiation` is returned, create an answer and send it to `calls/media/renegotiate`.
6. Close media and leave. A terminal room event or `revoked` event requires immediate local track stop and peer-connection teardown.
6. After initial negotiation, wait for the peer connection to connect and remote audio playback to start. Send `ready` on `callsRoom` with `{ "connectionId": "...", "generation": 1 }`. An empty listening room needs no transport or remote playback. If autoplay is blocked, wait for successful user-initiated playback before sending `ready`. The server validates the current bound media session, then emits the participant `joined` event and includes the participant in public snapshots and presence. Each new media generation must send `ready`; repeated commands for the same generation are idempotent.
7. Close media and leave. A terminal room event or `revoked` event requires immediate local track stop and peer-connection teardown.

If a heartbeat detects lost live state, the server emits `revoked` with reason `stale-generation`. Tear down the old peer connection and create a new media session/generation before republishing or resubscribing. Other revoke reasons are terminal for the current access decision and must not be retried without rejoining or refreshing authorization.

Expand Down Expand Up @@ -49,7 +50,9 @@ Retry network failures, HTTP 429, and provider 5xx responses with bounded expone

Canonical error categories are `invalid-request`, `authentication-required`, `access-denied`, `not-found`, `conflict`, `stale-revision`, `stale-generation`, `rate-limited`, `provider-unavailable`, and `terminal-room`. Error details must not contain secrets, raw SDP, ICE candidate addresses, or provider identifiers.

The independent browser example is in `packages/calls-reference-client`; it imports only the public `misskey-js` package.
Clients must implement the `ready` command to appear as participants. Reload existing clients after deploying this change.

The independent browser example is in `packages/calls-reference-client`; it imports only the public `misskey-js` package. Its browser integration must call `confirmReady()` after completing transport and audio playback setup, including after recovery.

## Instance operation

Expand Down
14 changes: 14 additions & 0 deletions packages/backend/src/core/calls/CallsLiveConnectionService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ export type CallsLiveConnection = {
sessionId: string | null;
createdAt: string;
lastSeenAt: string;
ready?: boolean;
};

export class StaleCallsConnectionError extends Error {}
Expand Down Expand Up @@ -46,6 +47,19 @@ export class CallsLiveConnectionService {
return value == null ? null : JSON.parse(value) as CallsLiveConnection;
}

public async isReady(participantId: string): Promise<boolean> {
return (await this.get(participantId))?.ready === true;
}

public async markReady(participantId: string, connectionId: string, generation: number): Promise<boolean> {
const connection = await this.assertCurrent(participantId, connectionId, generation);
if (connection.sessionId == null) throw new StaleCallsConnectionError();
if (connection.ready) return false;
connection.ready = true;
await this.redis.set(this.connectionKey(participantId), JSON.stringify(connection), 'EX', CallsLiveConnectionService.ttlSeconds);
return true;
}

public async touchHost(roomId: string, onlyIfMissing = false): Promise<void> {
const deadline = Date.now() + CallsLiveConnectionService.ttlSeconds * 1000;
if (onlyIfMissing) await this.redis.zadd('calls:host-deadlines', 'NX', deadline, roomId);
Expand Down
3 changes: 2 additions & 1 deletion packages/backend/src/core/calls/CallsMediaService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,8 @@ export class CallsMediaService {
}

public async reconcile(user: MiUser, roomId: string): Promise<{ roomRevision: number; publications: Array<{ id: string; participantId: string; mediaKind: 'audio' | 'video'; mediaSource: 'microphone' | 'camera' | 'screen'; screenPublicationId?: string }> }> {
const snapshot = await this.roomService.snapshot(user, roomId);
// Reserved participants still need publications to finish their media connection.
const snapshot = await this.roomService.snapshot(user, roomId, true);
const activeSpeakers = new Set(snapshot.participants.filter(p => p.role !== 'listener').map(p => p.id));
const publications = (await this.bindingService.listRoomPublications(roomId)).filter(binding => activeSpeakers.has(binding.participantId));
return { roomRevision: snapshot.room.revision, publications: publications.map(binding => ({ id: binding.id, participantId: binding.participantId, mediaKind: binding.mediaKind, mediaSource: binding.mediaSource ?? 'microphone', ...(binding.screenPublicationId == null ? {} : { screenPublicationId: binding.screenPublicationId }) })) };
Expand Down
33 changes: 26 additions & 7 deletions packages/backend/src/core/calls/CallsRoomService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -233,11 +233,16 @@ export class CallsRoomService implements OnModuleInit, OnApplicationShutdown {
}

@bindThis
public async snapshot(user: MiUser, roomId: string): Promise<{ room: MiCallsRoom; participants: MiCallsParticipant[] }> {
public async snapshot(user: MiUser, roomId: string, includeConnecting = false): Promise<{ room: MiCallsRoom; participants: MiCallsParticipant[] }> {
const room = await this.getRoom(roomId);
await this.assertCanAccess(user, room);
const participants = await this.callsParticipantsRepository.findBy({ roomId, state: 'active' });
return { room, participants };
return { room, participants: includeConnecting ? participants : await this.connectedParticipants(participants) };
}

private async connectedParticipants(participants: MiCallsParticipant[]): Promise<MiCallsParticipant[]> {
const ready = await Promise.all(participants.map(participant => this.callsLiveConnectionService.isReady(participant.id)));
return participants.filter((_, index) => ready[index]);
}

@bindThis
Expand All @@ -249,7 +254,7 @@ export class CallsRoomService implements OnModuleInit, OnApplicationShutdown {
if (following) {
followingUserIds = Object.keys(await this.cacheService.userFollowingsCache.fetch(user.id));
if (followingUserIds.length === 0) return [];
const participants = await this.callsParticipantsRepository.findBy({ userId: In(followingUserIds), state: 'active' });
const participants = await this.connectedParticipants(await this.callsParticipantsRepository.findBy({ userId: In(followingUserIds), state: 'active' }));
followingRoomIds = [...new Set(participants.map(participant => participant.roomId))];
}
const candidates = await this.callsRoomsRepository.find({
Expand All @@ -276,7 +281,7 @@ export class CallsRoomService implements OnModuleInit, OnApplicationShutdown {
@bindThis
public async listActiveRoomsForUsers(viewer: MiUser, userIds: MiUser['id'][]): Promise<Array<{ userId: MiUser['id']; roomId: MiCallsRoom['id'] }>> {
if (viewer.host !== null || userIds.length === 0) return [];
const participants = await this.callsParticipantsRepository.findBy({ userId: In(userIds), state: 'active' });
const participants = await this.connectedParticipants(await this.callsParticipantsRepository.findBy({ userId: In(userIds), state: 'active' }));
if (participants.length === 0) return [];
const rooms = await this.callsRoomsRepository.findBy({ id: In([...new Set(participants.map(participant => participant.roomId))]), state: 'open' });
const roomsById = new Map(rooms.map(room => [room.id, room]));
Expand Down Expand Up @@ -426,12 +431,26 @@ export class CallsRoomService implements OnModuleInit, OnApplicationShutdown {
updatedAt: now,
});
}
const revision = await this.bumpRevision(roomId);
await this.callsEventService.publish(roomId, revision, 'participant', { participantId: joined.id, action: 'joined' });
this.callsTelemetryService.lifecycle({ action: 'participant-joined', roomId, participantId: joined.id });
await this.bumpRevision(roomId);
return joined;
}

@bindThis
public async confirmReady(user: MiUser, roomId: string, identity: CallsConnectionIdentity): Promise<void> {
await this.callsLiveConnectionService.withRoomLock(roomId, async () => {
await this.assertCanJoin(user);
const room = await this.getRoom(roomId);
await this.assertCanAccess(user, room);
if (room.state !== 'open') throw new CallsRoomError('invalid-state');
const participant = await this.callsParticipantsRepository.findOneBy({ roomId, userId: user.id, state: 'active' });
if (participant == null) throw new CallsRoomError('participant-not-found');
if (!await this.callsLiveConnectionService.markReady(participant.id, identity.connectionId, identity.generation)) return;
const revision = await this.bumpRevision(roomId);
await this.callsEventService.publish(roomId, revision, 'participant', { participantId: participant.id, action: 'joined' });
this.callsTelemetryService.lifecycle({ action: 'participant-joined', roomId, participantId: participant.id });
});
}

@bindThis
public async leave(user: MiUser, roomId: string, identity?: CallsConnectionIdentity & { token?: string }): Promise<void> {
if (identity?.token != null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import { bindThis } from '@/decorators.js';
import type { GlobalEvents } from '@/core/GlobalEventService.js';
import { CallsFeatureDisabledError, CallsRoomError, CallsRoomService } from '@/core/calls/CallsRoomService.js';
import { CallsMediaService } from '@/core/calls/CallsMediaService.js';
import { CallsLiveConnectionService } from '@/core/calls/CallsLiveConnectionService.js';
import { CallsEntityService } from '@/core/entities/CallsEntityService.js';
import { DI } from '@/di-symbols.js';
import type { CallsParticipantsRepository } from '@/models/_.js';
Expand All @@ -30,6 +31,7 @@ export class CallsRoomChannel extends Channel {
@Inject(DI.callsParticipantsRepository)
private callsParticipantsRepository: CallsParticipantsRepository,
private callsEntityService: CallsEntityService,
private callsLiveConnectionService: CallsLiveConnectionService,
) { super(request); }

@bindThis
Expand Down Expand Up @@ -68,7 +70,7 @@ export class CallsRoomChannel extends Channel {
}
if (data.type === 'participant' && (data.body.action === 'joined' || data.body.action === 'updated')) {
const participant = await this.callsParticipantsRepository.findOneBy({ id: data.body.participantId, roomId: this.roomId });
if (participant != null) {
if (participant != null && await this.callsLiveConnectionService.isReady(participant.id)) {
const [packedParticipant] = await this.callsEntityService.packParticipants([participant], this.user);
this.send(data.type, data.body.action === 'updated'
? { ...data.body, participant: packedParticipant, moderatorUserIds: room.moderatorUserIds }
Expand All @@ -87,6 +89,9 @@ export class CallsRoomChannel extends Channel {
if (type === 'heartbeat' && typeof body === 'object' && body != null && !Array.isArray(body) && typeof body.connectionId === 'string' && typeof body.generation === 'number') {
void this.callsMediaService.heartbeat(this.user, this.roomId, body.connectionId, body.generation).catch(() => undefined);
}
if (type === 'ready' && typeof body === 'object' && body != null && !Array.isArray(body) && typeof body.connectionId === 'string' && typeof body.generation === 'number') {
void this.callsRoomService.confirmReady(this.user, this.roomId, { connectionId: body.connectionId, generation: body.generation }).catch(() => undefined);
}
}

@bindThis
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,20 @@ class FakeRedis {
}

describe('CallsLiveConnectionService', () => {
test('keeps a bound session hidden until ready and rejects readiness from an old generation', async () => {
const service = new CallsLiveConnectionService(new FakeRedis() as unknown as Redis.Redis);
await service.replace('participant', 'connection-1');
await service.bindSession('participant', 'connection-1', 1, 'session-1');
expect(await service.isReady('participant')).toBe(false);
expect(await service.markReady('participant', 'connection-1', 1)).toBe(true);
expect(await service.markReady('participant', 'connection-1', 1)).toBe(false);
await service.heartbeat('participant', 'connection-1', 1);
expect(await service.isReady('participant')).toBe(true);
await service.replace('participant', 'connection-2');
await expect(service.markReady('participant', 'connection-1', 1)).rejects.toBeInstanceOf(StaleCallsConnectionError);
expect(await service.isReady('participant')).toBe(false);
});

test('expires a host after 90 seconds and extends the deadline when they return', async () => {
vi.useFakeTimers();
try {
Expand Down
32 changes: 29 additions & 3 deletions packages/backend/test/unit/core/calls/CallsRoomChannel.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ test('leaves the room and revokes media when a stream event detects lost access'
const rooms = { getRoom: vi.fn().mockResolvedValue({ revision: 1 }), assertCanAccess: vi.fn().mockResolvedValue(undefined), leave: vi.fn().mockResolvedValue(undefined) };
const participants = { findOneBy: vi.fn() };
const entity = { packParticipants: vi.fn() };
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never);
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never, { isReady: async () => true } as never);
expect(await channel.init({ roomId: 'room-a' })).toBe(true);
rooms.assertCanAccess.mockRejectedValue(new CallsRoomError('access-denied'));
const handler = subscriber.on.mock.calls[0][1];
Expand All @@ -35,7 +35,7 @@ test.each(['joined', 'updated'] as const)('includes user data when a participant
const packedParticipant = { id: 'participant-a', user: { id: 'user-b', avatarUrl: 'avatar' } };
const participants = { findOneBy: vi.fn().mockResolvedValue(participant) };
const entity = { packParticipants: vi.fn().mockResolvedValue([packedParticipant]) };
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never);
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never, { isReady: async () => true } as never);
expect(await channel.init({ roomId: 'room-a' })).toBe(true);
const handler = subscriber.on.mock.calls[0][1];
const body = { sequence: 1, roomRevision: 1, occurredAt: new Date().toISOString(), participantId: 'participant-a', action };
Expand All @@ -53,7 +53,7 @@ test.each(['left', 'removed'] as const)('does not load user data when a particip
const rooms = { getRoom: vi.fn().mockResolvedValue({ revision: 1 }), assertCanAccess: vi.fn().mockResolvedValue(undefined) };
const participants = { findOneBy: vi.fn() };
const entity = { packParticipants: vi.fn() };
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never);
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never, { isReady: async () => true } as never);
expect(await channel.init({ roomId: 'room-a' })).toBe(true);
const handler = subscriber.on.mock.calls[0][1];
const body = { sequence: 1, roomRevision: 1, occurredAt: new Date().toISOString(), participantId: 'participant-a', action };
Expand All @@ -62,3 +62,29 @@ test.each(['left', 'removed'] as const)('does not load user data when a particip
expect(entity.packParticipants).not.toHaveBeenCalled();
expect(connection.sendMessageToWs).toHaveBeenCalledWith('channel', expect.objectContaining({ type: 'participant', body }));
});

test.each(['joined', 'updated'] as const)('does not expose a connecting participant on %s', async action => {
const subscriber = { on: vi.fn(), off: vi.fn() };
const connection = { user: { id: 'user-a' }, subscriber, sendMessageToWs: vi.fn() };
const rooms = { getRoom: vi.fn().mockResolvedValue({ revision: 1 }), assertCanAccess: vi.fn().mockResolvedValue(undefined) };
const participants = { findOneBy: vi.fn().mockResolvedValue({ id: 'participant-a' }) };
const entity = { packParticipants: vi.fn() };
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, participants as never, entity as never, { isReady: async () => false } as never);
await channel.init({ roomId: 'room-a' });
const body = { sequence: 1, roomRevision: 1, participantId: 'participant-a', action };
await subscriber.on.mock.calls[0][1]({ type: 'participant', body });
expect(entity.packParticipants).not.toHaveBeenCalled();
expect(connection.sendMessageToWs).toHaveBeenCalledWith('channel', expect.objectContaining({ type: 'participant', body }));
});

test.each([false, true])('ready requires Calls write permission (permitted: %s)', async permitted => {
const user = { id: 'user-a' };
const connection = { user, subscriber: { on: vi.fn() }, token: { permission: permitted ? ['read:calls', 'write:calls'] : ['read:calls'] } };
const rooms = { getRoom: vi.fn().mockResolvedValue({ revision: 1 }), assertCanAccess: vi.fn().mockResolvedValue(undefined), confirmReady: vi.fn().mockResolvedValue(undefined) };
const channel = new CallsRoomChannel({ id: 'channel-a', connection } as never, rooms as never, {} as never, {} as never, {} as never, {} as never);
await channel.init({ roomId: 'room-a' });
const identity = { connectionId: 'device-a', generation: 1 };
channel.onMessage('ready', identity);
if (permitted) expect(rooms.confirmReady).toHaveBeenCalledWith(user, 'room-a', identity);
else expect(rooms.confirmReady).not.toHaveBeenCalled();
});
Loading
Loading