sdk implementation

This commit is contained in:
Timo K.
2026-10-01 20:46:35 +02:00
parent 12bb87fa4c
commit 08abf1a387
31 changed files with 3015 additions and 348 deletions
+106 -27
View File
@@ -5,7 +5,12 @@ SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { expect, type Page } from "@playwright/test";
import {
type Browser,
expect,
type Locator,
type Page,
} from "@playwright/test";
import { SynapseAdmin } from "../utils/synapse-admin.ts";
@@ -18,35 +23,74 @@ export const SDK_HARNESS_URL = "https://localhost:3002";
const HOMESERVER_URL = "https://synapse.m.localhost";
const PASSWORD = "foobarbaz1!";
/** A registered user, with the session the registration opened. */
export interface User {
username: string;
displayName: string;
userId: string;
deviceId: string;
accessToken: string;
}
/**
* Registers two users through the Synapse admin API and has the first create a
* public room, so that the second can join it by id without an invite.
* public room, so that the second can join it by id without an invite. The
* registration logs each user in, which spares the harness a `/login` call
* that the homeserver rate-limits.
*/
export async function createUsersAndRoom(
name: string,
): Promise<{ usernames: [string, string]; roomId: string }> {
): Promise<{ users: [User, User]; roomId: string }> {
const admin = SynapseAdmin.forHomeserver(HOMESERVER_URL);
const usernames: [string, string] = [
`${name}_a_${Date.now()}`,
`${name}_b_${Date.now()}`,
];
const [{ access_token: accessToken }] = await Promise.all(
usernames.map(async (username, index) =>
admin.registerUser(username, PASSWORD, `${name} ${"AB"[index]}`),
),
);
const users = (await Promise.all(
["a", "b"].map(async (letter, index): Promise<User> => {
const username = `${name}_${letter}_${Date.now()}`;
const displayName = `${name} ${"AB"[index]}`;
const registration = await admin.registerUser(
username,
PASSWORD,
displayName,
);
return {
username,
displayName,
userId: registration.user_id,
deviceId: registration.device_id,
accessToken: registration.access_token,
};
}),
)) as [User, User];
const response = await fetch(
`${HOMESERVER_URL}/_matrix/client/v3/createRoom`,
{
method: "POST",
headers: {
Authorization: `Bearer ${accessToken}`,
Authorization: `Bearer ${users[0].accessToken}`,
"Content-Type": "application/json",
},
body: JSON.stringify({
name: `${name}'s session`,
preset: "public_chat",
// Per-participant media keys travel in encrypted to-device messages,
// which only reach devices the sender's crypto tracks, and it tracks
// the members of encrypted rooms
initial_state: [
{
type: "m.room.encryption",
state_key: "",
content: { algorithm: "m.megolm.v1.aes-sha2" },
},
],
// Every member has to be allowed to write its own membership, which
// a public room does not grant by default; the same levels Element
// Call gives a room it creates for a call
power_level_content_override: {
state_default: 0,
events_default: 0,
users_default: 0,
events: { "org.matrix.msc3401.call.member": 0 },
},
}),
},
);
@@ -56,33 +100,68 @@ export async function createUsersAndRoom(
);
const { room_id: roomId } = (await response.json()) as { room_id: string };
return { usernames, roomId };
return { users, roomId };
}
/**
* Opens the harness signed in as the given user and waits until the SDK
* reports the session as joined. The status line carries the SDK's error if
* it does not get there, so the failure says why.
* reports the session as connected. The status line carries the SDK's error
* if it does not get there, so the failure says why.
*/
export async function startHarness(
page: Page,
username: string,
user: User,
roomId: string,
): Promise<void> {
const query = new URLSearchParams({
homeserver: HOMESERVER_URL,
username,
password: PASSWORD,
accessToken: user.accessToken,
userId: user.userId,
deviceId: user.deviceId,
room: roomId,
});
await page.goto(`${SDK_HARNESS_URL}/?${query.toString()}`);
await page.getByRole("button", { name: "Start" }).click();
// A login, a crypto setup and an initial sync happen first. An error is
// final, so it is not worth waiting out the timeout for "Joined" after one.
const status = page.getByTestId("status");
await expect(status).toHaveText(/^(Joined|Error)/, { timeout: 120_000 });
const text = await status.textContent();
if (text !== "Joined")
throw new Error(`The harness did not join the session: ${text}`);
await waitForConnected(page);
}
/**
* A login, a crypto setup and an initial sync happen before the session can
* connect. An error is final, so it is not worth waiting out the timeout for
* "connected" after one.
*/
export async function waitForConnected(page: Page): Promise<void> {
const status = page.getByTestId("status");
await expect(status).toHaveText(/^(connected|Error)/, { timeout: 120_000 });
const text = await status.textContent();
if (text !== "connected")
throw new Error(`The harness did not connect to the session: ${text}`);
}
/**
* A context of its own per user, since a browser profile holds one login. No
* permissions to grant: each browser is launched with fake media that is
* handed out without asking (see playwright.config.ts).
*/
export async function newPage(browser: Browser): Promise<Page> {
const context = await browser.newContext({ ignoreHTTPSErrors: true });
return context.newPage();
}
/** Two users in one session, each on a page of their own. */
export async function startPair(
browser: Browser,
name: string,
): Promise<{ pages: [Page, Page]; users: [User, User]; roomId: string }> {
const { users, roomId } = await createUsersAndRoom(name);
const pages = await Promise.all([newPage(browser), newPage(browser)]);
await Promise.all(
pages.map(async (page, i) => startHarness(page, users[i], roomId)),
);
return { pages: pages as [Page, Page], users, roomId };
}
/** The tile a page shows for a user. */
export function tileOf(page: Page, user: User): Locator {
return page.locator(`[data-testid="member"][data-user-id="${user.userId}"]`);
}
+67
View File
@@ -0,0 +1,67 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { expect, test } from "@playwright/test";
import { startPair, tileOf, waitForConnected } from "./harness.ts";
/**
* Members coming and going: a member that leaves disappears for the peer, one
* that comes back is shown again, and the harness is left in the state its
* status line claims.
*/
test.describe.configure({ timeout: 300_000 });
test("a member that leaves disappears for the peer", async ({ browser }) => {
const {
pages: [leaver, stayer],
users: [leaving, staying],
} = await startPair(browser, "sdkleave");
await expect(stayer.getByTestId("member")).toHaveCount(2, {
timeout: 60_000,
});
await leaver.getByRole("button", { name: "Leave" }).click();
await expect(leaver.getByTestId("status")).toHaveText("Left");
await expect(leaver.getByTestId("member")).toHaveCount(0);
await expect(tileOf(stayer, leaving)).toHaveCount(0, { timeout: 60_000 });
await expect(tileOf(stayer, staying)).toHaveCount(1);
await expect(stayer.getByTestId("status")).toHaveText("connected");
});
test("a member that reloads is shown again with its media", async ({
browser,
}) => {
const {
pages: [reloader, watcher],
users: [reloading],
} = await startPair(browser, "sdkreload");
await expect(watcher.getByTestId("member")).toHaveCount(2, {
timeout: 60_000,
});
// The page comes back on the same device, so the new membership replaces
// the old one rather than sitting next to it
await reloader.reload();
await reloader.getByRole("button", { name: "Start" }).click();
await waitForConnected(reloader);
await expect(tileOf(watcher, reloading)).toHaveCount(1, {
timeout: 120_000,
});
const video = tileOf(watcher, reloading).locator("video");
await expect
.poll(async () => video.evaluate((v: HTMLVideoElement) => v.videoWidth), {
timeout: 60_000,
})
.toBeGreaterThan(0);
await expect(reloader.getByTestId("member")).toHaveCount(2, {
timeout: 60_000,
});
});
+158
View File
@@ -0,0 +1,158 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { expect, type Locator, type Page, test } from "@playwright/test";
import { startPair, tileOf, type User } from "./harness.ts";
/**
* What a member's media looks like from the other side: it plays, it is
* encrypted, it follows the mute switches, and the SFU sends only as much as
* the tile on screen can show. Each test reads the labels the harness puts on
* its media elements, which come straight from the SDK's track behaviors.
*
* One session serves every test here: logging two browsers in and connecting
* them takes most of a minute, and none of the tests leaves the session in a
* state the next cannot start from.
*/
test.describe.configure({ mode: "serial", timeout: 300_000 });
let pages: [Page, Page];
let users: [User, User];
test.beforeAll(async ({ browser }) => {
({ pages, users } = await startPair(browser, "sdkmedia"));
for (const page of pages)
await expect(page.getByTestId("member")).toHaveCount(2, {
timeout: 60_000,
});
});
test.afterAll(async () => {
for (const page of pages ?? []) await page.context().close();
});
test("the peer's camera and microphone play", async () => {
const peer = remoteTile(1, 0);
const video = peer.locator("video");
await expect
.poll(async () => video.evaluate((v: HTMLVideoElement) => v.videoWidth), {
timeout: 60_000,
})
.toBeGreaterThan(0);
await expect(video).toHaveJSProperty("paused", false);
const audio = peer.locator("audio");
await expect
.poll(async () =>
audio.evaluate((a: HTMLAudioElement) => a.srcObject !== null),
)
.toBe(true);
await expect(audio).toHaveJSProperty("paused", false);
});
test("the peer's tracks are encrypted", async () => {
const peer = remoteTile(1, 0);
await expect(peer.locator("video")).toHaveAttribute("data-encrypted", "true");
await expect(peer.locator("audio")).toHaveAttribute("data-encrypted", "true");
});
test("muting the microphone and the camera is seen by the peer", async () => {
const [page] = pages;
const peer = remoteTile(1, 0);
const video = peer.locator("video");
const audio = peer.locator("audio");
await expect(video).toHaveAttribute("data-muted", "false");
await expect(audio).toHaveAttribute("data-muted", "false");
await page.getByRole("button", { name: "Microphone" }).click();
await expect(audio).toHaveAttribute("data-muted", "true", {
timeout: 20_000,
});
await expect(video).toHaveAttribute("data-muted", "false");
await page.getByRole("button", { name: "Camera" }).click();
await expect(video).toHaveAttribute("data-muted", "true", {
timeout: 20_000,
});
await page.getByRole("button", { name: "Microphone" }).click();
await page.getByRole("button", { name: "Camera" }).click();
await expect(audio).toHaveAttribute("data-muted", "false", {
timeout: 20_000,
});
await expect(video).toHaveAttribute("data-muted", "false", {
timeout: 20_000,
});
});
test("a larger tile receives a higher resolution", async () => {
const peer = remoteTile(1, 0);
const video = peer.locator("video");
// The tile starts small enough for the lowest simulcast layer
const small = await receivedWidth(video, (width) => width > 0);
await resize(peer, 1280);
const large = await receivedWidth(video, (width) => width > small);
expect(large).toBeGreaterThan(small);
await resize(peer, 240);
await receivedWidth(video, (width) => width < large);
});
test("a hidden tile stops receiving video, and a shown one resumes", async () => {
const peer = remoteTile(1, 0);
const video = peer.locator("video");
await expect(frames(video)).resolves.toBeGreaterThan(0);
await peer.evaluate((tile) => (tile.style.display = "none"));
await expect
.poll(async () => framesStill(video), { timeout: 60_000 })
.toBe(true);
await peer.evaluate((tile) => (tile.style.display = ""));
await expect
.poll(async () => !(await framesStill(video)), { timeout: 60_000 })
.toBe(true);
});
/** The tile on one page for the user of the other page. */
function remoteTile(viewer: 0 | 1, shown: 0 | 1): Locator {
return tileOf(pages[viewer], users[shown]);
}
async function resize(tile: Locator, width: number): Promise<void> {
await tile.evaluate((t, w) => (t.style.width = `${w}px`), width);
}
/** The width of the video as decoded, once it satisfies the condition. */
async function receivedWidth(
video: Locator,
until: (width: number) => boolean,
): Promise<number> {
// Switching simulcast layers takes the SFU a few seconds
let width = 0;
await expect
.poll(
async () => {
width = await video.evaluate((v: HTMLVideoElement) => v.videoWidth);
return until(width);
},
{ timeout: 60_000 },
)
.toBe(true);
return width;
}
async function frames(video: Locator): Promise<number> {
return Number(await video.getAttribute("data-frames"));
}
/** Whether no frame arrived over a couple of seconds. */
async function framesStill(video: Locator): Promise<boolean> {
const before = await frames(video);
await new Promise((resolve) => setTimeout(resolve, 2500));
return (await frames(video)) === before;
}
+6 -30
View File
@@ -5,17 +5,13 @@ SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { type Browser, expect, type Page, test } from "@playwright/test";
import { expect, test } from "@playwright/test";
import { createUsersAndRoom, startHarness } from "./harness.ts";
import { startPair, tileOf } from "./harness.ts";
/**
* The MatrixRTC SDK driven through its harness in `sdk/dev`: no Element Call
* on the page, only the SDK's public API.
*
* This is the smoke test the implementation is built against. It fails until
* `createRtcSession` does something, and it fails on the status line first, so
* the SDK's own error is what the report shows.
*/
// Two browsers each log in, set up crypto and sync before anything is on
@@ -23,34 +19,14 @@ import { createUsersAndRoom, startHarness } from "./harness.ts";
test.describe.configure({ timeout: 300_000 });
test("two browsers see each other in one session", async ({ browser }) => {
const { usernames, roomId } = await createUsersAndRoom("sdksmoke");
const [pageA, pageB] = await Promise.all(
usernames.map(async () => newPage(browser)),
);
await Promise.all([
startHarness(pageA, usernames[0], roomId),
startHarness(pageB, usernames[1], roomId),
]);
const { pages, users } = await startPair(browser, "sdksmoke");
// Each page shows itself and the other, each tile named by its user
for (const page of [pageA, pageB]) {
for (const page of pages) {
await expect(page.getByTestId("member")).toHaveCount(2, {
timeout: 60_000,
});
for (const username of usernames)
await expect(
page.getByTestId("member").filter({ hasText: username }),
).toHaveCount(1);
for (const user of users)
await expect(tileOf(page, user)).toContainText(user.displayName);
}
});
/**
* A context of its own per user, since a browser profile holds one login. No
* permissions to grant: each browser is launched with fake media that is
* handed out without asking (see playwright.config.ts).
*/
async function newPage(browser: Browser): Promise<Page> {
const context = await browser.newContext({ ignoreHTTPSErrors: true });
return context.newPage();
}
+6 -4
View File
@@ -6,9 +6,11 @@ the media of every member, as observables. It has no UI. Element Call's own
`CallViewModel` is meant to become one consumer of it; the design is in
[`sdk-plan.md`](../sdk-plan.md).
**Status:** interface only. `createRtcSession` returns an object that does nothing.
The development harness and its e2e test exist so the implementation can be built
against them.
**Status:** first implementation. `createRtcSession` joins the session, connects to
the transport, publishes the local media and exposes every member's media; the
development harness and its e2e tests in `playwright/sdk` drive it. The
implementation still imports the building blocks it shares with Element Call from
`src/` (connections, memberships, key provider); moving them here is the next slice.
## Using it
@@ -25,7 +27,7 @@ const scope = new ObservableScope();
const session = createRtcSession(
scope,
client, // a matrix-js-sdk MatrixClient, logged in and syncing
room, // the matrix-js-sdk Room to hold the session in
room, // the matrix-js-sdk Room to hold the session in, from client.getRoom()
{
microphoneEnabled$: constant(true),
cameraEnabled$: constant(true),
+15 -2
View File
@@ -21,9 +21,12 @@ Please see LICENSE in the repository root for full details.
flex-wrap: wrap;
}
#members {
display: grid;
grid-template-columns: repeat(auto-fill, minmax(240px, 1fr));
display: flex;
gap: 0.5rem;
flex-wrap: wrap;
}
#members section {
width: 240px;
}
#members video {
width: 100%;
@@ -46,7 +49,17 @@ Please see LICENSE in the repository root for full details.
<input name="username" placeholder="username" />
<input name="password" type="password" placeholder="password" />
<input name="room" placeholder="room id or alias" />
<!-- A host that already holds a session passes it instead of a password -->
<input name="accessToken" type="hidden" />
<input name="userId" type="hidden" />
<input name="deviceId" type="hidden" />
<button type="submit">Start</button>
<button type="button" id="microphone" aria-pressed="true" hidden>
Microphone
</button>
<button type="button" id="camera" aria-pressed="true" hidden>
Camera
</button>
<button type="button" id="leave" hidden>Leave</button>
</form>
<p id="status" data-testid="status"></p>
+96 -21
View File
@@ -7,12 +7,18 @@ Please see LICENSE in the repository root for full details.
/**
* The smallest consumer of the SDK: log in, join a room, show every member's
* camera and play every remote member's microphone, leave. What a host writes,
* and nothing a host would not.
* camera and play every remote member's microphone, mute, leave. What a host
* writes, and nothing a host would not.
*/
import { logger } from "matrix-js-sdk/lib/logger";
import { combineLatest, type Observable, of, switchMap } from "rxjs";
import {
BehaviorSubject,
combineLatest,
type Observable,
of,
switchMap,
} from "rxjs";
import {
constant,
createRtcSession,
@@ -24,12 +30,16 @@ import {
type RtcSession,
} from "@element-hq/matrixrtc-sdk";
import { createSession } from "./session";
import { createSession, joinRoom, type Login } from "./session";
const form = document.querySelector("form")!;
const status = document.getElementById("status")!;
const members = document.getElementById("members")!;
const leaveButton = document.getElementById("leave") as HTMLButtonElement;
const buttons = {
microphone: document.getElementById("microphone") as HTMLButtonElement,
camera: document.getElementById("camera") as HTMLButtonElement,
leave: document.getElementById("leave") as HTMLButtonElement,
};
// So that a test, or a bookmark, can fill the form from the URL
const params = new URLSearchParams(location.search);
@@ -40,33 +50,36 @@ form.addEventListener("submit", (event) => {
event.preventDefault();
const fields = new FormData(form);
const field = (name: string): string => fields.get(name) as string;
void start(
field("homeserver"),
field("username"),
field("password"),
field("room"),
);
const login: Login = field("accessToken")
? {
accessToken: field("accessToken"),
userId: field("userId"),
deviceId: field("deviceId"),
}
: { username: field("username"), password: field("password") };
void start(field("homeserver"), login, field("room"));
});
async function start(
homeserver: string,
username: string,
password: string,
login: Login,
roomIdOrAlias: string,
): Promise<void> {
try {
status.textContent = "Logging in";
const client = await createSession(homeserver, username, password);
const room = await client.joinRoom(roomIdOrAlias);
const client = await createSession(homeserver, login);
const room = await joinRoom(client, roomIdOrAlias);
const scope = new ObservableScope();
const microphoneEnabled$ = new BehaviorSubject(true);
const cameraEnabled$ = new BehaviorSubject(true);
const session = createRtcSession(
scope,
client,
room,
{
microphoneEnabled$: constant(true),
cameraEnabled$: constant(true),
microphoneEnabled$,
cameraEnabled$,
audioInputDeviceId$: constant(undefined),
videoInputDeviceId$: constant(undefined),
videoProcessor$: constant(undefined),
@@ -76,16 +89,23 @@ async function start(
matrixRTCMode: MatrixRTCMode.Compatibility,
},
);
session.status$.pipe(scope.bind()).subscribe((s) => {
status.textContent = s;
});
session.fatalError$.pipe(scope.bind()).subscribe((error) => {
if (error !== null) status.textContent = `Error: ${error.message}`;
});
showMembers(scope, session);
session.join();
status.textContent = "Joined";
leaveButton.hidden = false;
leaveButton.onclick = (): void => {
toggle(buttons.microphone, microphoneEnabled$);
toggle(buttons.camera, cameraEnabled$);
buttons.leave.hidden = false;
buttons.leave.onclick = (): void => {
session.leave();
scope.end();
members.replaceChildren();
leaveButton.hidden = true;
for (const button of Object.values(buttons)) button.hidden = true;
status.textContent = "Left";
};
} catch (e) {
@@ -94,6 +114,17 @@ async function start(
}
}
function toggle(
button: HTMLButtonElement,
enabled$: BehaviorSubject<boolean>,
): void {
button.hidden = false;
button.onclick = (): void => {
enabled$.next(!enabled$.value);
button.ariaPressed = String(enabled$.value);
};
}
function showMembers(scope: ObservableScope, session: RtcSession): void {
const tiles = new Map<string, HTMLElement>();
combineLatest([session.localMember$, session.remoteMembers$])
@@ -118,6 +149,7 @@ function memberTile(scope: ObservableScope, member: RtcMember): HTMLElement {
const tile = document.createElement("section");
tile.dataset.testid = "member";
tile.dataset.userId = member.userId;
tile.dataset.local = String(member.local);
const name = tile.appendChild(document.createElement("h2"));
member.displayName$.pipe(scope.bind()).subscribe((n) => {
@@ -142,6 +174,10 @@ function memberTile(scope: ObservableScope, member: RtcMember): HTMLElement {
return tile;
}
/**
* Plays a track on an element and labels the element with the track's state,
* so that the page shows what the member is sending and how well it arrives.
*/
function render(
scope: ObservableScope,
track$: Observable<MediaTrack | undefined>,
@@ -154,4 +190,43 @@ function render(
track?.attach(element);
});
scope.onEnd(() => attached?.detach(element));
const label = (
key: string,
value$: Observable<string | number | boolean | undefined>,
): void => {
value$.pipe(scope.bind()).subscribe((value) => {
if (value === undefined) delete element.dataset[key];
else element.dataset[key] = String(value);
});
};
const of$ = <T>(
pick: (track: MediaTrack) => Observable<T>,
): Observable<T | undefined> =>
track$.pipe(switchMap((track) => (track ? pick(track) : of(undefined))));
label(
"muted",
of$((t) => t.muted$),
);
label(
"encrypted",
of$((t) => t.encrypted$),
);
const stats$ = of$((t) => t.stats$);
label("frameWidth", stats$.pipe(switchMap((s) => of(frames(s)?.frameWidth))));
label("frames", stats$.pipe(switchMap((s) => of(frames(s)?.count))));
}
/** The frame counters an RTP stream reports, from either end of it. */
function frames(
stats: RTCInboundRtpStreamStats | RTCOutboundRtpStreamStats | undefined,
): { frameWidth: number | undefined; count: number | undefined } | undefined {
if (stats === undefined) return undefined;
return {
frameWidth: stats.frameWidth,
count:
stats.type === "inbound-rtp"
? (stats as RTCInboundRtpStreamStats).framesDecoded
: (stats as RTCOutboundRtpStreamStats).framesEncoded,
};
}
+60 -10
View File
@@ -10,25 +10,29 @@ import {
createClient,
type MatrixClient,
MemoryStore,
type Room,
SyncState,
} from "matrix-js-sdk";
import { KnownMembership } from "matrix-js-sdk/lib/types";
/** Logs in with a password and returns a client that has finished its first sync. */
export type Login =
| { username: string; password: string }
| { accessToken: string; userId: string; deviceId: string };
/**
* Returns a client that has finished its first sync, logging in with a
* password unless the caller already holds a token.
*/
export async function createSession(
homeserver: string,
username: string,
password: string,
login: Login,
): Promise<MatrixClient> {
const login = await createClient({ baseUrl: homeserver }).login(
"m.login.password",
{ identifier: { type: "m.id.user", user: username }, password },
);
const credentials =
"accessToken" in login ? login : await logIn(homeserver, login);
const client = createClient({
baseUrl: homeserver,
accessToken: login.access_token,
userId: login.user_id,
deviceId: login.device_id,
...credentials,
store: new MemoryStore(),
useAuthorizationHeader: true,
fallbackICEServerAllowed: true,
@@ -47,3 +51,49 @@ export async function createSession(
return client;
}
/**
* Joins a room and returns it as the sync loop maintains it. The room object
* `joinRoom` itself returns for a room joined just now is a detached copy
* that never receives the state the sync delivers, and a session built on
* it would never see a member.
*/
export async function joinRoom(
client: MatrixClient,
roomIdOrAlias: string,
): Promise<Room> {
const { roomId } = await client.joinRoom(roomIdOrAlias);
const joinedRoom = (): Room | undefined => {
const room = client.getRoom(roomId);
return room?.hasMembershipState(client.getUserId()!, KnownMembership.Join)
? room
: undefined;
};
return (
joinedRoom() ??
new Promise<Room>((resolve) => {
const onSync = (): void => {
const room = joinedRoom();
if (room === undefined) return;
client.off(ClientEvent.Sync, onSync);
resolve(room);
};
client.on(ClientEvent.Sync, onSync);
})
);
}
async function logIn(
homeserver: string,
{ username, password }: { username: string; password: string },
): Promise<{ accessToken: string; userId: string; deviceId: string }> {
const login = await createClient({ baseUrl: homeserver }).login(
"m.login.password",
{ identifier: { type: "m.id.user", user: username }, password },
);
return {
accessToken: login.access_token,
userId: login.user_id,
deviceId: login.device_id,
};
}
+7 -251
View File
@@ -10,262 +10,18 @@ Please see LICENSE in the repository root for full details.
*
* MatrixRTC sessions with LiveKit media, without Element Call's UI: the call
* model Element Call's own view model is built on, for hosts that want to
* build a different one.
*
* Only the interface exists so far. `createRtcSession` returns an object that
* does nothing, so that the development harness in `sdk/dev` and its e2e test
* can be written against the contract before the implementation lands behind
* it. The design and the migration from Element Call's `CallViewModel` are in
* `sdk-plan.md` at the repository root.
* build a different one. The design and the migration from Element Call's
* `CallViewModel` are in `sdk-plan.md` at the repository root.
*/
import { type MatrixClient, type Room } from "matrix-js-sdk";
import {
type CallMembership,
type RTCCallIntent,
type RTCNotificationType,
type Transport,
} from "matrix-js-sdk/lib/matrixrtc";
import { type Track, type TrackProcessor } from "livekit-client";
import { type Observable } from "rxjs";
import { type Behavior } from "../src/state/Behavior";
import { type ObservableScope } from "../src/state/ObservableScope";
import { type EncryptionSystem } from "../src/e2ee/sharedKeyManagement";
import { type MatrixRTCMode } from "../src/config/ConfigOptions";
// Shared with Element Call. They live in `src` until the SDK has an
// implementation to move them with; a consumer gets them from here either way.
// Shared with Element Call. They live in `src` until the view model consumes
// the SDK, at which point they move here; a consumer gets them from this
// module either way.
export { type Behavior, constant } from "../src/state/Behavior";
export { ObservableScope } from "../src/state/ObservableScope";
export { E2eeType } from "../src/e2ee/e2eeType";
export { type EncryptionSystem } from "../src/e2ee/sharedKeyManagement";
export { MatrixRTCMode } from "../src/config/ConfigOptions";
// ---------------------------------------------------------------------------
// Session
export interface RtcSessionOptions {
encryptionSystem: EncryptionSystem;
/** Resolved by the host; the SDK reads neither config.json nor settings. */
matrixRTCMode: MatrixRTCMode;
/**
* MSC4075 notification sent with the join. Parameters of the MatrixRTC join
* itself, so they are here even though they are named after calls; reacting
* to a notification (ringing, timeouts, declines) is the application's job.
*/
sendNotificationType?: RTCNotificationType;
callIntent?: RTCCallIntent;
}
/** What the local member publishes. */
export interface LocalMediaInputs {
microphoneEnabled$: Behavior<boolean>;
cameraEnabled$: Behavior<boolean>;
audioInputDeviceId$: Behavior<string | undefined>;
videoInputDeviceId$: Behavior<string | undefined>;
/** Background blur and the like. */
videoProcessor$: Behavior<TrackProcessor<Track.Kind.Video> | undefined>;
}
export type SessionConnectionStatus =
| "waitingForTransport"
| "connecting"
| "connected"
| "reconnecting"
| "disconnected";
export class RtcSessionError extends Error {
public constructor(message: string, options?: ErrorOptions) {
super(message, options);
this.name = "RtcSessionError";
}
}
export interface RtcSession {
join(): void;
leave(): void;
/** Collapsed view of the local member's state machine. */
status$: Behavior<SessionConnectionStatus>;
connected$: Behavior<boolean>;
reconnecting$: Behavior<boolean>;
/** A transport, Matrix or connection error that stops the session. */
fatalError$: Behavior<RtcSessionError | null>;
localMember$: Behavior<LocalRtcMember | null>;
remoteMembers$: Behavior<RemoteRtcMember[]>;
/** `remoteMembers.length`, plus one for the local member once it exists. */
memberCount$: Behavior<number>;
/**
* Whether the session has grown large enough that MatrixRTC has stopped
* rotating the media encryption key.
*/
keyRotationSuppressed$: Behavior<boolean>;
/** Transports the session currently holds a live connection to. */
connectedTransports$: Behavior<TransportMetadata[]>;
}
// ---------------------------------------------------------------------------
// Transports
/**
* One transport advertised in a membership. Transport independent: `type` and
* `id` are all the SDK needs; `raw` and `resolved$` are there for a
* backend-specific developer panel and for connection diagnostics.
*/
export interface TransportMetadata {
/** `"livekit"` today. */
type: string;
/** Stable key, unique per transport in the session. For LiveKit, the service url. */
id: string;
/** The transport object as it appears in the membership. */
raw: Transport;
/**
* What the backend had to fetch before it could connect. Undefined until the
* connection has resolved it, and again after the connection stops.
*/
resolved$: Behavior<ResolvedTransport | undefined>;
}
export type ResolvedTransport =
| {
type: "livekit";
/** The SFU websocket url, as opposed to the JWT service url in `raw`. */
url: string;
/** A secret: fit for a developer panel, never for a log line. */
token: string;
roomAlias: string;
identity: string;
}
| { type: string; [key: string]: unknown };
// ---------------------------------------------------------------------------
// Members
export interface RtcMember {
local: boolean;
/** `${userId}:${deviceId}` before sticky events, a uuid in Matrix 2.0 mode. The backend identity. */
id: string;
userId: string;
deviceId: string;
membership$: Behavior<CallMembership>;
displayName$: Behavior<string>;
avatarUrl$: Behavior<string | undefined>;
/** Which transport this member is on; undefined when the membership has none. */
transport$: Behavior<TransportMetadata | undefined>;
/**
* Null while the member has a transport but no media has arrived on it yet
* ("waiting for media").
*/
media$: Behavior<MemberMedia | null>;
}
export interface RemoteRtcMember extends RtcMember {
local: false;
}
export interface LocalRtcMember extends RtcMember {
local: true;
media$: Behavior<LocalMemberMedia | null>;
sharingScreen$: Behavior<boolean>;
/** Null when the platform cannot share a screen. */
toggleScreenSharing: (() => void) | null;
screenShareError$: Behavior<Error | null>;
dismissScreenShareError(): void;
}
// ---------------------------------------------------------------------------
// Media
export type MediaSource =
| "microphone"
| "camera"
| "screenShare"
| "screenShareAudio";
export type MediaStreamStats =
| RTCInboundRtpStreamStats
| RTCOutboundRtpStreamStats
| undefined;
/** One published track of a member. */
export interface MediaTrack {
source: MediaSource;
kind: "audio" | "video";
/** Stable for the life of the track. */
id: string;
muted$: Behavior<boolean>;
/** False when the SFU reports the track as unencrypted. */
encrypted$: Behavior<boolean>;
/** Polled while subscribed. */
stats$: Behavior<MediaStreamStats>;
/**
* Rendering. The view hands its `<video>` or `<audio>` element over; the SDK
* sets its stream and, for video, registers the size and on-screen observers
* that pick a simulcast layer and pause the subscription while the element is
* hidden. `attach` is idempotent per element; `detach` has to be called
* before the element leaves the DOM so those observers are released.
*/
attach(element: HTMLMediaElement): void;
detach(element: HTMLMediaElement): void;
}
export interface AudioMediaTrack extends MediaTrack {
kind: "audio";
/** Route playback through Web Audio, for earpiece pan and gain. Undefined resets. */
setAudioContext(ctx: AudioContext | undefined, plugins?: AudioNode[]): void;
setVolume(volume: number): void;
}
export interface VideoMediaTrack extends MediaTrack {
kind: "video";
/** Undefined for remote tracks. */
facingMode$?: Behavior<"user" | "environment" | undefined>;
}
export type EncryptionError = "MissingKey" | "InvalidKey";
/**
* The media of one member, backed by a LiveKit participant inside the SDK.
* There is no identity field: the LiveKit identity is the member's `id`.
*/
export interface MemberMedia {
local: boolean;
speaking$: Behavior<boolean>;
screenShareEnabled$: Behavior<boolean>;
microphone$: Behavior<AudioMediaTrack | undefined>;
camera$: Behavior<VideoMediaTrack | undefined>;
screenShare$: Behavior<VideoMediaTrack | undefined>;
screenShareAudio$: Behavior<AudioMediaTrack | undefined>;
/** Emits when the SFU reports a key problem for this member. */
encryptionError$: Observable<EncryptionError>;
}
export interface LocalMemberMedia extends MemberMedia {
local: true;
}
// ---------------------------------------------------------------------------
// Factory
/**
* Takes the whole `MatrixClient` rather than a slice of it, and finds the
* MatrixRTC session itself. The js-sdk is the MatrixRTC implementation today;
* when the rust-rtc crate replaces it, the SDK has to bridge the crate to the
* js-sdk client, and only the SDK knows what that bridge needs. Holding the
* client keeps that change inside the SDK.
*/
export function createRtcSession(
scope: ObservableScope,
client: MatrixClient,
room: Room,
localMedia: LocalMediaInputs,
options: RtcSessionOptions,
): RtcSession {
// The interface ships ahead of the implementation; see the module comment.
return {} as RtcSession;
}
export * from "./src/api";
export { createRtcSession } from "./src/session/RtcSession";
+241
View File
@@ -0,0 +1,241 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
/**
* The public types of the SDK. `createRtcSession` in `session/RtcSession.ts`
* is the only way to obtain an implementation of them.
*/
import {
type CallMembership,
type RTCCallIntent,
type RTCNotificationType,
type Transport,
} from "matrix-js-sdk/lib/matrixrtc";
import { type Track, type TrackProcessor } from "livekit-client";
import { type Observable } from "rxjs";
import { type Behavior } from "../../src/state/Behavior";
import { type EncryptionSystem } from "../../src/e2ee/sharedKeyManagement";
import { type MatrixRTCMode } from "../../src/config/ConfigOptions";
// ---------------------------------------------------------------------------
// Session
export interface RtcSessionOptions {
encryptionSystem: EncryptionSystem;
/** Resolved by the host; the SDK reads neither config.json nor settings. */
matrixRTCMode: MatrixRTCMode;
/**
* MSC4075 notification sent with the join. Parameters of the MatrixRTC join
* itself, so they are here even though they are named after calls; reacting
* to a notification (ringing, timeouts, declines) is the application's job.
*/
sendNotificationType?: RTCNotificationType;
callIntent?: RTCCallIntent;
}
/** What the local member publishes. */
export interface LocalMediaInputs {
microphoneEnabled$: Behavior<boolean>;
cameraEnabled$: Behavior<boolean>;
audioInputDeviceId$: Behavior<string | undefined>;
videoInputDeviceId$: Behavior<string | undefined>;
/** Background blur and the like. */
videoProcessor$: Behavior<TrackProcessor<Track.Kind.Video> | undefined>;
}
export type SessionConnectionStatus =
| "waitingForTransport"
| "connecting"
| "connected"
| "reconnecting"
| "disconnected";
/**
* An error raised by the session. `cause` holds the underlying error, which
* lets a host that knows the backend tell the failures apart.
*/
export class RtcSessionError extends Error {
public constructor(message: string, options?: ErrorOptions) {
super(message, options);
this.name = "RtcSessionError";
}
}
/**
* A session in one room. The room has to be the one the client's sync loop
* maintains (`client.getRoom(roomId)`, once the join has synced): the
* detached copy `joinRoom` returns for a room joined just now never receives
* the state the members are read from.
*/
export interface RtcSession {
join(): void;
leave(): void;
/** Collapsed view of the local member's state machine. */
status$: Behavior<SessionConnectionStatus>;
connected$: Behavior<boolean>;
reconnecting$: Behavior<boolean>;
/** A transport, Matrix or connection error that stops the session. */
fatalError$: Behavior<RtcSessionError | null>;
localMember$: Behavior<LocalRtcMember | null>;
remoteMembers$: Behavior<RemoteRtcMember[]>;
/** `remoteMembers.length`, plus one for the local member once it exists. */
memberCount$: Behavior<number>;
/**
* Whether the session has grown large enough that MatrixRTC has stopped
* rotating the media encryption key.
*/
keyRotationSuppressed$: Behavior<boolean>;
/** Transports the session currently holds a live connection to. */
connectedTransports$: Behavior<TransportMetadata[]>;
}
// ---------------------------------------------------------------------------
// Transports
/**
* One transport advertised in a membership. Transport independent: `type` and
* `id` are all the SDK needs; `raw` and `resolved$` are there for a
* backend-specific developer panel and for connection diagnostics.
*/
export interface TransportMetadata {
/** `"livekit"` today. */
type: string;
/** Stable key, unique per transport in the session. For LiveKit, the service url. */
id: string;
/** The transport object as it appears in the membership. */
raw: Transport;
/**
* What the backend had to fetch before it could connect. Undefined until the
* connection has resolved it, and again after the connection stops.
*/
resolved$: Behavior<ResolvedTransport | undefined>;
}
export type ResolvedTransport =
| {
type: "livekit";
/** The SFU websocket url, as opposed to the JWT service url in `raw`. */
url: string;
/** A secret: fit for a developer panel, never for a log line. */
token: string;
roomAlias: string;
identity: string;
}
| { type: string; [key: string]: unknown };
// ---------------------------------------------------------------------------
// Members
export interface RtcMember {
local: boolean;
/** The identity the media backend knows this member by. */
id: string;
userId: string;
deviceId: string;
membership$: Behavior<CallMembership>;
displayName$: Behavior<string>;
avatarUrl$: Behavior<string | undefined>;
/** Which transport this member is on; undefined when the membership has none. */
transport$: Behavior<TransportMetadata | undefined>;
/**
* Null while the member has a transport but no media has arrived on it yet
* ("waiting for media").
*/
media$: Behavior<MemberMedia | null>;
}
export interface RemoteRtcMember extends RtcMember {
local: false;
}
export interface LocalRtcMember extends RtcMember {
local: true;
media$: Behavior<LocalMemberMedia | null>;
sharingScreen$: Behavior<boolean>;
/** Null when the platform cannot share a screen. */
toggleScreenSharing: (() => void) | null;
screenShareError$: Behavior<Error | null>;
dismissScreenShareError(): void;
}
// ---------------------------------------------------------------------------
// Media
export type MediaSource =
| "microphone"
| "camera"
| "screenShare"
| "screenShareAudio";
export type MediaStreamStats =
| RTCInboundRtpStreamStats
| RTCOutboundRtpStreamStats
| undefined;
/** One published track of a member. */
export interface MediaTrack {
source: MediaSource;
kind: "audio" | "video";
/** Stable for the life of the track. */
id: string;
muted$: Behavior<boolean>;
/** False when the SFU reports the track as unencrypted. */
encrypted$: Behavior<boolean>;
/** Polled while subscribed. */
stats$: Behavior<MediaStreamStats>;
/**
* Rendering. The view hands its `<video>` or `<audio>` element over; the SDK
* sets its stream and, for video, registers the size and on-screen observers
* that pick a simulcast layer and pause the subscription while the element is
* hidden. `attach` is idempotent per element; `detach` has to be called
* before the element leaves the DOM so those observers are released.
*/
attach(element: HTMLMediaElement): void;
detach(element: HTMLMediaElement): void;
}
export interface AudioMediaTrack extends MediaTrack {
kind: "audio";
/** Route playback through Web Audio, for earpiece pan and gain. Undefined resets. */
setAudioContext(ctx: AudioContext | undefined, plugins?: AudioNode[]): void;
setVolume(volume: number): void;
}
export interface VideoMediaTrack extends MediaTrack {
kind: "video";
/** Undefined for remote tracks. */
facingMode$?: Behavior<"user" | "environment" | undefined>;
}
export type EncryptionError = "MissingKey" | "InvalidKey";
/**
* The media of one member, backed by a LiveKit participant inside the SDK.
* There is no identity field: the LiveKit identity is the member's `id`.
*/
export interface MemberMedia {
local: boolean;
speaking$: Behavior<boolean>;
screenShareEnabled$: Behavior<boolean>;
microphone$: Behavior<AudioMediaTrack | undefined>;
camera$: Behavior<VideoMediaTrack | undefined>;
screenShare$: Behavior<VideoMediaTrack | undefined>;
screenShareAudio$: Behavior<AudioMediaTrack | undefined>;
/** Emits when the SFU reports a key problem for this member. */
encryptionError$: Observable<EncryptionError>;
}
export interface LocalMemberMedia extends MemberMedia {
local: true;
}
+160
View File
@@ -0,0 +1,160 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { EventEmitter } from "events";
import {
ParticipantEvent,
type RemoteParticipant,
type Room as LivekitRoom,
Track,
TrackEvent,
type TrackPublication,
} from "livekit-client";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { ObservableScope } from "../../../src/state/ObservableScope";
import { type AudioMediaTrack } from "../api";
import { createLivekitMediaTrack } from "./LivekitMediaTrack";
describe("createLivekitMediaTrack", () => {
let scope: ObservableScope;
beforeEach(() => {
scope = new ObservableScope();
});
afterEach(() => scope.end());
it("describes the publication", () => {
const { participant, publication, room } = fakes();
const track = createLivekitMediaTrack(
scope,
participant,
publication as unknown as TrackPublication,
room,
"camera",
);
expect(track.kind).toBe("video");
expect(track.source).toBe("camera");
expect(track.id).toBe("TR_1");
expect(track.encrypted$.value).toBe(true);
});
it("attaches elements to the track, even one that arrives later", () => {
const { participant, publication, room, emitter } = fakes({
track: undefined,
});
const element = document.createElement("video");
const track = createLivekitMediaTrack(
scope,
participant,
publication as unknown as TrackPublication,
room,
"camera",
);
track.attach(element);
track.attach(element);
const livekitTrack = fakeTrack();
publication.track = livekitTrack;
emitter.emit(TrackEvent.Subscribed, livekitTrack);
expect(livekitTrack.attach).toHaveBeenCalledTimes(1);
expect(livekitTrack.attach).toHaveBeenCalledWith(element);
track.detach(element);
expect(livekitTrack.detach).toHaveBeenCalledWith(element);
});
it("detaches everything when its scope ends", () => {
const { participant, publication, room } = fakes();
const element = document.createElement("video");
const scope = new ObservableScope();
const track = createLivekitMediaTrack(
scope,
participant,
publication as unknown as TrackPublication,
room,
"camera",
);
track.attach(element);
scope.end();
expect(publication.track!.detach).toHaveBeenCalledWith(element);
});
it("follows the mute state", () => {
const { participant, publication, room, emitter } = fakes();
const track = createLivekitMediaTrack(
scope,
participant,
publication as unknown as TrackPublication,
room,
"camera",
);
expect(track.muted$.value).toBe(false);
publication.isMuted = true;
emitter.emit(ParticipantEvent.TrackMuted, publication);
expect(track.muted$.value).toBe(true);
});
it("scales a remote member's volume on the participant", () => {
const { participant, publication, room } = fakes({
kind: Track.Kind.Audio,
});
const track = createLivekitMediaTrack(
scope,
participant,
publication as unknown as TrackPublication,
room,
"microphone",
) as AudioMediaTrack;
track.setVolume(0.5);
expect(participant.setVolume).toHaveBeenCalledWith(
0.5,
Track.Source.Microphone,
);
});
});
type FakeTrack = Track & {
attach: ReturnType<typeof vi.fn>;
detach: ReturnType<typeof vi.fn>;
};
type FakePublication = Omit<TrackPublication, "track" | "isMuted"> & {
track: FakeTrack | undefined;
isMuted: boolean;
};
function fakeTrack(): FakeTrack {
return { attach: vi.fn(), detach: vi.fn() } as unknown as FakeTrack;
}
function fakes({
track = fakeTrack(),
kind = Track.Kind.Video,
}: { track?: FakeTrack | undefined; kind?: Track.Kind } = {}): {
participant: RemoteParticipant & { setVolume: ReturnType<typeof vi.fn> };
publication: FakePublication;
room: LivekitRoom;
emitter: EventEmitter;
} {
const emitter = new EventEmitter();
const participant = Object.assign(emitter, {
isLocal: false,
identity: "@alice:example.org:DEVICE",
setVolume: vi.fn(),
getTrackPublication: (): TrackPublication =>
publication as unknown as TrackPublication,
}) as unknown as RemoteParticipant & { setVolume: ReturnType<typeof vi.fn> };
const publication = Object.assign(emitter, {
kind,
trackSid: "TR_1",
isMuted: false,
isEncrypted: true,
track,
}) as unknown as FakePublication;
const room = new EventEmitter() as unknown as LivekitRoom;
return { participant, publication, room, emitter };
}
+200
View File
@@ -0,0 +1,200 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
observeParticipantMedia,
roomEventSelector,
} from "@livekit/components-core";
import {
facingModeFromLocalTrack,
LocalTrack,
LocalVideoTrack,
type Participant,
RemoteAudioTrack,
type RemoteParticipant,
RemoteTrack,
type Room as LivekitRoom,
RoomEvent,
Track,
TrackEvent,
type TrackPublication,
} from "livekit-client";
import {
distinctUntilChanged,
fromEvent,
interval,
map,
merge,
of,
share,
startWith,
switchMap,
} from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import {
type AudioMediaTrack,
type MediaSource,
type MediaStreamStats,
type MediaTrack,
type VideoMediaTrack,
} from "../api";
import { LazyBehavior } from "../utils/LazyBehavior";
export const livekitSources: Record<MediaSource, Track.Source> = {
microphone: Track.Source.Microphone,
camera: Track.Source.Camera,
screenShare: Track.Source.ScreenShare,
screenShareAudio: Track.Source.ScreenShareAudio,
};
// One timer for every track so that a large session does not keep hundreds of
// them, each firing a statistics request, in the event loop.
const refreshStats$ = interval(1000).pipe(startWith(0), share());
/**
* One publication of a participant as a `MediaTrack`. The publication is
* fixed; the track behind it may come and go (a remote track arrives on
* subscription), and every attached element follows it.
*/
export function createLivekitMediaTrack(
scope: ObservableScope,
participant: Participant,
publication: TrackPublication,
room: LivekitRoom,
source: MediaSource,
): AudioMediaTrack | VideoMediaTrack {
const mediaChanged$ = observeParticipantMedia(participant);
const track$: Behavior<Track | undefined> = scope.behavior(
merge(
mediaChanged$,
fromEvent(publication, TrackEvent.Subscribed),
fromEvent(publication, TrackEvent.Unsubscribed),
).pipe(
map(() => publication.track),
startWith(publication.track),
distinctUntilChanged(),
),
);
const attached = new Set<HTMLMediaElement>();
let audioContext: AudioContext | undefined;
let audioPlugins: AudioNode[] = [];
const applyAudioContext = (track: Track | undefined): void => {
if (!(track instanceof RemoteAudioTrack)) return;
track.setAudioContext(audioContext);
track.setWebAudioPlugins(audioPlugins);
};
let current = track$.value;
track$.pipe(scope.bind()).subscribe((track) => {
for (const element of attached) current?.detach(element);
current = track;
applyAudioContext(track);
for (const element of attached) track?.attach(element);
});
scope.onEnd(() => {
for (const element of attached) current?.detach(element);
attached.clear();
});
const base: MediaTrack = {
source,
kind: publication.kind === Track.Kind.Audio ? "audio" : "video",
id: publication.trackSid,
muted$: scope.behavior(
mediaChanged$.pipe(map(() => publication.isMuted)),
publication.isMuted,
),
encrypted$: scope.behavior(
merge(
mediaChanged$,
roomEventSelector(room, RoomEvent.ParticipantEncryptionStatusChanged),
).pipe(map(() => publication.isEncrypted)),
publication.isEncrypted,
),
stats$: new LazyBehavior<MediaStreamStats>(
refreshStats$.pipe(
switchMap(async () => rtpStreamStats(publication, participant.isLocal)),
scope.bind(),
),
undefined,
),
attach: (element) => {
if (attached.has(element)) return;
attached.add(element);
current?.attach(element);
},
detach: (element) => {
if (!attached.delete(element)) return;
current?.detach(element);
},
};
if (base.kind === "audio")
return {
...base,
kind: "audio",
setAudioContext: (ctx, plugins = []) => {
audioContext = ctx;
audioPlugins = plugins;
applyAudioContext(current);
},
setVolume: (volume) => {
// Our own audio is never played back, so there is nothing to scale
if (participant.isLocal) return;
const remote = participant as RemoteParticipant;
if (source === "microphone")
remote.setVolume(volume, Track.Source.Microphone);
else if (source === "screenShareAudio")
remote.setVolume(volume, Track.Source.ScreenShareAudio);
},
};
return {
...base,
kind: "video",
...(participant.isLocal && { facingMode$: facingMode$(scope, track$) }),
};
}
async function rtpStreamStats(
publication: TrackPublication,
local: boolean,
): Promise<MediaStreamStats> {
const track = publication.track;
if (!(track instanceof RemoteTrack || track instanceof LocalTrack))
return undefined;
const report = await track.getRTCStatsReport();
if (!report) return undefined;
const type = local ? "outbound-rtp" : "inbound-rtp";
for (const stats of report.values()) if (stats.type === type) return stats;
return undefined;
}
function facingMode$(
scope: ObservableScope,
track$: Behavior<Track | undefined>,
): Behavior<"user" | "environment" | undefined> {
return scope.behavior(
track$.pipe(
switchMap((track) => {
if (!(track instanceof LocalVideoTrack)) return of(undefined);
return fromEvent(track, TrackEvent.Restarted).pipe(
startWith(null),
map(() => {
const { facingMode } = facingModeFromLocalTrack(track);
return facingMode === "user" || facingMode === "environment"
? facingMode
: undefined;
}),
);
}),
),
);
}
+148
View File
@@ -0,0 +1,148 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
observeParticipantEvents,
observeParticipantMedia,
roomEventSelector,
} from "@livekit/components-core";
import {
type LocalParticipant,
type Participant,
ParticipantEvent,
type Room as LivekitRoom,
RoomEvent,
type TrackPublication,
} from "livekit-client";
import { distinctUntilChanged, filter, map, type Observable } from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import { E2eeType } from "../../../src/e2ee/e2eeType";
import { type EncryptionSystem } from "../../../src/e2ee/sharedKeyManagement";
import {
type AudioMediaTrack,
type EncryptionError,
type LocalMemberMedia,
type MediaSource,
type MemberMedia,
type VideoMediaTrack,
} from "../api";
import { mapScoped } from "../utils/mapScoped";
import { createLivekitMediaTrack, livekitSources } from "./LivekitMediaTrack";
/**
* A participant as `MemberMedia`. The only place, with `LivekitMediaTrack`,
* that reads a LiveKit participant; everything above sees tracks and
* behaviors.
*/
export function createLivekitMemberMedia(
scope: ObservableScope,
participant: Participant,
room: LivekitRoom,
encryptionSystem: EncryptionSystem,
): MemberMedia {
const mediaChanged$ = observeParticipantMedia(participant);
const track$ = <T extends AudioMediaTrack | VideoMediaTrack>(
source: MediaSource,
): Behavior<T | undefined> =>
memberTrack$(scope, participant, room, source, mediaChanged$) as Behavior<
T | undefined
>;
return {
local: participant.isLocal,
speaking$: scope.behavior(
observeParticipantEvents(
participant,
ParticipantEvent.IsSpeakingChanged,
).pipe(map((p) => p.isSpeaking)),
participant.isSpeaking,
),
screenShareEnabled$: scope.behavior(
mediaChanged$.pipe(map((media) => media.isScreenShareEnabled)),
participant.isScreenShareEnabled,
),
microphone$: track$<AudioMediaTrack>("microphone"),
camera$: track$<VideoMediaTrack>("camera"),
screenShare$: track$<VideoMediaTrack>("screenShare"),
screenShareAudio$: track$<AudioMediaTrack>("screenShareAudio"),
encryptionError$: encryptionErrors$(
scope,
participant,
room,
encryptionSystem,
),
};
}
export function createLocalLivekitMemberMedia(
scope: ObservableScope,
participant: LocalParticipant,
room: LivekitRoom,
encryptionSystem: EncryptionSystem,
): LocalMemberMedia {
return {
...createLivekitMemberMedia(scope, participant, room, encryptionSystem),
local: true,
};
}
/** The member's track for a source, as long as the same publication is behind it. */
function memberTrack$(
scope: ObservableScope,
participant: Participant,
room: LivekitRoom,
source: MediaSource,
mediaChanged$: Observable<unknown>,
): Behavior<AudioMediaTrack | VideoMediaTrack | undefined> {
const publication = (): TrackPublication | undefined =>
participant.getTrackPublication(livekitSources[source]);
return mapScoped(
scope,
scope.behavior(
mediaChanged$.pipe(map(publication), distinctUntilChanged()),
publication(),
),
(trackScope, publication) =>
createLivekitMediaTrack(
trackScope,
participant,
publication,
room,
source,
),
);
}
function encryptionErrors$(
scope: ObservableScope,
participant: Participant,
room: LivekitRoom,
encryptionSystem: EncryptionSystem,
): Observable<EncryptionError> {
return roomEventSelector(room, RoomEvent.EncryptionError).pipe(
map(([error]) => error?.message ?? ""),
// The participant the error is about does not survive the trip from the
// worker, so the identity is matched in the message. A shared key is the
// same for everyone, so its errors are too.
filter(
(message) =>
encryptionSystem.kind === E2eeType.SHARED_KEY ||
message.includes(participant.identity),
),
map((message): EncryptionError | undefined =>
message.includes("MissingKey")
? "MissingKey"
: message.includes("InvalidKey")
? "InvalidKey"
: undefined,
),
filter((error) => error !== undefined),
scope.bind(),
);
}
+133
View File
@@ -0,0 +1,133 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
type BaseKeyProvider,
Room as LivekitRoom,
type RoomOptions,
} from "livekit-client";
// Inline so that the worker also loads when the SDK is served from another
// origin than the page
import E2EEWorker from "livekit-client/e2ee-worker?worker&inline";
import { type Logger } from "matrix-js-sdk/lib/logger";
import { type LivekitTransport } from "matrix-js-sdk/lib/matrixrtc";
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
import { BehaviorSubject, combineLatest, map } from "rxjs";
import { buildLiveKitOptions } from "../../../src/livekit/options";
import {
type OpenIDClientParts,
type SFUConfig,
} from "../../../src/livekit/openIDSFU";
import { type Behavior } from "../../../src/state/Behavior";
import {
Connection,
type ConnectionOpts,
ConnectionState,
} from "../../../src/state/CallViewModel/remoteMembers/Connection";
import { type ConnectionFactory } from "../../../src/state/CallViewModel/remoteMembers/ConnectionFactory";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import { type LocalMediaInputs, type ResolvedTransport } from "../api";
/**
* A connection that keeps what it fetched from the JWT service, so that the
* transport can report it as `resolved$`. Nothing while the config is still
* being fetched, and nothing again once the connection has stopped or failed.
*/
export class ResolvedConnection extends Connection {
private readonly sfuConfig$: BehaviorSubject<SFUConfig | undefined>;
public readonly resolved$: Behavior<ResolvedTransport | undefined>;
public constructor(opts: ConnectionOpts, logger: Logger) {
super(opts, logger);
this.sfuConfig$ = new BehaviorSubject(opts.existingSFUConfig);
this.resolved$ = opts.scope.behavior(
combineLatest([this.sfuConfig$, this.state$]).pipe(
map(([config, state]) =>
config === undefined ||
state === ConnectionState.Stopped ||
state instanceof Error
? undefined
: resolvedTransport(config),
),
),
);
}
protected override async getSFUConfigForRemoteConnection(): Promise<SFUConfig> {
const config = await super.getSFUConfigForRemoteConnection();
this.sfuConfig$.next(config);
return config;
}
}
function resolvedTransport(config: SFUConfig): ResolvedTransport {
return {
type: "livekit",
url: config.url,
token: config.jwt,
roomAlias: config.livekitAlias,
identity: config.livekitIdentity,
};
}
/**
* Creates a `ResolvedConnection` per transport, each with a LiveKit room of
* its own that captures from the devices the host selected.
*/
export class LivekitConnectionFactory implements ConnectionFactory {
public constructor(
private readonly client: OpenIDClientParts,
private readonly roomId: string,
private readonly localMedia: LocalMediaInputs,
private readonly keyProvider: BaseKeyProvider | undefined,
) {}
public createConnection(
scope: ObservableScope,
transport: LivekitTransport,
ownMembershipIdentity: CallMembershipIdentityParts,
logger: Logger,
sfuConfig?: SFUConfig,
): Connection {
return new ResolvedConnection(
{
existingSFUConfig: sfuConfig,
roomId: this.roomId,
transport,
client: this.client,
scope,
livekitRoomFactory: () =>
new LivekitRoom(roomOptions(this.localMedia, this.keyProvider)),
ownMembershipIdentity,
},
logger,
);
}
}
function roomOptions(
localMedia: LocalMediaInputs,
keyProvider: BaseKeyProvider | undefined,
): RoomOptions {
const base = buildLiveKitOptions();
return {
...base,
videoCaptureDefaults: {
...base.videoCaptureDefaults,
deviceId: localMedia.videoInputDeviceId$.value,
processor: localMedia.videoProcessor$.value,
},
audioCaptureDefaults: {
...base.audioCaptureDefaults,
deviceId: localMedia.audioInputDeviceId$.value,
},
// Every room needs a worker of its own: one gets confused by streams from
// several rooms
e2ee: keyProvider && { keyProvider, worker: new E2EEWorker() },
};
}
+38
View File
@@ -0,0 +1,38 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { type BaseKeyProvider, ExternalE2EEKeyProvider } from "livekit-client";
import { type Logger } from "matrix-js-sdk/lib/logger";
import { type MatrixRTCSession as JsSdkRtcSession } from "matrix-js-sdk/lib/matrixrtc";
import { E2eeType } from "../../../src/e2ee/e2eeType";
import { MatrixKeyProvider } from "../../../src/e2ee/matrixKeyProvider";
import { type EncryptionSystem } from "../../../src/e2ee/sharedKeyManagement";
/** The LiveKit key provider for the encryption the host asked for, or none. */
export function createKeyProvider(
encryptionSystem: EncryptionSystem,
session: JsSdkRtcSession,
logger: Logger,
): BaseKeyProvider | undefined {
switch (encryptionSystem.kind) {
case E2eeType.NONE:
return undefined;
case E2eeType.PER_PARTICIPANT: {
const keyProvider = new MatrixKeyProvider();
keyProvider.setRTCSession(session);
return keyProvider;
}
case E2eeType.SHARED_KEY: {
const keyProvider = new ExternalE2EEKeyProvider();
keyProvider
.setKey(encryptionSystem.secret)
.catch((e) => logger.error("Failed to set the shared E2EE key", e));
return keyProvider;
}
}
}
+359
View File
@@ -0,0 +1,359 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
type LocalParticipant,
type ScreenShareCaptureOptions,
} from "livekit-client";
import { type Logger } from "matrix-js-sdk/lib/logger";
import {
type LivekitTransport,
type MatrixRTCSession as JsSdkRtcSession,
type Status as RTCSessionStatus,
} from "matrix-js-sdk/lib/matrixrtc";
import { deepCompare } from "matrix-js-sdk/lib/utils";
import {
BehaviorSubject,
catchError,
combineLatest,
concat,
distinctUntilChanged,
map,
NEVER,
type Observable,
of,
pairwise,
race,
Subject,
switchMap,
} from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import { type HomeserverConnected } from "../../../src/state/CallViewModel/localMember/HomeserverConnected";
import { observeSharingScreen$ } from "../../../src/state/CallViewModel/localMember/LocalMember";
import {
type Connection,
ConnectionState,
} from "../../../src/state/CallViewModel/remoteMembers/Connection";
import { type IConnectionManager } from "../../../src/state/CallViewModel/remoteMembers/ConnectionManager";
import { type RtcSessionError } from "../api";
import { toRtcSessionError } from "../utils/errors";
import { type LocalTransport } from "./LocalTransport";
import { type Publisher } from "./Publisher";
export enum TransportState {
Waiting = "transport_waiting",
}
export enum PublishState {
WaitingForUser = "publish_waiting_for_user",
Publishing = "publish_publishing",
}
export type LocalMemberMediaState =
| { connection: ConnectionState | RtcSessionError }
| PublishState
| RtcSessionError;
export type LocalMemberState =
| RtcSessionError
| TransportState.Waiting
| {
media: LocalMemberMediaState;
matrix: RtcSessionError | RTCSessionStatus;
};
interface Props {
scope: ObservableScope;
connectionManager: IConnectionManager;
localTransport$: Observable<LocalTransport>;
homeserverConnected: HomeserverConnected;
createPublisher: (connection: Connection) => Publisher;
joinMatrixRTC: (transport: LivekitTransport) => void;
/** The membership manager giving up on keeping the membership alive. */
membershipManagerError$: Observable<unknown>;
matrixRTCSession: Pick<
JsSdkRtcSession,
"updateCallIntent" | "leaveRoomSession"
>;
cameraEnabled$: Behavior<boolean>;
logger: Logger;
}
export interface LocalMembership {
requestJoinAndPublish: () => void;
requestDisconnect: () => void;
joinRequested$: Behavior<boolean>;
state$: Behavior<LocalMemberState>;
participant$: Behavior<LocalParticipant | null>;
connection$: Behavior<Connection | null>;
/** Fully connected: to the homeserver, the session and the transport. */
connected$: Behavior<boolean>;
/** Connected once, and currently not. */
reconnecting$: Behavior<boolean>;
sharingScreen$: Behavior<boolean>;
toggleScreenSharing: (() => void) | null;
screenShareError$: Behavior<Error | null>;
dismissScreenShareError: () => void;
}
/**
* The local member's state machine: waits for the transport, publishes on
* its connection once asked to join, and enters and leaves the MatrixRTC
* session in step.
*/
export function createLocalMembership$({
scope,
connectionManager,
localTransport$: localTransportWithErrors$,
homeserverConnected,
createPublisher,
joinMatrixRTC,
membershipManagerError$,
matrixRTCSession,
cameraEnabled$,
logger: parentLogger,
}: Props): LocalMembership {
const logger = parentLogger.getChild("[LocalMember]");
const fatalTransportError$ = new Subject<RtcSessionError>();
const localTransport$ = localTransportWithErrors$.pipe(
catchError((e: unknown) => {
fatalTransportError$.next(toRtcSessionError(e));
return NEVER;
}),
);
const transport$ = scope.behavior<LocalTransport | null>(
localTransport$,
null,
);
const connection$ = scope.behavior(
combineLatest([
connectionManager.connectionManagerData$,
localTransport$,
]).pipe(
map(([{ value: connections }, { transport }]) =>
connections.getConnectionForTransport(transport),
),
),
null,
);
const joinRequested$ = new BehaviorSubject(false);
const publisher$ = new BehaviorSubject<Publisher | null>(null);
const publishError$ = new BehaviorSubject<RtcSessionError | null>(null);
const matrixError$ = new BehaviorSubject<RtcSessionError | null>(null);
scope.reconcile(connection$, async (connection) => {
if (connection === null) return;
const publisher = createPublisher(connection);
publisher$.next(publisher);
return Promise.resolve(async (): Promise<void> => {
publisher$.next(null);
await publisher.destroy();
});
});
scope.reconcile(
scope.behavior(combineLatest([publisher$, joinRequested$])),
async ([publisher, shouldPublish]) => {
if (publisher === null) return;
try {
if (shouldPublish) {
publisher.createAndSetupTracks();
await publisher.startPublishing();
} else if (publisher.shouldPublish) await publisher.stopPublishing();
} catch (e) {
if (publishError$.value === null)
publishError$.next(toRtcSessionError(e));
else logger.error("Another publish error", e);
}
},
);
scope.reconcile(
scope.behavior(combineLatest([transport$, joinRequested$])),
async ([transport, shouldJoin]) => {
if (transport === null || !shouldJoin) return;
try {
joinMatrixRTC(transport.transport);
} catch (e) {
logger.error("Failed to enter the session", e);
if (matrixError$.value === null)
matrixError$.next(toRtcSessionError(e));
}
return Promise.resolve(async (): Promise<void> => {
try {
await matrixRTCSession.leaveRoomSession(1000);
} catch (e) {
logger.error("Failed to leave the session", e);
}
});
},
);
membershipManagerError$.pipe(scope.bind()).subscribe((e) => {
logger.error("The membership manager stopped", e);
if (matrixError$.value === null) matrixError$.next(toRtcSessionError(e));
});
cameraEnabled$.pipe(scope.bind()).subscribe((videoEnabled) => {
matrixRTCSession
.updateCallIntent(videoEnabled ? "video" : "audio")
.catch((e) => {
// Expected before the join: the intent is sent with it instead
if (e instanceof Error && e.message === "Not connected yet") return;
logger.error("Failed to update the call intent", e);
});
});
const connectionState$ = connection$.pipe(
switchMap((connection) => connection?.state$ ?? of(null)),
);
const mediaState$ = scope.behavior<LocalMemberMediaState>(
combineLatest([connectionState$, joinRequested$]).pipe(
map(([connectionState, shouldPublish]) => {
if (connectionState !== ConnectionState.LivekitConnected)
return {
connection:
connectionState instanceof Error
? toRtcSessionError(connectionState)
: (connectionState ?? ConnectionState.Initialized),
};
return shouldPublish
? PublishState.Publishing
: PublishState.WaitingForUser;
}),
distinctUntilChanged(deepCompare),
),
);
const state$ = scope.behavior<LocalMemberState>(
concat(
of(TransportState.Waiting),
race(
fatalTransportError$,
localTransport$.pipe(
switchMap(() =>
combineLatest(
[
mediaState$,
homeserverConnected.rtsSession$,
matrixError$,
publishError$,
],
(media, sessionStatus, matrixError, publishError) => ({
matrix: matrixError ?? sessionStatus,
media: publishError ?? media,
}),
),
),
),
),
),
);
const connected$ = scope.behavior(
combineLatest(
[homeserverConnected.combined$, connectionState$],
([homeserverConnected], connectionState) =>
homeserverConnected &&
connectionState === ConnectionState.LivekitConnected,
),
);
const reconnecting$ = scope.behavior(
connected$.pipe(
pairwise(),
map(([was, is]) => was && !is),
),
false,
);
const participant$ = scope.behavior(
connection$.pipe(map((c) => c?.livekitRoom.localParticipant ?? null)),
);
// Nothing leaves this device while it may already have been dropped from
// the session: the member would show as away while still being heard
combineLatest([participant$, homeserverConnected.combined$])
.pipe(scope.bind())
.subscribe(([participant, [connected]]) => {
if (participant === null) return;
for (const { track } of participant.trackPublications.values()) {
if (!track) continue;
if (connected && track.isUpstreamPaused)
track.resumeUpstream().catch((e) => {
logger.error(`Failed to resume the ${track.kind} track`, e);
});
else if (!connected && !track.isUpstreamPaused)
track.pauseUpstream().catch((e) => {
logger.error(`Failed to pause the ${track.kind} track`, e);
});
}
});
const sharingScreen$ = scope.behavior(
participant$.pipe(
switchMap((p) => (p === null ? of(false) : observeSharingScreen$(p))),
),
);
const screenShareError$ = new BehaviorSubject<Error | null>(null);
const toggleScreenSharing =
"getDisplayMedia" in (navigator.mediaDevices ?? {})
? (): void => {
const participant = participant$.value;
if (participant === null) return;
const enable = !sharingScreen$.value;
participant
.setScreenShareEnabled(enable, screenShareCaptureOptions)
.catch((e: unknown) => {
logger.error(
`Screen share ${enable ? "start" : "stop"} failed`,
e,
);
// The user closing the picker is not an error worth showing
if (e instanceof DOMException && e.name === "NotAllowedError")
return;
screenShareError$.next(
e instanceof Error ? e : new Error(String(e)),
);
});
}
: null;
return {
requestJoinAndPublish: () => joinRequested$.next(true),
requestDisconnect: () => joinRequested$.next(false),
joinRequested$,
state$,
participant$,
connection$,
connected$,
reconnecting$,
sharingScreen$,
toggleScreenSharing,
screenShareError$,
dismissScreenShareError: () => screenShareError$.next(null),
};
}
const screenShareCaptureOptions: ScreenShareCaptureOptions = {
// No echo cancellation: it would cancel the other members' voices out of
// the shared audio
audio: {
autoGainControl: false,
noiseSuppression: false,
voiceIsolation: false,
},
selfBrowserSurface: "include",
surfaceSwitching: "include",
systemAudio: "include",
};
+66
View File
@@ -0,0 +1,66 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { type MatrixClient } from "matrix-js-sdk";
import { type Logger } from "matrix-js-sdk/lib/logger";
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
import {
DEFAULT_CONFIG,
type MatrixRTCMode,
} from "../../../src/config/ConfigOptions";
import { getSFUConfigWithOpenID } from "../../../src/livekit/openIDSFU";
import { type LocalTransport } from "../../../src/state/CallViewModel/localMember/LocalTransport";
import { RtcTransportAutoDiscovery } from "../../../src/state/CallViewModel/localMember/RtcTransportAutoDiscovery";
import { MatrixRTCTransportMissingError } from "../../../src/utils/errors";
export type { LocalTransport };
interface Props {
client: Pick<
MatrixClient,
| "getDomain"
| "_unstable_getRTCTransports"
| "getOpenIdToken"
| "getDeviceId"
>;
ownMembershipIdentity: CallMembershipIdentityParts;
roomId: string;
matrixRTCMode: MatrixRTCMode;
logger: Logger;
}
/**
* The transport the local member publishes on: the homeserver's preferred
* one, authenticated with so that the session can be joined with a token in
* hand. Only the homeserver is asked; the SDK has no configuration of its own
* to fall back on.
*/
export async function getLocalTransport({
client,
ownMembershipIdentity,
roomId,
matrixRTCMode,
logger,
}: Props): Promise<LocalTransport> {
const transport = await new RtcTransportAutoDiscovery({
client,
resolvedConfig: DEFAULT_CONFIG,
logger,
}).discoverPreferredTransport();
if (transport === null)
throw new MatrixRTCTransportMissingError(client.getDomain() ?? "");
const sfuConfig = await getSFUConfigWithOpenID(
client,
ownMembershipIdentity,
transport.livekit_service_url,
roomId,
{ matrixRTCMode },
logger,
);
return { transport, sfuConfig };
}
+162
View File
@@ -0,0 +1,162 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
type LocalParticipant,
type Participant,
type Room as LivekitRoom,
} from "livekit-client";
import { type CallMembership } from "matrix-js-sdk/lib/matrixrtc";
import { combineLatest, distinctUntilChanged, map } from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import { type Connection } from "../../../src/state/CallViewModel/remoteMembers/Connection";
import { type createMatrixMemberMetadata$ } from "../../../src/state/CallViewModel/remoteMembers/MatrixMemberMetadata";
import { type RemoteMatrixLivekitMember } from "../../../src/state/CallViewModel/remoteMembers/MatrixLivekitMembers";
import { type EncryptionSystem } from "../../../src/e2ee/sharedKeyManagement";
import {
type LocalRtcMember,
type MemberMedia,
type RemoteRtcMember,
type RtcMember,
} from "../api";
import {
createLivekitMemberMedia,
createLocalLivekitMemberMedia,
} from "../media/LivekitMemberMedia";
import { mapScoped } from "../utils/mapScoped";
import { type LocalMembership } from "./LocalMember";
import { type TransportRegistry } from "./Transports";
export interface MemberContext {
metadata: ReturnType<typeof createMatrixMemberMetadata$>;
transports: TransportRegistry;
encryptionSystem: EncryptionSystem;
}
/** Everything that tells one membership from another, for keying items. */
export function membershipKeys(
membership: CallMembership,
): [string, string, string, string] {
return [
membership.userId,
membership.deviceId,
membership.memberId,
membership.rtcBackendIdentity,
];
}
export function createRemoteRtcMember(
scope: ObservableScope,
member: RemoteMatrixLivekitMember,
context: MemberContext,
): RemoteRtcMember {
return {
...createRtcMember(scope, member.membership$, context),
local: false,
media$: mediaFor(
scope,
member.participant.value$,
member.connection$,
(mediaScope, participant, room) =>
createLivekitMemberMedia(
mediaScope,
participant,
room,
context.encryptionSystem,
),
),
};
}
export function createLocalRtcMember(
scope: ObservableScope,
membership$: Behavior<CallMembership>,
localMembership: LocalMembership,
context: MemberContext,
): LocalRtcMember {
return {
...createRtcMember(scope, membership$, context),
local: true,
media$: mediaFor(
scope,
localMembership.participant$,
localMembership.connection$,
(mediaScope, participant, room) =>
createLocalLivekitMemberMedia(
mediaScope,
participant,
room,
context.encryptionSystem,
),
),
sharingScreen$: localMembership.sharingScreen$,
toggleScreenSharing: localMembership.toggleScreenSharing,
screenShareError$: localMembership.screenShareError$,
dismissScreenShareError: localMembership.dismissScreenShareError,
};
}
function createRtcMember(
scope: ObservableScope,
membership$: Behavior<CallMembership>,
{ metadata, transports }: MemberContext,
): Omit<RtcMember, "local" | "media$"> {
const { userId, deviceId, rtcBackendIdentity } = membership$.value;
return {
id: rtcBackendIdentity,
userId,
deviceId,
membership$,
displayName$: scope.behavior(
metadata
.createDisplayNameBehavior$(scope, userId)
.pipe(map((name) => name ?? userId)),
),
avatarUrl$: metadata.createAvatarUrlBehavior$(scope, userId),
transport$: scope.behavior(
membership$.pipe(
map((membership) => membership.getTransport()),
distinctUntilChanged(),
map((transport) => transport && transports.get(transport)),
),
),
};
}
/**
* The member's media once its participant is known, in a scope that ends when
* the participant or the connection changes.
*/
function mediaFor<
P extends Participant | LocalParticipant,
M extends MemberMedia,
>(
scope: ObservableScope,
participant$: Behavior<P | null>,
connection$: Behavior<Connection | null>,
factory: (scope: ObservableScope, participant: P, room: LivekitRoom) => M,
): Behavior<M | null> {
const source$ = scope.behavior(
combineLatest([participant$, connection$]).pipe(
map(([participant, connection]) =>
participant && connection
? { participant, room: connection.livekitRoom }
: null,
),
distinctUntilChanged(
(a, b) => a?.participant === b?.participant && a?.room === b?.room,
),
),
);
return scope.behavior(
mapScoped(scope, source$, (mediaScope, { participant, room }) =>
factory(mediaScope, participant, room),
).pipe(map((media) => media ?? null)),
);
}
+226
View File
@@ -0,0 +1,226 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { observeParticipantMedia } from "@livekit/components-core";
import {
ConnectionState as LivekitConnectionState,
type LocalTrackPublication,
LocalVideoTrack,
ParticipantEvent,
type Room as LivekitRoom,
Track,
} from "livekit-client";
import { type Logger } from "matrix-js-sdk/lib/logger";
import { combineLatest, distinctUntilChanged, map, skip } from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { ObservableScope } from "../../../src/state/ObservableScope";
import { type LocalMediaInputs } from "../api";
/**
* Publishes the local media on one LiveKit room, following `LocalMediaInputs`.
*
* LiveKit publishes a track the moment it is created, but a member must not be
* heard before it has joined the MatrixRTC session, nor after it has left. The
* tracks are therefore kept, and their upstream paused, while `shouldPublish`
* is false: the local preview stays live while nothing reaches the room.
*/
export class Publisher {
public shouldPublish = false;
private tracksRequested = false;
private readonly scope = new ObservableScope();
private readonly room: LivekitRoom;
public constructor(
room: LivekitRoom,
private readonly inputs: LocalMediaInputs,
private readonly logger: Logger,
) {
this.room = room;
room.setE2EEEnabled(room.options.e2ee !== undefined)?.catch((e: Error) => {
this.logger.error("Failed to enable E2EE on the room", e);
});
this.followVideoProcessor();
this.followDevices();
this.onLocalTrackPublished = this.onLocalTrackPublished.bind(this);
room.localParticipant.on(
ParticipantEvent.LocalTrackPublished,
this.onLocalTrackPublished,
);
}
public async destroy(): Promise<void> {
this.scope.end();
this.room.localParticipant.off(
ParticipantEvent.LocalTrackPublished,
this.onLocalTrackPublished,
);
try {
await this.stopTracks();
} catch (e) {
this.logger.error("Failed to stop the local tracks", e);
}
}
/**
* Creates the microphone and camera tracks the inputs ask for, and keeps
* them in step with the inputs from then on. Both are enabled in one call so
* that the browser asks for permission once. Safe to call more than once.
*/
public createAndSetupTracks(): void {
if (this.tracksRequested) return;
this.tracksRequested = true;
const participant = this.room.localParticipant;
const audio = this.inputs.microphoneEnabled$.value;
const video = this.inputs.cameraEnabled$.value;
// LiveKit resolves these once the track is published, which may block on
// the connection; LocalTrackPublished is what tells us a track exists.
if (audio && video) void participant.enableCameraAndMicrophone();
else if (audio) void participant.setMicrophoneEnabled(true);
else if (video) void participant.setCameraEnabled(true);
this.follow(this.inputs.microphoneEnabled$, Track.Source.Microphone);
this.follow(this.inputs.cameraEnabled$, Track.Source.Camera);
}
public async startPublishing(): Promise<void> {
if (this.shouldPublish) return;
this.shouldPublish = true;
// Enabling a track does not resume an upstream that was paused while it
// was already enabled, so it is done explicitly
await this.resumeUpstreams([Track.Source.Microphone, Track.Source.Camera]);
}
public async stopPublishing(): Promise<void> {
this.shouldPublish = false;
await this.pauseUpstreams([
Track.Source.Microphone,
Track.Source.Camera,
Track.Source.ScreenShare,
]);
}
private async stopTracks(): Promise<void> {
const participant = this.room.localParticipant;
for (const source of [
Track.Source.Microphone,
Track.Source.Camera,
Track.Source.ScreenShare,
]) {
const track = participant.getTrackPublication(source)?.track;
if (track) await participant.unpublishTrack(track, true);
}
}
private onLocalTrackPublished(publication: LocalTrackPublication): void {
this.logger.info(`Local ${publication.source} track published`);
if (!this.shouldPublish)
this.pauseUpstreams([publication.source]).catch((e) => {
this.logger.error("Failed to pause the upstream", e);
});
// The input may have changed while the track was being created
const enabled =
publication.source === Track.Source.Microphone
? this.inputs.microphoneEnabled$.value
: publication.source === Track.Source.Camera
? this.inputs.cameraEnabled$.value
: undefined;
if (enabled === false) this.setEnabled(publication.source, false);
}
private follow(
enabled$: Behavior<boolean>,
source: Track.Source.Microphone | Track.Source.Camera,
): void {
enabled$
.pipe(skip(1), distinctUntilChanged(), this.scope.bind())
.subscribe((enabled) => this.setEnabled(source, enabled));
}
private setEnabled(source: Track.Source, enabled: boolean): void {
const participant = this.room.localParticipant;
const toggle =
source === Track.Source.Microphone
? participant.setMicrophoneEnabled(enabled)
: participant.setCameraEnabled(enabled);
toggle
.then(async () => {
// Unmuting restarts the upstream; until the member has joined, it
// has to stay paused
if (enabled && !this.shouldPublish) await this.pauseUpstreams([source]);
})
.catch((e) => {
this.logger.error(`Failed to set ${source} enabled=${enabled}`, e);
});
}
private async pauseUpstreams(sources: Track.Source[]): Promise<void> {
for (const source of sources) {
const track =
this.room.localParticipant.getTrackPublication(source)?.track;
if (track && !track.isUpstreamPaused) await track.pauseUpstream();
}
}
private async resumeUpstreams(sources: Track.Source[]): Promise<void> {
for (const source of sources) {
const track =
this.room.localParticipant.getTrackPublication(source)?.track;
if (track?.isUpstreamPaused) await track.resumeUpstream();
}
}
private followDevices(): void {
const sync = (
kind: MediaDeviceKind,
deviceId$: Behavior<string | undefined>,
): void => {
deviceId$.pipe(this.scope.bind()).subscribe((deviceId) => {
if (
deviceId === undefined ||
this.room.state !== LivekitConnectionState.Connected ||
this.room.getActiveDevice(kind) === deviceId
)
return;
this.room
.switchActiveDevice(kind, deviceId)
.catch((e) => this.logger.error(`Failed to switch ${kind}`, e));
});
};
sync("audioinput", this.inputs.audioInputDeviceId$);
sync("videoinput", this.inputs.videoInputDeviceId$);
}
private followVideoProcessor(): void {
const participant = this.room.localParticipant;
const cameraTrack$ = observeParticipantMedia(participant).pipe(
map(() => {
const track = participant.getTrackPublication(
Track.Source.Camera,
)?.track;
return track instanceof LocalVideoTrack ? track : undefined;
}),
distinctUntilChanged(),
);
combineLatest([cameraTrack$, this.inputs.videoProcessor$])
.pipe(this.scope.bind())
.subscribe(([track, processor]) => {
if (!track) return;
if (processor && !track.getProcessor()) {
// A processor cannot be built on a track that has already ended
if (track.mediaStreamTrack.readyState === "ended") return;
track.setProcessor(processor).catch((e) => {
this.logger.warn("Failed to attach the video processor", e);
});
} else if (!processor && track.getProcessor()) {
track.stopProcessor().catch((e) => {
this.logger.warn("Failed to stop the video processor", e);
});
}
});
}
}
+235
View File
@@ -0,0 +1,235 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { type MatrixClient, type Room } from "matrix-js-sdk";
import { logger as rootLogger } from "matrix-js-sdk/lib/logger";
import { MatrixRTCSessionEvent } from "matrix-js-sdk/lib/matrixrtc";
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
import { v4 as uuidv4 } from "uuid";
import { combineLatest, from, fromEvent, map } from "rxjs";
import { MatrixRTCMode } from "../../../src/config/ConfigOptions";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import {
createKeyRotationSuppressed$,
createMemberships$,
membershipsAndTransports$,
} from "../../../src/state/SessionBehaviors";
import { createHomeserverConnected$ } from "../../../src/state/CallViewModel/localMember/HomeserverConnected";
import { createConnectionManager$ } from "../../../src/state/CallViewModel/remoteMembers/ConnectionManager";
import { createRemoteMatrixLivekitMembers$ } from "../../../src/state/CallViewModel/remoteMembers/MatrixLivekitMembers";
import {
createMatrixMemberMetadata$,
createRoomMembers$,
} from "../../../src/state/CallViewModel/remoteMembers/MatrixMemberMetadata";
import { filterBehavior, generateItems } from "../../../src/utils/observable";
import {
type LocalMediaInputs,
type LocalRtcMember,
type RemoteRtcMember,
type RtcSession,
RtcSessionError,
type RtcSessionOptions,
} from "../api";
import { LivekitConnectionFactory } from "./ConnectionFactory";
import { joinRtcSession, sessionTimings } from "./joinRtcSession";
import { createKeyProvider } from "./KeyProvider";
import { createLocalMembership$ } from "./LocalMember";
import { getLocalTransport } from "./LocalTransport";
import {
createLocalRtcMember,
createRemoteRtcMember,
membershipKeys,
} from "./Members";
import { Publisher } from "./Publisher";
import { mapScoped } from "../utils/mapScoped";
import { fatalError, sessionStatus } from "./status";
import { createTransportRegistry } from "./Transports";
/**
* Takes the whole `MatrixClient` rather than a slice of it, and finds the
* MatrixRTC session itself. The js-sdk is the MatrixRTC implementation today;
* when the rust-rtc crate replaces it, the SDK has to bridge the crate to the
* js-sdk client, and only the SDK knows what that bridge needs. Holding the
* client keeps that change inside the SDK.
*/
export function createRtcSession(
scope: ObservableScope,
client: MatrixClient,
room: Room,
localMedia: LocalMediaInputs,
options: RtcSessionOptions,
): RtcSession {
const logger = rootLogger.getChild("[RtcSession]");
const userId = client.getUserId();
const deviceId = client.getDeviceId();
if (!(userId && deviceId))
throw new RtcSessionError("The client has to be logged in");
const { encryptionSystem, matrixRTCMode } = options;
const jsSdkSession = client.matrixRTC.getRoomSession(room);
const keyProvider = createKeyProvider(encryptionSystem, jsSdkSession, logger);
const ownMembershipIdentity: CallMembershipIdentityParts = {
userId,
deviceId,
// A pre-sticky membership names itself `${userId}:${deviceId}`, and the
// key transport stamps this id into every key event, so a uuid there
// would name a member no peer can resolve
memberId:
matrixRTCMode === MatrixRTCMode.Matrix_2_0
? uuidv4()
: `${userId}:${deviceId}`,
};
const memberships$ = createMemberships$(scope, jsSdkSession);
const { membershipsWithTransport$, transports$ } = membershipsAndTransports$(
scope,
memberships$,
);
const localTransport$ = from(
getLocalTransport({
client,
ownMembershipIdentity,
roomId: room.roomId,
matrixRTCMode,
logger,
}),
);
const connectionManager = createConnectionManager$({
scope,
connectionFactory: new LivekitConnectionFactory(
client,
room.roomId,
localMedia,
keyProvider,
),
localTransport$,
remoteTransports$: transports$,
logger,
ownMembershipIdentity,
});
const transports = createTransportRegistry(scope, connectionManager);
const localMembership = createLocalMembership$({
scope,
connectionManager,
localTransport$,
homeserverConnected: createHomeserverConnected$(
scope,
client,
jsSdkSession,
sessionTimings.syncDisconnectGracePeriodMs,
),
createPublisher: (connection) =>
new Publisher(
connection.livekitRoom,
localMedia,
logger.getChild(
`[Publisher ${connection.transport.livekit_service_url}]`,
),
),
joinMatrixRTC: (transport) =>
joinRtcSession(jsSdkSession, ownMembershipIdentity, transport, {
encryptMedia: keyProvider !== undefined,
matrixRTCMode,
sendNotificationType: options.sendNotificationType,
callIntent: options.callIntent,
}),
membershipManagerError$: fromEvent(
jsSdkSession,
MatrixRTCSessionEvent.MembershipManagerError,
),
matrixRTCSession: jsSdkSession,
cameraEnabled$: localMedia.cameraEnabled$,
logger,
});
const context = {
metadata: createMatrixMemberMetadata$(
scope,
scope.behavior(memberships$.pipe(map(({ value }) => value))),
createRoomMembers$(scope, room),
),
transports,
encryptionSystem,
};
const remoteMembers$ = scope.behavior<RemoteRtcMember[]>(
createRemoteMatrixLivekitMembers$({
scope,
membershipsWithTransport$,
connectionManager,
localUser: { userId, deviceId },
}).pipe(
map(({ value }) => value),
generateItems(
"RtcSession remoteMembers",
function* (members) {
for (const member of members)
yield {
keys: membershipKeys(member.membership$.value),
data: member,
};
},
(memberScope, member$) =>
createRemoteRtcMember(memberScope, member$.value, context),
),
),
);
const localMembership$ = scope.behavior(
memberships$.pipe(
map(
({ value }) =>
value.find(
(membership) =>
membership.userId === userId && membership.deviceId === deviceId,
) ?? null,
),
filterBehavior((membership) => membership !== null),
),
);
const localMember$ = scope.behavior<LocalRtcMember | null>(
mapScoped(scope, localMembership$, (memberScope, membership$) =>
createLocalRtcMember(memberScope, membership$, localMembership, context),
).pipe(map((member) => member ?? null)),
);
const status$ = scope.behavior(
combineLatest(
[
localMembership.state$,
localMembership.joinRequested$,
localMembership.connected$,
localMembership.reconnecting$,
],
sessionStatus,
),
);
return {
join: localMembership.requestJoinAndPublish,
leave: localMembership.requestDisconnect,
status$,
connected$: localMembership.connected$,
reconnecting$: localMembership.reconnecting$,
fatalError$: scope.behavior(localMembership.state$.pipe(map(fatalError))),
localMember$,
remoteMembers$,
memberCount$: scope.behavior(
combineLatest(
[localMember$, remoteMembers$],
(local, remote) => remote.length + (local === null ? 0 : 1),
),
),
keyRotationSuppressed$: createKeyRotationSuppressed$(scope, jsSdkSession),
connectedTransports$: transports.connected$,
};
}
+93
View File
@@ -0,0 +1,93 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
isLivekitTransport,
type Transport,
} from "matrix-js-sdk/lib/matrixrtc";
import { combineLatest, distinctUntilChanged, map, of, switchMap } from "rxjs";
import { constant, type Behavior } from "../../../src/state/Behavior";
import { type ObservableScope } from "../../../src/state/ObservableScope";
import { ConnectionState } from "../../../src/state/CallViewModel/remoteMembers/Connection";
import { type IConnectionManager } from "../../../src/state/CallViewModel/remoteMembers/ConnectionManager";
import { type TransportMetadata } from "../api";
import { ResolvedConnection } from "./ConnectionFactory";
export interface TransportRegistry {
/** The one `TransportMetadata` for a transport, however many memberships name it. */
get(raw: Transport): TransportMetadata;
/** Transports with a live connection. */
connected$: Behavior<TransportMetadata[]>;
}
export function createTransportRegistry(
scope: ObservableScope,
connectionManager: IConnectionManager,
): TransportRegistry {
const transports = new Map<string, TransportMetadata>();
const get = (raw: Transport): TransportMetadata => {
const id = transportId(raw);
let transport = transports.get(id);
if (transport === undefined) {
transport = createTransportMetadata(scope, connectionManager, raw, id);
transports.set(id, transport);
}
return transport;
};
const connected$ = scope.behavior(
connectionManager.connectionManagerData$.pipe(
switchMap(({ value }) => {
const connections = value.getConnections();
if (connections.length === 0) return of([]);
return combineLatest(
connections.map((connection) =>
connection.state$.pipe(
map((state) =>
state === ConnectionState.LivekitConnected
? get(connection.transport)
: null,
),
),
),
).pipe(map((transports) => transports.filter((t) => t !== null)));
}),
),
);
return { get, connected$ };
}
function createTransportMetadata(
scope: ObservableScope,
connectionManager: IConnectionManager,
raw: Transport,
id: string,
): TransportMetadata {
const resolved$ = isLivekitTransport(raw)
? scope.behavior(
connectionManager.connectionManagerData$.pipe(
map(({ value }) => value.getConnectionForTransport(raw)),
distinctUntilChanged(),
switchMap((connection) =>
connection instanceof ResolvedConnection
? connection.resolved$
: of(undefined),
),
),
)
: constant(undefined);
return { type: raw.type, id, raw, resolved$ };
}
function transportId(raw: Transport): string {
return isLivekitTransport(raw)
? raw.livekit_service_url
: JSON.stringify(raw);
}
+68
View File
@@ -0,0 +1,68 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
type LivekitTransport,
type MatrixRTCSession as JsSdkRtcSession,
type RTCCallIntent,
type RTCNotificationType,
} from "matrix-js-sdk/lib/matrixrtc";
import { type CallMembershipIdentityParts } from "matrix-js-sdk/lib/matrixrtc/EncryptionManager";
import {
DEFAULT_CONFIG,
MatrixRTCMode,
} from "../../../src/config/ConfigOptions";
interface Options {
encryptMedia: boolean;
matrixRTCMode: MatrixRTCMode;
sendNotificationType?: RTCNotificationType;
callIntent?: RTCCallIntent;
}
/** The session's timing defaults; a later slice makes them an input. */
export const sessionTimings = {
syncDisconnectGracePeriodMs: DEFAULT_CONFIG.sync_disconnect_grace_period_ms,
...DEFAULT_CONFIG.matrix_rtc_session,
};
/**
* Sends the membership and starts the membership manager, which keeps it
* alive and retries on failure.
* @throws If the join could not be started.
*/
export function joinRtcSession(
session: JsSdkRtcSession,
ownMembershipIdentity: CallMembershipIdentityParts,
transport: LivekitTransport,
{ encryptMedia, matrixRTCMode, sendNotificationType, callIntent }: Options,
): void {
const timings = sessionTimings.delayed_leave;
// Give up on the network as soon as either the sync has been down for the
// grace period or the delayed leave has probably been sent, whichever is
// sooner
const maxWaitMs = Math.min(
sessionTimings.syncDisconnectGracePeriodMs,
timings.delay_ms,
);
session.joinRTCSession(ownMembershipIdentity, [transport], {
notificationType: sendNotificationType,
callIntent,
manageMediaKeys: encryptMedia,
delayedLeaveEventRestartMs: timings.restart_ms,
delayedLeaveEventDelayMs: timings.delay_ms,
delayedLeaveEventRestartLocalTimeoutMs: timings.restart_timeout_ms,
networkErrorRetryMs: sessionTimings.network_error_retry_ms,
makeKeyDelay: sessionTimings.wait_for_key_rotation_ms,
membershipEventExpiryMs: sessionTimings.membership_event_expiry_ms,
keyRotationParticipantLimit: sessionTimings.key_rotation_participant_limit,
unstableSendStickyEvents: matrixRTCMode === MatrixRTCMode.Matrix_2_0,
maximumNetworkErrorRetryCount:
Math.ceil(maxWaitMs / sessionTimings.network_error_retry_ms) + 1,
});
}
+77
View File
@@ -0,0 +1,77 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { describe, expect, it } from "vitest";
import { Status } from "matrix-js-sdk/lib/matrixrtc";
import { ConnectionState } from "../../../src/state/CallViewModel/remoteMembers/Connection";
import { RtcSessionError } from "../api";
import {
type LocalMemberState,
PublishState,
TransportState,
} from "./LocalMember";
import { fatalError, sessionStatus } from "./status";
const publishing: LocalMemberState = {
media: PublishState.Publishing,
matrix: Status.Connected,
};
const connecting: LocalMemberState = {
media: { connection: ConnectionState.LivekitConnecting },
matrix: Status.Connecting,
};
const error = new RtcSessionError("gone");
describe("sessionStatus", () => {
it.each<[string, LocalMemberState, boolean, boolean, boolean, string]>([
[
"no transport yet",
TransportState.Waiting,
true,
false,
false,
"waitingForTransport",
],
["not asked to join", connecting, false, false, false, "disconnected"],
["joining", connecting, true, false, false, "connecting"],
["connected", publishing, true, true, false, "connected"],
["dropped after connecting", connecting, true, false, true, "reconnecting"],
["a transport error", error, true, false, false, "disconnected"],
[
"a matrix error",
{ ...publishing, matrix: error },
true,
true,
false,
"disconnected",
],
])("%s", (_name, state, joinRequested, connected, reconnecting, expected) => {
expect(sessionStatus(state, joinRequested, connected, reconnecting)).toBe(
expected,
);
});
});
describe("fatalError", () => {
it("is null while nothing is wrong", () => {
expect(fatalError(TransportState.Waiting)).toBeNull();
expect(fatalError(publishing)).toBeNull();
});
it("finds the error wherever the state holds it", () => {
expect(fatalError(error)).toBe(error);
expect(fatalError({ ...publishing, matrix: error })).toBe(error);
expect(
fatalError({ media: { connection: error }, matrix: Status.Connected }),
).toBe(error);
});
it("does not treat a publish error as fatal", () => {
expect(fatalError({ media: error, matrix: Status.Connected })).toBeNull();
});
});
+37
View File
@@ -0,0 +1,37 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { type RtcSessionError, type SessionConnectionStatus } from "../api";
import { type LocalMemberState, TransportState } from "./LocalMember";
/** The local member's state, collapsed to what a host shows. */
export function sessionStatus(
state: LocalMemberState,
joinRequested: boolean,
connected: boolean,
reconnecting: boolean,
): SessionConnectionStatus {
if (fatalError(state) !== null) return "disconnected";
if (state === TransportState.Waiting) return "waitingForTransport";
if (connected) return "connected";
if (reconnecting) return "reconnecting";
return joinRequested ? "connecting" : "disconnected";
}
/** The error that stops the session, if the state holds one. */
export function fatalError(state: LocalMemberState): RtcSessionError | null {
if (state === TransportState.Waiting) return null;
if (state instanceof Error) return state;
if (state.matrix instanceof Error) return state.matrix;
if (
typeof state.media === "object" &&
"connection" in state.media &&
state.media.connection instanceof Error
)
return state.media.connection;
return null;
}
+52
View File
@@ -0,0 +1,52 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { describe, expect, it, vi } from "vitest";
import { Observable, Subject } from "rxjs";
import { LazyBehavior } from "./LazyBehavior";
describe("LazyBehavior", () => {
it("runs its source only while subscribed", () => {
const teardown = vi.fn();
let starts = 0;
const values$ = new Subject<number>();
const source$ = new Observable<number>((subscriber) => {
starts++;
const subscription = values$.subscribe(subscriber);
return (): void => {
teardown();
subscription.unsubscribe();
};
});
const lazy$ = new LazyBehavior(source$, 0);
expect(starts).toBe(0);
expect(lazy$.value).toBe(0);
const seenByFirst: number[] = [];
const first = lazy$.subscribe((v) => seenByFirst.push(v));
expect(starts).toBe(1);
values$.next(1);
expect(seenByFirst).toEqual([0, 1]);
expect(lazy$.value).toBe(1);
const seenBySecond: number[] = [];
const second = lazy$.subscribe((v) => seenBySecond.push(v));
expect(starts).toBe(1);
expect(seenBySecond).toEqual([1]);
first.unsubscribe();
expect(teardown).not.toHaveBeenCalled();
second.unsubscribe();
expect(teardown).toHaveBeenCalledTimes(1);
values$.next(2);
expect(lazy$.value).toBe(1);
lazy$.subscribe().unsubscribe();
expect(starts).toBe(2);
});
});
+53
View File
@@ -0,0 +1,53 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import {
BehaviorSubject,
type Observable,
type Subscriber,
type Subscription,
} from "rxjs";
/**
* A behavior that only runs its source while someone is subscribed. The source
* is subscribed with the first subscriber and unsubscribed with the last, so a
* value that is expensive to keep current, such as polled statistics, costs
* nothing while nobody looks at it. `value` is whatever the source last
* produced, or the initial value.
*/
export class LazyBehavior<T> extends BehaviorSubject<T> {
private subscribers = 0;
private upstream: Subscription | undefined;
public constructor(
private readonly source$: Observable<T>,
initialValue: T,
) {
super(initialValue);
}
// Every subscription of an Observable goes through this hook, which rxjs
// declares as internal and so leaves out of its types. The base class's
// version replays the current value and registers the subscriber; it is
// wrapped rather than replaced.
protected _subscribe(subscriber: Subscriber<T>): Subscription {
if (this.subscribers++ === 0)
this.upstream = this.source$.subscribe((value) => this.next(value));
const subscription = (
BehaviorSubject.prototype as unknown as {
_subscribe(this: BehaviorSubject<T>, s: Subscriber<T>): Subscription;
}
)._subscribe.call(this, subscriber);
subscription.add(() => {
if (--this.subscribers === 0) {
this.upstream?.unsubscribe();
this.upstream = undefined;
}
});
return subscription;
}
}
+35
View File
@@ -0,0 +1,35 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { describe, expect, it } from "vitest";
import { RtcSessionError } from "../api";
import { toRtcSessionError } from "./errors";
describe("toRtcSessionError", () => {
it("keeps an RtcSessionError as it is", () => {
const error = new RtcSessionError("gone");
expect(toRtcSessionError(error)).toBe(error);
});
it("wraps another error as the cause", () => {
const cause = new Error("boom");
const error = toRtcSessionError(cause);
expect(error).toBeInstanceOf(RtcSessionError);
expect(error.message).toBe("boom");
expect(error.cause).toBe(cause);
});
it("falls back to the code of an error without a message", () => {
const cause = Object.assign(new Error(""), { code: "SFU_ERROR" });
expect(toRtcSessionError(cause).message).toBe("SFU_ERROR");
});
it("stringifies anything else", () => {
expect(toRtcSessionError("nope").message).toBe("nope");
});
});
+27
View File
@@ -0,0 +1,27 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { RtcSessionError } from "../api";
/** Wraps whatever a backend threw into the SDK's error type, keeping it as the cause. */
export function toRtcSessionError(error: unknown): RtcSessionError {
if (error instanceof RtcSessionError) return error;
if (error instanceof Error)
return new RtcSessionError(messageOf(error), { cause: error });
return new RtcSessionError(String(error));
}
/**
* Element Call's errors carry a translated title as their message, which is
* empty where no translations are loaded; their code says what went wrong
* regardless.
*/
function messageOf(error: Error): string {
if (error.message !== "") return error.message;
if ("code" in error && typeof error.code === "string") return error.code;
return error.name;
}
+37
View File
@@ -0,0 +1,37 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { describe, expect, it, vi } from "vitest";
import { BehaviorSubject } from "rxjs";
import { ObservableScope } from "../../../src/state/ObservableScope";
import { mapScoped } from "./mapScoped";
describe("mapScoped", () => {
it("gives each value an item with a scope that ends with the value", () => {
const scope = new ObservableScope();
const source$ = new BehaviorSubject<string | null>("a");
const ended = vi.fn();
const items$ = mapScoped(scope, source$, (itemScope, value) => {
itemScope.onEnd(() => ended(value));
return value.toUpperCase();
});
expect(items$.value).toBe("A");
source$.next("b");
expect(items$.value).toBe("B");
expect(ended).toHaveBeenCalledWith("a");
source$.next(null);
expect(items$.value).toBeUndefined();
expect(ended).toHaveBeenCalledWith("b");
source$.next("c");
scope.end();
expect(ended).toHaveBeenCalledWith("c");
});
});
+35
View File
@@ -0,0 +1,35 @@
/*
Copyright 2026 Element Creations Ltd.
SPDX-License-Identifier: AGPL-3.0-only OR LicenseRef-Element-Commercial
Please see LICENSE in the repository root for full details.
*/
import { Observable, of, switchMap } from "rxjs";
import { type Behavior } from "../../../src/state/Behavior";
import { ObservableScope } from "../../../src/state/ObservableScope";
/**
* Builds one item per present value of a behavior, each in a scope of its own
* that ends when the value changes or the outer scope ends. The item for a
* null or undefined value is undefined.
*/
export function mapScoped<T, Item>(
scope: ObservableScope,
source$: Behavior<T | null | undefined>,
factory: (scope: ObservableScope, value: T) => Item,
): Behavior<Item | undefined> {
return scope.behavior(
source$.pipe(
switchMap((value) => {
if (value === null || value === undefined) return of(undefined);
return new Observable<Item>((subscriber) => {
const itemScope = new ObservableScope();
subscriber.next(factory(itemScope, value));
return (): void => itemScope.end();
});
}),
),
);
}
+12 -3
View File
@@ -90,7 +90,18 @@ export const vitePluginsConfig = ({
);
}
return { plugins };
return {
plugins,
resolve: {
alias: {
// The SDK by its package name, resolved to its source, for every
// config built on this one: see the matching entry in tsconfig.json
"@element-hq/matrixrtc-sdk": fileURLToPath(
new URL("./sdk/index.ts", import.meta.url),
),
},
},
};
};
// https://vitejs.dev/config/
// Modified type helper from defineConfig to allow for packageType (see defineConfig from vite)
@@ -162,8 +173,6 @@ export default ({
// which Vite for some reason refuses to work with, so we point it to
// src/index.ts instead
"matrix-widget-api": "matrix-widget-api/src/index.ts",
// The SDK by its package name, resolved to its source: see the matching
// entry in tsconfig.json
"@element-hq/matrixrtc-sdk": fileURLToPath(
new URL("./sdk/index.ts", import.meta.url),
),