mirror of
https://github.com/vector-im/element-call.git
synced 2026-08-29 21:15:19 +00:00
Delegate delayed events using dedicated endpoint
This commit is contained in:
@@ -18,7 +18,6 @@ import {
|
|||||||
NoMatrix2AuthorizationService,
|
NoMatrix2AuthorizationService,
|
||||||
} from "../utils/errors";
|
} from "../utils/errors";
|
||||||
import { doNetworkOperationWithRetry } from "../utils/matrix";
|
import { doNetworkOperationWithRetry } from "../utils/matrix";
|
||||||
import { Config } from "../config/Config";
|
|
||||||
import { JwtEndpointVersion } from "../state/CallViewModel/localMember/LocalTransport";
|
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.
|
* 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 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.
|
* 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.
|
* @param logger optional logger.
|
||||||
* @returns Object containing the token information
|
* @returns Object containing the token information
|
||||||
* @throws FailToGetOpenIdToken
|
* @throws FailToGetOpenIdToken
|
||||||
@@ -98,8 +95,6 @@ export async function getSFUConfigWithOpenID(
|
|||||||
roomId: string,
|
roomId: string,
|
||||||
opts?: {
|
opts?: {
|
||||||
forceJwtEndpoint?: JwtEndpointVersion;
|
forceJwtEndpoint?: JwtEndpointVersion;
|
||||||
delayEndpointBaseUrl?: string;
|
|
||||||
delayId?: string;
|
|
||||||
},
|
},
|
||||||
logger?: Logger,
|
logger?: Logger,
|
||||||
): Promise<SFUConfig> {
|
): Promise<SFUConfig> {
|
||||||
@@ -121,20 +116,18 @@ export async function getSFUConfigWithOpenID(
|
|||||||
const forceMatrix2Jwt =
|
const forceMatrix2Jwt =
|
||||||
opts?.forceJwtEndpoint === JwtEndpointVersion.Matrix_2_0;
|
opts?.forceJwtEndpoint === JwtEndpointVersion.Matrix_2_0;
|
||||||
|
|
||||||
// We want to start using the new endpoint (with optional delay delegation)
|
// We want to start using the new endpoint if we can use both or if we are
|
||||||
// if we can use both or if we are forced to use the new one.
|
// forced to use the new one.
|
||||||
if (tryBothJwtEndpoints || forceMatrix2Jwt) {
|
if (tryBothJwtEndpoints || forceMatrix2Jwt) {
|
||||||
try {
|
try {
|
||||||
logger?.info(
|
logger?.info(
|
||||||
`Trying to get JWT with delegation for focus ${serviceUrl}...`,
|
`Trying to get JWT with delegation for focus ${serviceUrl}...`,
|
||||||
);
|
);
|
||||||
const sfuConfig = await getLiveKitJWTWithDelayDelegation(
|
const sfuConfig = await getLiveKitJWT(
|
||||||
membership,
|
membership,
|
||||||
serviceUrl,
|
serviceUrl,
|
||||||
roomId,
|
roomId,
|
||||||
openIdToken,
|
openIdToken,
|
||||||
opts?.delayEndpointBaseUrl,
|
|
||||||
opts?.delayId,
|
|
||||||
);
|
);
|
||||||
|
|
||||||
return extractFullConfigFromToken(sfuConfig);
|
return extractFullConfigFromToken(sfuConfig);
|
||||||
@@ -154,13 +147,11 @@ export async function getSFUConfigWithOpenID(
|
|||||||
logger?.info(
|
logger?.info(
|
||||||
`Trying to get JWT with legacy endpoint for focus ${serviceUrl}...`,
|
`Trying to get JWT with legacy endpoint for focus ${serviceUrl}...`,
|
||||||
);
|
);
|
||||||
sfuConfig = await getLiveKitJWT(
|
sfuConfig = await getLiveKitJWTLegacy(
|
||||||
membership.deviceId,
|
membership.deviceId,
|
||||||
serviceUrl,
|
serviceUrl,
|
||||||
roomId,
|
roomId,
|
||||||
openIdToken,
|
openIdToken,
|
||||||
opts?.delayEndpointBaseUrl,
|
|
||||||
opts?.delayId,
|
|
||||||
);
|
);
|
||||||
logger?.info(`Got JWT from call's active focus URL.`);
|
logger?.info(`Got JWT from call's active focus URL.`);
|
||||||
return extractFullConfigFromToken(sfuConfig);
|
return extractFullConfigFromToken(sfuConfig);
|
||||||
@@ -188,33 +179,14 @@ function extractFullConfigFromToken(sfuConfig: {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
async function getLiveKitJWT(
|
async function getLiveKitJWTLegacy(
|
||||||
deviceId: string,
|
deviceId: string,
|
||||||
livekitServiceURL: string,
|
livekitServiceURL: string,
|
||||||
matrixRoomId: string,
|
matrixRoomId: string,
|
||||||
openIDToken: IOpenIDToken,
|
openIDToken: IOpenIDToken,
|
||||||
delayEndpointBaseUrl?: string,
|
|
||||||
delayId?: string,
|
|
||||||
): Promise<{ url: string; jwt: string }> {
|
): Promise<{ url: string; jwt: string }> {
|
||||||
interface IDelayParams {
|
const res = await doNetworkOperationWithRetry(async () =>
|
||||||
delay_id?: string;
|
fetch(livekitServiceURL + "/sfu/get", {
|
||||||
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", {
|
|
||||||
method: "POST",
|
method: "POST",
|
||||||
headers: {
|
headers: {
|
||||||
"Content-Type": "application/json",
|
"Content-Type": "application/json",
|
||||||
@@ -225,31 +197,9 @@ async function getLiveKitJWT(
|
|||||||
room: matrixRoomId,
|
room: matrixRoomId,
|
||||||
openid_token: openIDToken,
|
openid_token: openIDToken,
|
||||||
device_id: deviceId,
|
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) {
|
if (!res.ok) {
|
||||||
throw parseErrorResponse(res, await res.text());
|
throw parseErrorResponse(res, await res.text());
|
||||||
@@ -264,17 +214,15 @@ class NotSupportedError extends Error {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export async function getLiveKitJWTWithDelayDelegation(
|
export async function getLiveKitJWT(
|
||||||
membership: CallMembershipIdentityParts,
|
membership: CallMembershipIdentityParts,
|
||||||
livekitServiceURL: string,
|
livekitServiceURL: string,
|
||||||
matrixRoomId: string,
|
matrixRoomId: string,
|
||||||
openIDToken: IOpenIDToken,
|
openIDToken: IOpenIDToken,
|
||||||
delayEndpointBaseUrl?: string,
|
|
||||||
delayId?: string,
|
|
||||||
): Promise<{ url: string; jwt: string }> {
|
): Promise<{ url: string; jwt: string }> {
|
||||||
const { userId, deviceId, memberId } = membership;
|
const { userId, deviceId, memberId } = membership;
|
||||||
|
|
||||||
const body = {
|
const body = JSON.stringify({
|
||||||
room_id: matrixRoomId,
|
room_id: matrixRoomId,
|
||||||
slot_id: "m.call#ROOM",
|
slot_id: "m.call#ROOM",
|
||||||
openid_token: openIDToken,
|
openid_token: openIDToken,
|
||||||
@@ -283,27 +231,13 @@ export async function getLiveKitJWTWithDelayDelegation(
|
|||||||
claimed_user_id: userId,
|
claimed_user_id: userId,
|
||||||
claimed_device_id: deviceId,
|
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 () => {
|
const res = await doNetworkOperationWithRetry(async () => {
|
||||||
return await fetch(livekitServiceURL + "/get_token", {
|
return await fetch(livekitServiceURL + "/get_token", {
|
||||||
method: "POST",
|
method: "POST",
|
||||||
headers: {
|
headers: { "Content-Type": "application/json" },
|
||||||
"Content-Type": "application/json",
|
body,
|
||||||
},
|
|
||||||
body: JSON.stringify({ ...body, ...bodyDalayParts }),
|
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -49,7 +49,6 @@ import {
|
|||||||
import { type IWidgetApiRequest } from "matrix-widget-api";
|
import { type IWidgetApiRequest } from "matrix-widget-api";
|
||||||
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
|
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
|
||||||
import { v4 as uuidv4 } from "uuid";
|
import { v4 as uuidv4 } from "uuid";
|
||||||
import { type IMembershipManager } from "matrix-js-sdk/lib/matrixrtc/IMembershipManager";
|
|
||||||
|
|
||||||
import {
|
import {
|
||||||
createToggle$,
|
createToggle$,
|
||||||
@@ -488,16 +487,6 @@ export function createCallViewModel$(
|
|||||||
memberships$: memberships$,
|
memberships$: memberships$,
|
||||||
ownMembershipIdentity,
|
ownMembershipIdentity,
|
||||||
client,
|
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,
|
roomId: matrixRoom.roomId,
|
||||||
forceJwtEndpoint:
|
forceJwtEndpoint:
|
||||||
mode === MatrixRTCMode.Matrix_2_0
|
mode === MatrixRTCMode.Matrix_2_0
|
||||||
@@ -588,9 +577,20 @@ export function createCallViewModel$(
|
|||||||
);
|
);
|
||||||
},
|
},
|
||||||
connectionManager,
|
connectionManager,
|
||||||
|
client,
|
||||||
matrixRTCSession,
|
matrixRTCSession,
|
||||||
localTransport$,
|
localTransport$,
|
||||||
roomId: matrixRoom.roomId,
|
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()}]`),
|
logger: logger.getChild(`[${Date.now()}]`),
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -15,6 +15,11 @@ import {
|
|||||||
MediaDeviceFailure,
|
MediaDeviceFailure,
|
||||||
} from "livekit-client";
|
} from "livekit-client";
|
||||||
import { observeParticipantEvents } from "@livekit/components-core";
|
import { observeParticipantEvents } from "@livekit/components-core";
|
||||||
|
import {
|
||||||
|
type IOpenIDToken,
|
||||||
|
parseErrorResponse,
|
||||||
|
type MatrixClient,
|
||||||
|
} from "matrix-js-sdk";
|
||||||
import {
|
import {
|
||||||
Status as RTCSessionStatus,
|
Status as RTCSessionStatus,
|
||||||
type LivekitTransport,
|
type LivekitTransport,
|
||||||
@@ -72,6 +77,7 @@ import {
|
|||||||
import { type HomeserverConnected } from "./HomeserverConnected.ts";
|
import { type HomeserverConnected } from "./HomeserverConnected.ts";
|
||||||
import { type LocalTransport } from "./LocalTransport.ts";
|
import { type LocalTransport } from "./LocalTransport.ts";
|
||||||
import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts";
|
import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts";
|
||||||
|
import { doNetworkOperationWithRetry } from "../../../utils/matrix.ts";
|
||||||
|
|
||||||
export enum TransportState {
|
export enum TransportState {
|
||||||
/** Not even a transport is available to the LocalMembership */
|
/** Not even a transport is available to the LocalMembership */
|
||||||
@@ -136,11 +142,14 @@ interface Props {
|
|||||||
joinMatrixRTC: (transport: LivekitTransportConfig) => void;
|
joinMatrixRTC: (transport: LivekitTransportConfig) => void;
|
||||||
homeserverConnected: HomeserverConnected;
|
homeserverConnected: HomeserverConnected;
|
||||||
roomId: string;
|
roomId: string;
|
||||||
|
ownMembershipIdentity: CallMembershipIdentityParts;
|
||||||
localTransport$: Behavior<LocalTransport>;
|
localTransport$: Behavior<LocalTransport>;
|
||||||
|
client: Pick<MatrixClient, "getOpenIdToken" | "baseUrl">;
|
||||||
matrixRTCSession: Pick<
|
matrixRTCSession: Pick<
|
||||||
MatrixRTCSession,
|
MatrixRTCSession,
|
||||||
"updateCallIntent" | "leaveRoomSession"
|
"slotId" | "updateCallIntent" | "leaveRoomSession"
|
||||||
>;
|
>;
|
||||||
|
delayId$: Behavior<string | null>;
|
||||||
logger: Logger;
|
logger: Logger;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -176,8 +185,11 @@ export const createLocalMembership$ = ({
|
|||||||
joinMatrixRTC,
|
joinMatrixRTC,
|
||||||
logger: parentLogger,
|
logger: parentLogger,
|
||||||
muteStates,
|
muteStates,
|
||||||
|
client,
|
||||||
matrixRTCSession,
|
matrixRTCSession,
|
||||||
roomId,
|
roomId,
|
||||||
|
ownMembershipIdentity,
|
||||||
|
delayId$,
|
||||||
}: Props): {
|
}: Props): {
|
||||||
/**
|
/**
|
||||||
* This request to start audio and video tracks.
|
* 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(
|
const participant$ = scope.behavior(
|
||||||
localConnection$.pipe(
|
localConnection$.pipe(
|
||||||
map((c) => c?.livekitRoom?.localParticipant ?? null),
|
map((c) => c?.livekitRoom?.localParticipant ?? null),
|
||||||
|
|||||||
@@ -14,11 +14,8 @@ import {
|
|||||||
type MockedObject,
|
type MockedObject,
|
||||||
vi,
|
vi,
|
||||||
} from "vitest";
|
} from "vitest";
|
||||||
import {
|
import { type CallMembership } from "matrix-js-sdk/lib/matrixrtc";
|
||||||
type CallMembership,
|
import { lastValueFrom } from "rxjs";
|
||||||
type LivekitTransportConfig,
|
|
||||||
} from "matrix-js-sdk/lib/matrixrtc";
|
|
||||||
import { BehaviorSubject, filter, lastValueFrom } from "rxjs";
|
|
||||||
import fetchMock from "fetch-mock";
|
import fetchMock from "fetch-mock";
|
||||||
|
|
||||||
import {
|
import {
|
||||||
@@ -27,11 +24,7 @@ import {
|
|||||||
ownMemberMock,
|
ownMemberMock,
|
||||||
testScope,
|
testScope,
|
||||||
} from "../../../utils/test";
|
} from "../../../utils/test";
|
||||||
import {
|
import { createLocalTransport$, JwtEndpointVersion } from "./LocalTransport";
|
||||||
createLocalTransport$,
|
|
||||||
JwtEndpointVersion,
|
|
||||||
type LocalTransportWithSFUConfig,
|
|
||||||
} from "./LocalTransport";
|
|
||||||
import { constant } from "../../Behavior";
|
import { constant } from "../../Behavior";
|
||||||
import { Epoch, ObservableScope } from "../../ObservableScope";
|
import { Epoch, ObservableScope } from "../../ObservableScope";
|
||||||
import {
|
import {
|
||||||
@@ -68,7 +61,6 @@ describe("LocalTransport", () => {
|
|||||||
},
|
},
|
||||||
ownMembershipIdentity: ownMemberMock,
|
ownMembershipIdentity: ownMemberMock,
|
||||||
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
||||||
delayId$: constant("delay_id_mock"),
|
|
||||||
});
|
});
|
||||||
await flushPromises();
|
await flushPromises();
|
||||||
|
|
||||||
@@ -109,7 +101,6 @@ describe("LocalTransport", () => {
|
|||||||
},
|
},
|
||||||
ownMembershipIdentity: ownMemberMock,
|
ownMembershipIdentity: ownMemberMock,
|
||||||
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
||||||
delayId$: constant("delay_id_mock"),
|
|
||||||
});
|
});
|
||||||
active$.subscribe(
|
active$.subscribe(
|
||||||
(o) => observations.push(o),
|
(o) => observations.push(o),
|
||||||
@@ -151,7 +142,6 @@ describe("LocalTransport", () => {
|
|||||||
},
|
},
|
||||||
ownMembershipIdentity: ownMemberMock,
|
ownMembershipIdentity: ownMemberMock,
|
||||||
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
||||||
delayId$: constant("delay_id_mock"),
|
|
||||||
});
|
});
|
||||||
|
|
||||||
openIdResolver.resolve?.({
|
openIdResolver.resolve?.({
|
||||||
@@ -195,7 +185,6 @@ describe("LocalTransport", () => {
|
|||||||
scope: testScope(),
|
scope: testScope(),
|
||||||
roomId: "!example_room_id",
|
roomId: "!example_room_id",
|
||||||
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
||||||
delayId$: constant(null),
|
|
||||||
memberships$: constant(new Epoch<CallMembership[]>([])),
|
memberships$: constant(new Epoch<CallMembership[]>([])),
|
||||||
client: {
|
client: {
|
||||||
baseUrl: "https://example.org",
|
baseUrl: "https://example.org",
|
||||||
@@ -307,7 +296,6 @@ describe("LocalTransport", () => {
|
|||||||
ownMembershipIdentity: ownMemberMock,
|
ownMembershipIdentity: ownMemberMock,
|
||||||
roomId: "!example_room_id",
|
roomId: "!example_room_id",
|
||||||
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
forceJwtEndpoint: JwtEndpointVersion.Legacy,
|
||||||
delayId$: constant(null),
|
|
||||||
memberships$: constant(new Epoch<CallMembership[]>([])),
|
memberships$: constant(new Epoch<CallMembership[]>([])),
|
||||||
client: {
|
client: {
|
||||||
getDomain: () => "example.org",
|
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(),
|
|
||||||
);
|
|
||||||
});
|
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -10,14 +10,7 @@ import {
|
|||||||
type LivekitTransportConfig,
|
type LivekitTransportConfig,
|
||||||
} from "matrix-js-sdk/lib/matrixrtc";
|
} from "matrix-js-sdk/lib/matrixrtc";
|
||||||
import { type MatrixClient } from "matrix-js-sdk";
|
import { type MatrixClient } from "matrix-js-sdk";
|
||||||
import {
|
import { distinctUntilChanged, from, map, of, switchMap } from "rxjs";
|
||||||
combineLatest,
|
|
||||||
distinctUntilChanged,
|
|
||||||
from,
|
|
||||||
map,
|
|
||||||
of,
|
|
||||||
switchMap,
|
|
||||||
} from "rxjs";
|
|
||||||
import { logger as rootLogger, type Logger } from "matrix-js-sdk/lib/logger";
|
import { logger as rootLogger, type Logger } from "matrix-js-sdk/lib/logger";
|
||||||
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
|
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.
|
// Used by the jwt service to create the livekit room and compute the livekit alias.
|
||||||
roomId: string;
|
roomId: string;
|
||||||
forceJwtEndpoint: JwtEndpointVersion;
|
forceJwtEndpoint: JwtEndpointVersion;
|
||||||
delayId$: Behavior<string | null>;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
export enum JwtEndpointVersion {
|
export enum JwtEndpointVersion {
|
||||||
@@ -123,7 +115,6 @@ export const createLocalTransport$ = ({
|
|||||||
client,
|
client,
|
||||||
roomId,
|
roomId,
|
||||||
forceJwtEndpoint,
|
forceJwtEndpoint,
|
||||||
delayId$,
|
|
||||||
}: Props): LocalTransport => {
|
}: Props): LocalTransport => {
|
||||||
const logger = rootLogger.getChild("[LocalTransport]");
|
const logger = rootLogger.getChild("[LocalTransport]");
|
||||||
|
|
||||||
@@ -162,8 +153,8 @@ export const createLocalTransport$ = ({
|
|||||||
distinctUntilChanged(areLivekitTransportsEqual),
|
distinctUntilChanged(areLivekitTransportsEqual),
|
||||||
);
|
);
|
||||||
|
|
||||||
const preferredTransport$ = combineLatest([preferredConfig$, delayId$]).pipe(
|
const preferredTransport$ = preferredConfig$.pipe(
|
||||||
switchMap(async ([transport, delayId]) => {
|
switchMap(async (transport) => {
|
||||||
try {
|
try {
|
||||||
return await doOpenIdAndJWTFromUrl(
|
return await doOpenIdAndJWTFromUrl(
|
||||||
transport,
|
transport,
|
||||||
@@ -171,7 +162,6 @@ export const createLocalTransport$ = ({
|
|||||||
ownMembershipIdentity,
|
ownMembershipIdentity,
|
||||||
roomId,
|
roomId,
|
||||||
client,
|
client,
|
||||||
delayId ?? undefined,
|
|
||||||
logger,
|
logger,
|
||||||
);
|
);
|
||||||
} catch (e) {
|
} catch (e) {
|
||||||
@@ -223,7 +213,6 @@ export const createLocalTransport$ = ({
|
|||||||
* @param membership The identity of the local member.
|
* @param membership The identity of the local member.
|
||||||
* @param roomId The room ID to use for the JWT.
|
* @param roomId The room ID to use for the JWT.
|
||||||
* @param client The client to use for the OpenID token.
|
* @param client The client to use for the OpenID token.
|
||||||
* @param delayId The delayId to use for the JWT.
|
|
||||||
*
|
*
|
||||||
* @throws FailToGetOpenIdToken, NoMatrix2AuthorizationService
|
* @throws FailToGetOpenIdToken, NoMatrix2AuthorizationService
|
||||||
*/
|
*/
|
||||||
@@ -232,12 +221,7 @@ async function doOpenIdAndJWTFromUrl(
|
|||||||
forceJwtEndpoint: JwtEndpointVersion,
|
forceJwtEndpoint: JwtEndpointVersion,
|
||||||
membership: CallMembershipIdentityParts,
|
membership: CallMembershipIdentityParts,
|
||||||
roomId: string,
|
roomId: string,
|
||||||
client: Pick<
|
client: Pick<MatrixClient, "_unstable_getRTCTransports"> & OpenIDClientParts,
|
||||||
MatrixClient,
|
|
||||||
"getDomain" | "baseUrl" | "_unstable_getRTCTransports"
|
|
||||||
> &
|
|
||||||
OpenIDClientParts,
|
|
||||||
delayId?: string,
|
|
||||||
logger?: Logger,
|
logger?: Logger,
|
||||||
): Promise<LocalTransportWithSFUConfig> {
|
): Promise<LocalTransportWithSFUConfig> {
|
||||||
const sfuConfig = await getSFUConfigWithOpenID(
|
const sfuConfig = await getSFUConfigWithOpenID(
|
||||||
@@ -247,8 +231,6 @@ async function doOpenIdAndJWTFromUrl(
|
|||||||
roomId,
|
roomId,
|
||||||
{
|
{
|
||||||
forceJwtEndpoint: forceJwtEndpoint,
|
forceJwtEndpoint: forceJwtEndpoint,
|
||||||
delayEndpointBaseUrl: client.baseUrl,
|
|
||||||
delayId,
|
|
||||||
},
|
},
|
||||||
logger,
|
logger,
|
||||||
);
|
);
|
||||||
|
|||||||
Reference in New Issue
Block a user