mirror of
https://github.com/element-hq/element-call.git
synced 2026-10-01 04:38:11 +00:00
Attach background effects through one switchable pipeline
- One pipeline, shared by the pre-join preview and the call, attached the first time an effect is wanted and never detached after; turning blur off switches it to a disabled mode in place. Main dropped the processor whenever blur went off, and a pipeline taken off the camera is destroyed and re-primed, so each time blur came back one unprocessed frame went out. Built once, the priming frame is spent once. - Still driven by the existing blur setting: nothing a user sees changes. - Its logic is BackgroundEffects in src/state/, which the provider only builds, so the tests drive it without rendering anything. - Disabled drops the blur radius too. The SDK's disabled mode keeps it, and with a radius every frame is still segmented and thrown away; with neither a radius nor a picture, frames pass untouched. - Switches run one at a time and skip any overtaken, so an earlier, slower switch can't land after a later choice. - Builds, rebuilds and teardowns of the pipeline run one at a time (OneStepPipeline). The preview's and the call's tracks attach and stop it independently, and the SDK's wrapper holds one track's streams: two builds at once overwrite each other's, and a stop during a build does nothing. - Every camera track starts with the pipeline on: the room's capture defaults follow it, and a camera turned off before the first effect is replaced on turning it on, since a stopped track can't take one. One attached after publishing lets the room through while it is built. - BlurBackgroundTransformer becomes BackgroundEffectTransformer: it also sits disabled, and its init applies that state as the SDK's would. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,24 @@
|
||||
/*
|
||||
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 { BackgroundEffectTransformer } from "./BackgroundEffectTransformer";
|
||||
|
||||
describe("BackgroundEffectTransformer", () => {
|
||||
it("passes frames on untouched once disabled, even after blurring", async () => {
|
||||
const transformer = new BackgroundEffectTransformer({});
|
||||
await transformer.update({ blurRadius: 15, backgroundDisabled: false });
|
||||
await transformer.update({
|
||||
imagePath: undefined,
|
||||
backgroundDisabled: true,
|
||||
});
|
||||
// With neither, the library passes a frame on without segmenting it.
|
||||
expect(transformer.options.blurRadius).toBeUndefined();
|
||||
expect(transformer.options.imagePath).toBeUndefined();
|
||||
});
|
||||
});
|
||||
+16
-1
@@ -1,5 +1,6 @@
|
||||
/*
|
||||
Copyright 2024-2025 New Vector Ltd.
|
||||
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.
|
||||
@@ -9,6 +10,7 @@ import {
|
||||
BackgroundTransformer,
|
||||
VideoTransformer,
|
||||
type VideoTransformerInitOptions,
|
||||
type BackgroundOptions,
|
||||
} from "@livekit/track-processors";
|
||||
import { ImageSegmenter } from "@mediapipe/tasks-vision";
|
||||
|
||||
@@ -49,7 +51,18 @@ const wasmFileset: WasmFileset = {
|
||||
* loads the segmentation models from our own bundle rather than as an external
|
||||
* resource fetched from the public internet.
|
||||
*/
|
||||
export class BlurBackgroundTransformer extends BackgroundTransformer {
|
||||
export class BackgroundEffectTransformer extends BackgroundTransformer {
|
||||
/**
|
||||
* As the library's, except that disabling also drops the blur radius: kept,
|
||||
* it has every frame segmented and the result thrown away, where with
|
||||
* neither a radius nor a picture frames are passed on untouched.
|
||||
*/
|
||||
public override async update(opts: BackgroundOptions): Promise<void> {
|
||||
await super.update(
|
||||
opts.backgroundDisabled ? { ...opts, blurRadius: undefined } : opts,
|
||||
);
|
||||
}
|
||||
|
||||
public async init({
|
||||
outputCanvas,
|
||||
inputElement: inputVideo,
|
||||
@@ -73,8 +86,10 @@ export class BlurBackgroundTransformer extends BackgroundTransformer {
|
||||
outputConfidenceMasks: false,
|
||||
});
|
||||
|
||||
// BackgroundTransformer's own init applies these, and this one replaces it.
|
||||
if (this.options.blurRadius) {
|
||||
this.gl?.setBlurRadius(this.options.blurRadius);
|
||||
}
|
||||
this.gl?.setBackgroundDisabled(this.options.backgroundDisabled ?? false);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
/*
|
||||
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 { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { ProcessorWrapper } from "@livekit/track-processors";
|
||||
|
||||
import { OneStepPipeline } from "./OneStepPipeline";
|
||||
import { BackgroundEffectTransformer } from "./BackgroundEffectTransformer";
|
||||
import { flushPromises } from "../utils/test";
|
||||
|
||||
const opts = {} as Parameters<OneStepPipeline["init"]>[0];
|
||||
|
||||
describe("OneStepPipeline", () => {
|
||||
afterEach(() => vi.restoreAllMocks());
|
||||
|
||||
it("stops and builds again only once a build has finished", async () => {
|
||||
let finishBuild!: () => void;
|
||||
const init = vi
|
||||
.spyOn(ProcessorWrapper.prototype, "init")
|
||||
.mockImplementationOnce(
|
||||
async () => new Promise((resolve) => (finishBuild = resolve)),
|
||||
)
|
||||
.mockResolvedValue();
|
||||
const destroy = vi
|
||||
.spyOn(ProcessorWrapper.prototype, "destroy")
|
||||
.mockResolvedValue();
|
||||
const pipeline = new OneStepPipeline(
|
||||
new BackgroundEffectTransformer({}),
|
||||
"test",
|
||||
);
|
||||
|
||||
const first = pipeline.init(opts);
|
||||
const stopped = pipeline.destroy();
|
||||
const second = pipeline.init(opts);
|
||||
await flushPromises();
|
||||
expect(init).toHaveBeenCalledTimes(1);
|
||||
expect(destroy).not.toHaveBeenCalled();
|
||||
|
||||
finishBuild();
|
||||
await Promise.all([first, stopped, second]);
|
||||
expect(init).toHaveBeenCalledTimes(2);
|
||||
expect(destroy.mock.invocationCallOrder[0]).toBeLessThan(
|
||||
init.mock.invocationCallOrder[1],
|
||||
);
|
||||
});
|
||||
|
||||
it("restarts in one step, without waiting on itself", async () => {
|
||||
const init = vi
|
||||
.spyOn(ProcessorWrapper.prototype, "init")
|
||||
.mockResolvedValue();
|
||||
const destroy = vi
|
||||
.spyOn(ProcessorWrapper.prototype, "destroy")
|
||||
.mockResolvedValue();
|
||||
const pipeline = new OneStepPipeline(
|
||||
new BackgroundEffectTransformer({}),
|
||||
"test",
|
||||
);
|
||||
|
||||
await pipeline.restart(opts);
|
||||
expect(destroy).toHaveBeenCalledWith({ willProcessorRestart: true });
|
||||
expect(init).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("carries on after a step that failed", async () => {
|
||||
vi.spyOn(ProcessorWrapper.prototype, "init")
|
||||
.mockRejectedValueOnce(new Error("no GPU"))
|
||||
.mockResolvedValue();
|
||||
const pipeline = new OneStepPipeline(
|
||||
new BackgroundEffectTransformer({}),
|
||||
"test",
|
||||
);
|
||||
|
||||
await expect(pipeline.init(opts)).rejects.toThrow("no GPU");
|
||||
await expect(pipeline.init(opts)).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,43 @@
|
||||
/*
|
||||
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 { BackgroundProcessorWrapper } from "@livekit/track-processors";
|
||||
|
||||
type InitOptions = Parameters<BackgroundProcessorWrapper["init"]>[0];
|
||||
type DestroyOptions = Parameters<BackgroundProcessorWrapper["destroy"]>[0];
|
||||
|
||||
/**
|
||||
* A background pipeline that is built, rebuilt and destroyed one step at a
|
||||
* time. The camera tracks sharing it attach and stop it independently, and
|
||||
* the SDK's wrapper holds a single track's streams: two builds at once
|
||||
* overwrite each other's, and a destroy during a build does nothing.
|
||||
*/
|
||||
export class OneStepPipeline extends BackgroundProcessorWrapper {
|
||||
private steps: Promise<void> = Promise.resolve();
|
||||
|
||||
public override async init(opts: InitOptions): Promise<void> {
|
||||
return this.inTurn(async () => super.init(opts));
|
||||
}
|
||||
|
||||
// The SDK's restart calls init and destroy, which would wait on itself.
|
||||
public override async restart(opts: InitOptions): Promise<void> {
|
||||
return this.inTurn(async () => {
|
||||
await super.destroy({ willProcessorRestart: true });
|
||||
await super.init(opts);
|
||||
});
|
||||
}
|
||||
|
||||
public override async destroy(options?: DestroyOptions): Promise<void> {
|
||||
return this.inTurn(async () => super.destroy(options));
|
||||
}
|
||||
|
||||
private async inTurn(step: () => Promise<void>): Promise<void> {
|
||||
const turn = this.steps.then(step);
|
||||
this.steps = turn.catch(() => {});
|
||||
return turn;
|
||||
}
|
||||
}
|
||||
@@ -5,17 +5,54 @@ 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 { act, createElement, type FC } from "react";
|
||||
import { render } from "@testing-library/react";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { type LocalVideoTrack } from "livekit-client";
|
||||
import {
|
||||
type BackgroundOptions,
|
||||
type ProcessorWrapper,
|
||||
} from "@livekit/track-processors";
|
||||
|
||||
import { applyProcessor, trackProcessorSync } from "./TrackProcessorContext";
|
||||
import {
|
||||
applyProcessor,
|
||||
ProcessorProvider,
|
||||
type ProcessorState,
|
||||
trackProcessorSync,
|
||||
useTrackProcessor,
|
||||
useTrackProcessorSync,
|
||||
} from "./TrackProcessorContext";
|
||||
import { backgroundBlur } from "../settings/settings";
|
||||
import { constant } from "../state/Behavior";
|
||||
import { flushPromises, testScope } from "../utils/test";
|
||||
|
||||
const pipelines = vi.hoisted(() => ({
|
||||
built: 0,
|
||||
destroyed: 0,
|
||||
switches: [] as unknown[],
|
||||
}));
|
||||
|
||||
// The pipeline itself needs WebGL and MediaPipe; what these tests are about is
|
||||
// what the provider asks of it.
|
||||
vi.mock("@livekit/track-processors", () => ({
|
||||
BackgroundProcessorWrapper: vi.fn(function (this: Record<string, unknown>) {
|
||||
pipelines.built += 1;
|
||||
this.switchTo = vi.fn(async (options: unknown) => {
|
||||
pipelines.switches.push(options);
|
||||
return Promise.resolve();
|
||||
});
|
||||
this.destroy = vi.fn(async () => {
|
||||
pipelines.destroyed += 1;
|
||||
return Promise.resolve();
|
||||
});
|
||||
}),
|
||||
supportsBackgroundProcessors: (): boolean => true,
|
||||
}));
|
||||
vi.mock("./BackgroundEffectTransformer", () => ({
|
||||
BackgroundEffectTransformer: vi.fn(),
|
||||
}));
|
||||
vi.mock("../Platform", () => ({ platform: "desktop" }));
|
||||
|
||||
const processor = {} as ProcessorWrapper<BackgroundOptions>;
|
||||
|
||||
function mockTrack(
|
||||
@@ -85,3 +122,105 @@ describe("trackProcessorSync", () => {
|
||||
expect(track.setProcessor).toHaveBeenCalledWith(processor);
|
||||
});
|
||||
});
|
||||
|
||||
describe("ProcessorProvider", () => {
|
||||
/** A camera track that holds whatever processor it is given. */
|
||||
function cameraTrack(): LocalVideoTrack {
|
||||
let current: unknown;
|
||||
return {
|
||||
mediaStreamTrack: { readyState: "live" },
|
||||
getProcessor: vi.fn(() => current),
|
||||
setProcessor: vi.fn(async (next: unknown) => {
|
||||
current = next;
|
||||
return Promise.resolve();
|
||||
}),
|
||||
stopProcessor: vi.fn(async () => {
|
||||
current = undefined;
|
||||
return Promise.resolve();
|
||||
}),
|
||||
} as unknown as LocalVideoTrack;
|
||||
}
|
||||
|
||||
let seen: ProcessorState[];
|
||||
const Surface: FC<{ track: LocalVideoTrack | null }> = ({ track }) => {
|
||||
seen.push(useTrackProcessor());
|
||||
useTrackProcessorSync(track);
|
||||
return null;
|
||||
};
|
||||
// One component for every render, so a rerender updates the tree rather than
|
||||
// mounting a second provider with a pipeline of its own.
|
||||
const Surfaces: FC<{ tracks: (LocalVideoTrack | null)[] }> = ({ tracks }) =>
|
||||
createElement(
|
||||
ProcessorProvider,
|
||||
null,
|
||||
createElement(
|
||||
"div",
|
||||
null,
|
||||
...tracks.map((track, key) => createElement(Surface, { key, track })),
|
||||
),
|
||||
);
|
||||
const surfaces = (
|
||||
...tracks: (LocalVideoTrack | null)[]
|
||||
): ReturnType<typeof createElement> => createElement(Surfaces, { tracks });
|
||||
const blur = async (on: boolean): Promise<void> => {
|
||||
await act(async () => {
|
||||
backgroundBlur.setValue(on);
|
||||
await flushPromises();
|
||||
});
|
||||
};
|
||||
const latest = (): ProcessorState => seen[seen.length - 1];
|
||||
|
||||
beforeEach(() => {
|
||||
seen = [];
|
||||
pipelines.built = 0;
|
||||
pipelines.destroyed = 0;
|
||||
pipelines.switches = [];
|
||||
backgroundBlur.setValue(false);
|
||||
});
|
||||
afterEach(() => backgroundBlur.setValue(false));
|
||||
|
||||
// A pipeline primes itself when it is built and when it is destroyed, and a
|
||||
// primed pipeline lets one frame through untouched: so the frame is spent
|
||||
// once, as long as it is built once and never taken off the camera.
|
||||
it("spends one priming frame on first use and none afterwards", async () => {
|
||||
const track = cameraTrack();
|
||||
render(surfaces(track));
|
||||
for (const on of [true, false, true, false, true]) await blur(on);
|
||||
|
||||
expect(pipelines.built).toBe(1);
|
||||
expect(pipelines.destroyed).toBe(0);
|
||||
expect(track.setProcessor).toHaveBeenCalledTimes(1);
|
||||
expect(track.stopProcessor).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("preview and call share one pipeline", async () => {
|
||||
const preview = cameraTrack();
|
||||
const call = cameraTrack();
|
||||
render(surfaces(preview, call));
|
||||
await blur(true);
|
||||
|
||||
const pipeline = latest().processor;
|
||||
expect(preview.setProcessor).toHaveBeenCalledWith(pipeline);
|
||||
expect(call.setProcessor).toHaveBeenCalledWith(pipeline);
|
||||
expect(pipelines.built).toBe(1);
|
||||
});
|
||||
|
||||
it("applies an effect chosen while the camera is off", async () => {
|
||||
const view = render(surfaces(null));
|
||||
await blur(true);
|
||||
// Already switched to blur, with nothing yet to attach it to.
|
||||
expect(pipelines.switches).toEqual([
|
||||
{ mode: "background-blur", blurRadius: 15 },
|
||||
]);
|
||||
|
||||
const track = cameraTrack();
|
||||
await act(async () => {
|
||||
view.rerender(surfaces(track));
|
||||
await flushPromises();
|
||||
});
|
||||
expect(track.setProcessor).toHaveBeenCalledTimes(1);
|
||||
expect(track.setProcessor).toHaveBeenCalledWith(latest().processor);
|
||||
expect(pipelines.built).toBe(1);
|
||||
expect(pipelines.switches).toHaveLength(1);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
/*
|
||||
Copyright 2024-2025 New Vector Ltd.
|
||||
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 {
|
||||
ProcessorWrapper,
|
||||
type ProcessorWrapper,
|
||||
supportsBackgroundProcessors as supportsBackgroundProcessorsLivekitSdk,
|
||||
type BackgroundOptions,
|
||||
} from "@livekit/track-processors";
|
||||
@@ -16,20 +17,19 @@ import {
|
||||
type JSX,
|
||||
use,
|
||||
useEffect,
|
||||
useMemo,
|
||||
useState,
|
||||
} from "react";
|
||||
import { type LocalVideoTrack } from "livekit-client";
|
||||
import { logger } from "matrix-js-sdk/lib/logger";
|
||||
import { combineLatest, map, type Observable } from "rxjs";
|
||||
import { useObservable } from "observable-hooks";
|
||||
import { combineLatest } from "rxjs";
|
||||
|
||||
import {
|
||||
backgroundBlur as backgroundBlurSettings,
|
||||
useSetting,
|
||||
} from "../settings/settings";
|
||||
import { BlurBackgroundTransformer } from "./BlurBackgroundTransformer";
|
||||
import { backgroundBlur as backgroundBlurSettings } from "../settings/settings";
|
||||
import { BackgroundEffectTransformer } from "./BackgroundEffectTransformer";
|
||||
import { OneStepPipeline } from "./OneStepPipeline";
|
||||
import { type Behavior } from "../state/Behavior";
|
||||
import { type ObservableScope } from "../state/ObservableScope";
|
||||
import { ObservableScope } from "../state/ObservableScope";
|
||||
import { BackgroundEffects } from "../state/BackgroundEffects";
|
||||
import { useBehavior } from "../useBehavior";
|
||||
import { platform } from "../Platform";
|
||||
|
||||
//TODO-MULTI-SFU: This is not yet fully there.
|
||||
@@ -41,29 +41,21 @@ export type ProcessorState = {
|
||||
processor: undefined | ProcessorWrapper<BackgroundOptions>;
|
||||
};
|
||||
|
||||
const ProcessorContext = createContext<ProcessorState | undefined>(undefined);
|
||||
const ProcessorContext = createContext<BackgroundEffects | undefined>(
|
||||
undefined,
|
||||
);
|
||||
|
||||
export function useTrackProcessor(): ProcessorState {
|
||||
const state = use(ProcessorContext);
|
||||
if (state === undefined)
|
||||
throw new Error(
|
||||
"useTrackProcessor must be used within a ProcessorProvider",
|
||||
);
|
||||
return state;
|
||||
return useBehavior(useTrackProcessorObservable$());
|
||||
}
|
||||
|
||||
export function useTrackProcessorObservable$(): Observable<ProcessorState> {
|
||||
const state = use(ProcessorContext);
|
||||
if (state === undefined)
|
||||
export function useTrackProcessorObservable$(): Behavior<ProcessorState> {
|
||||
const effects = use(ProcessorContext);
|
||||
if (effects === undefined)
|
||||
throw new Error(
|
||||
"useTrackProcessor must be used within a ProcessorProvider",
|
||||
);
|
||||
const state$ = useObservable(
|
||||
(init$) => init$.pipe(map(([init]) => init)),
|
||||
[state],
|
||||
);
|
||||
|
||||
return state$;
|
||||
return effects.state$;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -83,8 +75,13 @@ export function applyProcessor(
|
||||
logger.debug("Not attaching video processor to an ended track");
|
||||
return;
|
||||
}
|
||||
const track = videoTrack.mediaStreamTrack;
|
||||
videoTrack.setProcessor(processor).catch((e) => {
|
||||
logger.warn("Failed to attach video processor", e);
|
||||
// The pipeline builds one track at a time, so a track can end while its
|
||||
// build waits a turn. The camera's next track attaches it again.
|
||||
if (track.readyState === "ended")
|
||||
logger.debug("Video processor not attached: the track ended first");
|
||||
else logger.warn("Failed to attach video processor", e);
|
||||
});
|
||||
}
|
||||
if (!processor && videoTrack.getProcessor()) {
|
||||
@@ -130,26 +127,22 @@ function supportsBackgroundProcessors(): boolean {
|
||||
}
|
||||
|
||||
export const ProcessorProvider: FC<Props> = ({ children }) => {
|
||||
// The setting the user wants to have
|
||||
const [blurActivated] = useSetting(backgroundBlurSettings);
|
||||
const supported = useMemo(() => supportsBackgroundProcessors(), []);
|
||||
const blur = useMemo(
|
||||
() =>
|
||||
new ProcessorWrapper(
|
||||
new BlurBackgroundTransformer({ blurRadius: 15 }),
|
||||
"background-blur",
|
||||
),
|
||||
[],
|
||||
);
|
||||
const [effects, setEffects] = useState<BackgroundEffects | null>(null);
|
||||
useEffect(() => {
|
||||
const scope = new ObservableScope();
|
||||
setEffects(
|
||||
new BackgroundEffects(scope, {
|
||||
supported: supportsBackgroundProcessors(),
|
||||
blur$: backgroundBlurSettings.value$,
|
||||
pipeline: new OneStepPipeline(
|
||||
new BackgroundEffectTransformer({ backgroundDisabled: true }),
|
||||
"background-effect",
|
||||
),
|
||||
}),
|
||||
);
|
||||
return (): void => scope.end();
|
||||
}, []);
|
||||
|
||||
// This is the actual state exposed through the context
|
||||
const processorState = useMemo(
|
||||
() => ({
|
||||
supported,
|
||||
processor: supported && blurActivated ? blur : undefined,
|
||||
}),
|
||||
[supported, blurActivated, blur],
|
||||
);
|
||||
|
||||
return <ProcessorContext value={processorState}>{children}</ProcessorContext>;
|
||||
if (effects === null) return null;
|
||||
return <ProcessorContext value={effects}>{children}</ProcessorContext>;
|
||||
};
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
/*
|
||||
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 { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { BehaviorSubject, distinctUntilChanged, map } from "rxjs";
|
||||
import { type BackgroundProcessorWrapper } from "@livekit/track-processors";
|
||||
|
||||
import {
|
||||
BackgroundEffects,
|
||||
type BackgroundEffectsOptions,
|
||||
} from "./BackgroundEffects";
|
||||
import { type ProcessorState } from "../livekit/TrackProcessorContext";
|
||||
import { flushPromises, testScope, withTestScheduler } from "../utils/test";
|
||||
|
||||
/** A pipeline that records what it is switched to. */
|
||||
function fakePipeline(): {
|
||||
pipeline: BackgroundProcessorWrapper;
|
||||
switches: unknown[];
|
||||
/** Holds the next switch until it resolves, as a picture loading does. */
|
||||
holdNext: () => () => void;
|
||||
} {
|
||||
const switches: unknown[] = [];
|
||||
let slow: Promise<void> | undefined;
|
||||
const pipeline = {
|
||||
switchTo: vi.fn(async (options: unknown) => {
|
||||
switches.push(options);
|
||||
const held = slow;
|
||||
slow = undefined;
|
||||
await held;
|
||||
}),
|
||||
} as unknown as BackgroundProcessorWrapper;
|
||||
return {
|
||||
pipeline,
|
||||
switches,
|
||||
holdNext: () => {
|
||||
let release!: () => void;
|
||||
slow = new Promise((resolve) => (release = resolve));
|
||||
return release;
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
/** One letter per state: idle, attached. */
|
||||
function letter(state: ProcessorState): string {
|
||||
return state.processor === undefined ? "i" : "a";
|
||||
}
|
||||
|
||||
describe("the pipeline's state", () => {
|
||||
function testState({
|
||||
effect,
|
||||
expected,
|
||||
}: {
|
||||
effect: string;
|
||||
expected: string;
|
||||
}): void {
|
||||
withTestScheduler(({ behavior, expectObservable }) => {
|
||||
const effects = new BackgroundEffects(testScope(), {
|
||||
supported: true,
|
||||
blur$: behavior(effect, { n: false, b: true }),
|
||||
pipeline: fakePipeline().pipeline,
|
||||
});
|
||||
expectObservable(
|
||||
effects.state$.pipe(map(letter), distinctUntilChanged()),
|
||||
).toBe(expected);
|
||||
});
|
||||
}
|
||||
|
||||
it("attaches on first use and stays attached", () =>
|
||||
testState({ effect: "nbn", expected: "ia" }));
|
||||
});
|
||||
|
||||
describe("background effects", () => {
|
||||
let blur$: BehaviorSubject<boolean>;
|
||||
let fake: ReturnType<typeof fakePipeline>;
|
||||
|
||||
function build(
|
||||
options: Partial<BackgroundEffectsOptions> = {},
|
||||
): BackgroundEffects {
|
||||
return new BackgroundEffects(testScope(), {
|
||||
supported: true,
|
||||
blur$,
|
||||
pipeline: fake.pipeline,
|
||||
...options,
|
||||
});
|
||||
}
|
||||
const blur = async (on: boolean): Promise<void> => {
|
||||
blur$.next(on);
|
||||
await flushPromises();
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
blur$ = new BehaviorSubject(false);
|
||||
fake = fakePipeline();
|
||||
});
|
||||
|
||||
it("switches in place rather than reattaching", async () => {
|
||||
const effects = build();
|
||||
await blur(true);
|
||||
const pipeline = effects.state$.value.processor;
|
||||
await blur(false);
|
||||
await blur(true);
|
||||
|
||||
expect(effects.state$.value.processor).toBe(pipeline);
|
||||
expect(fake.switches).toEqual([
|
||||
{ mode: "background-blur", blurRadius: 15 },
|
||||
{ mode: "disabled" },
|
||||
{ mode: "background-blur", blurRadius: 15 },
|
||||
]);
|
||||
});
|
||||
|
||||
it("switches one at a time, skipping those overtaken", async () => {
|
||||
const finished = fake.holdNext();
|
||||
build();
|
||||
await blur(true);
|
||||
await blur(false);
|
||||
await blur(true);
|
||||
await blur(false);
|
||||
expect(fake.switches).toHaveLength(1);
|
||||
|
||||
finished();
|
||||
await flushPromises();
|
||||
expect(fake.switches).toEqual([
|
||||
{ mode: "background-blur", blurRadius: 15 },
|
||||
{ mode: "disabled" },
|
||||
]);
|
||||
});
|
||||
});
|
||||
@@ -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 { combineLatest, distinctUntilChanged, filter, map, scan } from "rxjs";
|
||||
import {
|
||||
type BackgroundProcessorWrapper,
|
||||
type SwitchBackgroundProcessorOptions,
|
||||
} from "@livekit/track-processors";
|
||||
import { logger } from "matrix-js-sdk/lib/logger";
|
||||
|
||||
import { type Behavior } from "./Behavior";
|
||||
import { type ObservableScope } from "./ObservableScope";
|
||||
import { type ProcessorState } from "../livekit/TrackProcessorContext";
|
||||
|
||||
const blurRadius = 15;
|
||||
|
||||
export interface BackgroundEffectsOptions {
|
||||
/** Whether this browser can run a pipeline at all. */
|
||||
supported: boolean;
|
||||
/** Whether blur is chosen. */
|
||||
blur$: Behavior<boolean>;
|
||||
/**
|
||||
* Shared by the pre-join preview and the call. Building or destroying it is
|
||||
* what primes it, and a primed pipeline lets its next frame through
|
||||
* untouched, so it is switched rather than rebuilt.
|
||||
*/
|
||||
pipeline: BackgroundProcessorWrapper;
|
||||
}
|
||||
|
||||
/** The background effect pipeline, as the camera tracks and the menus see it. */
|
||||
export class BackgroundEffects {
|
||||
public readonly state$: Behavior<ProcessorState>;
|
||||
|
||||
public constructor(
|
||||
scope: ObservableScope,
|
||||
{ supported, blur$, pipeline }: BackgroundEffectsOptions,
|
||||
) {
|
||||
this.state$ = scope.behavior(
|
||||
blur$.pipe(
|
||||
scan<boolean, ProcessorState>(
|
||||
(previous, wanted) => {
|
||||
// Attached the first time an effect is wanted and never detached
|
||||
// after, so someone who never turns one on pays for none of it.
|
||||
const attached =
|
||||
previous.processor !== undefined || (supported && wanted);
|
||||
return { supported, processor: attached ? pipeline : undefined };
|
||||
},
|
||||
{ supported, processor: undefined },
|
||||
),
|
||||
),
|
||||
);
|
||||
|
||||
const switchTo = oneSwitchAtATime(pipeline);
|
||||
combineLatest([this.state$, blur$])
|
||||
.pipe(
|
||||
filter(([{ processor }]) => processor !== undefined),
|
||||
map(([, blur]): SwitchBackgroundProcessorOptions =>
|
||||
blur ? { mode: "background-blur", blurRadius } : { mode: "disabled" },
|
||||
),
|
||||
distinctUntilChanged((a, b) => JSON.stringify(a) === JSON.stringify(b)),
|
||||
scope.bind(),
|
||||
)
|
||||
.subscribe((options) => {
|
||||
switchTo(options).catch((e) =>
|
||||
logger.warn("Failed to switch background effect", e),
|
||||
);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Switches the pipeline one choice at a time, skipping those overtaken while
|
||||
* they waited. A switch to a picture ends only once it has loaded, so a
|
||||
* slower, earlier choice would otherwise land after a later one.
|
||||
*/
|
||||
function oneSwitchAtATime(
|
||||
pipeline: BackgroundProcessorWrapper,
|
||||
): (options: SwitchBackgroundProcessorOptions) => Promise<void> {
|
||||
let latest: SwitchBackgroundProcessorOptions | undefined;
|
||||
let queue = Promise.resolve();
|
||||
return async (options) => {
|
||||
latest = options;
|
||||
const turn = queue.then(async () => {
|
||||
if (options === latest) await pipeline.switchTo(options);
|
||||
});
|
||||
queue = turn.catch(() => {});
|
||||
return turn;
|
||||
};
|
||||
}
|
||||
@@ -11,11 +11,16 @@ import {
|
||||
LocalParticipant,
|
||||
type LocalTrack,
|
||||
type LocalTrackPublication,
|
||||
LocalVideoTrack,
|
||||
ParticipantEvent,
|
||||
Track,
|
||||
} from "livekit-client";
|
||||
import { BehaviorSubject } from "rxjs";
|
||||
import { logger } from "matrix-js-sdk/lib/logger";
|
||||
import {
|
||||
type BackgroundOptions,
|
||||
type ProcessorWrapper,
|
||||
} from "@livekit/track-processors";
|
||||
|
||||
import { ObservableScope } from "../../ObservableScope";
|
||||
import { constant } from "../../Behavior";
|
||||
@@ -25,6 +30,7 @@ import {
|
||||
mockMediaDevices,
|
||||
} from "../../../utils/test";
|
||||
import { Publisher } from "./Publisher";
|
||||
import { type ProcessorState } from "../../../livekit/TrackProcessorContext";
|
||||
import { type Connection } from "../remoteMembers/Connection";
|
||||
import { type MuteStates } from "../../MuteStates";
|
||||
|
||||
@@ -400,3 +406,67 @@ describe("Bug fix", () => {
|
||||
await publisher.destroy();
|
||||
});
|
||||
});
|
||||
|
||||
describe("turning the camera on with an effect chosen", () => {
|
||||
const processor = {} as ProcessorWrapper<BackgroundOptions>;
|
||||
let state$: BehaviorSubject<ProcessorState>;
|
||||
let publisher: Publisher;
|
||||
/** The processor the SDK would give each camera track it made. */
|
||||
let madeWith: unknown[];
|
||||
|
||||
beforeEach(() => {
|
||||
state$ = new BehaviorSubject<ProcessorState>({
|
||||
supported: true,
|
||||
processor: undefined,
|
||||
});
|
||||
publisher = new Publisher(
|
||||
connection,
|
||||
mockMediaDevices({}),
|
||||
muteStates,
|
||||
state$,
|
||||
logger,
|
||||
false,
|
||||
);
|
||||
madeWith = [];
|
||||
vi.spyOn(localParticipant, "setCameraEnabled").mockImplementation(
|
||||
async () => {
|
||||
madeWith.push(
|
||||
connection.livekitRoom.options.videoCaptureDefaults?.processor,
|
||||
);
|
||||
return Promise.resolve(undefined);
|
||||
},
|
||||
);
|
||||
vi.spyOn(localParticipant, "unpublishTrack").mockResolvedValue(undefined);
|
||||
});
|
||||
afterEach(async () => {
|
||||
await publisher.destroy();
|
||||
});
|
||||
|
||||
it("makes each camera track with the pipeline attached since", async () => {
|
||||
await publisher.createAndSetupTracks();
|
||||
state$.next({ supported: true, processor });
|
||||
videoEnabled$.next(true);
|
||||
await flushPromises();
|
||||
|
||||
expect(madeWith).toEqual([processor]);
|
||||
});
|
||||
|
||||
it("replaces a camera track turned off before the effect was chosen", async () => {
|
||||
state$.next({ supported: true, processor });
|
||||
const stopped = Object.assign(Object.create(LocalVideoTrack.prototype), {
|
||||
source: Track.Source.Camera,
|
||||
getProcessor: (): undefined => undefined,
|
||||
}) as LocalVideoTrack;
|
||||
trackPublications.push({
|
||||
track: stopped,
|
||||
source: Track.Source.Camera,
|
||||
} as Partial<LocalTrackPublication> as LocalTrackPublication);
|
||||
|
||||
await publisher.createAndSetupTracks();
|
||||
videoEnabled$.next(true);
|
||||
await flushPromises();
|
||||
|
||||
expect(localParticipant.unpublishTrack).toHaveBeenCalledWith(stopped);
|
||||
expect(madeWith).toEqual([processor]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -63,7 +63,7 @@ export class Publisher {
|
||||
private connection: Pick<Connection, "livekitRoom" | "state$">, //setE2EEEnabled,
|
||||
devices: MediaDevices,
|
||||
private readonly muteStates: MuteStates,
|
||||
trackerProcessorState$: Behavior<ProcessorState>,
|
||||
private readonly trackerProcessorState$: Behavior<ProcessorState>,
|
||||
private logger: Logger,
|
||||
controlledAudioDevices: boolean,
|
||||
) {
|
||||
@@ -74,7 +74,7 @@ export class Publisher {
|
||||
});
|
||||
|
||||
// Setup track processor syncing (blur)
|
||||
this.observeTrackProcessors(this.scope, room, trackerProcessorState$);
|
||||
this.observeTrackProcessors(this.scope, room, this.trackerProcessorState$);
|
||||
// Observe media device changes and update LiveKit active devices accordingly
|
||||
this.observeMediaDevices(this.scope, devices, controlledAudioDevices);
|
||||
|
||||
@@ -412,6 +412,8 @@ export class Publisher {
|
||||
this.muteStates.video.setHandler(async (enable) => {
|
||||
try {
|
||||
this.logger.debug(`handler: Setting LiveKit camera enabled: ${enable}`);
|
||||
const { processor } = this.trackerProcessorState$.value;
|
||||
if (enable && processor) await this.dropCameraWithoutEffect(lkRoom);
|
||||
await lkRoom.localParticipant.setCameraEnabled(enable);
|
||||
// Unmute will restart the track if it was paused upstream,
|
||||
// but until explicitly requested, we want to keep it paused.
|
||||
@@ -426,6 +428,19 @@ export class Publisher {
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Unpublishes a camera track that has no effect on, so that turning the
|
||||
* camera on creates one with it. A camera turned off is stopped, and a
|
||||
* stopped track can't take an effect chosen while it was off.
|
||||
*/
|
||||
private async dropCameraWithoutEffect(lkRoom: LivekitRoom): Promise<void> {
|
||||
const track = lkRoom.localParticipant.getTrackPublication(
|
||||
Track.Source.Camera,
|
||||
)?.track;
|
||||
if (track instanceof LocalVideoTrack && !track.getProcessor())
|
||||
await lkRoom.localParticipant.unpublishTrack(track);
|
||||
}
|
||||
|
||||
private observeTrackProcessors(
|
||||
scope: ObservableScope,
|
||||
room: LivekitRoom,
|
||||
@@ -441,5 +456,14 @@ export class Publisher {
|
||||
null,
|
||||
);
|
||||
trackProcessorSync(scope, track$, trackerProcessorState$);
|
||||
// Every camera track the SDK makes, on joining or on turning the camera
|
||||
// on, then starts with the pipeline on. One attached after publishing lets
|
||||
// the room through while it is built.
|
||||
trackerProcessorState$.pipe(scope.bind()).subscribe(({ processor }) => {
|
||||
room.options.videoCaptureDefaults = {
|
||||
...room.options.videoCaptureDefaults,
|
||||
processor,
|
||||
};
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user