Compare commits

..

1 Commits

Author SHA1 Message Date
Robin
85d3cc33d9 Delegate delayed events using dedicated endpoint 2026-08-26 20:27:50 +02:00
5 changed files with 100 additions and 209 deletions

View File

@@ -18,7 +18,6 @@ import {
NoMatrix2AuthorizationService,
} from "../utils/errors";
import { doNetworkOperationWithRetry } from "../utils/matrix";
import { Config } from "../config/Config";
import { JwtEndpointVersion } from "../state/CallViewModel/localMember/LocalTransport";
/**
@@ -85,8 +84,6 @@ export type OpenIDClientParts = Pick<
* This function by default uses whatever is possible with the current jwt service installed next to the SFU.
* For remote connections this does not matter, since we will not publish there we can rely on the newest option.
* For our own connection we can only use the hashed version if we also send the new matrix2.0 sticky events.
* @param opts.delayEndpointBaseUrl The URL of the matrix homeserver.
* @param opts.delayId The delay id used for the jwt service to manage.
* @param logger optional logger.
* @returns Object containing the token information
* @throws FailToGetOpenIdToken
@@ -98,8 +95,6 @@ export async function getSFUConfigWithOpenID(
roomId: string,
opts?: {
forceJwtEndpoint?: JwtEndpointVersion;
delayEndpointBaseUrl?: string;
delayId?: string;
},
logger?: Logger,
): Promise<SFUConfig> {
@@ -121,20 +116,18 @@ export async function getSFUConfigWithOpenID(
const forceMatrix2Jwt =
opts?.forceJwtEndpoint === JwtEndpointVersion.Matrix_2_0;
// We want to start using the new endpoint (with optional delay delegation)
// if we can use both or if we are forced to use the new one.
// We want to start using the new endpoint if we can use both or if we are
// forced to use the new one.
if (tryBothJwtEndpoints || forceMatrix2Jwt) {
try {
logger?.info(
`Trying to get JWT with delegation for focus ${serviceUrl}...`,
);
const sfuConfig = await getLiveKitJWTWithDelayDelegation(
const sfuConfig = await getLiveKitJWT(
membership,
serviceUrl,
roomId,
openIdToken,
opts?.delayEndpointBaseUrl,
opts?.delayId,
);
return extractFullConfigFromToken(sfuConfig);
@@ -154,13 +147,11 @@ export async function getSFUConfigWithOpenID(
logger?.info(
`Trying to get JWT with legacy endpoint for focus ${serviceUrl}...`,
);
sfuConfig = await getLiveKitJWT(
sfuConfig = await getLiveKitJWTLegacy(
membership.deviceId,
serviceUrl,
roomId,
openIdToken,
opts?.delayEndpointBaseUrl,
opts?.delayId,
);
logger?.info(`Got JWT from call's active focus URL.`);
return extractFullConfigFromToken(sfuConfig);
@@ -188,33 +179,14 @@ function extractFullConfigFromToken(sfuConfig: {
};
}
async function getLiveKitJWT(
async function getLiveKitJWTLegacy(
deviceId: string,
livekitServiceURL: string,
matrixRoomId: string,
openIDToken: IOpenIDToken,
delayEndpointBaseUrl?: string,
delayId?: string,
): Promise<{ url: string; jwt: string }> {
interface IDelayParams {
delay_id?: string;
delay_timeout?: number;
delay_cs_api_url?: string;
}
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_cs_api_url: delayEndpointBaseUrl,
};
}
const makeRequest = async (delayParts: IDelayParams): Promise<Response> => {
return await fetch(livekitServiceURL + "/sfu/get", {
const res = await doNetworkOperationWithRetry(async () =>
fetch(livekitServiceURL + "/sfu/get", {
method: "POST",
headers: {
"Content-Type": "application/json",
@@ -225,31 +197,9 @@ async function getLiveKitJWT(
room: matrixRoomId,
openid_token: openIDToken,
device_id: deviceId,
...delayParts,
}),
});
};
const res = await doNetworkOperationWithRetry(async () => {
let response = await makeRequest(bodyDalayParts);
// Old service compatibility check
const oldServiceDoesNotSupportDelayParts =
response.status === 400 && Object.keys(bodyDalayParts).length > 0;
// If http status 400 with M_BAD_JSON and we sent delay parts, retry without them
if (oldServiceDoesNotSupportDelayParts) {
try {
const errorBody = await response.json();
if (errorBody.errcode === "M_BAD_JSON") {
response = await makeRequest({});
}
} catch {
// If we can't parse the error, treat as real error
}
}
return response;
});
}),
);
if (!res.ok) {
throw parseErrorResponse(res, await res.text());
@@ -264,17 +214,15 @@ class NotSupportedError extends Error {
}
}
export async function getLiveKitJWTWithDelayDelegation(
export async function getLiveKitJWT(
membership: CallMembershipIdentityParts,
livekitServiceURL: string,
matrixRoomId: string,
openIDToken: IOpenIDToken,
delayEndpointBaseUrl?: string,
delayId?: string,
): Promise<{ url: string; jwt: string }> {
const { userId, deviceId, memberId } = membership;
const body = {
const body = JSON.stringify({
room_id: matrixRoomId,
slot_id: "m.call#ROOM",
openid_token: openIDToken,
@@ -283,27 +231,13 @@ export async function getLiveKitJWTWithDelayDelegation(
claimed_user_id: userId,
claimed_device_id: deviceId,
},
};
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_cs_api_url: delayEndpointBaseUrl,
};
}
});
const res = await doNetworkOperationWithRetry(async () => {
return await fetch(livekitServiceURL + "/get_token", {
method: "POST",
headers: {
"Content-Type": "application/json",
},
body: JSON.stringify({ ...body, ...bodyDalayParts }),
headers: { "Content-Type": "application/json" },
body,
});
});

View File

@@ -49,7 +49,6 @@ import {
import { type IWidgetApiRequest } from "matrix-widget-api";
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
import { v4 as uuidv4 } from "uuid";
import { type IMembershipManager } from "matrix-js-sdk/lib/matrixrtc/IMembershipManager";
import {
createToggle$,
@@ -488,16 +487,6 @@ export function createCallViewModel$(
memberships$: memberships$,
ownMembershipIdentity,
client,
delayId$: scope.behavior(
(
fromEvent(
matrixRTCSession,
MembershipManagerEvent.DelayIdChanged,
// The type of reemitted event includes the original emitted as the second arg.
) as Observable<[string | undefined, IMembershipManager]>
).pipe(map(([delayId]) => delayId ?? null)),
matrixRTCSession.delayId ?? null,
),
roomId: matrixRoom.roomId,
forceJwtEndpoint:
mode === MatrixRTCMode.Matrix_2_0
@@ -588,9 +577,20 @@ export function createCallViewModel$(
);
},
connectionManager,
client,
matrixRTCSession,
localTransport$,
roomId: matrixRoom.roomId,
ownMembershipIdentity,
delayId$: scope.behavior(
(
fromEvent(
matrixRTCSession,
MembershipManagerEvent.DelayIdChanged,
) as Observable<[string | undefined]>
).pipe(map(([delayId]) => delayId ?? null)),
matrixRTCSession.delayId ?? null,
),
logger: logger.getChild(`[${Date.now()}]`),
});

View File

@@ -15,6 +15,11 @@ import {
MediaDeviceFailure,
} from "livekit-client";
import { observeParticipantEvents } from "@livekit/components-core";
import {
type IOpenIDToken,
parseErrorResponse,
type MatrixClient,
} from "matrix-js-sdk";
import {
Status as RTCSessionStatus,
type LivekitTransport,
@@ -72,6 +77,7 @@ 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";
export enum TransportState {
/** Not even a transport is available to the LocalMembership */
@@ -136,11 +142,14 @@ interface Props {
joinMatrixRTC: (transport: LivekitTransportConfig) => void;
homeserverConnected: HomeserverConnected;
roomId: string;
ownMembershipIdentity: CallMembershipIdentityParts;
localTransport$: Behavior<LocalTransport>;
client: Pick<MatrixClient, "getOpenIdToken" | "baseUrl">;
matrixRTCSession: Pick<
MatrixRTCSession,
"updateCallIntent" | "leaveRoomSession"
"slotId" | "updateCallIntent" | "leaveRoomSession"
>;
delayId$: Behavior<string | null>;
logger: Logger;
}
@@ -176,8 +185,11 @@ export const createLocalMembership$ = ({
joinMatrixRTC,
logger: parentLogger,
muteStates,
client,
matrixRTCSession,
roomId,
ownMembershipIdentity,
delayId$,
}: Props): {
/**
* This request to start audio and video tracks.
@@ -639,6 +651,61 @@ export const createLocalMembership$ = ({
},
);
scope.reconcile(
scope.behavior(combineLatest([delayId$, activeTransport$])),
async ([delayId, transport]) => {
if (delayId !== null && transport !== null) {
logger.info(
`Delegating delayed event ${delayId} to ${transport.livekit_service_url}`,
);
let openIdToken: IOpenIDToken;
try {
openIdToken = await doNetworkOperationWithRetry(async () =>
client.getOpenIdToken(),
);
} catch (e) {
logger.error("Failed to get OpenID token for delegation", e);
return;
}
const body = JSON.stringify({
room_id: roomId,
slot_id: matrixRTCSession.slotId,
member: {
id: ownMembershipIdentity.memberId,
claimed_user_id: ownMembershipIdentity.userId,
claimed_device_id: ownMembershipIdentity.deviceId,
},
delay_id: delayId,
delay_timeout:
Config.get().matrix_rtc_session?.delayed_leave_event_delay_ms,
delay_cs_api_url: client.baseUrl,
openid_token: openIdToken,
});
try {
const res = await doNetworkOperationWithRetry(async () =>
fetch(transport.livekit_service_url + "/delegate_delayed_leave", {
method: "POST",
headers: { "Content-Type": "application/json" },
body,
}),
);
if (!res.ok) {
if (res.status === 404)
logger.warn(
`Delayed event delegation not available on ${transport.livekit_service_url}`,
);
else throw parseErrorResponse(res, await res.text());
}
} catch (e) {
logger.error(
`Failed to delegate ${delayId} to ${transport.livekit_service_url}`,
e,
);
}
}
},
);
const participant$ = scope.behavior(
localConnection$.pipe(
map((c) => c?.livekitRoom?.localParticipant ?? null),

View File

@@ -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$, JwtEndpointVersion } from "./LocalTransport";
import { constant } from "../../Behavior";
import { Epoch, ObservableScope } from "../../ObservableScope";
import {
@@ -68,7 +61,6 @@ describe("LocalTransport", () => {
},
ownMembershipIdentity: ownMemberMock,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
});
await flushPromises();
@@ -109,7 +101,6 @@ describe("LocalTransport", () => {
},
ownMembershipIdentity: ownMemberMock,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
});
active$.subscribe(
(o) => observations.push(o),
@@ -151,7 +142,6 @@ describe("LocalTransport", () => {
},
ownMembershipIdentity: ownMemberMock,
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant("delay_id_mock"),
});
openIdResolver.resolve?.({
@@ -195,7 +185,6 @@ describe("LocalTransport", () => {
scope: testScope(),
roomId: "!example_room_id",
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant(null),
memberships$: constant(new Epoch<CallMembership[]>([])),
client: {
baseUrl: "https://example.org",
@@ -307,7 +296,6 @@ describe("LocalTransport", () => {
ownMembershipIdentity: ownMemberMock,
roomId: "!example_room_id",
forceJwtEndpoint: JwtEndpointVersion.Legacy,
delayId$: constant(null),
memberships$: constant(new Epoch<CallMembership[]>([])),
client: {
getDomain: () => "example.org",
@@ -329,84 +317,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(),
);
});
});

View File

@@ -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";
@@ -54,7 +47,6 @@ interface Props {
// 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 {
@@ -123,7 +115,6 @@ export const createLocalTransport$ = ({
client,
roomId,
forceJwtEndpoint,
delayId$,
}: Props): LocalTransport => {
const logger = rootLogger.getChild("[LocalTransport]");
@@ -162,8 +153,8 @@ export const createLocalTransport$ = ({
distinctUntilChanged(areLivekitTransportsEqual),
);
const preferredTransport$ = combineLatest([preferredConfig$, delayId$]).pipe(
switchMap(async ([transport, delayId]) => {
const preferredTransport$ = preferredConfig$.pipe(
switchMap(async (transport) => {
try {
return await doOpenIdAndJWTFromUrl(
transport,
@@ -171,7 +162,6 @@ export const createLocalTransport$ = ({
ownMembershipIdentity,
roomId,
client,
delayId ?? undefined,
logger,
);
} catch (e) {
@@ -223,7 +213,6 @@ export const createLocalTransport$ = ({
* @param membership The identity of the local member.
* @param roomId The room ID to use for the JWT.
* @param client The client to use for the OpenID token.
* @param delayId The delayId to use for the JWT.
*
* @throws FailToGetOpenIdToken, NoMatrix2AuthorizationService
*/
@@ -232,12 +221,7 @@ async function doOpenIdAndJWTFromUrl(
forceJwtEndpoint: JwtEndpointVersion,
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(
@@ -247,8 +231,6 @@ async function doOpenIdAndJWTFromUrl(
roomId,
{
forceJwtEndpoint: forceJwtEndpoint,
delayEndpointBaseUrl: client.baseUrl,
delayId,
},
logger,
);