Merge pull request #4242 from element-hq/delayed-leave-timings

Adapt delayed leave timings when delegation is available
This commit is contained in:
Robin
2026-09-09 08:19:58 +02:00
committed by GitHub
18 changed files with 570 additions and 527 deletions
+42 -55
View File
@@ -54,7 +54,6 @@ import { type IMembershipManager } from "matrix-js-sdk/lib/matrixrtc/IMembership
import {
createToggle$,
filterBehavior,
generateItem,
generateItems,
pauseWhen,
} from "../../utils/observable";
@@ -64,7 +63,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";
@@ -110,7 +112,6 @@ import {
} from "./localMember/LocalMember.ts";
import {
createLocalTransport$,
JwtEndpointVersion,
type LocalTransport,
} from "./localMember/LocalTransport.ts";
import {
@@ -188,7 +189,7 @@ export interface CallViewModelOptions {
/** Optional value overriding the connection factory, for testing purposes. */
connectionFactory?: ConnectionFactory;
/** The version & compatibility mode of MatrixRTC that we should use. */
matrixRTCMode$?: Behavior<MatrixRTCMode>;
matrixRTCMode?: MatrixRTCMode;
/** Optional behavior overriding for the screensharing, for testing */
toggleScreensharing?: () => void;
}
@@ -450,10 +451,8 @@ export function createCallViewModel$(
const configMatrixRTCMode = Config.get().matrix_rtc_mode as
| MatrixRTCMode
| undefined;
const matrixRTCMode$ =
configMatrixRTCMode !== undefined
? constant(configMatrixRTCMode)
: (options.matrixRTCMode$ ?? constant(MatrixRTCMode.Compatibility));
const matrixRTCMode =
configMatrixRTCMode ?? options.matrixRTCMode ?? MatrixRTCMode.Compatibility;
// Each hbar seperates a block of input variables required for the CallViewModel to function.
// The outputs of this block is written under the hbar.
@@ -487,38 +486,16 @@ export function createCallViewModel$(
memberId: uuidv4(),
};
const localTransport$ = scope.behavior(
matrixRTCMode$.pipe(
generateItem(
"CallViewModel localTransport$",
// Re-create LocalTransport whenever the mode changes
(mode) => ({ keys: [mode], data: undefined }),
(scope, _data$, mode) =>
options.localTransport ??
createLocalTransport$({
scope: scope,
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,
forceJwtEndpoint:
mode === MatrixRTCMode.Matrix_2_0
? JwtEndpointVersion.Matrix_2_0
: JwtEndpointVersion.Legacy,
}),
),
),
);
const localTransport =
options.localTransport ??
createLocalTransport$({
scope: scope,
memberships$: memberships$,
ownMembershipIdentity,
client,
roomId: matrixRoom.roomId,
matrixRTCMode,
});
const connectionFactory =
options.connectionFactory ??
@@ -536,8 +513,7 @@ export function createCallViewModel$(
scope: scope,
connectionFactory: connectionFactory,
localTransport$: scope.behavior(
localTransport$.pipe(
switchMap((t) => t.active$),
localTransport.active$.pipe(
catchError((e: unknown) => {
logger.info(
"could not pass local transport to createConnectionManager$. localTransport$ threw an error",
@@ -562,16 +538,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 +546,19 @@ export function createCallViewModel$(
matrixRTCSession,
),
muteStates,
joinMatrixRTC: (transport: LivekitTransportConfig) => {
joinMatrixRTC: (
transport: LivekitTransportConfig,
delayedLeaveTimings: ResolvedDelayedLeaveTimings,
) => {
return enterRTCSession(
matrixRTCSession,
ownMembershipIdentity,
transport,
connectOptions$.value,
{
encryptMedia: livekitKeyProvider !== undefined,
matrixRTCMode,
delayedLeaveTimings,
},
);
},
createPublisherFactory: (connection: Connection) => {
@@ -600,9 +573,23 @@ 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()}]`),
});
@@ -238,7 +238,7 @@ export function withCallViewModel(mode: MatrixRTCMode) {
);
},
},
matrixRTCMode$: constant(mode),
matrixRTCMode: mode,
...options,
},
raisedHands$,
@@ -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,
@@ -58,6 +64,7 @@ import {
type LocalTransport,
type LocalTransportWithSFUConfig,
} from "./LocalTransport";
import * as openIDSFU from "../../../livekit/openIDSFU";
initializeWidget();
@@ -95,112 +102,112 @@ 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,
},
);
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,
}),
};
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 }),
);
});
describe("enterRTCSession", () => {
const transport: LivekitTransportConfig = {
livekit_alias: "roomId",
livekit_service_url: "http://my-livekit-service-url.com",
type: "livekit",
};
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 options = {
encryptMedia: true,
matrixRTCMode: MATRIX_RTC_MODE,
delayedLeaveTimings: timings,
};
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 mockedSession = vi.mocked({
room: {
roomId: "roomId",
client: mockedClient,
},
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,
},
);
beforeEach(() =>
mockConfig({
livekit: { livekit_service_url: "http://my-default-service-url.com" },
}),
);
expect(mockedSession.joinRTCSession).toHaveBeenLastCalledWith(
expect.any(Object),
[],
expect.any(Object),
expect.objectContaining({
keyRotationParticipantLimit: 50,
}),
);
});
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 +233,30 @@ describe("LocalMembership", () => {
rtsSession$: constant(RTCMemberStatus.Connected),
},
roomId: "!test-room-id:example.org",
baseUrl: "https://matrix.example.org",
ownMembershipIdentity: ownMemberMock,
client: mockedClient,
delayId$: constant(null),
matrixRTCMode: MATRIX_RTC_MODE,
};
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>(
@@ -254,9 +283,8 @@ describe("LocalMembership", () => {
scope,
...defaultCreateLocalMemberValues,
connectionManager: mockConnectionManager,
localTransport$: behavior("a", { a: aLocalTransport }),
localTransport: aLocalTransport,
});
localMembership.requestJoinAndPublish();
expectObservable(localMembership.localMemberState$).toBe("ne", {
n: TransportState.Waiting,
@@ -299,9 +327,8 @@ describe("LocalMembership", () => {
scope,
...defaultCreateLocalMemberValues,
connectionManager: mockConnectionManager,
localTransport$: behavior("a", { a: aLocalTransport }),
localTransport: aLocalTransport,
});
localMembership.requestJoinAndPublish();
expectObservable(localMembership.localMemberState$).toBe("n-e", {
n: TransportState.Waiting,
@@ -317,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 = {
@@ -336,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)";
@@ -397,12 +424,90 @@ 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],
[
"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 joinMatrixRTC = vi.fn();
const delayId$ = new BehaviorSubject<string | null>(null);
if (delegationUrl !== null)
fetchMock.post(delegationUrl, () => ({ status: 401, body: {} }));
const localMembership = createLocalMembership$({
scope,
...defaultCreateLocalMemberValues,
connectionManager: {
connectionManagerData$: constant(
new Epoch(new ConnectionManagerData()),
),
},
joinMatrixRTC,
localTransport: {
advertised$: constant(aTransport),
active$: constant(aTransportWithSFUConfig),
},
delayId$,
});
localMembership.requestJoinAndPublish();
void (await fetchMock.flush());
await flushPromises();
// Joins with timings appropriate for the level of delegation support
expect(joinMatrixRTC).toHaveBeenCalledWith(
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");
}
},
);
it("recreates publisher if new connection is used, always unpublish and end tracks", async () => {
const scope = new ObservableScope();
const activeTransport$ = new BehaviorSubject(aTransportWithSFUConfig);
const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
advertised$: constant(aTransport),
active$: activeTransport$,
};
@@ -440,7 +545,7 @@ describe("LocalMembership", () => {
connectionManager: {
connectionManagerData$: constant(new Epoch(connectionManagerData)),
},
localTransport$: new BehaviorSubject(aLocalTransport),
localTransport: aLocalTransport,
});
await flushPromises();
activeTransport$.next({
@@ -471,7 +576,7 @@ describe("LocalMembership", () => {
const publishers: Publisher[] = [];
const tracks$ = new BehaviorSubject<LocalTrack[]>([]);
const publishing$ = new BehaviorSubject<boolean>(false);
const publishing$ = constant<boolean>(false);
defaultCreateLocalMemberValues.createPublisherFactory.mockImplementation(
() => {
const p = {
@@ -495,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();
@@ -508,7 +613,7 @@ describe("LocalMembership", () => {
connectionManager: {
connectionManagerData$: constant(new Epoch(connectionManagerData)),
},
localTransport$: new BehaviorSubject(aLocalTransport),
localTransport: aLocalTransport,
});
await flushPromises();
expect(publisherFactory).toHaveBeenCalledOnce();
@@ -536,7 +641,7 @@ describe("LocalMembership", () => {
new BehaviorSubject<null | LocalTransportWithSFUConfig>(null);
const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
advertised$: constant(aTransport),
active$: activeTransport$,
};
@@ -579,7 +684,7 @@ describe("LocalMembership", () => {
connectionManager: {
connectionManagerData$,
},
localTransport$: new BehaviorSubject(aLocalTransport),
localTransport: aLocalTransport,
});
await flushPromises();
@@ -714,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();
@@ -754,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();
@@ -796,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();
@@ -848,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();
@@ -917,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 };
};
@@ -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,
@@ -62,7 +63,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 +76,8 @@ import {
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 */
@@ -133,14 +139,22 @@ 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>;
ownMembershipIdentity: CallMembershipIdentityParts;
localTransport: LocalTransport;
client: Pick<MatrixClient, "getDeviceId" | "getOpenIdToken">;
matrixRTCSession: Pick<
MatrixRTCSession,
"updateCallIntent" | "leaveRoomSession"
>;
baseUrl: string;
delayId$: Behavior<string | null>;
matrixRTCMode: MatrixRTCMode;
logger: Logger;
}
@@ -155,10 +169,12 @@ interface Props {
* @param props.createPublisherFactory Factory to create a publisher once we have a connection.
* @param props.joinMatrixRTC Callback to join the matrix RTC session once we have a transport.
* @param props.homeserverConnected The homeserver connected state.
* @param props.localTransport$ The transport to advertise in our membership.
* @param props.localTransport The transport to advertise in our membership.
* @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.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.
@@ -170,14 +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.
@@ -236,25 +257,69 @@ export const createLocalMembership$ = ({
return of(null);
};
// This is the transport that we will advertise in our membership.
const advertisedTransport$ = localTransport$.pipe(
switchMap((lt) => lt.advertised$),
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. Unfortunately
// we can't wrap this in a retry loop, as many servers don't just disable
// delegation support, but in fact are from a time before the endpoint
// existed at all, therefore we can hit CORS errors which would just gum
// up the retry loop. (May be revisited after Matrix 2.0.)
const res = await 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.advertised$.pipe(
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.
const activeTransport$ = scope.behavior(
localTransport$.pipe(
switchMap((lt) => {
return combineLatest([lt.active$, lt.advertised$]).pipe(
map(([active, advertised]) => {
// Our policy is to not publish to another transport if our prefered transport is miss-configured
if (advertised == null) return null;
combineLatest([localTransport.active$, localTransport.advertised$]).pipe(
map(([active, advertised]) => {
// Our policy is to not publish to another transport if our prefered transport is miss-configured
if (advertised == null) return null;
return active?.transport ?? null;
}),
);
return active?.transport ?? null;
}),
catchError(handleTransportError),
distinctUntilChanged(areLivekitTransportsEqual),
@@ -617,16 +682,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)
@@ -653,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.
@@ -864,6 +961,7 @@ export function observeSharingScreen$(p: Participant): Observable<boolean> {
interface EnterRTCSessionOptions {
encryptMedia: boolean;
matrixRTCMode: MatrixRTCMode;
delayedLeaveTimings: ResolvedDelayedLeaveTimings;
}
/**
@@ -876,6 +974,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 +983,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 +992,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 +1014,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 +1033,14 @@ 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,
keyRotationParticipantLimit:
matrixRtcSessionConfig?.key_rotation_participant_limit,
delayedLeaveTimings.restart_timeout_ms,
networkErrorRetryMs: sessionConfig.network_error_retry_ms,
makeKeyDelay: sessionConfig.wait_for_key_rotation_ms,
membershipEventExpiryMs: sessionConfig.membership_event_expiry_ms,
keyRotationParticipantLimit: sessionConfig.key_rotation_participant_limit,
unstableSendStickyEvents: matrixRTCMode === MatrixRTCMode.Matrix_2_0,
maximumNetworkErrorRetryCount: maximumNetworkErrorRetryCount,
},
@@ -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 {
@@ -41,6 +34,7 @@ import {
import * as openIDSFU from "../../../livekit/openIDSFU";
import { customLivekitUrl } from "../../../settings/settings";
import { testJWTToken } from "../../../utils/test-fixtures";
import { MatrixRTCMode } from "../../../config/ConfigOptions";
describe("LocalTransport", () => {
const openIdResponse: openIDSFU.SFUConfig = {
@@ -61,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,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
matrixRTCMode: MatrixRTCMode.Compatibility,
});
await flushPromises();
@@ -100,7 +92,6 @@ describe("LocalTransport", () => {
roomId: "!example_room_id",
memberships$: constant(new Epoch<CallMembership[]>([])),
client: {
baseUrl: "https://example.org",
getDomain: () => "example.org",
// eslint-disable-next-line @typescript-eslint/naming-convention
_unstable_getRTCTransports: async () => Promise.resolve([]),
@@ -108,8 +99,7 @@ describe("LocalTransport", () => {
getDeviceId: vi.fn(),
},
ownMembershipIdentity: ownMemberMock,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
matrixRTCMode: MatrixRTCMode.Compatibility,
});
active$.subscribe(
(o) => observations.push(o),
@@ -147,11 +137,9 @@ describe("LocalTransport", () => {
getDomain: () => "example.org",
getOpenIdToken: vi.fn(),
getDeviceId: vi.fn(),
baseUrl: "https://example.org",
},
ownMembershipIdentity: ownMemberMock,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
matrixRTCMode: MatrixRTCMode.Compatibility,
});
openIdResolver.resolve?.({
@@ -194,11 +182,9 @@ describe("LocalTransport", () => {
ownMembershipIdentity: ownMemberMock,
scope: testScope(),
roomId: "!example_room_id",
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant(null),
matrixRTCMode: MatrixRTCMode.Compatibility,
memberships$: constant(new Epoch<CallMembership[]>([])),
client: {
baseUrl: "https://example.org",
getDomain: vi.fn().mockReturnValue("example.org"),
// eslint-disable-next-line @typescript-eslint/naming-convention
_unstable_getRTCTransports: vi.fn().mockResolvedValue([]),
@@ -306,12 +292,10 @@ describe("LocalTransport", () => {
scope: testScope(),
ownMembershipIdentity: ownMemberMock,
roomId: "!example_room_id",
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant(null),
matrixRTCMode: MatrixRTCMode.Compatibility,
memberships$: constant(new Epoch<CallMembership[]>([])),
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
@@ -329,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<string | null>(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<CallMembership[]>([])),
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(),
);
});
});
@@ -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";
@@ -37,6 +30,7 @@ import {
import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts";
import { customLivekitUrl } from "../../../settings/settings.ts";
import { RtcTransportAutoDiscovery } from "./RtcTransportAutoDiscovery.ts";
import { type MatrixRTCMode } from "../../../config/ConfigOptions.ts";
/*
* It figures out “which LiveKit focus URL/alias the local user should use,”
@@ -46,20 +40,11 @@ interface Props {
scope: ObservableScope;
ownMembershipIdentity: CallMembershipIdentityParts;
memberships$: Behavior<Epoch<CallMembership[]>>;
client: Pick<
MatrixClient,
"getDomain" | "baseUrl" | "_unstable_getRTCTransports"
> &
client: Pick<MatrixClient, "getDomain" | "_unstable_getRTCTransports"> &
OpenIDClientParts;
// Used by the jwt service to create the livekit room and compute the livekit alias.
roomId: string;
forceJwtEndpoint: JwtEndpointVersion;
delayId$: Behavior<string | null>;
}
export enum JwtEndpointVersion {
Legacy = "legacy",
Matrix_2_0 = "matrix_2_0",
matrixRTCMode: MatrixRTCMode;
}
// TODO livekit_alias-cleanup
@@ -122,8 +107,7 @@ export const createLocalTransport$ = ({
ownMembershipIdentity,
client,
roomId,
forceJwtEndpoint,
delayId$,
matrixRTCMode,
}: Props): LocalTransport => {
const logger = rootLogger.getChild("[LocalTransport]");
@@ -138,40 +122,36 @@ 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,
forceJwtEndpoint,
matrixRTCMode,
ownMembershipIdentity,
roomId,
client,
delayId ?? undefined,
logger,
);
} catch (e) {
@@ -193,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,25 +185,19 @@ export const createLocalTransport$ = ({
* use we don't want to risk any issues by re-using a token.
*
* @param transport The transport to authenticate with.
* @param forceJwtEndpoint Whether to force the JWT endpoint to be used.
* @param matrixRTCMode Whether to force the JWT endpoint to be used.
* @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
*/
async function doOpenIdAndJWTFromUrl(
transport: LivekitTransportConfig,
forceJwtEndpoint: JwtEndpointVersion,
matrixRTCMode: MatrixRTCMode,
membership: CallMembershipIdentityParts,
roomId: string,
client: Pick<
MatrixClient,
"getDomain" | "baseUrl" | "_unstable_getRTCTransports"
> &
OpenIDClientParts,
delayId?: string,
client: Pick<MatrixClient, "_unstable_getRTCTransports"> & OpenIDClientParts,
logger?: Logger,
): Promise<LocalTransportWithSFUConfig> {
const sfuConfig = await getSFUConfigWithOpenID(
@@ -245,11 +205,7 @@ async function doOpenIdAndJWTFromUrl(
membership,
transport.livekit_service_url,
roomId,
{
forceJwtEndpoint: forceJwtEndpoint,
delayEndpointBaseUrl: client.baseUrl,
delayId,
},
{ matrixRTCMode },
logger,
);
return {