From cff8766aa618f8350eb98a71aa49df6c19719847 Mon Sep 17 00:00:00 2001 From: Robin Date: Tue, 8 Sep 2026 20:54:27 +0200 Subject: [PATCH] Delegate delayed leaves in LocalMember rather than LocalTransport IMO this is where the delegation calls should have lived all along, since the leave event is part of the membership lifecycle, and we otherwise end up with an awkward hack to ignore transport updates. Doing this now ensures that the client won't send any delegation requests if delegation is unsupported, and prepares the code for a future change in which we use the dedicated delegation endpoint from the CS API. --- src/state/CallViewModel/CallViewModel.ts | 24 +-- .../localMember/LocalMember.test.ts | 149 +++++++++++------- .../CallViewModel/localMember/LocalMember.ts | 41 ++++- .../localMember/LocalTransport.test.ts | 103 +----------- .../localMember/LocalTransport.ts | 94 ++++------- 5 files changed, 178 insertions(+), 233 deletions(-) diff --git a/src/state/CallViewModel/CallViewModel.ts b/src/state/CallViewModel/CallViewModel.ts index 0bd33f213..aa88e6115 100644 --- a/src/state/CallViewModel/CallViewModel.ts +++ b/src/state/CallViewModel/CallViewModel.ts @@ -493,16 +493,6 @@ export function createCallViewModel$( memberships$: memberships$, ownMembershipIdentity, client, - delayId$: scope.behavior( - ( - fromEvent( - matrixRTCSession, - MembershipManagerEvent.DelayIdChanged, - // The type of reemitted event includes the original emitted as the second arg. - ) as Observable<[string | undefined, IMembershipManager]> - ).pipe(map(([delayId]) => delayId ?? null)), - matrixRTCSession.delayId ?? null, - ), roomId: matrixRoom.roomId, matrixRTCMode, }); @@ -583,10 +573,22 @@ export function createCallViewModel$( ); }, connectionManager, + client, matrixRTCSession, - localTransport$, + localTransport, roomId: matrixRoom.roomId, baseUrl: client.baseUrl, + ownMembershipIdentity, + delayId$: scope.behavior( + ( + fromEvent( + matrixRTCSession, + MembershipManagerEvent.DelayIdChanged, + // The type of reemitted event includes the original emitted as the second arg. + ) as Observable<[string | undefined, IMembershipManager]> + ).pipe(map(([delayId]) => delayId ?? null)), + matrixRTCSession.delayId ?? null, + ), matrixRTCMode, logger: logger.getChild(`[${Date.now()}]`), }); diff --git a/src/state/CallViewModel/localMember/LocalMember.test.ts b/src/state/CallViewModel/localMember/LocalMember.test.ts index 6fe2155bd..e7c5b1e45 100644 --- a/src/state/CallViewModel/localMember/LocalMember.test.ts +++ b/src/state/CallViewModel/localMember/LocalMember.test.ts @@ -64,6 +64,7 @@ import { type LocalTransport, type LocalTransportWithSFUConfig, } from "./LocalTransport"; +import * as openIDSFU from "../../../livekit/openIDSFU"; initializeWidget(); @@ -113,6 +114,17 @@ const delegatedTimings: ResolvedDelayedLeaveTimings = { restart_timeout_ms: timings.restart_timeout_ms! * 10, }; +const mockedClient = { + getDomain: vi.fn().mockReturnValue("example.org"), + getDeviceId: vi.fn().mockReturnValue("AAAA"), + getOpenIdToken: vi.fn().mockResolvedValue({ + access_token: "ACCCESS_TOKEN", + token_type: "Bearer", + matrix_server_name: "localhost", + expires_in: 10000, + }), +}; + describe("enterRTCSession", () => { const transport: LivekitTransportConfig = { livekit_alias: "roomId", @@ -129,15 +141,7 @@ describe("enterRTCSession", () => { const mockedSession = vi.mocked({ room: { roomId: "roomId", - client: { - getDomain: vi.fn().mockReturnValue("example.org"), - getOpenIdToken: vi.fn().mockResolvedValue({ - access_token: "ACCCESS_TOKEN", - token_type: "Bearer", - matrix_server_name: "localhost", - expires_in: 10000, - }), - }, + client: mockedClient, }, memberships: [], joinRTCSession: vi.fn(), @@ -230,6 +234,9 @@ describe("LocalMembership", () => { }, roomId: "!test-room-id:example.org", baseUrl: "https://matrix.example.org", + ownMembershipIdentity: ownMemberMock, + client: mockedClient, + delayId$: constant(null), matrixRTCMode: MATRIX_RTC_MODE, }; @@ -276,7 +283,7 @@ describe("LocalMembership", () => { scope, ...defaultCreateLocalMemberValues, connectionManager: mockConnectionManager, - localTransport$: behavior("a", { a: aLocalTransport }), + localTransport: aLocalTransport, }); expectObservable(localMembership.localMemberState$).toBe("ne", { @@ -320,7 +327,7 @@ describe("LocalMembership", () => { scope, ...defaultCreateLocalMemberValues, connectionManager: mockConnectionManager, - localTransport$: constant(aLocalTransport), + localTransport: aLocalTransport, }); expectObservable(localMembership.localMemberState$).toBe("n-e", { @@ -337,8 +344,8 @@ describe("LocalMembership", () => { const scope = new ObservableScope(); const aLocalTransport: LocalTransport = { - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), }; const mockConnectionManager = { @@ -356,7 +363,7 @@ describe("LocalMembership", () => { leaveRoomSession: vi.fn(), }, connectionManager: mockConnectionManager, - localTransport$: new BehaviorSubject(aLocalTransport), + localTransport: aLocalTransport, }); const expextedLog = "'not connected yet' while updating the call intent (this is expected on startup)"; @@ -417,6 +424,11 @@ describe("LocalMembership", () => { livekitRoom: mockLivekitRoom({}), } as unknown as Connection; + const authCallSpy = vi + .spyOn(openIDSFU, "getSFUConfigWithOpenID") + .mockImplementation(() => mockedClient.getOpenIdToken()); + afterEach(() => authCallSpy.mockClear()); + it.each([ ["no", null, timings], [ @@ -429,14 +441,8 @@ describe("LocalMembership", () => { "joins session with %s delegation support", async (_serviceName, delegationUrl, delayedLeaveTimings) => { const scope = testScope(); - - const activeTransport$ = constant(aTransportWithSFUConfig); - const aLocalTransport: LocalTransport = { - advertised$: constant(aTransport), - active$: activeTransport$, - }; - const connectionManagerData = new ConnectionManagerData(); const joinMatrixRTC = vi.fn(); + const delayId$ = new BehaviorSubject(null); if (delegationUrl !== null) fetchMock.post(delegationUrl, () => ({ status: 401, body: {} })); @@ -445,10 +451,16 @@ describe("LocalMembership", () => { scope, ...defaultCreateLocalMemberValues, connectionManager: { - connectionManagerData$: constant(new Epoch(connectionManagerData)), + connectionManagerData$: constant( + new Epoch(new ConnectionManagerData()), + ), }, - localTransport$: constant(aLocalTransport), joinMatrixRTC, + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, + delayId$, }); localMembership.requestJoinAndPublish(); @@ -459,6 +471,34 @@ describe("LocalMembership", () => { aTransport, delayedLeaveTimings, ); + + expect(authCallSpy).not.toHaveBeenCalled(); + delayId$.next("leave1"); + await flushPromises(); + if (delegationUrl === null) { + expect(authCallSpy).not.toHaveBeenCalled(); + } else { + // Delegation is supported in this test case, so go on to check that + // LocalMember actually performs delegation + const expectDelegation = (delayId: string) => + expect(authCallSpy).toHaveBeenLastCalledWith( + mockedClient, + ownMemberMock, + "a", + "!test-room-id:example.org", + { + matrixRTCMode: MATRIX_RTC_MODE, + delayEndpointBaseUrl: "https://matrix.example.org", + delayId, + }, + expect.anything(), + ); + + expectDelegation("leave1"); + delayId$.next("leave2"); // Can change delegated leaves + await flushPromises(); + expectDelegation("leave2"); + } }, ); @@ -467,7 +507,7 @@ describe("LocalMembership", () => { const activeTransport$ = new BehaviorSubject(aTransportWithSFUConfig); const aLocalTransport: LocalTransport = { - advertised$: new BehaviorSubject(aTransport), + advertised$: constant(aTransport), active$: activeTransport$, }; @@ -505,7 +545,7 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject(aLocalTransport), + localTransport: aLocalTransport, }); await flushPromises(); activeTransport$.next({ @@ -536,7 +576,7 @@ describe("LocalMembership", () => { const publishers: Publisher[] = []; const tracks$ = new BehaviorSubject([]); - const publishing$ = new BehaviorSubject(false); + const publishing$ = constant(false); defaultCreateLocalMemberValues.createPublisherFactory.mockImplementation( () => { const p = { @@ -560,8 +600,8 @@ describe("LocalMembership", () => { >; const aLocalTransport: LocalTransport = { - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), }; const connectionManagerData = new ConnectionManagerData(); @@ -573,7 +613,7 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject(aLocalTransport), + localTransport: aLocalTransport, }); await flushPromises(); expect(publisherFactory).toHaveBeenCalledOnce(); @@ -601,7 +641,7 @@ describe("LocalMembership", () => { new BehaviorSubject(null); const aLocalTransport: LocalTransport = { - advertised$: new BehaviorSubject(aTransport), + advertised$: constant(aTransport), active$: activeTransport$, }; @@ -644,7 +684,7 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$, }, - localTransport$: new BehaviorSubject(aLocalTransport), + localTransport: aLocalTransport, }); await flushPromises(); @@ -779,10 +819,10 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject({ - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), - }), + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, }); await flushPromises(); @@ -819,10 +859,10 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject({ - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), - }), + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, }); await flushPromises(); @@ -861,18 +901,19 @@ describe("LocalMembership", () => { scope, ...defaultCreateLocalMemberValues, homeserverConnected: { - combined$: new BehaviorSubject< - [boolean, HomeserverDisconnectReason | null] - >([true, null]), + combined$: constant<[boolean, HomeserverDisconnectReason | null]>([ + true, + null, + ]), rtsSession$: constant(RTCMemberStatus.Connected), }, connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject({ - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), - }), + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, }); await flushPromises(); @@ -913,10 +954,10 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject({ - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), - }), + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, }); await flushPromises(); @@ -982,10 +1023,10 @@ describe("LocalMembership", () => { connectionManager: { connectionManagerData$: constant(new Epoch(connectionManagerData)), }, - localTransport$: new BehaviorSubject({ - advertised$: new BehaviorSubject(aTransport), - active$: new BehaviorSubject(aTransportWithSFUConfig), - }), + localTransport: { + advertised$: constant(aTransport), + active$: constant(aTransportWithSFUConfig), + }, }); return { scope, localMembership }; }; diff --git a/src/state/CallViewModel/localMember/LocalMember.ts b/src/state/CallViewModel/localMember/LocalMember.ts index 9b9218c54..7805805b8 100644 --- a/src/state/CallViewModel/localMember/LocalMember.ts +++ b/src/state/CallViewModel/localMember/LocalMember.ts @@ -15,6 +15,7 @@ import { MediaDeviceFailure, } from "livekit-client"; import { observeParticipantEvents } from "@livekit/components-core"; +import { type MatrixClient } from "matrix-js-sdk"; import { Status as RTCSessionStatus, type LivekitTransport, @@ -76,6 +77,7 @@ import { type HomeserverConnected } from "./HomeserverConnected.ts"; import { type LocalTransport } from "./LocalTransport.ts"; import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts"; import { or$ } from "../../../utils/observable.ts"; +import { getSFUConfigWithOpenID } from "../../../livekit/openIDSFU.ts"; export enum TransportState { /** Not even a transport is available to the LocalMembership */ @@ -143,12 +145,16 @@ interface Props { ) => void; homeserverConnected: HomeserverConnected; roomId: string; + ownMembershipIdentity: CallMembershipIdentityParts; localTransport: LocalTransport; + client: Pick; matrixRTCSession: Pick< MatrixRTCSession, "updateCallIntent" | "leaveRoomSession" >; baseUrl: string; + delayId$: Behavior; + matrixRTCMode: MatrixRTCMode; logger: Logger; } @@ -168,6 +174,7 @@ interface Props { * @param props.muteStates The mute states for video and audio. * @param props.matrixRTCSession The matrix RTC session to join. * @param props.baseUrl Base URL of the homeserver. + * @param props.delayId$ ID of the delayed leave event to delegate to the SFU. * @param props.roomId The room ID used as the call identifier in analytics events. * @returns * - publisher: The handle to create tracks and publish them to the room. @@ -179,15 +186,19 @@ interface Props { export const createLocalMembership$ = ({ scope, connectionManager, - localTransport$, + localTransport, homeserverConnected, createPublisherFactory, joinMatrixRTC, logger: parentLogger, muteStates, + client, matrixRTCSession, baseUrl, roomId, + ownMembershipIdentity, + delayId$, + matrixRTCMode, }: Props): { /** * This request to start audio and video tracks. @@ -711,6 +722,34 @@ export const createLocalMembership$ = ({ ), ); + // Delegate delayed leaves to the SFU + scope.reconcile( + scope.behavior(combineLatest([joinParams$, delayId$])), + async ([joinParams, delayId]) => { + if (joinParams?.delegationSupported && delayId !== null) { + try { + // This will technically cause the service to issue a new JWT token, + // but it's safe to discard. We're only interested in triggering + // delegation. + await getSFUConfigWithOpenID( + client, + ownMembershipIdentity, + joinParams.transport.livekit_service_url, + roomId, + { matrixRTCMode, delayEndpointBaseUrl: baseUrl, delayId }, + logger, + ); + } catch (e) { + // TODO: Surface this to the user as a service interruption? + logger.error( + `Failed to delegate leave to ${joinParams.transport.livekit_service_url}`, + e, + ); + } + } + }, + ); + // Pause upstream of all local media tracks when we're disconnected from // MatrixRTC, because it can be an unpleasant surprise for the app to say // 'reconnecting' and yet still be transmitting your media to others. diff --git a/src/state/CallViewModel/localMember/LocalTransport.test.ts b/src/state/CallViewModel/localMember/LocalTransport.test.ts index 7f9ecb98a..a77c7a16d 100644 --- a/src/state/CallViewModel/localMember/LocalTransport.test.ts +++ b/src/state/CallViewModel/localMember/LocalTransport.test.ts @@ -14,11 +14,8 @@ import { type MockedObject, vi, } from "vitest"; -import { - type CallMembership, - type LivekitTransportConfig, -} from "matrix-js-sdk/lib/matrixrtc"; -import { BehaviorSubject, filter, lastValueFrom } from "rxjs"; +import { type CallMembership } from "matrix-js-sdk/lib/matrixrtc"; +import { lastValueFrom } from "rxjs"; import fetchMock from "fetch-mock"; import { @@ -27,11 +24,7 @@ import { ownMemberMock, testScope, } from "../../../utils/test"; -import { - createLocalTransport$, - JwtEndpointVersion, - type LocalTransportWithSFUConfig, -} from "./LocalTransport"; +import { createLocalTransport$ } from "./LocalTransport"; import { constant } from "../../Behavior"; import { Epoch, ObservableScope } from "../../ObservableScope"; import { @@ -62,14 +55,12 @@ describe("LocalTransport", () => { // eslint-disable-next-line @typescript-eslint/naming-convention _unstable_getRTCTransports: async () => Promise.resolve([]), getDomain: () => "example.org", - baseUrl: "example.org", // These won't be called in this error path but satisfy the type getOpenIdToken: vi.fn(), getDeviceId: vi.fn(), }, ownMembershipIdentity: ownMemberMock, matrixRTCMode: MatrixRTCMode.Compatibility, - delayId$: constant("delay_id_mock"), }); await flushPromises(); @@ -101,7 +92,6 @@ describe("LocalTransport", () => { roomId: "!example_room_id", memberships$: constant(new Epoch([])), client: { - baseUrl: "https://example.org", getDomain: () => "example.org", // eslint-disable-next-line @typescript-eslint/naming-convention _unstable_getRTCTransports: async () => Promise.resolve([]), @@ -110,7 +100,6 @@ describe("LocalTransport", () => { }, ownMembershipIdentity: ownMemberMock, matrixRTCMode: MatrixRTCMode.Compatibility, - delayId$: constant("delay_id_mock"), }); active$.subscribe( (o) => observations.push(o), @@ -148,11 +137,9 @@ describe("LocalTransport", () => { getDomain: () => "example.org", getOpenIdToken: vi.fn(), getDeviceId: vi.fn(), - baseUrl: "https://example.org", }, ownMembershipIdentity: ownMemberMock, matrixRTCMode: MatrixRTCMode.Compatibility, - delayId$: constant("delay_id_mock"), }); openIdResolver.resolve?.({ @@ -196,10 +183,8 @@ describe("LocalTransport", () => { scope: testScope(), roomId: "!example_room_id", matrixRTCMode: MatrixRTCMode.Compatibility, - delayId$: constant(null), memberships$: constant(new Epoch([])), client: { - baseUrl: "https://example.org", getDomain: vi.fn().mockReturnValue("example.org"), // eslint-disable-next-line @typescript-eslint/naming-convention _unstable_getRTCTransports: vi.fn().mockResolvedValue([]), @@ -308,11 +293,9 @@ describe("LocalTransport", () => { ownMembershipIdentity: ownMemberMock, roomId: "!example_room_id", matrixRTCMode: MatrixRTCMode.Compatibility, - delayId$: constant(null), memberships$: constant(new Epoch([])), client: { getDomain: () => "example.org", - baseUrl: "https://example.org", // eslint-disable-next-line @typescript-eslint/naming-convention _unstable_getRTCTransports: async () => Promise.resolve([]), // These won't be called in this error path but satisfy the type @@ -330,84 +313,4 @@ describe("LocalTransport", () => { ); }); }); - - it("should not update advertised/active transport on delayID changes, but delay Id delegation should be called", async () => { - // For simplicity, we'll just use the config livekit - customLivekitUrl.setValue("https://lk.example.org"); - - const authCallSpy = vi - .spyOn(openIDSFU, "getSFUConfigWithOpenID") - .mockResolvedValue(openIdResponse); - - const delayId$ = new BehaviorSubject(null); - - const { advertised$, active$ } = createLocalTransport$({ - scope: testScope(), - ownMembershipIdentity: ownMemberMock, - roomId: "!example_room_id", - // We want multi-sdu - forceJwtEndpoint: JwtEndpointVersion.Legacy, - delayId$: delayId$, - memberships$: constant(new Epoch([])), - client: { - getDomain: () => "example.org", - baseUrl: "https://example.org", - // eslint-disable-next-line @typescript-eslint/naming-convention - _unstable_getRTCTransports: async () => Promise.resolve([]), - // These won't be called in this error path but satisfy the type - getOpenIdToken: vi.fn(), - getDeviceId: vi.fn(), - }, - }); - - const advertisedValues: LivekitTransportConfig[] = []; - const activeValues: LocalTransportWithSFUConfig[] = []; - advertised$ - .pipe(filter((v) => v !== null)) - .subscribe((t) => advertisedValues.push(t)); - active$ - .pipe(filter((v) => v !== null)) - .subscribe((t) => activeValues.push(t)); - - await flushPromises(); - - // we have now an active and an advertised - expect(advertisedValues.length).toEqual(1); - expect(activeValues.length).toEqual(1); - expect(advertisedValues[0]!.livekit_service_url).toEqual( - "https://lk.example.org", - ); - expect(activeValues[0]!.transport.livekit_service_url).toEqual( - "https://lk.example.org", - ); - - expect(authCallSpy).toHaveBeenCalledTimes(2); - // Now emits 3 new delays id - delayId$.next("delay_id_1"); - await flushPromises(); - delayId$.next("delay_id_2"); - await flushPromises(); - delayId$.next("delay_id_3"); - await flushPromises(); - - // No new emissions should've happened, it is the same transport. - expect(advertisedValues.length).toEqual(1); - expect(activeValues.length).toEqual(1); - - // Still we should have updated the delayID to auth - expect(authCallSpy).toHaveBeenCalledTimes( - 4 * 2 /* 2 calls for each delayId ?? why */, - ); - - expect(authCallSpy).toHaveBeenLastCalledWith( - expect.anything(), - expect.anything(), - expect.anything(), - expect.anything(), - expect.objectContaining({ - delayId: "delay_id_3", - }), - expect.anything(), - ); - }); }); diff --git a/src/state/CallViewModel/localMember/LocalTransport.ts b/src/state/CallViewModel/localMember/LocalTransport.ts index 66be511a4..c5aa39590 100644 --- a/src/state/CallViewModel/localMember/LocalTransport.ts +++ b/src/state/CallViewModel/localMember/LocalTransport.ts @@ -10,14 +10,7 @@ import { type LivekitTransportConfig, } from "matrix-js-sdk/lib/matrixrtc"; import { type MatrixClient } from "matrix-js-sdk"; -import { - combineLatest, - distinctUntilChanged, - from, - map, - of, - switchMap, -} from "rxjs"; +import { distinctUntilChanged, from, map, of, switchMap } from "rxjs"; import { logger as rootLogger, type Logger } from "matrix-js-sdk/lib/logger"; import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager"; @@ -47,15 +40,11 @@ interface Props { scope: ObservableScope; ownMembershipIdentity: CallMembershipIdentityParts; memberships$: Behavior>; - client: Pick< - MatrixClient, - "getDomain" | "baseUrl" | "_unstable_getRTCTransports" - > & + client: Pick & OpenIDClientParts; // Used by the jwt service to create the livekit room and compute the livekit alias. roomId: string; matrixRTCMode: MatrixRTCMode; - delayId$: Behavior; } // TODO livekit_alias-cleanup @@ -119,7 +108,6 @@ export const createLocalTransport$ = ({ client, roomId, matrixRTCMode, - delayId$, }: Props): LocalTransport => { const logger = rootLogger.getChild("[LocalTransport]"); @@ -134,32 +122,29 @@ export const createLocalTransport$ = ({ transportDiscovery.discoverPreferredTransport(), ); - const preferredConfig$ = customLivekitUrl.value$ - .pipe( - switchMap((customUrl) => { - if (customUrl) { - return of({ - type: "livekit", - livekit_service_url: customUrl, - } as LivekitTransportConfig); - } else { - return discoveredTransport$; - } - }), - ) - .pipe( - map((config) => { - if (!config) { - // Bubbled up from the preferredConfig$ observable. - throw new MatrixRTCTransportMissingError(client.getDomain() ?? ""); - } - return config; - }), - distinctUntilChanged(areLivekitTransportsEqual), - ); + const preferredConfig$ = customLivekitUrl.value$.pipe( + switchMap((customUrl) => { + if (customUrl) { + return of({ + type: "livekit", + livekit_service_url: customUrl, + } as LivekitTransportConfig); + } else { + return discoveredTransport$; + } + }), + map((config) => { + if (!config) { + // Bubbled up from the preferredConfig$ observable. + throw new MatrixRTCTransportMissingError(client.getDomain() ?? ""); + } + return config; + }), + distinctUntilChanged(areLivekitTransportsEqual), + ); - const preferredTransport$ = combineLatest([preferredConfig$, delayId$]).pipe( - switchMap(async ([transport, delayId]) => { + const preferredTransport$ = preferredConfig$.pipe( + switchMap(async (transport) => { try { return await doOpenIdAndJWTFromUrl( transport, @@ -167,7 +152,6 @@ export const createLocalTransport$ = ({ ownMembershipIdentity, roomId, client, - delayId ?? undefined, logger, ); } catch (e) { @@ -189,21 +173,7 @@ export const createLocalTransport$ = ({ ), null, ), - active$: scope.behavior( - preferredTransport$.pipe( - // XXX: WORK AROUND due to a reconnection glitch. - // To remove when we have a proper way to refresh the delegation event ID without refreshing - // the whole credentials. - // We deliberately hide any changes to the SFU config because we - // do not want the app to reconnect whenever the JWT - // token changes due to us delegating a new delayed event. The - // initial SFU config for the transport is all the app needs. - distinctUntilChanged((prev, next) => - areLivekitTransportsEqual(prev.transport, next.transport), - ), - ), - null, - ), + active$: scope.behavior(preferredTransport$, null), }; }; @@ -219,7 +189,6 @@ export const createLocalTransport$ = ({ * @param membership The identity of the local member. * @param roomId The room ID to use for the JWT. * @param client The client to use for the OpenID token. - * @param delayId The delayId to use for the JWT. * * @throws FailToGetOpenIdToken, NoMatrix2AuthorizationService */ @@ -228,12 +197,7 @@ async function doOpenIdAndJWTFromUrl( matrixRTCMode: MatrixRTCMode, membership: CallMembershipIdentityParts, roomId: string, - client: Pick< - MatrixClient, - "getDomain" | "baseUrl" | "_unstable_getRTCTransports" - > & - OpenIDClientParts, - delayId?: string, + client: Pick & OpenIDClientParts, logger?: Logger, ): Promise { const sfuConfig = await getSFUConfigWithOpenID( @@ -241,11 +205,7 @@ async function doOpenIdAndJWTFromUrl( membership, transport.livekit_service_url, roomId, - { - matrixRTCMode, - delayEndpointBaseUrl: client.baseUrl, - delayId, - }, + { matrixRTCMode }, logger, ); return {