review: Move the all advertised/active down to the LocalMember

And let the local member use it properly to send membership event and publish media
This commit is contained in:
Valere
2026-04-09 15:22:10 +02:00
parent 08006d640a
commit 40dacd523b
4 changed files with 94 additions and 60 deletions

View File

@@ -120,7 +120,6 @@ import {
type LocalMatrixLivekitMember, type LocalMatrixLivekitMember,
type RemoteMatrixLivekitMember, type RemoteMatrixLivekitMember,
type MatrixLivekitMember, type MatrixLivekitMember,
areLivekitTransportsEqual,
} from "./remoteMembers/MatrixLivekitMembers.ts"; } from "./remoteMembers/MatrixLivekitMembers.ts";
import { import {
type AutoLeaveReason, type AutoLeaveReason,
@@ -518,26 +517,6 @@ export function createCallViewModel$(
), ),
); );
// Observe the transport we should publish
const publishingTransport$ = localTransport$.pipe(
// observe the active$ transport
switchMap((t) => {
return combineLatest([t.active$, t.advertised$]).pipe(
map(([active, advertised]) => {
if (active?.transport) {
// use the active one (oldest member transport)
return active.transport;
} else {
// There is no active transport, we might just be the first member in the call
// so use the advertised to start
return advertised;
}
}),
);
}),
distinctUntilChanged(areLivekitTransportsEqual),
);
const localMembership = createLocalMembership$({ const localMembership = createLocalMembership$({
scope, scope,
homeserverConnected: createHomeserverConnected$( homeserverConnected: createHomeserverConnected$(
@@ -567,7 +546,7 @@ export function createCallViewModel$(
}, },
connectionManager, connectionManager,
matrixRTCSession, matrixRTCSession,
localTransport$: scope.behavior(publishingTransport$), localTransport$,
logger: logger.getChild(`[${Date.now()}]`), logger: logger.getChild(`[${Date.now()}]`),
}); });

View File

@@ -40,6 +40,10 @@ import { ConnectionManagerData } from "../remoteMembers/ConnectionManager";
import { ConnectionState, type Connection } from "../remoteMembers/Connection"; import { ConnectionState, type Connection } from "../remoteMembers/Connection";
import { type Publisher } from "./Publisher"; import { type Publisher } from "./Publisher";
import { initializeWidget } from "../../../widget"; import { initializeWidget } from "../../../widget";
import {
type LocalTransport,
type LocalTransportWithSFUConfig,
} from "./LocalTransport";
initializeWidget(); initializeWidget();
@@ -214,7 +218,7 @@ describe("LocalMembership", () => {
}; };
it("throws error on missing RTC config error", () => { it("throws error on missing RTC config error", () => {
withTestScheduler(({ scope, hot, expectObservable }) => { withTestScheduler(({ scope, hot, behavior, expectObservable }) => {
const localTransport$ = scope.behavior<null | LivekitTransportConfig>( const localTransport$ = scope.behavior<null | LivekitTransportConfig>(
hot("1ms #", {}, new MatrixRTCTransportMissingError("domain.com")), hot("1ms #", {}, new MatrixRTCTransportMissingError("domain.com")),
null, null,
@@ -230,11 +234,16 @@ describe("LocalMembership", () => {
), ),
}; };
const aLocalTransport: LocalTransport = {
advertised$: localTransport$,
active$: behavior("a", { a: null }),
};
const localMembership = createLocalMembership$({ const localMembership = createLocalMembership$({
scope, scope,
...defaultCreateLocalMemberValues, ...defaultCreateLocalMemberValues,
connectionManager: mockConnectionManager, connectionManager: mockConnectionManager,
localTransport$, localTransport$: behavior("a", { a: aLocalTransport }),
}); });
localMembership.requestJoinAndPublish(); localMembership.requestJoinAndPublish();
@@ -248,7 +257,11 @@ describe("LocalMembership", () => {
it("logs if callIntent cannot be updated", async () => { it("logs if callIntent cannot be updated", async () => {
const scope = new ObservableScope(); const scope = new ObservableScope();
const localTransport$ = new BehaviorSubject(aTransport); const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
active$: new BehaviorSubject(aTransportWithSFUConfig),
};
const mockConnectionManager = { const mockConnectionManager = {
transports$: constant(new Epoch([])), transports$: constant(new Epoch([])),
connectionManagerData$: constant(new Epoch(new ConnectionManagerData())), connectionManagerData$: constant(new Epoch(new ConnectionManagerData())),
@@ -264,7 +277,7 @@ describe("LocalMembership", () => {
leaveRoomSession: vi.fn(), leaveRoomSession: vi.fn(),
}, },
connectionManager: mockConnectionManager, connectionManager: mockConnectionManager,
localTransport$, localTransport$: new BehaviorSubject(aLocalTransport),
}); });
const expextedLog = const expextedLog =
"'not connected yet' while updating the call intent (this is expected on startup)"; "'not connected yet' while updating the call intent (this is expected on startup)";
@@ -279,6 +292,17 @@ describe("LocalMembership", () => {
const aTransport = { const aTransport = {
livekit_service_url: "a", livekit_service_url: "a",
} as LivekitTransportConfig; } as LivekitTransportConfig;
const aTransportWithSFUConfig = {
transport: aTransport,
sfuConfig: {
jwt: "foo",
livekitAlias: "bar",
livekitIdentity: "baz",
url: "bro",
},
} as LocalTransportWithSFUConfig;
const bTransport = { const bTransport = {
livekit_service_url: "b", livekit_service_url: "b",
} as LivekitTransportConfig; } as LivekitTransportConfig;
@@ -307,7 +331,11 @@ describe("LocalMembership", () => {
it("recreates publisher if new connection is used, always unpublish and end tracks", async () => { it("recreates publisher if new connection is used, always unpublish and end tracks", async () => {
const scope = new ObservableScope(); const scope = new ObservableScope();
const localTransport$ = new BehaviorSubject(aTransport); const activeTransport$ = new BehaviorSubject(aTransportWithSFUConfig);
const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
active$: activeTransport$,
};
const publishers: Publisher[] = []; const publishers: Publisher[] = [];
let seed = 0; let seed = 0;
@@ -343,10 +371,13 @@ describe("LocalMembership", () => {
connectionManager: { connectionManager: {
connectionManagerData$: constant(new Epoch(connectionManagerData)), connectionManagerData$: constant(new Epoch(connectionManagerData)),
}, },
localTransport$, localTransport$: new BehaviorSubject(aLocalTransport),
}); });
await flushPromises(); await flushPromises();
localTransport$.next(bTransport); activeTransport$.next({
...aTransportWithSFUConfig,
transport: bTransport,
});
await flushPromises(); await flushPromises();
expect(publisherFactory).toHaveBeenCalledTimes(2); expect(publisherFactory).toHaveBeenCalledTimes(2);
@@ -368,8 +399,6 @@ describe("LocalMembership", () => {
it("only start tracks if requested", async () => { it("only start tracks if requested", async () => {
const scope = new ObservableScope(); const scope = new ObservableScope();
const localTransport$ = new BehaviorSubject(aTransport);
const publishers: Publisher[] = []; const publishers: Publisher[] = [];
const tracks$ = new BehaviorSubject<LocalTrack[]>([]); const tracks$ = new BehaviorSubject<LocalTrack[]>([]);
@@ -396,6 +425,11 @@ describe("LocalMembership", () => {
typeof vi.fn typeof vi.fn
>; >;
const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
active$: new BehaviorSubject(aTransportWithSFUConfig),
};
const connectionManagerData = new ConnectionManagerData(); const connectionManagerData = new ConnectionManagerData();
connectionManagerData.add(connectionTransportAConnected, []); connectionManagerData.add(connectionTransportAConnected, []);
// connectionManagerData.add(connectionTransportB, []); // connectionManagerData.add(connectionTransportB, []);
@@ -405,7 +439,7 @@ describe("LocalMembership", () => {
connectionManager: { connectionManager: {
connectionManagerData$: constant(new Epoch(connectionManagerData)), connectionManagerData$: constant(new Epoch(connectionManagerData)),
}, },
localTransport$, localTransport$: new BehaviorSubject(aLocalTransport),
}); });
await flushPromises(); await flushPromises();
expect(publisherFactory).toHaveBeenCalledOnce(); expect(publisherFactory).toHaveBeenCalledOnce();
@@ -428,9 +462,15 @@ describe("LocalMembership", () => {
const scope = new ObservableScope(); const scope = new ObservableScope();
const connectionManagerData = new ConnectionManagerData(); const connectionManagerData = new ConnectionManagerData();
const localTransport$ = new BehaviorSubject<null | LivekitTransportConfig>(
null, const activeTransport$ =
); new BehaviorSubject<null | LocalTransportWithSFUConfig>(null);
const aLocalTransport: LocalTransport = {
advertised$: new BehaviorSubject(aTransport),
active$: activeTransport$,
};
const connectionManagerData$ = new BehaviorSubject( const connectionManagerData$ = new BehaviorSubject(
new Epoch(connectionManagerData), new Epoch(connectionManagerData),
); );
@@ -470,14 +510,14 @@ describe("LocalMembership", () => {
connectionManager: { connectionManager: {
connectionManagerData$, connectionManagerData$,
}, },
localTransport$, localTransport$: new BehaviorSubject(aLocalTransport),
}); });
await flushPromises(); await flushPromises();
expect(localMembership.localMemberState$.value).toStrictEqual( expect(localMembership.localMemberState$.value).toStrictEqual(
TransportState.Waiting, TransportState.Waiting,
); );
localTransport$.next(aTransport); activeTransport$.next(aTransportWithSFUConfig);
await flushPromises(); await flushPromises();
expect(localMembership.localMemberState$.value).toStrictEqual({ expect(localMembership.localMemberState$.value).toStrictEqual({
matrix: RTCMemberStatus.Connected, matrix: RTCMemberStatus.Connected,

View File

@@ -62,6 +62,8 @@ import {
} from "../remoteMembers/Connection.ts"; } from "../remoteMembers/Connection.ts";
import { type HomeserverConnected } from "./HomeserverConnected.ts"; import { type HomeserverConnected } from "./HomeserverConnected.ts";
import { and$ } from "../../../utils/observable.ts"; import { and$ } from "../../../utils/observable.ts";
import { type LocalTransport } from "./LocalTransport.ts";
import { areLivekitTransportsEqual } from "../remoteMembers/MatrixLivekitMembers.ts";
export enum TransportState { export enum TransportState {
/** Not even a transport is available to the LocalMembership */ /** Not even a transport is available to the LocalMembership */
@@ -127,7 +129,7 @@ interface Props {
createPublisherFactory: (connection: Connection) => Publisher; createPublisherFactory: (connection: Connection) => Publisher;
joinMatrixRTC: (transport: LivekitTransportConfig) => void; joinMatrixRTC: (transport: LivekitTransportConfig) => void;
homeserverConnected: HomeserverConnected; homeserverConnected: HomeserverConnected;
localTransport$: Behavior<LivekitTransportConfig | null>; localTransport$: Behavior<LocalTransport>;
matrixRTCSession: Pick< matrixRTCSession: Pick<
MatrixRTCSession, MatrixRTCSession,
"updateCallIntent" | "leaveRoomSession" "updateCallIntent" | "leaveRoomSession"
@@ -160,7 +162,7 @@ interface Props {
export const createLocalMembership$ = ({ export const createLocalMembership$ = ({
scope, scope,
connectionManager, connectionManager,
localTransport$: localTransportCanThrow$, localTransport$,
homeserverConnected, homeserverConnected,
createPublisherFactory, createPublisherFactory,
joinMatrixRTC, joinMatrixRTC,
@@ -205,23 +207,34 @@ export const createLocalMembership$ = ({
const logger = parentLogger.getChild("[LocalMembership]"); const logger = parentLogger.getChild("[LocalMembership]");
logger.debug(`Creating local membership..`); logger.debug(`Creating local membership..`);
// We consider error on the transport as fatal.
// Whether it is the active transport or the preferred transport.
const handleTransportError = (e: unknown): Observable<null> => {
let error: ElementCallError;
if (e instanceof ElementCallError) {
error = e;
} else {
error = new UnknownCallError(
e instanceof Error ? e : new Error("Unknown error from localTransport"),
);
}
setTransportError(error);
return of(null);
};
// This is the transport that we will advertise in our membership.
const advertisedTransport$ = localTransport$.pipe(
switchMap((lt) => lt.advertised$),
catchError(handleTransportError),
distinctUntilChanged(areLivekitTransportsEqual),
);
// Unwrap the local transport and set the state of the LocalMembership to error in case the transport is an error. // Unwrap the local transport and set the state of the LocalMembership to error in case the transport is an error.
const localTransport$ = scope.behavior( const activeTransport$ = scope.behavior(
localTransportCanThrow$.pipe( localTransport$.pipe(
catchError((e: unknown) => { switchMap((lt) => lt.active$.pipe(map((t) => t?.transport ?? null))),
let error: ElementCallError; catchError(handleTransportError),
if (e instanceof ElementCallError) { distinctUntilChanged(areLivekitTransportsEqual),
error = e;
} else {
error = new UnknownCallError(
e instanceof Error
? e
: new Error("Unknown error from localTransport"),
);
}
setTransportError(error);
return of(null);
}),
), ),
); );
@@ -229,7 +242,7 @@ export const createLocalMembership$ = ({
const localConnection$ = scope.behavior( const localConnection$ = scope.behavior(
combineLatest([ combineLatest([
connectionManager.connectionManagerData$, connectionManager.connectionManagerData$,
localTransport$, activeTransport$,
]).pipe( ]).pipe(
map(([{ value: connectionData }, localTransport]) => { map(([{ value: connectionData }, localTransport]) => {
if (localTransport === null) { if (localTransport === null) {
@@ -398,7 +411,7 @@ export const createLocalMembership$ = ({
const mediaState$: Behavior<LocalMemberMediaState> = scope.behavior( const mediaState$: Behavior<LocalMemberMediaState> = scope.behavior(
combineLatest([ combineLatest([
localConnectionState$, localConnectionState$,
localTransport$, activeTransport$,
joinAndPublishRequested$, joinAndPublishRequested$,
from(trackStartRequested.promise).pipe( from(trackStartRequested.promise).pipe(
map(() => true), map(() => true),
@@ -537,9 +550,11 @@ export const createLocalMembership$ = ({
}); });
}); });
// Keep matrix rtc session in sync with localTransport$, connectRequested$ // Keep matrix rtc session in sync with advertisedTransport$, connectRequested$
scope.reconcile( scope.reconcile(
scope.behavior(combineLatest([localTransport$, joinAndPublishRequested$])), scope.behavior(
combineLatest([advertisedTransport$, joinAndPublishRequested$]),
),
async ([transport, shouldConnect]) => { async ([transport, shouldConnect]) => {
if (!transport) return; if (!transport) return;
// if shouldConnect=false we will do the disconnect as the cleanup from the previous reconcile iteration. // if shouldConnect=false we will do the disconnect as the cleanup from the previous reconcile iteration.

View File

@@ -106,7 +106,7 @@ export function isLocalTransportWithSFUConfig(
return "transport" in obj && "sfuConfig" in obj; return "transport" in obj && "sfuConfig" in obj;
} }
interface LocalTransport { export interface LocalTransport {
/** /**
* The transport to be advertised in our MatrixRTC membership. `null` when not * The transport to be advertised in our MatrixRTC membership. `null` when not
* yet fetched/validated. * yet fetched/validated.