diff --git a/js/publish/src/audio/encoder.test.ts b/js/publish/src/audio/encoder.test.ts index 076d29d4fd..76be253266 100644 --- a/js/publish/src/audio/encoder.test.ts +++ b/js/publish/src/audio/encoder.test.ts @@ -1,8 +1,10 @@ import { describe, expect, mock, test } from "bun:test"; import * as Catalog from "@moq/hang/catalog"; +import * as Moq from "@moq/net"; import { Time } from "@moq/net"; -import type { Format } from "./capture"; -import { resolve } from "./encoder"; +import { Signal } from "@moq/signals"; +import type { AudioFrame, Format } from "./capture"; +import { Encoder, resolve } from "./encoder"; // Bun does not load Vite's worklet URL imports from the public audio entrypoint. mock.module("./capture-worklet.ts?worklet", () => ({ default: "blob:fake-capture" })); @@ -52,3 +54,198 @@ describe("resolve", () => { expect(resolved.catalog.jitter).toBe(Catalog.u53(Math.ceil((1024 / 48_000) * 1000))); }); }); + +// Like Chrome's Opus encoder, it holds the newest chunks until later input pushes them out. +class LaggingAudioEncoder { + static readonly LAG = 2; + + // Called on configure; the encoder publishes its pipeline synchronously right after. + static onConfigure: (() => void) | undefined; + + state: CodecState = "unconfigured"; + #output: EncodedAudioChunkOutputCallback; + #held: { timestamp: number; duration: number }[] = []; + + constructor(init: AudioEncoderInit) { + this.#output = init.output; + } + + configure(): void { + this.state = "configured"; + LaggingAudioEncoder.onConfigure?.(); + } + + encode(data: AudioData): void { + const duration = Math.round((data.numberOfFrames / data.sampleRate) * 1_000_000); + this.#held.push({ timestamp: data.timestamp, duration }); + while (this.#held.length > LaggingAudioEncoder.LAG) { + const { timestamp, duration } = this.#held.shift() as { timestamp: number; duration: number }; + const chunk = { + type: "key", + timestamp, + duration, + byteLength: 1, + copyTo: (dest: Uint8Array) => dest.set([1]), + }; + this.#output(chunk as unknown as EncodedAudioChunk); + } + } + + close(): void { + this.state = "closed"; + } +} + +class FakeAudioData { + readonly timestamp: number; + readonly numberOfFrames: number; + readonly sampleRate: number; + + constructor(init: AudioDataInit) { + // WebIDL's `long long` conversion truncates a fractional timestamp. + this.timestamp = Math.trunc(init.timestamp); + this.numberOfFrames = init.numberOfFrames; + this.sampleRate = init.sampleRate; + } + + close(): void {} +} + +function installFakeWebCodecs() { + const names = ["AudioEncoder", "AudioDecoder", "AudioData"] as const; + const originals = names.map((name) => Object.getOwnPropertyDescriptor(globalThis, name)); + const fakes = [LaggingAudioEncoder, class {}, FakeAudioData]; + names.forEach((name, i) => { + Object.defineProperty(globalThis, name, { configurable: true, writable: true, value: fakes[i] }); + }); + + return { + [Symbol.dispose]() { + names.forEach((name, i) => { + const original = originals[i]; + if (original) Object.defineProperty(globalThis, name, original); + else Reflect.deleteProperty(globalThis, name); + }); + }, + }; +} + +// A capture stream that hands over one frame per read. The reader pushes each frame through the +// pipeline before reading again, so a pending read proves the previous frame was fully processed. +class Feed { + readonly stream: ReadableStream; + #deliver: ((frame: AudioFrame) => void) | undefined; + #requested!: () => void; + #request = this.#next(); + + constructor() { + this.stream = new ReadableStream( + { + pull: (controller) => + new Promise((resolve) => { + this.#deliver = (frame) => { + controller.enqueue(frame); + resolve(); + }; + this.#requested(); + }), + }, + { highWaterMark: 0 }, + ); + } + + #next(): Promise { + return new Promise((resolve) => { + this.#requested = resolve; + }); + } + + // Resolves once every frame pushed so far has been processed. + async drain(): Promise { + await this.#request; + } + + async push(frame: AudioFrame): Promise { + await this.drain(); + this.#request = this.#next(); + this.#deliver?.(frame); + } +} + +// The encoder outlives a demand gap, so chunks it held when demand disappeared surface after the +// resume. Written after the marker, they would put pre-gap media on the live edge, and a rounding +// step below the marker aborts every subscriber. +test("a demand gap marks where submitted audio ends and drops the chunks held across it", async () => { + using _webcodecs = installFakeWebCodecs(); + const configured = new Promise((resolve) => { + LaggingAudioEncoder.onConfigure = resolve; + }); + + const track = new Moq.Track.Producer("audio").accept(); + const written: [number, number][] = []; + let onWrite: (() => void) | undefined; + const writeFrame = track.writeFrame.bind(track); + track.writeFrame = (frame) => { + const [timestamp, payload] = Moq.Varint.decode(frame.payload); + written.push([timestamp, payload.byteLength]); + writeFrame(frame); + onWrite?.(); + }; + + const rendition = { + config: new Signal(undefined), + track: new Signal(track), + close: () => track.close(), + }; + + const feed = new Feed(); + const capture = { + in: { source: new Signal(undefined) }, + out: { + root: new Signal(undefined), + format: new Signal({ sampleRate: 48_000, channelCount: 1 }), + frames: new Signal({ subscribe: () => feed.stream }), + }, + }; + + const encoder = new Encoder("audio", { + broadcast: { audio: () => rendition } as never, + capture: capture as never, + }); + + // One 20ms Opus frame per push, on a clock with a fractional microsecond origin. + let index = 0; + const push = async (count: number) => { + for (let i = 0; i < count; i++, index++) { + await feed.push({ timestamp: Time.Micro(18_699.6 + index * 20_000), channels: [new Float32Array(960)] }); + } + await feed.drain(); + }; + + try { + await configured; + await push(4); // two written, two held + + const marked = new Promise((resolve) => { + onWrite = resolve; + }); + rendition.track.set(undefined); + await marked; + onWrite = undefined; + + await push(2); // gated + rendition.track.set(track); + await push(4); // releases the two held pre-gap chunks, then two resumed ones + + expect(written).toEqual([ + [18_700, 1], + [38_700, 1], + [98_700, 0], + [138_700, 1], + [158_700, 1], + ]); + } finally { + LaggingAudioEncoder.onConfigure = undefined; + encoder.close(); + } +}); diff --git a/js/publish/src/audio/encoder.ts b/js/publish/src/audio/encoder.ts index 9bfabb30f1..d86a973bee 100644 --- a/js/publish/src/audio/encoder.ts +++ b/js/publish/src/audio/encoder.ts @@ -155,9 +155,13 @@ export class Encoder { // discontinuity and re-anchors on. #pipeline: Pipeline | undefined; - // The exclusive end of the newest frame written to the live track, where a demand gap's - // discontinuity marker goes. Cleared once the marker is written. - #end: Time.Micro | undefined; + // Where the next frame submitted to the AudioEncoder starts, i.e. the exclusive end of the + // newest one, where a demand gap's discontinuity marker goes. Cleared once the marker is written. + #next: Time.Micro | undefined; + + // The newest demand gap's marker. The AudioEncoder outlives the gap, so chunks it still held + // when demand disappeared surface after the resume; they sit below the marker and are dropped. + #floor: Time.Micro | undefined; // The fatal error an AudioEncoder reported, if any. That instance can never encode again and // reconfiguring it would be a retry, so the rendition stays down for the life of this encoder. @@ -264,14 +268,15 @@ export class Encoder { // When demand disappears, end the epoch with a discontinuity marker (see // Container.Legacy.Producer.cut) so a later subscriber resumes on the same track without the - // pre-gap frames reading as live. Its empty payload marks where the source media ends. + // pre-gap frames reading as live. Its empty payload marks where the submitted media ends. effect.run((effect) => { const track = effect.get(rendition.track); if (!track) return; effect.cleanup(() => { - const end = this.#end; - this.#end = undefined; + const end = this.#next; + this.#next = undefined; if (end === undefined || track.closed.peek() !== undefined) return; + this.#floor = end; track.writeFrame({ payload: Container.Legacy.encodeFrame(new Uint8Array(), end), timestamp: Time.Timestamp.fromMicros(end), @@ -407,11 +412,11 @@ export class Encoder { // waiting for a group boundary. Loss is handled by the codec's PLC. const live = track.peek(); if (!live) return; + if (this.#floor !== undefined && frame.timestamp < this.#floor) return; live.writeFrame({ payload: Container.Legacy.encodeFrame(frame, frame.timestamp as Time.Micro), timestamp: Time.Timestamp.fromMicros(frame.timestamp as Time.Micro), }); - this.#end = (frame.timestamp + (frame.duration ?? 0)) as Time.Micro; }, error: (err) => { console.error("encoder error", err); @@ -439,6 +444,10 @@ export class Encoder { // on the capture clock, but there is nowhere to send a chunk with no subscriber. if (!track.peek()) continue; + // Round to whole microseconds once, here, so a chunk's timestamp and the marker + // placed at the next frame's start agree exactly. + const timestamp = Math.round(data.timestamp) as Time.Micro; + const joinedLength = data.channels.reduce((total, channel) => total + channel.length, 0); const joined = new Float32Array(joinedLength); @@ -452,13 +461,14 @@ export class Encoder { sampleRate: config.sampleRate, numberOfFrames: data.channels[0].length, numberOfChannels: data.channels.length, - timestamp: data.timestamp, + timestamp, data: joined, transfer: [joined.buffer], }); encoder.encode(frame); frame.close(); + this.#next = Math.round(framer.next) as Time.Micro; } }, }; diff --git a/js/publish/src/audio/framer.ts b/js/publish/src/audio/framer.ts index b1ca2ef6c3..ad2352361b 100644 --- a/js/publish/src/audio/framer.ts +++ b/js/publish/src/audio/framer.ts @@ -80,6 +80,11 @@ export class Framer { return output; } + /** Where the next frame will start. Throws until the first input sets the origin. */ + get next(): Time.Micro { + return this.#timestamp(); + } + // Whether this chunk starts somewhere other than where the previous one left off. #discontinuous(timestamp: Time.Micro): boolean { if (this.#origin === undefined) return false;