From a4f8a42622c5b23d0a34bffe9112ad43696a19bb Mon Sep 17 00:00:00 2001 From: mattyatea Date: Sun, 4 Oct 2026 22:07:32 +0000 Subject: [PATCH 1/2] =?UTF-8?q?fix(calls):=20=E9=9F=B3=E5=A3=B0=E6=8E=A5?= =?UTF-8?q?=E7=B6=9A=E3=81=8C=E5=AE=8C=E4=BA=86=E3=81=97=E3=81=A6=E3=81=8B?= =?UTF-8?q?=E3=82=89=E5=8F=82=E5=8A=A0=E8=80=85=E3=82=92=E8=A1=A8=E7=A4=BA?= =?UTF-8?q?=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/calls-protocol.md | 9 ++- .../core/calls/CallsLiveConnectionService.ts | 14 ++++ .../src/core/calls/CallsMediaService.ts | 3 +- .../src/core/calls/CallsRoomService.ts | 33 ++++++++-- .../server/api/stream/channels/calls-room.ts | 7 +- .../core/calls/CallsLiveConnectionService.ts | 14 ++++ .../test/unit/core/calls/CallsRoomChannel.ts | 32 ++++++++- .../test/unit/core/calls/CallsRoomService.ts | 32 ++++++++- packages/calls-reference-client/src/index.ts | 6 ++ .../test/conformance.test.ts | 3 + .../src/composables/use-calls-room.ts | 2 + packages/frontend/src/utility/calls-media.ts | 4 +- .../frontend/src/utility/calls-session.ts | 42 +++++++++--- .../test/unit/calls-media-controller.test.ts | 35 +++++++++- .../frontend/test/unit/calls-session.test.ts | 65 +++++++++++++++++-- .../frontend/test/unit/use-calls-room.test.ts | 15 +++++ packages/misskey-js/src/streaming.types.ts | 1 + 17 files changed, 286 insertions(+), 31 deletions(-) diff --git a/docs/calls-protocol.md b/docs/calls-protocol.md index 13ba6ba8eb4..43ff58eae93 100644 --- a/docs/calls-protocol.md +++ b/docs/calls-protocol.md @@ -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. @@ -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 diff --git a/packages/backend/src/core/calls/CallsLiveConnectionService.ts b/packages/backend/src/core/calls/CallsLiveConnectionService.ts index 3349ba6f87f..8f79c979d75 100644 --- a/packages/backend/src/core/calls/CallsLiveConnectionService.ts +++ b/packages/backend/src/core/calls/CallsLiveConnectionService.ts @@ -16,6 +16,7 @@ export type CallsLiveConnection = { sessionId: string | null; createdAt: string; lastSeenAt: string; + ready?: boolean; }; export class StaleCallsConnectionError extends Error {} @@ -46,6 +47,19 @@ export class CallsLiveConnectionService { return value == null ? null : JSON.parse(value) as CallsLiveConnection; } + public async isReady(participantId: string): Promise { + return (await this.get(participantId))?.ready === true; + } + + public async markReady(participantId: string, connectionId: string, generation: number): Promise { + 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 { const deadline = Date.now() + CallsLiveConnectionService.ttlSeconds * 1000; if (onlyIfMissing) await this.redis.zadd('calls:host-deadlines', 'NX', deadline, roomId); diff --git a/packages/backend/src/core/calls/CallsMediaService.ts b/packages/backend/src/core/calls/CallsMediaService.ts index cdc8357af2f..dd90ef90c00 100644 --- a/packages/backend/src/core/calls/CallsMediaService.ts +++ b/packages/backend/src/core/calls/CallsMediaService.ts @@ -207,7 +207,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' }> }> { - 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' })) }; diff --git a/packages/backend/src/core/calls/CallsRoomService.ts b/packages/backend/src/core/calls/CallsRoomService.ts index 877afc9afe7..d9dde1ca0f5 100644 --- a/packages/backend/src/core/calls/CallsRoomService.ts +++ b/packages/backend/src/core/calls/CallsRoomService.ts @@ -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 { + const ready = await Promise.all(participants.map(participant => this.callsLiveConnectionService.isReady(participant.id))); + return participants.filter((_, index) => ready[index]); } @bindThis @@ -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({ @@ -276,7 +281,7 @@ export class CallsRoomService implements OnModuleInit, OnApplicationShutdown { @bindThis public async listActiveRoomsForUsers(viewer: MiUser, userIds: MiUser['id'][]): Promise> { 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])); @@ -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 { + 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 { if (identity?.token != null) { diff --git a/packages/backend/src/server/api/stream/channels/calls-room.ts b/packages/backend/src/server/api/stream/channels/calls-room.ts index a75020d9133..930213cba9d 100644 --- a/packages/backend/src/server/api/stream/channels/calls-room.ts +++ b/packages/backend/src/server/api/stream/channels/calls-room.ts @@ -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'; @@ -30,6 +31,7 @@ export class CallsRoomChannel extends Channel { @Inject(DI.callsParticipantsRepository) private callsParticipantsRepository: CallsParticipantsRepository, private callsEntityService: CallsEntityService, + private callsLiveConnectionService: CallsLiveConnectionService, ) { super(request); } @bindThis @@ -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 } @@ -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 diff --git a/packages/backend/test/unit/core/calls/CallsLiveConnectionService.ts b/packages/backend/test/unit/core/calls/CallsLiveConnectionService.ts index 9b6a271be90..4d3bb415b45 100644 --- a/packages/backend/test/unit/core/calls/CallsLiveConnectionService.ts +++ b/packages/backend/test/unit/core/calls/CallsLiveConnectionService.ts @@ -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 { diff --git a/packages/backend/test/unit/core/calls/CallsRoomChannel.ts b/packages/backend/test/unit/core/calls/CallsRoomChannel.ts index 1a0d69f550a..f55bef49570 100644 --- a/packages/backend/test/unit/core/calls/CallsRoomChannel.ts +++ b/packages/backend/test/unit/core/calls/CallsRoomChannel.ts @@ -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]; @@ -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 }; @@ -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 }; @@ -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(); +}); diff --git a/packages/backend/test/unit/core/calls/CallsRoomService.ts b/packages/backend/test/unit/core/calls/CallsRoomService.ts index c94276ea5f6..213d3f10ba0 100644 --- a/packages/backend/test/unit/core/calls/CallsRoomService.ts +++ b/packages/backend/test/unit/core/calls/CallsRoomService.ts @@ -50,6 +50,36 @@ function createConnectionFixture(role: 'host' | 'listener') { } describe('CallsRoomService lifecycle', () => { + test('hides connecting participants from snapshots and user presence while retaining them for negotiation', async () => { + const service = createAccessFixture(); + const room = { ...baseRoom, state: 'open', visibility: 'public' }; + const participants = [{ id: 'ready', roomId: room.id, userId: 'ready-user' }, { id: 'connecting', roomId: room.id, userId: viewer.id }]; + Object.assign(service, { + callsRoomsRepository: { findOneBy: async () => room, findBy: async () => [room] }, + callsParticipantsRepository: { findBy: async () => participants }, + callsLiveConnectionService: { isReady: async (id: string) => id === 'ready' }, + }); + expect((await service.snapshot(viewer, room.id)).participants).toEqual([participants[0]]); + expect((await service.snapshot(viewer, room.id, true)).participants).toEqual(participants); + expect(await service.listActiveRoomsForUsers(viewer, ['ready-user', viewer.id])).toEqual([{ userId: 'ready-user', roomId: room.id }]); + }); + + test('announces participation once after media readiness is confirmed', async () => { + const fixture = createConnectionFixture('host'); + const markReady = vi.fn().mockResolvedValueOnce(true).mockResolvedValue(false); + const publish = vi.fn(); + Object.assign(fixture.service, { + callsLiveConnectionService: { ...fixture.live, markReady }, + callsEventService: { publish }, + }); + const user = { id: 'owner-a', host: null } as MiUser; + const identity = { connectionId: 'new-device', generation: 2 }; + await fixture.service.confirmReady(user, fixture.room.id, identity); + await fixture.service.confirmReady(user, fixture.room.id, identity); + expect(markReady).toHaveBeenCalledWith('participant-a', 'new-device', 2); + expect(publish).toHaveBeenCalledExactlyOnceWith(fixture.room.id, 2, 'participant', { participantId: 'participant-a', action: 'joined' }); + }); + test('filters followed hosts and active participants before limiting rooms, preserving access checks', async () => { const service = createAccessFixture({ followings: { 'followed-host': {}, 'followed-listener': {} } }); const hosted = { ...baseRoom, id: 'hosted', ownerUserId: 'followed-host', state: 'open', visibility: 'public' }; @@ -57,7 +87,7 @@ describe('CallsRoomService lifecycle', () => { const privateRoom = { ...baseRoom, id: 'private', state: 'open' }; const find = vi.fn().mockResolvedValue([hosted, attended, privateRoom]); const findBy = vi.fn().mockResolvedValue([{ roomId: 'attended' }, { roomId: 'private' }]); - Object.assign(service, { callsRoomsRepository: { find }, callsParticipantsRepository: { findBy } }); + Object.assign(service, { callsRoomsRepository: { find }, callsParticipantsRepository: { findBy }, callsLiveConnectionService: { isReady: async () => true } }); await expect(service.listDiscoverable(viewer, 10, undefined, ['open'], true)).resolves.toEqual([hosted, attended]); expect(findBy).toHaveBeenCalledWith({ userId: expect.objectContaining({ _value: ['followed-host', 'followed-listener'] }), state: 'active' }); expect(find).toHaveBeenCalledWith(expect.objectContaining({ where: [ diff --git a/packages/calls-reference-client/src/index.ts b/packages/calls-reference-client/src/index.ts index 0037a2f31e6..7762878b4d9 100644 --- a/packages/calls-reference-client/src/index.ts +++ b/packages/calls-reference-client/src/index.ts @@ -62,6 +62,12 @@ export class CallsReferenceClient { await this.createMediaSession(roomId); } + // The browser integration calls this after transport connection and successful audio playback. + public confirmReady(): void { + if (this.peer == null || this.generation === 0) return; + this.channel?.send('ready', { connectionId: this.connectionId, generation: this.generation }); + } + private async createMediaSession(roomId: string): Promise { this.peer?.close(); this.subscriptions.clear(); diff --git a/packages/calls-reference-client/test/conformance.test.ts b/packages/calls-reference-client/test/conformance.test.ts index 9fb5db1a36d..1e22f6df3df 100644 --- a/packages/calls-reference-client/test/conformance.test.ts +++ b/packages/calls-reference-client/test/conformance.test.ts @@ -78,6 +78,9 @@ describe('third-party Calls protocol conformance', () => { test('negotiates capability, joins, creates media, heartbeats, reconciles a sequence gap, and honors revoke', async () => { const client = new CallsReferenceClient('https://misskey.example', 'token'); await client.join('room-a'); + expect(fixture.send).not.toHaveBeenCalledWith('ready', expect.anything()); + client.confirmReady(); + expect(fixture.send).toHaveBeenCalledWith('ready', { connectionId: expect.any(String), generation: 1 }); expect(fixture.request.mock.calls.map(call => call[0])).toEqual(expect.arrayContaining([ 'calls/capabilities', 'calls/rooms/join', 'calls/media/session/create', 'calls/media/reconcile', diff --git a/packages/frontend/src/composables/use-calls-room.ts b/packages/frontend/src/composables/use-calls-room.ts index 32130397df1..4281e4c228f 100644 --- a/packages/frontend/src/composables/use-calls-room.ts +++ b/packages/frontend/src/composables/use-calls-room.ts @@ -104,6 +104,8 @@ export function createCallsRoomConnection(roomId: string) { return { room, endReason, participants, connected, speakingParticipantIds, refresh, + identifyParticipant(participantId: string) { ownParticipantId = participantId; }, + ready(connectionId: string, generation: number) { channel.send('ready', { connectionId, generation }); }, setMuted(isMuted: boolean) { channel.send('mute', isMuted); }, setSpeaking(speaking: boolean) { channel.send('speaking', speaking); }, heartbeat(connectionId: string, generation: number) { channel.send('heartbeat', { connectionId, generation }); }, diff --git a/packages/frontend/src/utility/calls-media.ts b/packages/frontend/src/utility/calls-media.ts index 78a9b80ea0d..82f4c463fbc 100644 --- a/packages/frontend/src/utility/calls-media.ts +++ b/packages/frontend/src/utility/calls-media.ts @@ -71,6 +71,7 @@ export class CallsMediaController { previousConnection?: { connectionId: string; generation: number }, private replaceExisting = false, private videoCallbacks?: { + ready?: () => void; localTrack: (source: CallsVideoSource, track: MediaStreamTrack | null) => void; microphoneTrack?: (track: MediaStreamTrack | null) => void; noiseSuppressionChanged?: (enabled: boolean) => void; @@ -256,11 +257,12 @@ export class CallsMediaController { } const subscribed = await this.reconcileNow(); // SDP exchanges stay serialized, but receiving need not wait for the sender's transport. - if (this.localTrack != null) { + if (this.localTrack != null || subscribed) { await this.waitUntilConnected(peer); this.setState('connected'); } if (this.localTrack == null && !subscribed && this.peer === peer) this.setState('connected'); + if (this.peer === peer) this.videoCallbacks?.ready?.(); } public reconcile(): Promise { diff --git a/packages/frontend/src/utility/calls-session.ts b/packages/frontend/src/utility/calls-session.ts index 66a09b7a188..39d0dfcfdd1 100644 --- a/packages/frontend/src/utility/calls-session.ts +++ b/packages/frontend/src/utility/calls-session.ts @@ -4,6 +4,7 @@ */ import { computed, ref, shallowRef, watch } from 'vue'; +import type * as Misskey from 'misskey-js'; import type { CallsMediaFailure, CallsMediaState, CallsRemotePublication, CallsVideoQuality, CallsVideoSource } from '@/utility/calls-media.js'; import type { MenuItem } from '@/types/menu.js'; import { createCallsRoomConnection } from '@/composables/use-calls-room.js'; @@ -51,6 +52,8 @@ const connection = shallowRef(null); const media = shallowRef(null); const mediaState = ref('idle'); const mediaFailure = ref(null); +const mediaReady = ref(false); +const sessionParticipant = shallowRef(null); const muted = ref(false); const joining = ref(false); const replacedRoomId = ref(null); @@ -64,7 +67,7 @@ const reconnectCandidate = ref(null); const reconnectRoomState = ref<'checking' | 'open' | 'unavailable'>('checking'); const reconnectSecondsRemaining = ref(0); const speakerRequestResult = ref<'rejected' | null>(null); -const remoteAudio = new Map(); +const remoteAudio = new Map }>(); const participantVolumes = shallowRef(new Map()); const localVideos = shallowRef(new Map()); const remoteVideos = shallowRef(new Map()); @@ -132,7 +135,12 @@ const room = computed(() => connection.value?.room.value ?? null); const participants = computed(() => connection.value?.participants.value ?? []); const speakingParticipantIds = computed(() => connection.value?.speakingParticipantIds.value ?? new Set()); const connected = computed(() => connection.value?.connected.value ?? false); -const myParticipant = computed(() => participants.value.find(participant => participant.userId === $i?.id) ?? null); +const myParticipant = computed(() => participants.value.find(participant => participant.userId === $i?.id) ?? sessionParticipant.value); + +watch(participants, current => { + const own = current.find(participant => participant.userId === $i?.id); + if (own != null) sessionParticipant.value = own; +}); const isActive = computed(() => currentRoomId.value != null && myParticipant.value != null && room.value?.state === 'open'); watch(isActive, (active, _, onCleanup) => { @@ -224,10 +232,10 @@ function addRemoteTrack(track: MediaStreamTrack, publication: CallsRemotePublica audio.autoplay = true; audio.hidden = true; audio.srcObject = new MediaStream([track]); - remoteAudio.set(publication.id, { participantId: publication.participantId, element: audio }); + const playback = audio.play().catch(() => { needsAudioResume.value = true; }); + remoteAudio.set(publication.id, { participantId: publication.participantId, element: audio, playback }); applyParticipantVolumes(); window.document.body.append(audio); - void audio.play().catch(() => { needsAudioResume.value = true; }); track.addEventListener('ended', () => { removeRemoteTrack(publication.id); }, { once: true }); @@ -256,6 +264,7 @@ async function connectMedia(generation: number, previousConnection?: { connectio if (participant == null || targetRoomId == null) return; await media.value?.close().catch(() => undefined); if (generation !== sessionGeneration) return; + mediaReady.value = false; const controller = new CallsMediaController( targetRoomId, participant.role, @@ -263,12 +272,18 @@ async function connectMedia(generation: number, previousConnection?: { connectio if (generation !== sessionGeneration || media.value !== controller) return; mediaState.value = state; mediaFailure.value = failure; + if (state !== 'connected') mediaReady.value = false; }, (track, publication) => { if (generation === sessionGeneration && media.value === controller) addRemoteTrack(track, publication); }, (_stats, speaking) => connection.value?.setSpeaking(speaking), previousConnection, replaceExisting, { + ready() { + if (generation !== sessionGeneration || media.value !== controller) return; + mediaReady.value = true; + void announceMediaReady(controller, generation); + }, noiseSuppressionChanged(enabled) { if (generation === sessionGeneration && media.value === controller) noiseSuppression.value = enabled; }, @@ -311,6 +326,8 @@ async function clearSession(): Promise { sessionGeneration += 1; const controller = media.value; media.value = null; + mediaReady.value = false; + sessionParticipant.value = null; disposeConnection(); currentRoomId.value = null; mediaState.value = 'idle'; @@ -371,10 +388,10 @@ async function join(roomId: string, alreadyParticipant: boolean, reconnectToken? generation = ++sessionGeneration; currentRoomId.value = roomId; const next = attachConnection(roomId); - if (!alreadyParticipant) { - await misskeyApi('calls/rooms/join', { roomId, reconnectToken }); - joinedNow = true; - } + const participant = await misskeyApi('calls/rooms/join', { roomId, reconnectToken }); + joinedNow = !alreadyParticipant; + sessionParticipant.value = participant; + next.identifyParticipant(participant.id); await next.refresh(); if (myParticipant.value == null) throw new Error('Calls participant state was not created'); muted.value = myParticipant.value.isMuted; @@ -645,6 +662,15 @@ function openScreenSettings(event: MouseEvent): void { async function resumeAudio(): Promise { await Promise.all([...remoteAudio.values()].map(({ element }) => element.play())); needsAudioResume.value = false; + if (media.value != null) await announceMediaReady(media.value, sessionGeneration); +} + +async function announceMediaReady(controller: CallsMediaController, generation: number): Promise { + const identity = controller.connectionIdentity; + await Promise.all([...remoteAudio.values()].map(({ playback }) => playback)); + if (generation !== sessionGeneration || media.value !== controller || !mediaReady.value || needsAudioResume.value || identity == null) return; + if (controller.connectionIdentity?.generation !== identity.generation) return; + connection.value?.ready(identity.connectionId, identity.generation); } watch(videos, current => { diff --git a/packages/frontend/test/unit/calls-media-controller.test.ts b/packages/frontend/test/unit/calls-media-controller.test.ts index 649eade66db..ed319485675 100644 --- a/packages/frontend/test/unit/calls-media-controller.test.ts +++ b/packages/frontend/test/unit/calls-media-controller.test.ts @@ -117,6 +117,37 @@ describe('CallsMediaController', () => { } }); + test('a listener announces readiness after subscription negotiation and transport connection', async () => { + installBrowserMedia(vi.fn()); + const respond = apiMock.getMockImplementation()!; + let finishNegotiation!: (value: unknown) => void; + apiMock.mockImplementation(async (endpoint: string) => { + if (endpoint === 'calls/media/reconcile') { + FakePeerConnection.instances[0].connectionState = 'connecting'; + return { roomRevision: 1, publications: [{ id: 'remote-audio', participantId: 'participant-b', mediaKind: 'audio', mediaSource: 'microphone' }] }; + } + if (endpoint === 'calls/media/tracks/subscribe') return { subscriptions: [{ publicationId: 'remote-audio', mid: '1' }], requiresImmediateRenegotiation: false, sessionDescription: { type: 'offer', sdp: 'subscribe-offer' } }; + if (endpoint === 'calls/media/renegotiate') return new Promise(resolve => { finishNegotiation = resolve; }); + return respond(endpoint); + }); + const ready = vi.fn(); + const controller = new CallsMediaController('room-a', 'listener', undefined, undefined, undefined, undefined, false, { ready, localTrack: vi.fn(), remoteRemoved: vi.fn(), error: vi.fn() }); + const connecting = controller.connect(); + try { + await vi.waitFor(() => expect(finishNegotiation).toBeDefined()); + expect(ready).not.toHaveBeenCalled(); + finishNegotiation({ requiresImmediateRenegotiation: false, sessionDescription: null }); + await new Promise(resolve => setTimeout(resolve, 0)); + expect(ready).not.toHaveBeenCalled(); + FakePeerConnection.instances[0].connectionState = 'connected'; + FakePeerConnection.instances[0].dispatchEvent(new Event('connectionstatechange')); + await connecting; + expect(ready).toHaveBeenCalledOnce(); + } finally { + await controller.close(); + } + }); + test('a muted host connects without waiting for outbound audio statistics', async () => { vi.useFakeTimers(); installBrowserMedia(vi.fn().mockResolvedValue(stream(makeTrack('audio')))); @@ -601,7 +632,7 @@ describe('CallsMediaController', () => { expect(apiMock.mock.calls.filter(([endpoint]) => endpoint === 'calls/media/session/create')).toHaveLength(sessionCreates); }); - test.each(['failed', 'disconnected'])('recovers a listener when a replacement peer becomes %s before connecting', async state => { + test.each(['failed', 'disconnected'])('recovers a listener when a replacement peer becomes %s', async state => { installBrowserMedia(vi.fn()); const respond = apiMock.getMockImplementation()!; apiMock.mockImplementation(async (endpoint: string) => { @@ -617,7 +648,7 @@ describe('CallsMediaController', () => { peer.connectionState = 'failed'; peer.dispatchEvent(new Event('connectionstatechange')); await controller.reconcile(); - expect(controller.state).toBe('reconnecting'); + expect(controller.state).toBe('connected'); const replacement = FakePeerConnection.instances[1]!; replacement.connectionState = state; diff --git a/packages/frontend/test/unit/calls-session.test.ts b/packages/frontend/test/unit/calls-session.test.ts index 341ddfe43c6..dd39800777b 100644 --- a/packages/frontend/test/unit/calls-session.test.ts +++ b/packages/frontend/test/unit/calls-session.test.ts @@ -19,10 +19,14 @@ const fixture = vi.hoisted(() => ({ captureCamera: vi.fn(), keepalive: vi.fn(), connectionExists: false, + connectingParticipant: false, + connectMedia: vi.fn(), participantMuted: true, microphoneAvailable: true, setMuted: vi.fn(), heartbeat: vi.fn(), + ready: vi.fn(), + identifyParticipant: vi.fn(), role: 'listener' as 'listener' | 'host', revoked: [] as Array<(event: { reason: string; connectionId?: string; generation?: number }) => void>, connections: [] as Array<{ room: { value: { id: string; title: string; state: string; revision: number } }; endReason: { value: 'host-timeout' | null }; participants: { value: Array<{ id: string; userId: string; role: string; isMuted: boolean; joinedAt?: string }> }; refresh: ReturnType }>, @@ -47,10 +51,10 @@ vi.mock('@/composables/use-calls-room.js', async () => { const connection = { room: ref({ id: roomId, title: 'Room', state: 'open', revision: 1 }), endReason: ref<'host-timeout' | null>(null), - participants: ref([{ id: 'participant-a', userId: 'user-a', role: fixture.role, isMuted: fixture.participantMuted, joinedAt: new Date().toISOString() }]), + participants: ref(fixture.connectingParticipant ? [] : [{ id: 'participant-a', userId: 'user-a', role: fixture.role, isMuted: fixture.participantMuted, joinedAt: new Date().toISOString() }]), speakingParticipantIds: ref(new Set()), connected: ref(true), - refresh: vi.fn(), dispose: vi.fn(), setMuted: fixture.setMuted, setSpeaking: vi.fn(), heartbeat: fixture.heartbeat, + refresh: vi.fn(), dispose: vi.fn(), setMuted: fixture.setMuted, setSpeaking: vi.fn(), heartbeat: fixture.heartbeat, ready: fixture.ready, identifyParticipant: fixture.identifyParticipant, onTrackChange: () => vi.fn(), onRevoked: (callback: typeof fixture.revoked[number]) => { fixture.revoked.push(callback); return vi.fn(); }, }; @@ -74,10 +78,12 @@ vi.mock('@/utility/calls-media.js', () => ({ public reconcile = vi.fn().mockResolvedValue(undefined); public connect = vi.fn(async () => { if (fixture.connectionExists && !this.replaceExisting) throw Object.assign(new Error('Connection exists'), { code: 'CALLS_CONNECTION_EXISTS' }); + await fixture.connectMedia(); this.onState('connected'); + this.callbacks?.ready?.(); }); public setMuted = vi.fn(); - constructor(_roomId: string, _role: string, private onState: (state: string) => void, onRemoteTrack: typeof fixture.remoteTrackCallbacks[number], _onStats: unknown, previousConnection?: { connectionId: string; generation: number }, public replaceExisting = false) { + constructor(_roomId: string, _role: string, private onState: (state: string) => void, onRemoteTrack: typeof fixture.remoteTrackCallbacks[number], _onStats: unknown, previousConnection?: { connectionId: string; generation: number }, public replaceExisting = false, private callbacks?: { ready?: () => void }) { this.connectionIdentity = previousConnection ?? { connectionId: `device-${fixture.controllers.length}`, generation: 1 }; fixture.controllers.push(this); fixture.remoteTrackCallbacks.push(onRemoteTrack); @@ -91,7 +97,9 @@ describe('Calls session device handoff', () => { beforeEach(async () => { vi.resetModules(); vi.useFakeTimers(); - fixture.api.mockReset().mockResolvedValue({}); + fixture.api.mockReset().mockImplementation(async (endpoint: string) => endpoint === 'calls/rooms/join' ? { id: 'participant-a', roomId: 'room-a', userId: 'user-a', role: fixture.role, isMuted: fixture.participantMuted } : {}); + fixture.ready.mockClear(); + fixture.identifyParticipant.mockClear(); fixture.toast.mockClear(); fixture.playSound.mockClear(); fixture.alert.mockClear(); @@ -101,6 +109,8 @@ describe('Calls session device handoff', () => { fixture.keepalive.mockClear(); fixture.confirm.mockReset().mockResolvedValue({ canceled: true }); fixture.connectionExists = false; + fixture.connectingParticipant = false; + fixture.connectMedia.mockReset(); fixture.participantMuted = true; fixture.microphoneAvailable = true; fixture.setMuted.mockClear(); @@ -126,6 +136,50 @@ describe('Calls session device handoff', () => { } }); + test('keeps provisional membership out of the participant list until media is ready', async () => { + fixture.connectingParticipant = true; + let complete!: () => void; + fixture.connectMedia.mockImplementationOnce(() => new Promise(resolve => { complete = resolve; })); + const joining = session.join('room-a', false); + await vi.waitFor(() => expect(fixture.connectMedia).toHaveBeenCalled()); + expect(session.participants.value).toEqual([]); + expect(session.myParticipant.value?.id).toBe('participant-a'); + expect(fixture.identifyParticipant).toHaveBeenCalledWith('participant-a'); + expect(fixture.ready).not.toHaveBeenCalled(); + complete(); + await joining; + await vi.waitFor(() => expect(fixture.ready).toHaveBeenCalledWith('device-0', 1)); + expect(session.participants.value).toEqual([]); + }); + + test.each([false, true])('waits for remote playback before announcing readiness (autoplay blocked: %s)', async blocked => { + let startPlayback!: () => void; + const play = vi.spyOn(HTMLMediaElement.prototype, 'play').mockImplementationOnce(() => blocked + ? Promise.reject(new DOMException('Autoplay blocked', 'NotAllowedError')) + : new Promise(resolve => { startPlayback = resolve; })); + const pause = vi.spyOn(HTMLMediaElement.prototype, 'pause').mockImplementation(() => {}); + fixture.connectMedia.mockImplementationOnce(() => { + fixture.remoteTrackCallbacks[0](Object.assign(new EventTarget(), { kind: 'audio' }) as MediaStreamTrack, + { id: 'remote-audio', participantId: 'participant-b', mediaKind: 'audio', mediaSource: 'microphone' }); + }); + try { + await session.join('room-a', false); + expect(fixture.ready).not.toHaveBeenCalled(); + if (blocked) { + expect(session.needsAudioResume.value).toBe(true); + play.mockResolvedValue(); + await session.resumeAudio(); + } else { + startPlayback(); + } + await vi.waitFor(() => expect(fixture.ready).toHaveBeenCalledWith('device-0', 1)); + await session.leave(); + } finally { + play.mockRestore(); + pause.mockRestore(); + } + }); + test('plays join, participant changes and leave sounds without sounding the initial snapshot', async () => { fixture.api.mockImplementation(async (endpoint: string) => { if (endpoint === 'calls/rooms/join') fixture.connections[0].participants.value.push({ id: 'participant-c', userId: 'user-c', role: 'listener', isMuted: true }); @@ -589,6 +643,9 @@ describe('Calls session device handoff', () => { fixture.revoked[0]({ reason: 'stale-generation', ...identity }); await vi.waitFor(() => expect(fixture.controllers).toHaveLength(2)); expect(fixture.controllers[1].connectionIdentity).toEqual(identity); + fixture.connections[0].participants.value = []; + await nextTick(); + expect(session.myParticipant.value?.id).toBe('participant-a'); expect(fixture.playSound.mock.calls).toEqual([['callsJoin']]); expect(session.isActive.value).toBe(true); expect(session.replacedRoomId.value).toBeNull(); diff --git a/packages/frontend/test/unit/use-calls-room.test.ts b/packages/frontend/test/unit/use-calls-room.test.ts index b6acb4b9c6c..78b10f90186 100644 --- a/packages/frontend/test/unit/use-calls-room.test.ts +++ b/packages/frontend/test/unit/use-calls-room.test.ts @@ -47,6 +47,21 @@ describe('useCallsRoom streaming updates', () => { fixture.api.mockImplementation(async (endpoint: string) => endpoint === 'calls/rooms/show' ? structuredClone(snapshot) : { roomRevision: 1, publications: [] }); }); + test('recognizes a provisional participant for revocation without displaying it', async () => { + fixture.api.mockResolvedValueOnce({ ...snapshot, participants: [] }); + const calls = useCallsRoom('room-a'); + calls.identifyParticipant('participant-a'); + await calls.refresh(); + const revoked = vi.fn(); + calls.onRevoked(revoked); + fixture.channelHandlers.get('revoked')?.({ sequence: 1, roomRevision: 2, participantId: 'participant-a', reason: 'access' }); + expect(revoked).toHaveBeenCalledOnce(); + expect(calls.participants.value).toEqual([]); + calls.ready('device-a', 1); + expect(fixture.send).toHaveBeenCalledWith('ready', { connectionId: 'device-a', generation: 1 }); + calls.dispose(); + }); + test('shares room loading and updates across five cards until the last card is disposed', async () => { const cards = Array.from({ length: 5 }, () => retainCallsRoomConnection('shared-room')); await Promise.all(cards.map(card => card.load())); diff --git a/packages/misskey-js/src/streaming.types.ts b/packages/misskey-js/src/streaming.types.ts index 91d69b4824a..eeae2a2c023 100644 --- a/packages/misskey-js/src/streaming.types.ts +++ b/packages/misskey-js/src/streaming.types.ts @@ -326,6 +326,7 @@ export type Channels = { mute: boolean; speaking: boolean; heartbeat: { connectionId: string; generation: number }; + ready: { connectionId: string; generation: number }; }; }; callsRooms: { From 75cdbf5a801f2402e5723da816142e3781a7c10b Mon Sep 17 00:00:00 2001 From: mattyatea Date: Mon, 5 Oct 2026 08:35:35 +0000 Subject: [PATCH 2/2] =?UTF-8?q?fix(calls):=20CI=20=E3=81=AE=E5=9E=8B?= =?UTF-8?q?=E6=A4=9C=E6=9F=BB=E3=81=A8=20API=20=E3=83=AC=E3=83=9D=E3=83=BC?= =?UTF-8?q?=E3=83=88=E3=82=92=E4=BF=AE=E6=AD=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- packages/frontend/test/unit/calls-session.test.ts | 4 ++-- packages/misskey-js/etc/misskey-js.api.md | 4 ++++ 2 files changed, 6 insertions(+), 2 deletions(-) diff --git a/packages/frontend/test/unit/calls-session.test.ts b/packages/frontend/test/unit/calls-session.test.ts index dd39800777b..88dab07a9a9 100644 --- a/packages/frontend/test/unit/calls-session.test.ts +++ b/packages/frontend/test/unit/calls-session.test.ts @@ -534,8 +534,8 @@ describe('Calls session device handoff', () => { test.each(['host', 'listener'] as const)('warns before reloading while participating as %s without disconnecting', async role => { fixture.role = role; - const addListener = vi.spyOn(window, 'addEventListener'); - const removeListener = vi.spyOn(window, 'removeEventListener'); + const addListener = vi.spyOn(window as Window, 'addEventListener'); + const removeListener = vi.spyOn(window as Window, 'removeEventListener'); expect(addListener.mock.calls.some(([type]) => type === 'beforeunload')).toBe(false); await session.join('room-a', true); const listener = addListener.mock.calls.find(([type]) => type === 'beforeunload')![1] as EventListener; diff --git a/packages/misskey-js/etc/misskey-js.api.md b/packages/misskey-js/etc/misskey-js.api.md index b73fc2621e2..d92b22ea75e 100644 --- a/packages/misskey-js/etc/misskey-js.api.md +++ b/packages/misskey-js/etc/misskey-js.api.md @@ -1279,6 +1279,10 @@ export type Channels = { connectionId: string; generation: number; }; + ready: { + connectionId: string; + generation: number; + }; }; }; callsRooms: {