From 02654d57f97e54cd375733460421fb30da7e4dfd Mon Sep 17 00:00:00 2001 From: fkwp Date: Fri, 25 Sep 2026 14:42:11 +0200 Subject: [PATCH] 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) --- .../BackgroundEffectTransformer.test.ts | 24 +++ ...rmer.ts => BackgroundEffectTransformer.ts} | 17 ++- src/livekit/OneStepPipeline.test.ts | 80 ++++++++++ src/livekit/OneStepPipeline.ts | 43 ++++++ src/livekit/TrackProcessorContext.test.ts | 143 +++++++++++++++++- src/livekit/TrackProcessorContext.tsx | 89 +++++------ src/state/BackgroundEffects.test.ts | 131 ++++++++++++++++ src/state/BackgroundEffects.ts | 93 ++++++++++++ .../localMember/Publisher.test.ts | 70 +++++++++ .../CallViewModel/localMember/Publisher.ts | 28 +++- 10 files changed, 665 insertions(+), 53 deletions(-) create mode 100644 src/livekit/BackgroundEffectTransformer.test.ts rename src/livekit/{BlurBackgroundTransformer.ts => BackgroundEffectTransformer.ts} (78%) create mode 100644 src/livekit/OneStepPipeline.test.ts create mode 100644 src/livekit/OneStepPipeline.ts create mode 100644 src/state/BackgroundEffects.test.ts create mode 100644 src/state/BackgroundEffects.ts diff --git a/src/livekit/BackgroundEffectTransformer.test.ts b/src/livekit/BackgroundEffectTransformer.test.ts new file mode 100644 index 000000000..b8a26b86b --- /dev/null +++ b/src/livekit/BackgroundEffectTransformer.test.ts @@ -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(); + }); +}); diff --git a/src/livekit/BlurBackgroundTransformer.ts b/src/livekit/BackgroundEffectTransformer.ts similarity index 78% rename from src/livekit/BlurBackgroundTransformer.ts rename to src/livekit/BackgroundEffectTransformer.ts index f86120d33..965359211 100644 --- a/src/livekit/BlurBackgroundTransformer.ts +++ b/src/livekit/BackgroundEffectTransformer.ts @@ -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 { + 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); } } diff --git a/src/livekit/OneStepPipeline.test.ts b/src/livekit/OneStepPipeline.test.ts new file mode 100644 index 000000000..12aa94773 --- /dev/null +++ b/src/livekit/OneStepPipeline.test.ts @@ -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[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(); + }); +}); diff --git a/src/livekit/OneStepPipeline.ts b/src/livekit/OneStepPipeline.ts new file mode 100644 index 000000000..e2060a0ea --- /dev/null +++ b/src/livekit/OneStepPipeline.ts @@ -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[0]; +type DestroyOptions = Parameters[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 = Promise.resolve(); + + public override async init(opts: InitOptions): Promise { + 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 { + return this.inTurn(async () => { + await super.destroy({ willProcessorRestart: true }); + await super.init(opts); + }); + } + + public override async destroy(options?: DestroyOptions): Promise { + return this.inTurn(async () => super.destroy(options)); + } + + private async inTurn(step: () => Promise): Promise { + const turn = this.steps.then(step); + this.steps = turn.catch(() => {}); + return turn; + } +} diff --git a/src/livekit/TrackProcessorContext.test.ts b/src/livekit/TrackProcessorContext.test.ts index 5b977bd4d..c873276e0 100644 --- a/src/livekit/TrackProcessorContext.test.ts +++ b/src/livekit/TrackProcessorContext.test.ts @@ -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) { + 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; 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 => createElement(Surfaces, { tracks }); + const blur = async (on: boolean): Promise => { + 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); + }); +}); diff --git a/src/livekit/TrackProcessorContext.tsx b/src/livekit/TrackProcessorContext.tsx index 96897929f..7ee7cf562 100644 --- a/src/livekit/TrackProcessorContext.tsx +++ b/src/livekit/TrackProcessorContext.tsx @@ -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; }; -const ProcessorContext = createContext(undefined); +const ProcessorContext = createContext( + 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 { - const state = use(ProcessorContext); - if (state === undefined) +export function useTrackProcessorObservable$(): Behavior { + 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 = ({ 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(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 {children}; + if (effects === null) return null; + return {children}; }; diff --git a/src/state/BackgroundEffects.test.ts b/src/state/BackgroundEffects.test.ts new file mode 100644 index 000000000..6be4fe8e8 --- /dev/null +++ b/src/state/BackgroundEffects.test.ts @@ -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 | 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; + let fake: ReturnType; + + function build( + options: Partial = {}, + ): BackgroundEffects { + return new BackgroundEffects(testScope(), { + supported: true, + blur$, + pipeline: fake.pipeline, + ...options, + }); + } + const blur = async (on: boolean): Promise => { + 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" }, + ]); + }); +}); diff --git a/src/state/BackgroundEffects.ts b/src/state/BackgroundEffects.ts new file mode 100644 index 000000000..5c50acc7d --- /dev/null +++ b/src/state/BackgroundEffects.ts @@ -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; + /** + * 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; + + public constructor( + scope: ObservableScope, + { supported, blur$, pipeline }: BackgroundEffectsOptions, + ) { + this.state$ = scope.behavior( + blur$.pipe( + scan( + (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 { + 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; + }; +} diff --git a/src/state/CallViewModel/localMember/Publisher.test.ts b/src/state/CallViewModel/localMember/Publisher.test.ts index 5d9f03784..6884ce0e3 100644 --- a/src/state/CallViewModel/localMember/Publisher.test.ts +++ b/src/state/CallViewModel/localMember/Publisher.test.ts @@ -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; + let state$: BehaviorSubject; + let publisher: Publisher; + /** The processor the SDK would give each camera track it made. */ + let madeWith: unknown[]; + + beforeEach(() => { + state$ = new BehaviorSubject({ + 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 as LocalTrackPublication); + + await publisher.createAndSetupTracks(); + videoEnabled$.next(true); + await flushPromises(); + + expect(localParticipant.unpublishTrack).toHaveBeenCalledWith(stopped); + expect(madeWith).toEqual([processor]); + }); +}); diff --git a/src/state/CallViewModel/localMember/Publisher.ts b/src/state/CallViewModel/localMember/Publisher.ts index 5353a98d4..176e3c2e0 100644 --- a/src/state/CallViewModel/localMember/Publisher.ts +++ b/src/state/CallViewModel/localMember/Publisher.ts @@ -63,7 +63,7 @@ export class Publisher { private connection: Pick, //setE2EEEnabled, devices: MediaDevices, private readonly muteStates: MuteStates, - trackerProcessorState$: Behavior, + private readonly trackerProcessorState$: Behavior, 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 { + 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, + }; + }); } }