Adapt delayed leave timings when delegation is available

Splits the config options for the timings of a delayed leave event into two sets: one for when delegation is available (as you can relax the timings and get more stable calls this way), and another for when it's unavailable (as we must continue to gracefully downgrade even after Matrix 2.0 is fully rolled out).

This works by bluntly hitting the delegation endpoints without auth before joining to check for a 404.
This commit is contained in:
Robin
2026-09-07 13:10:54 +02:00
parent a376280c91
commit 6a9583a2ce
13 changed files with 363 additions and 208 deletions
+16 -13
View File
@@ -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()}]`),
});
@@ -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<null | LivekitTransportConfig>(
@@ -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();
@@ -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<LocalTransport>;
@@ -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<boolean> {
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<boolean> {
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,
},