Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
201 changes: 199 additions & 2 deletions js/publish/src/audio/encoder.test.ts
Original file line number Diff line number Diff line change
@@ -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" }));
Expand Down Expand Up @@ -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<AudioFrame>;
#deliver: ((frame: AudioFrame) => void) | undefined;
#requested!: () => void;
#request = this.#next();

constructor() {
this.stream = new ReadableStream<AudioFrame>(
{
pull: (controller) =>
new Promise<void>((resolve) => {
this.#deliver = (frame) => {
controller.enqueue(frame);
resolve();
};
this.#requested();
}),
},
{ highWaterMark: 0 },
);
}

#next(): Promise<void> {
return new Promise((resolve) => {
this.#requested = resolve;
});
}

// Resolves once every frame pushed so far has been processed.
async drain(): Promise<void> {
await this.#request;
}

async push(frame: AudioFrame): Promise<void> {
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<void>((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<Moq.Track.Producer | undefined>(track),
close: () => track.close(),
};

const feed = new Feed();
const capture = {
in: { source: new Signal(undefined) },
out: {
root: new Signal(undefined),
format: new Signal<Format>({ 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<void>((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();
}
});
26 changes: 18 additions & 8 deletions js/publish/src/audio/encoder.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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);

Expand All @@ -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;
}
},
};
Expand Down
5 changes: 5 additions & 0 deletions js/publish/src/audio/framer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading