diff --git a/config/config.devenv.json b/config/config.devenv.json index 48602406b..294209a67 100644 --- a/config/config.devenv.json +++ b/config/config.devenv.json @@ -9,8 +9,6 @@ "matrix_rtc_session": { "wait_for_key_rotation_ms": 3000, "membership_event_expiry_ms": 180000000, - "delayed_leave_event_delay_ms": 18000, - "delayed_leave_event_restart_ms": 4000, "network_error_retry_ms": 100 } } diff --git a/config/config.sample.json b/config/config.sample.json index 0baa522e4..9f7ab61ac 100644 --- a/config/config.sample.json +++ b/config/config.sample.json @@ -13,8 +13,6 @@ "matrix_rtc_session": { "wait_for_key_rotation_ms": 3000, "membership_event_expiry_ms": 180000000, - "delayed_leave_event_delay_ms": 18000, - "delayed_leave_event_restart_ms": 4000, "network_error_retry_ms": 100 } } diff --git a/config/config_netlify_preview.json b/config/config_netlify_preview.json index 313f0d02c..179177bac 100644 --- a/config/config_netlify_preview.json +++ b/config/config_netlify_preview.json @@ -9,8 +9,6 @@ "matrix_rtc_session": { "wait_for_key_rotation_ms": 3000, "membership_event_expiry_ms": 180000000, - "delayed_leave_event_delay_ms": 18000, - "delayed_leave_event_restart_ms": 4000, "network_error_retry_ms": 100 }, "posthog": { diff --git a/config/config_netlify_preview_sdk.json b/config/config_netlify_preview_sdk.json index 784f0c7ef..fa60fe8d4 100644 --- a/config/config_netlify_preview_sdk.json +++ b/config/config_netlify_preview_sdk.json @@ -9,8 +9,6 @@ "matrix_rtc_session": { "wait_for_key_rotation_ms": 3000, "membership_event_expiry_ms": 180000000, - "delayed_leave_event_delay_ms": 18000, - "delayed_leave_event_restart_ms": 4000, "network_error_retry_ms": 100 } } diff --git a/src/config/ConfigOptions.ts b/src/config/ConfigOptions.ts index 536966e00..872e35168 100644 --- a/src/config/ConfigOptions.ts +++ b/src/config/ConfigOptions.ts @@ -24,6 +24,35 @@ export enum MatrixRTCMode { Matrix_2_0 = "matrix_2_0", } +export interface DelayedLeaveTimings { + /** + * The delay (in milliseconds) with which delayed leave events are sent. + * + * If the server receives no keep-alives from the client for any longer than + * this duration, it will send the leave event, automatically removing the + * user from the call. + */ + delay_ms?: number; + + /** + * How frequently (in milliseconds) the client sends keep-alives to the server + * to restart the timer for a delayed leave event. Should be less than + * {@link DelayedLeaveTimings.delay_ms}. + */ + restart_ms?: number; + + /** + * The time (in milliseconds) after which we consider a delayed event restart HTTP request to have failed. + * Setting this to a lower value will result in more frequent retries, but then we will also give up earlier. + * + * In the presence of network packet loss (hurting TCP connections), the custom delayedEventRestartLocalTimeoutMs + * helps by keeping more delayed event reset candidates in flight, + * improving the chances of a successful reset. (its is equivalent to the js-sdk `localTimeout` configuration, + * but only applies to calls to the `_unstable_updateDelayedEvent` endpoint with a body of `{action:"restart"}`.) + */ + restart_timeout_ms?: number; +} + export interface ConfigOptions { /** * The Posthog endpoint to which analytics data will be sent. @@ -184,29 +213,6 @@ export interface ConfigOptions { */ wait_for_key_rotation_ms?: number; - /** - * The duration (in milliseconds) after the most recent keep-alive (delayed leave event restart) - * that the server waits before sending the leave MatrixRTC membership event. - */ - delayed_leave_event_delay_ms?: number; - - /** - * The time (in milliseconds) after which we consider a delayed event restart http request to have failed. - * Setting this to a lower value will result in more frequent retries but also a higher chance of failiour. - * - * In the presence of network packet loss (hurting TCP connections), the custom delayedEventRestartLocalTimeoutMs - * helps by keeping more delayed event reset candidates in flight, - * improving the chances of a successful reset. (its is equivalent to the js-sdk `localTimeout` configuration, - * but only applies to calls to the `_unstable_updateDelayedEvent` endpoint with a body of `{action:"restart"}`.) - */ - delayed_leave_event_restart_local_timeout_ms?: number; - - /** - * The time interval (in milliseconds) at which the client sends membership keep-alive - * messages to the server by restarting the timer for the delayed leave event. - */ - delayed_leave_event_restart_ms?: number; - /** * How long we wait before retrying after a network error on any of the requests. */ @@ -231,9 +237,28 @@ export interface ConfigOptions { * Defaults to the js-sdk default (undefined). Which means that rotation will always happen. */ key_rotation_participant_limit?: number; + + /** + * Timing options for delayed leave events, which are used to remove a user + * from a call when they lose connection. + */ + delayed_leave?: DelayedLeaveTimings; + + /** + * Timing options for delayed leave events, in cases where the ability to + * send the event can be delegated to the SFU. + * + * We recommend setting {@link DelayedLeaveTimings.delay_ms} >> + * {@link sync_disconnect_grace_period_ms} here. + */ + delegated_delayed_leave?: DelayedLeaveTimings; }; } +export interface ResolvedDelayedLeaveTimings extends DelayedLeaveTimings { + delay_ms: number; // Required +} + // Overrides members from ConfigOptions that are always provided by the // default config and are therefore non-optional. export interface ResolvedConfigOptions extends ConfigOptions { @@ -257,19 +282,15 @@ export interface ResolvedConfigOptions extends ConfigOptions { > >; }; - matrix_rtc_session: { - wait_for_key_rotation_ms?: number; - delayed_leave_event_delay_ms: number; - delayed_leave_event_restart_local_timeout_ms?: number; - delayed_leave_event_restart_ms?: number; + matrix_rtc_session: ConfigOptions["matrix_rtc_session"] & { network_error_retry_ms: number; - membership_event_expiry_ms?: number; - key_rotation_participant_limit?: number; + delayed_leave: ResolvedDelayedLeaveTimings; + delegated_delayed_leave: ResolvedDelayedLeaveTimings; }; } export const DEFAULT_CONFIG: ResolvedConfigOptions = { - sync_disconnect_grace_period_ms: 10000, + sync_disconnect_grace_period_ms: 10_000, ssla: "https://static.element.io/legal/element-software-and-services-license-agreement-uk-1.pdf", media_quality: { video_codec: "vp8", @@ -285,7 +306,8 @@ export const DEFAULT_CONFIG: ResolvedConfigOptions = { }, }, matrix_rtc_session: { - delayed_leave_event_delay_ms: 10000, - network_error_retry_ms: 1000, + network_error_retry_ms: 1_000, + delayed_leave: { delay_ms: 18_000, restart_ms: 4_000 }, + delegated_delayed_leave: { delay_ms: 3_600_000, restart_ms: 300_000 }, }, }; diff --git a/src/livekit/openIDSFU.test.ts b/src/livekit/openIDSFU.test.ts index 2ddb6c95c..f759a3aac 100644 --- a/src/livekit/openIDSFU.test.ts +++ b/src/livekit/openIDSFU.test.ts @@ -148,7 +148,7 @@ describe("getSFUConfigWithOpenID", () => { // Verify, that the request contains the expected delay parameters if ( body.delay_id === "mock_delay_id" && - body.delay_timeout === 10000 && + body.delay_timeout === 3600000 && body.delay_cs_api_url === "https://homeserverserver.org/cs_api" ) { return { @@ -229,7 +229,7 @@ describe("getSFUConfigWithOpenID", () => { expect(calls[0][0]).toStrictEqual("https://sfu.example.org/get_token"); expect(calls[0][1]).toStrictEqual({ // check if it uses correct delayID! - body: '{"room_id":"!example_room_id","slot_id":"m.call#ROOM","member":{"id":"@alice:example.org:DEVICE","claimed_user_id":"@alice:example.org","claimed_device_id":"DEVICE"},"delay_id":"mock_delay_id","delay_timeout":10000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', + body: '{"room_id":"!example_room_id","slot_id":"m.call#ROOM","member":{"id":"@alice:example.org:DEVICE","claimed_user_id":"@alice:example.org","claimed_device_id":"DEVICE"},"delay_id":"mock_delay_id","delay_timeout":3600000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', method: "POST", headers: { "Content-Type": "application/json", @@ -239,7 +239,7 @@ describe("getSFUConfigWithOpenID", () => { expect(calls[1][0]).toStrictEqual("https://sfu.example.org/sfu/get"); expect(calls[1][1]).toStrictEqual({ - body: '{"room":"!example_room_id","device_id":"DEVICE","delay_id":"mock_delay_id","delay_timeout":10000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', + body: '{"room":"!example_room_id","device_id":"DEVICE","delay_id":"mock_delay_id","delay_timeout":3600000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', headers: { "Content-Type": "application/json", }, @@ -284,7 +284,7 @@ describe("getSFUConfigWithOpenID", () => { expect(calls[0][0]).toStrictEqual("https://sfu.example.org/get_token"); expect(calls[0][1]).toStrictEqual({ // check if it uses correct delayID! - body: '{"room_id":"!example_room_id","slot_id":"m.call#ROOM","member":{"id":"@alice:example.org:DEVICE","claimed_user_id":"@alice:example.org","claimed_device_id":"DEVICE"},"delay_id":"mock_delay_id","delay_timeout":10000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', + body: '{"room_id":"!example_room_id","slot_id":"m.call#ROOM","member":{"id":"@alice:example.org:DEVICE","claimed_user_id":"@alice:example.org","claimed_device_id":"DEVICE"},"delay_id":"mock_delay_id","delay_timeout":3600000,"delay_cs_api_url":"https://matrix.homeserverserver.org"}', method: "POST", headers: { "Content-Type": "application/json", diff --git a/src/livekit/openIDSFU.ts b/src/livekit/openIDSFU.ts index 00cf69b1a..67dd5e9e7 100644 --- a/src/livekit/openIDSFU.ts +++ b/src/livekit/openIDSFU.ts @@ -204,11 +204,10 @@ async function getLiveKitJWT( let bodyDalayParts: IDelayParams = {}; // Also check for empty string if (delayId && delayEndpointBaseUrl) { - const delayTimeoutMs = - Config.get().matrix_rtc_session?.delayed_leave_event_delay_ms; bodyDalayParts = { delay_id: delayId, - delay_timeout: delayTimeoutMs, + delay_timeout: + Config.get().matrix_rtc_session.delegated_delayed_leave.delay_ms, delay_cs_api_url: delayEndpointBaseUrl, }; } @@ -288,11 +287,10 @@ export async function getLiveKitJWTWithDelayDelegation( let bodyDalayParts = {}; // Also check for empty string if (delayId && delayEndpointBaseUrl) { - const delayTimeoutMs = - Config.get().matrix_rtc_session?.delayed_leave_event_delay_ms; bodyDalayParts = { delay_id: delayId, - delay_timeout: delayTimeoutMs, + delay_timeout: + Config.get().matrix_rtc_session.delegated_delayed_leave.delay_ms, delay_cs_api_url: delayEndpointBaseUrl, }; } diff --git a/src/state/CallViewModel/CallViewModel.ts b/src/state/CallViewModel/CallViewModel.ts index c854363be..62f05a476 100644 --- a/src/state/CallViewModel/CallViewModel.ts +++ b/src/state/CallViewModel/CallViewModel.ts @@ -64,7 +64,10 @@ import { showReactions, } from "../../settings/settings"; import { Config } from "../../config/Config"; -import { MatrixRTCMode } from "../../config/ConfigOptions"; +import { + MatrixRTCMode, + type ResolvedDelayedLeaveTimings, +} from "../../config/ConfigOptions"; import { isFirefox, platform } from "../../Platform"; import { setPipEnabled$ } from "../../controls"; import { TileStore } from "../TileStore"; @@ -562,16 +565,6 @@ export function createCallViewModel$( localUser: { userId, deviceId }, }); - const connectOptions$ = scope.behavior( - matrixRTCMode$.pipe( - map((mode) => ({ - encryptMedia: livekitKeyProvider !== undefined, - // TODO. This might need to get called again on each change of matrixRTCMode... - matrixRTCMode: mode, - })), - ), - ); - const localMembership = createLocalMembership$({ scope, homeserverConnected: createHomeserverConnected$( @@ -580,12 +573,21 @@ export function createCallViewModel$( matrixRTCSession, ), muteStates, - joinMatrixRTC: (transport: LivekitTransportConfig) => { + joinMatrixRTC: ( + transport: LivekitTransportConfig, + delayedLeaveTimings: ResolvedDelayedLeaveTimings, + ) => { return enterRTCSession( matrixRTCSession, ownMembershipIdentity, transport, - connectOptions$.value, + { + encryptMedia: livekitKeyProvider !== undefined, + // We merely sample the current mode here, so the user would need to + // manually rejoin to switch to a different one + matrixRTCMode: matrixRTCMode$.value, + delayedLeaveTimings, + }, ); }, createPublisherFactory: (connection: Connection) => { @@ -603,6 +605,7 @@ export function createCallViewModel$( matrixRTCSession, localTransport$, roomId: matrixRoom.roomId, + baseUrl: client.baseUrl, logger: logger.getChild(`[${Date.now()}]`), }); diff --git a/src/state/CallViewModel/localMember/LocalMember.test.ts b/src/state/CallViewModel/localMember/LocalMember.test.ts index aec641776..66cd8cda0 100644 --- a/src/state/CallViewModel/localMember/LocalMember.test.ts +++ b/src/state/CallViewModel/localMember/LocalMember.test.ts @@ -19,13 +19,18 @@ import { beforeAll, afterAll, beforeEach, + afterEach, } from "vitest"; import { BehaviorSubject, map, of } from "rxjs"; import { logger } from "matrix-js-sdk/lib/logger"; import { type LocalParticipant, type LocalTrack } from "livekit-client"; +import fetchMock from "fetch-mock"; import { PosthogAnalytics } from "../../../analytics/PosthogAnalytics"; -import { MatrixRTCMode } from "../../../config/ConfigOptions"; +import { + MatrixRTCMode, + type ResolvedDelayedLeaveTimings, +} from "../../../config/ConfigOptions"; import { type HomeserverDisconnectReason } from "./HomeserverConnected"; import { flushPromises, @@ -35,6 +40,7 @@ import { mockMuteStates, withTestScheduler, ownMemberMock, + testScope, } from "../../../utils/test"; import { TransportState, @@ -95,112 +101,109 @@ describe("watchScreenShareToggle", () => { }); }); -describe("LocalMembership", () => { - describe("enterRTCSession", () => { - it("It joins the correct Session", () => { - mockConfig({ - livekit: { livekit_service_url: "http://my-default-service-url.com" }, - }); +const timings: ResolvedDelayedLeaveTimings = { + delay_ms: 10000, + restart_ms: 4000, + restart_timeout_ms: 1000, +}; - 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, - }), - }, - }, - memberships: [], - joinRTCSession: vi.fn(), - }) as unknown as MatrixRTCSession; +const delegatedTimings: ResolvedDelayedLeaveTimings = { + delay_ms: timings.delay_ms * 10, + restart_ms: timings.restart_ms! * 10, + restart_timeout_ms: timings.restart_timeout_ms! * 10, +}; - enterRTCSession( - mockedSession, - ownMemberMock, - { - livekit_alias: "roomId", - livekit_service_url: "http://my-livekit-service-url.com", - type: "livekit", - }, - { - encryptMedia: true, - matrixRTCMode: MATRIX_RTC_MODE, - }, - ); +describe("enterRTCSession", () => { + const transport: LivekitTransportConfig = { + livekit_alias: "roomId", + livekit_service_url: "http://my-livekit-service-url.com", + type: "livekit", + }; - expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith( - { - deviceId: "DEVICE", - memberId: "@alice:example.org:DEVICE", - userId: "@alice:example.org", - }, - [], - { - livekit_alias: "roomId", - livekit_service_url: "http://my-livekit-service-url.com", - type: "livekit", - }, - expect.objectContaining({ manageMediaKeys: true }), - ); - }); + const options = { + encryptMedia: true, + matrixRTCMode: MATRIX_RTC_MODE, + delayedLeaveTimings: timings, + }; - it("passes keyRotationParticipantLimit from config to joinRTCSession", () => { - mockConfig({ - livekit: { livekit_service_url: "http://my-default-service-url.com" }, - matrix_rtc_session: { - delayed_leave_event_delay_ms: 0, - network_error_retry_ms: 0, - key_rotation_participant_limit: 50, - }, - }); - - 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, - }), - }, - }, - memberships: [], - joinRTCSession: vi.fn(), - }) as unknown as MatrixRTCSession; - - enterRTCSession( - mockedSession, - ownMemberMock, - { - livekit_alias: "roomId", - livekit_service_url: "http://my-livekit-service-url.com", - type: "livekit", - }, - { - encryptMedia: true, - matrixRTCMode: MATRIX_RTC_MODE, - }, - ); - - expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith( - expect.any(Object), - [], - expect.any(Object), - expect.objectContaining({ - keyRotationParticipantLimit: 50, + 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, }), - ); - }); + }, + }, + memberships: [], + joinRTCSession: vi.fn(), + }) as unknown as MatrixRTCSession; + + beforeEach(() => + mockConfig({ + livekit: { livekit_service_url: "http://my-default-service-url.com" }, + }), + ); + + it("It joins the correct Session", () => { + enterRTCSession(mockedSession, ownMemberMock, transport, options); + + expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith( + { + deviceId: "DEVICE", + memberId: "@alice:example.org:DEVICE", + userId: "@alice:example.org", + }, + [], + transport, + expect.objectContaining({ manageMediaKeys: true }), + ); }); + it("passes keyRotationParticipantLimit from config to joinRTCSession", () => { + mockConfig({ + livekit: { livekit_service_url: "http://my-default-service-url.com" }, + matrix_rtc_session: { + network_error_retry_ms: 0, + key_rotation_participant_limit: 50, + delayed_leave: timings, + delegated_delayed_leave: timings, + }, + }); + + enterRTCSession(mockedSession, ownMemberMock, transport, options); + + expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith( + expect.any(Object), + [], + expect.any(Object), + expect.objectContaining({ + keyRotationParticipantLimit: 50, + }), + ); + }); + + it("uses the specified delayed leave timings", () => { + enterRTCSession(mockedSession, ownMemberMock, transport, options); + + expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith( + expect.anything(), + expect.anything(), + expect.anything(), + expect.objectContaining({ + delayedLeaveEventRestartMs: timings.restart_ms, + delayedLeaveEventDelayMs: timings.delay_ms, + delayedLeaveEventRestartLocalTimeoutMs: timings.restart_timeout_ms, + }), + ); + }); +}); + +describe("LocalMembership", () => { const defaultCreateLocalMemberValues = { options: constant({ encryptMedia: false, @@ -226,8 +229,26 @@ describe("LocalMembership", () => { rtsSession$: constant(RTCMemberStatus.Connected), }, roomId: "!test-room-id:example.org", + baseUrl: "https://matrix.example.org", }; + beforeEach(() => { + mockConfig({ + livekit: { livekit_service_url: "http://my-default-service-url.com" }, + matrix_rtc_session: { + network_error_retry_ms: 1000, + delayed_leave: timings, + delegated_delayed_leave: delegatedTimings, + }, + }); + fetchMock.catch(404); + }); + + afterEach(async () => { + void (await fetchMock.flush()); + fetchMock.reset(); + }); + it("throws error on missing RTC config error", () => { withTestScheduler(({ scope, hot, behavior, expectObservable }) => { const localTransport$ = scope.behavior( @@ -256,7 +277,6 @@ describe("LocalMembership", () => { connectionManager: mockConnectionManager, localTransport$: behavior("a", { a: aLocalTransport }), }); - localMembership.requestJoinAndPublish(); expectObservable(localMembership.localMemberState$).toBe("ne", { n: TransportState.Waiting, @@ -299,9 +319,8 @@ describe("LocalMembership", () => { scope, ...defaultCreateLocalMemberValues, connectionManager: mockConnectionManager, - localTransport$: behavior("a", { a: aLocalTransport }), + localTransport$: constant(aLocalTransport), }); - localMembership.requestJoinAndPublish(); expectObservable(localMembership.localMemberState$).toBe("n-e", { n: TransportState.Waiting, @@ -397,6 +416,51 @@ describe("LocalMembership", () => { livekitRoom: mockLivekitRoom({}), } as unknown as Connection; + it.each([ + ["no", null, timings], + [ + "homeserver", + "https://matrix.example.org/_matrix/client/unstable/io.element.msc4195/rtc/livekit/delegate_delayed_leave", + delegatedTimings, + ], + ["transport", "/a/delegate_delayed_leave", delegatedTimings], + ])( + "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(); + + if (delegationUrl !== null) + fetchMock.post(delegationUrl, () => ({ status: 401, body: {} })); + + const localMembership = createLocalMembership$({ + scope, + ...defaultCreateLocalMemberValues, + connectionManager: { + connectionManagerData$: constant(new Epoch(connectionManagerData)), + }, + localTransport$: constant(aLocalTransport), + joinMatrixRTC, + }); + + localMembership.requestJoinAndPublish(); + void (await fetchMock.flush()); + await flushPromises(); + // Joins with timings appropriate for the level of delegation support + expect(joinMatrixRTC).toHaveBeenCalledWith( + aTransport, + delayedLeaveTimings, + ); + }, + ); + it("recreates publisher if new connection is used, always unpublish and end tracks", async () => { const scope = new ObservableScope(); diff --git a/src/state/CallViewModel/localMember/LocalMember.ts b/src/state/CallViewModel/localMember/LocalMember.ts index dfe8936a7..838b14ee2 100644 --- a/src/state/CallViewModel/localMember/LocalMember.ts +++ b/src/state/CallViewModel/localMember/LocalMember.ts @@ -62,7 +62,10 @@ import { screenShareCodec, parseResolution, } from "../../../settings/settings.ts"; -import { MatrixRTCMode } from "../../../config/ConfigOptions.ts"; +import { + MatrixRTCMode, + type ResolvedDelayedLeaveTimings, +} from "../../../config/ConfigOptions.ts"; import { Config } from "../../../config/Config.ts"; import { ConnectionState, @@ -72,6 +75,8 @@ import { import { type HomeserverConnected } from "./HomeserverConnected.ts"; import { type LocalTransport } from "./LocalTransport.ts"; import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts"; +import { doNetworkOperationWithRetry } from "../../../utils/matrix.ts"; +import { or$ } from "../../../utils/observable.ts"; export enum TransportState { /** Not even a transport is available to the LocalMembership */ @@ -133,7 +138,10 @@ interface Props { muteStates: MuteStates; connectionManager: IConnectionManager; createPublisherFactory: (connection: Connection) => Publisher; - joinMatrixRTC: (transport: LivekitTransportConfig) => void; + joinMatrixRTC: ( + transport: LivekitTransportConfig, + delayedLeaveTimings: ResolvedDelayedLeaveTimings, + ) => void; homeserverConnected: HomeserverConnected; roomId: string; localTransport$: Behavior; @@ -141,6 +149,7 @@ interface Props { MatrixRTCSession, "updateCallIntent" | "leaveRoomSession" >; + baseUrl: string; logger: Logger; } @@ -159,6 +168,7 @@ interface Props { * @param props.logger The logger to use. * @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.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. @@ -177,6 +187,7 @@ export const createLocalMembership$ = ({ logger: parentLogger, muteStates, matrixRTCSession, + baseUrl, roomId, }: Props): { /** @@ -236,11 +247,58 @@ export const createLocalMembership$ = ({ return of(null); }; - // This is the transport that we will advertise in our membership. - const advertisedTransport$ = localTransport$.pipe( + async function checkDelegationSupport( + endpointUrl: string, + serviceName: string, + ): Promise { + logger.info(`Checking whether ${serviceName} supports delegation…`); + try { + // Bluntly hit the endpoint without auth to check for a 404 + const res = await doNetworkOperationWithRetry(async () => + fetch(endpointUrl, { method: "POST" }), + ); + if (res.status === 404) { + logger.warn(`${serviceName} does not support delegation`); + return false; + } else { + logger.info(`${serviceName} supports delegation`); + return true; + } + } catch (e) { + logger.warn( + `Failed to determine whether ${serviceName} supports delegation, assuming no support`, + e, + ); + return false; + } + } + + const homeserverSupportsDelegation = checkDelegationSupport( + baseUrl + + "/_matrix/client/unstable/io.element.msc4195/rtc/livekit/delegate_delayed_leave", + "homeserver", + ); + + // The transport that we will advertise in our membership, paired with info as + // to whether delayed event delegation is supported + const joinParams$ = localTransport$.pipe( switchMap((lt) => lt.advertised$), catchError(handleTransportError), distinctUntilChanged(areLivekitTransportsEqual), + switchMap((transport) => { + if (transport === null) return of(null); + const transportSupportsDelegation = checkDelegationSupport( + transport.livekit_service_url + "/delegate_delayed_leave", + `transport ${transport.livekit_service_url}`, + ); + return or$( + from(homeserverSupportsDelegation), + from(transportSupportsDelegation), + ).pipe( + map((delegationSupported) => ({ transport, delegationSupported })), + startWith(null), + ); + }), ); // Unwrap the local transport and set the state of the LocalMembership to error in case the transport is an error. @@ -617,16 +675,20 @@ export const createLocalMembership$ = ({ // Keep matrix rtc session in sync with advertisedTransport$, connectRequested$ scope.reconcile( - scope.behavior( - combineLatest([advertisedTransport$, joinAndPublishRequested$]), - ), - async ([transport, shouldConnect]) => { - if (!transport) return; + scope.behavior(combineLatest([joinParams$, joinAndPublishRequested$])), + async ([joinParams, shouldConnect]) => { + if (!joinParams) return; // if shouldConnect=false we will do the disconnect as the cleanup from the previous reconcile iteration. if (!shouldConnect) return; + const sessionConfig = Config.get().matrix_rtc_session; try { - joinMatrixRTC(transport); + joinMatrixRTC( + joinParams.transport, + joinParams.delegationSupported + ? sessionConfig.delegated_delayed_leave + : sessionConfig.delayed_leave, + ); } catch (error) { logger.error("Error entering RTC session", error); if (error instanceof Error) @@ -864,6 +926,7 @@ export function observeSharingScreen$(p: Participant): Observable { interface EnterRTCSessionOptions { encryptMedia: boolean; matrixRTCMode: MatrixRTCMode; + delayedLeaveTimings: ResolvedDelayedLeaveTimings; } /** @@ -876,6 +939,7 @@ interface EnterRTCSessionOptions { * @param rtcSession - The MatrixRTCSession to join. * @param ownMembershipIdentity - Options for entering the RTC session. * @param transport - The LivekitTransport to use for this session. + * @param delayedLeaveTimings - The preferred timings for delayed leave events. * @param options - `encryptMedia`: Whether to encrypt media `matrixRTCMode`: The Matrix RTC mode to use. * @throws If the widget could not send ElementWidgetActions.JoinCall action. */ @@ -884,9 +948,8 @@ export function enterRTCSession( rtcSession: MatrixRTCSession, ownMembershipIdentity: CallMembershipIdentityParts, transport: LivekitTransportConfig, - options: EnterRTCSessionOptions, + { encryptMedia, matrixRTCMode, delayedLeaveTimings }: EnterRTCSessionOptions, ): void { - const { encryptMedia, matrixRTCMode } = options; PosthogAnalytics.instance.eventCallEnded.cacheStartCall(new Date()); PosthogAnalytics.instance.eventCallStarted.track(rtcSession.room.roomId); @@ -894,7 +957,11 @@ export function enterRTCSession( // have started tracking by the time calls start getting created. // groupCallOTelMembership?.onJoinCall(); - const { matrix_rtc_session: matrixRtcSessionConfig } = Config.get(); + const { + sync_disconnect_grace_period_ms: gracePeriod, + matrix_rtc_session: sessionConfig, + } = Config.get(); + const retryInterval = sessionConfig.network_error_retry_ms; const { sendNotificationType: notificationType, callIntent } = getUrlParams(); const multiSFU = matrixRTCMode === MatrixRTCMode.Compatibility || @@ -912,16 +979,11 @@ export function enterRTCSession( }; } - // Calculates `maximumNetworkErrorRetryCount`. The connection is failed if EITHER: - // - The /sync loop is unresponsive for > `gracePeriod` ms, or - // - A delayed leave event is emitted (after `leaveDelay` ms period). - // Note: Use leaveDelay >> gracePeriod for delegated leave events. - const gracePeriod = Config.get().sync_disconnect_grace_period_ms; - const leaveDelay = matrixRtcSessionConfig?.delayed_leave_event_delay_ms; - const retryInterval = matrixRtcSessionConfig?.network_error_retry_ms; - + // Set maximumNetworkErrorRetryCount such that we will consider the client + // disconnected as soon as either it fails to sync for longer than the grace + // period, or it is likely that a delayed leave event has been sent. // Math.min is used to account for the respective worst case: /sync not available or leave event emitted. - const maxWaitTime = Math.min(gracePeriod, leaveDelay); + const maxWaitTime = Math.min(gracePeriod, delayedLeaveTimings.delay_ms); const maximumNetworkErrorRetryCount = Math.ceil(maxWaitTime / retryInterval) + 1; @@ -936,18 +998,15 @@ export function enterRTCSession( notificationType, callIntent, manageMediaKeys: encryptMedia, - delayedLeaveEventRestartMs: - matrixRtcSessionConfig?.delayed_leave_event_restart_ms, - delayedLeaveEventDelayMs: - matrixRtcSessionConfig?.delayed_leave_event_delay_ms, + delayedLeaveEventRestartMs: delayedLeaveTimings.restart_ms, + delayedLeaveEventDelayMs: delayedLeaveTimings.delay_ms, delayedLeaveEventRestartLocalTimeoutMs: - matrixRtcSessionConfig?.delayed_leave_event_restart_local_timeout_ms, - networkErrorRetryMs: matrixRtcSessionConfig?.network_error_retry_ms, - makeKeyDelay: matrixRtcSessionConfig?.wait_for_key_rotation_ms, - membershipEventExpiryMs: - matrixRtcSessionConfig?.membership_event_expiry_ms, + delayedLeaveTimings.restart_timeout_ms, + networkErrorRetryMs: sessionConfig?.network_error_retry_ms, + makeKeyDelay: sessionConfig?.wait_for_key_rotation_ms, + membershipEventExpiryMs: sessionConfig?.membership_event_expiry_ms, keyRotationParticipantLimit: - matrixRtcSessionConfig?.key_rotation_participant_limit, + sessionConfig?.key_rotation_participant_limit, unstableSendStickyEvents: matrixRTCMode === MatrixRTCMode.Matrix_2_0, maximumNetworkErrorRetryCount: maximumNetworkErrorRetryCount, }, diff --git a/src/utils/observable.test.ts b/src/utils/observable.test.ts index 80cbb3c83..392f742ab 100644 --- a/src/utils/observable.test.ts +++ b/src/utils/observable.test.ts @@ -9,9 +9,30 @@ import { expect, test } from "vitest"; import { type Observable, of, Subject, switchMap } from "rxjs"; import { withTestScheduler } from "./test"; -import { filterBehavior, generateItems, pauseWhen } from "./observable"; +import { or$, filterBehavior, generateItems, pauseWhen } from "./observable"; import { type Behavior } from "../state/Behavior"; +const yesNo = { + y: true, + n: false, +}; + +test("or$", () => { + withTestScheduler(({ behavior, expectObservable }) => { + const input1Marbles = "ny--n--"; + const input2Marbles = "n-y--n-"; + const input3Marbles = "n--y--n"; + const outputMarbles = "nyyyyyn"; + expectObservable( + or$( + behavior(input1Marbles, yesNo), + behavior(input2Marbles, yesNo), + behavior(input3Marbles, yesNo), + ), + ).toBe(outputMarbles, yesNo); + }); +}); + test("pauseWhen", () => { withTestScheduler(({ behavior, expectObservable }) => { const inputMarbles = " abcdefgh-i-jk-"; diff --git a/src/utils/observable.ts b/src/utils/observable.ts index c32254db3..d4ea90b22 100644 --- a/src/utils/observable.ts +++ b/src/utils/observable.ts @@ -114,13 +114,11 @@ export function getValue(state$: Observable): T { } /** - * Creates an Observable that has a value of true whenever all its inputs are - * true. - * - * @public + * Creates an Observable that has a value of true whenever some of its inputs + * are true. */ -export function and$(...inputs: Observable[]): Observable { - return combineLatest(inputs, (...flags) => flags.every((flag) => flag)); +export function or$(...inputs: Observable[]): Observable { + return combineLatest(inputs, (...flags) => flags.some((flag) => flag)); } /** diff --git a/vite-embedded.config.ts b/vite-embedded.config.ts index 27a42fbbf..228d2cf8f 100644 --- a/vite-embedded.config.ts +++ b/vite-embedded.config.ts @@ -27,8 +27,6 @@ export default defineConfig((env) => data: { matrix_rtc_session: { wait_for_key_rotation_ms: 5000, - delayed_leave_event_restart_ms: 4000, - delayed_leave_event_delay_ms: 18000, }, }, },