diff --git a/Cargo.lock b/Cargo.lock index d83041b540..a4871f50bd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4494,6 +4494,7 @@ dependencies = [ "getrandom 0.4.3", "kio 0.6.1", "loom", + "moq-net-sim", "moq-pattern", "num_enum", "rand 0.10.3", @@ -4516,6 +4517,19 @@ dependencies = [ "moq-net", ] +[[package]] +name = "moq-net-sim" +version = "0.0.0" +dependencies = [ + "futures", + "kio 0.6.1", + "moq-net-sim-macros", +] + +[[package]] +name = "moq-net-sim-macros" +version = "0.0.0" + [[package]] name = "moq-noq" version = "2.0.0" diff --git a/Cargo.toml b/Cargo.toml index e58b93783a..4065deb5dc 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -22,6 +22,8 @@ members = [ "rs/moq-native", "rs/moq-net", "rs/moq-net/fuzz", + "rs/moq-net/sim", + "rs/moq-net/sim/macros", "rs/moq-nvenc", "rs/moq-pattern", "rs/moq-relay", diff --git a/js/net/src/lite/group.ts b/js/net/src/lite/group.ts index 2bb058051f..dd85e2de6e 100644 --- a/js/net/src/lite/group.ts +++ b/js/net/src/lite/group.ts @@ -64,7 +64,7 @@ export class Group { } } -/** Decode an unsigned zigzag varint back to a signed delta (mirrors Rust `VarInt::to_zigzag`). */ +/** Decode an unsigned zigzag varint back to a signed delta (mirrors Rust `varint::unzigzag`). */ function unzigzag(v: bigint): bigint { return (v >> 1n) ^ -(v & 1n); } diff --git a/js/net/src/lite/publisher.ts b/js/net/src/lite/publisher.ts index 7ed9622222..5cb11cfc55 100644 --- a/js/net/src/lite/publisher.ts +++ b/js/net/src/lite/publisher.ts @@ -43,7 +43,7 @@ const PROBE_MAX_AGE = 10_000; // ms const PROBE_MAX_DELTA = 0.25; const PROBE_RTT_DELTA = 0.25; -/** Map a signed delta to an unsigned zigzag varint value (mirrors Rust `VarInt::from_zigzag`). */ +/** Map a signed delta to an unsigned zigzag varint value (mirrors Rust `varint::zigzag`). */ function zigzag(delta: bigint): bigint { return delta >= 0n ? delta << 1n : (-delta << 1n) - 1n; } diff --git a/quest/m1/rs2ts/README.md b/quest/m1/rs2ts/README.md index e37edf9eb0..4547ca0a01 100644 --- a/quest/m1/rs2ts/README.md +++ b/quest/m1/rs2ts/README.md @@ -30,9 +30,12 @@ Decided in planning (2026-09-27), with the spike data in and bytes out, no runtime. The async helper methods move behind an `async` cargo feature; rs2ts reads the crate without it and JS reimplements the helpers with Promises. No second crate. -- Varints are not bounded to 2^53 on the wire: 62 bits in QUIC form, 64 in - leading-ones form. JS holds any `u64` as a `U64` with checked conversion to - and from `number`; varints are only its wire encoding. +- Values are plain `u64` in Rust, and varint is a wire encoding in the codec, + not a type. The spec is not bounded to 2^53: the leading-ones form + (moq-transport draft-17+) carries all 64 bits, and the QUIC form (moq-lite, + drafts 14-16) refuses anything past 2^62 - 1 rather than truncating. Rust + `u64` maps to a TypeScript `U64` (two `u32` halves), generically, with + checked conversion to and from `number`. - The generated TypeScript is committed and a CI lane regenerates it and fails on drift, so JS contributors and npm publishing never need the nightly toolchain Charon pins. It lives inside js/net and `@moq/net` stays @@ -41,7 +44,8 @@ Decided in planning (2026-09-27), with the spike data in is no worse to use; watch, publish, hang, and the demos update in the same change. - Parity: `just test interop --all`, plus moq-net's own tests translated with - the code once they run on a mock clock instead of tokio. + the code. They run on simulated time with no runtime (`moq-net-sim`), so the + async-free ones translate as they stand. - The line lands on `dev`: the Rust refactors break moq-net's published API, and the translator and generated code build on them. Only the additive JS `U64` (`js/net/src/util/u64.ts`) is on `main`, package-internal. @@ -54,13 +58,13 @@ js/net it replaces, measured with the [browser benchmarks](/quest/m1/browser-ben ## Required -- [VarInt codec](/quest/m1/rs2ts/varint-codec.md) - moq-net encodes through a `VarInt` newtype and a concrete slice-based codec, not generic traits on primitives - [rs2ts](/quest/m1/rs2ts/translator.md) - a Charon-based translator emits readable TypeScript for moq-net's lite codec, committed and checked for drift in CI - [Sans-IO moq-net](/quest/m1/rs2ts/sans-io/README.md) - moq-net builds and runs without a runtime; async helpers sit behind an `async` feature -- [Mock-clock tests](/quest/m1/rs2ts/mock-clock.md) - moq-net's tests run on the sans-IO clock instead of tokio, so they translate with the code - [Generated lite](/quest/m1/rs2ts/lite.md) - @moq/net's lite session and model layer are generated from moq-net +- [IETF parameters](/quest/m1/rs2ts/ietf-params.md) - the IETF codec drops its `Param` trait on primitives, so it translates like lite - [Generated IETF](/quest/m1/rs2ts/ietf.md) - @moq/net's moq-transport session is generated too - [Remove moq-wasm](/quest/m1/rs2ts/remove-wasm.md) - the WASM experiment is deleted once generated lite ships +- [Browser benchmarks](/quest/m1/browser-benchmarks.md) - the harness the no-downgrade report uses ## Closes @@ -68,10 +72,6 @@ js/net it replaces, measured with the [browser benchmarks](/quest/m1/browser-ben - [#2822](https://github.com/moq-dev/moq/issues/2822) - close this issue when the quest finishes - [#2835](https://github.com/moq-dev/moq/issues/2835) - close this issue when the quest finishes -## Required - -- [Browser benchmarks](/quest/m1/browser-benchmarks.md) - the harness the no-downgrade report uses - ## Related - [#2850](/quest/m1/2850-js-net-give-reader-a-synchronous-decode-so-the-publisher.md) - the same synchronous decode shape, in hand-written js/net today diff --git a/quest/m1/rs2ts/ietf-params.md b/quest/m1/rs2ts/ietf-params.md new file mode 100644 index 0000000000..d31886dd59 --- /dev/null +++ b/quest/m1/rs2ts/ietf-params.md @@ -0,0 +1,25 @@ +# [S] IETF parameters without primitive traits + +## Goal + +moq-net's IETF message parameters encode and decode through concrete +per-kind methods, with no `Param` impls on `u8`, `bool`, `u64`, or +`Option`, so the IETF codec has the same concrete shape as the lite codec +the translator targets. + +## Plan + +The VarInt codec (#4463) removed Encode/Decode on +primitives, but `ietf/parameters.rs` still implements its own `Param` trait on +them, which rs2ts can only translate with dictionary passing. Replace the +impls with methods on the concrete `Decoder`/`Encoder` (or on the parameter +kinds), keeping the draft-14..16 varint cast and the draft-17+ forms byte for +byte. Benchmark IETF message encode and decode before and after with +`--bench codec`. + +Public API: none (`ietf` parameters are crate-private). Lands on `dev` with +the line. Wire: none. + +## Required + +- #4463 (VarInt codec) merged into the questline, since this builds on its `Decoder`/`Encoder` diff --git a/quest/m1/rs2ts/ietf.md b/quest/m1/rs2ts/ietf.md index b72b735753..3de2f1ec70 100644 --- a/quest/m1/rs2ts/ietf.md +++ b/quest/m1/rs2ts/ietf.md @@ -18,3 +18,4 @@ Public API: breaks `@moq/net`; retargets to `dev`. Wire: none. - [Generated lite](/quest/m1/rs2ts/lite.md) - the pipeline this reuses - [Sans-IO IETF session](/quest/m1/rs2ts/sans-io/ietf.md) - the session shape it translates +- [IETF parameters](/quest/m1/rs2ts/ietf-params.md) - the concrete parameter codec it translates diff --git a/quest/m1/rs2ts/lite.md b/quest/m1/rs2ts/lite.md index f42bccd742..d7234978df 100644 --- a/quest/m1/rs2ts/lite.md +++ b/quest/m1/rs2ts/lite.md @@ -27,6 +27,4 @@ Public API: breaks `@moq/net`; retargets to `dev`. Wire: none. - [rs2ts](/quest/m1/rs2ts/translator.md) - the translator - [Sans-IO lite session](/quest/m1/rs2ts/sans-io/lite.md) - the session shape it translates -- [Sans-IO model](/quest/m1/rs2ts/sans-io/model.md) - the model shape it translates - [The async feature](/quest/m1/rs2ts/sans-io/async-feature.md) - rs2ts reads moq-net without it -- [Mock-clock tests](/quest/m1/rs2ts/mock-clock.md) - the tests that prove parity diff --git a/quest/m1/rs2ts/mock-clock.md b/quest/m1/rs2ts/mock-clock.md deleted file mode 100644 index 3f1ea344c9..0000000000 --- a/quest/m1/rs2ts/mock-clock.md +++ /dev/null @@ -1,26 +0,0 @@ -# [L] Mock-clock tests - -## Goal - -moq-net's tests run on the sans-IO clock instead of tokio, so they need no -runtime and rs2ts translates them alongside the code. The generated -TypeScript runs the same tests under bun. - -## Plan - -tokio is in moq-net's tests only for paused, advanceable time -(`start_paused`, `advance`): 591 `tokio::test`s on dev, 86 of them paused. -Drive them from the model's injectable clock and a small synchronous -executor instead. - -Guidance: - -- Port mechanically where possible; keep each test's assertions unchanged. -- Tests that exercise the `async` helpers stay behind that feature and are - not translated. - -Public API: none. Wire: none. - -## Required - -- [Sans-IO model](/quest/m1/rs2ts/sans-io/model.md) - supplies the clock seam diff --git a/quest/m1/rs2ts/sans-io/README.md b/quest/m1/rs2ts/sans-io/README.md index 03a4cfb82b..6ee16d7ab6 100644 --- a/quest/m1/rs2ts/sans-io/README.md +++ b/quest/m1/rs2ts/sans-io/README.md @@ -18,6 +18,5 @@ The line has no work of its own beyond its children. ## Required - [Sans-IO lite session](/quest/m1/rs2ts/sans-io/lite.md) - the lite session is driven by bytes, stream events, and `tick(now)` -- [Sans-IO model](/quest/m1/rs2ts/sans-io/model.md) - origin, broadcast, track, and group handles run without a runtime, with time supplied by the caller - [The async feature](/quest/m1/rs2ts/sans-io/async-feature.md) - the async helpers sit behind an `async` feature and a CI lane builds and tests moq-net without it - [Sans-IO IETF session](/quest/m1/rs2ts/sans-io/ietf.md) - the moq-transport session is driven the same way as lite diff --git a/quest/m1/rs2ts/sans-io/async-feature.md b/quest/m1/rs2ts/sans-io/async-feature.md index f7b7e35046..eabcfbc21b 100644 --- a/quest/m1/rs2ts/sans-io/async-feature.md +++ b/quest/m1/rs2ts/sans-io/async-feature.md @@ -24,4 +24,3 @@ Wire: none. ## Required - [Sans-IO lite session](/quest/m1/rs2ts/sans-io/lite.md) - the session builds without a runtime -- [Sans-IO model](/quest/m1/rs2ts/sans-io/model.md) - the model builds without a runtime diff --git a/quest/m1/rs2ts/sans-io/ietf.md b/quest/m1/rs2ts/sans-io/ietf.md index 542c1965a1..40a07401f4 100644 --- a/quest/m1/rs2ts/sans-io/ietf.md +++ b/quest/m1/rs2ts/sans-io/ietf.md @@ -13,6 +13,8 @@ twice (for `Session` and `ControlStreamAdapter`); collapse that while here. If the [async feature](/quest/m1/rs2ts/sans-io/async-feature.md) landed first with the IETF session behind it, move the session out. +Turn the session's async test bodies into synchronous `poll_*` tests with +explicit instants as it is rewritten, like lite. Public API: breaks moq-net's IETF session API; retargets to `dev`. Wire: none. diff --git a/quest/m1/rs2ts/sans-io/lite.md b/quest/m1/rs2ts/sans-io/lite.md index 51e426c403..60da5b369d 100644 --- a/quest/m1/rs2ts/sans-io/lite.md +++ b/quest/m1/rs2ts/sans-io/lite.md @@ -22,5 +22,7 @@ Guidance: - Stream handles, write backpressure, and close codes become explicit events or return values the driver acts on. - Keep maps as Vec slabs where the key space is small. +- Turn the session's async test bodies into synchronous `poll_*` tests with + explicit instants as it is rewritten, so they translate with the code. Public API: breaks moq-net's session API; retargets to `dev`. Wire: none. diff --git a/quest/m1/rs2ts/sans-io/model.md b/quest/m1/rs2ts/sans-io/model.md deleted file mode 100644 index 9060397847..0000000000 --- a/quest/m1/rs2ts/sans-io/model.md +++ /dev/null @@ -1,18 +0,0 @@ -# [L] Sans-IO model - -## Goal - -The origin, broadcast, track, group, and frame producers and consumers run -without an async runtime. Their waiting is poll-based on kio, and anything -time-based reads a clock the caller supplies, so the model translates to -TypeScript and its tests can run on a mock clock. - -## Plan - -The model is already poll-based on kio waiters; what remains is every place -that reaches a runtime or wall clock directly (`runtime::Deadline`, the cache -pool's expiry, stats timers). Route time through one injectable clock, the -seam the [mock-clock tests](/quest/m1/rs2ts/mock-clock.md) use. - -Public API: may break moq-net's model constructors; retargets to `dev`. -Wire: none. diff --git a/quest/m1/rs2ts/translator.md b/quest/m1/rs2ts/translator.md index 2f9f222427..327215aaa2 100644 --- a/quest/m1/rs2ts/translator.md +++ b/quest/m1/rs2ts/translator.md @@ -29,17 +29,23 @@ Mapping decided in planning: `[Symbol.dispose]`; `Arc`/`Rc` of a type with drop glue become an explicit refcount. JS is single-threaded, so `Mutex` and atomics become plain access. -- Rust `u64` (and `VarInt` while it lasts) maps to js/net's `U64` - (`js/net/src/util/u64.ts`), read and written as a varint by - `Cursor.varint()` and `Writer.varint()`. - Integers up to 32 bits and `usize` map to `number` with checked arithmetic - that throws on overflow; never wrap silently. A `u64` or `i64` never maps - to a lossy `number`: the model accepts `u64::MAX` (e.g. - `model/subscription.rs`), so each one either becomes `U64` or an - `Option` in the source, or maps to a full-width 64-bit TypeScript type. +- Rust `u64` maps to js/net's `U64` (`js/net/src/util/u64.ts`, two `u32` + halves), generically, read and written as a varint by `Cursor.varint()` and + `Writer.varint()`. Integers up to 32 bits and `usize` map to `number` with + checked arithmetic that throws on overflow; never wrap silently. A `u64` or + `i64` never maps to a lossy `number`: the model accepts `u64::MAX` (e.g. + `model/subscription.rs`). Varint is a wire encoding in the codec, not a + type, so nothing maps by that name. Guidance: +- Keep rs2ts a generic translator for a Rust subset, not a moq-net tool. It + maps by Rust type and construct (e.g. `u64` to one 64-bit TypeScript type), + never by moq-net names; anything project-specific lives in moq-net's source + or a small config, so another crate could use rs2ts unchanged. +- `varint::zigzag` and `unzigzag` (lite per-frame timestamps) still use + 64-bit bit math, which the subset forbids. Rewrite them on two `u32` halves, + or give the 64-bit TypeScript type the operations they need. - Pin Charon and its nightly in the nix shell for the regeneration lane only. - Borrow rust-js's MIT oxc printer for formatting and source maps. - Readability pass: inline single-use temporaries and keep source branch @@ -54,7 +60,3 @@ Guidance: Public API: none (internal tool). Lands on `dev` with the codec it translates. Wire: none. - -## Required - -- [VarInt codec](/quest/m1/rs2ts/varint-codec.md) - the codec shape the translator targets diff --git a/quest/m1/rs2ts/varint-codec.md b/quest/m1/rs2ts/varint-codec.md deleted file mode 100644 index 39557cb31b..0000000000 --- a/quest/m1/rs2ts/varint-codec.md +++ /dev/null @@ -1,31 +0,0 @@ -# [M] VarInt codec - -## Goal - -Every varint moq-net puts on the wire is a `VarInt`, and `VarInt` is the only -integer type with Encode/Decode. Messages encode and decode through a -concrete, slice-based codec instead of generic traits implemented on `u64`, -`usize`, `bool`, `String`, `Option`, and `Vec`. - -## Plan - -Two reasons, one refactor. Not every `u64` is a valid varint, so the type -should say which fields are. And the rs2ts translator cannot map generic -traits on primitives without dictionary passing, its most expensive feature; -the same generics are most of why moq-net's lite codec built to 47 KB gzip in -WASM against 6 KB for a hand-carved one. - -Guidance: - -- Keep the 62-bit range; the wire is not bounded to 2^53. -- Message types keep a local trait; the primitives become inherent methods - on concrete reader and writer types (`varint`, `string`, `bytes`, ...). - Make the version a concrete type rather than a generic `V` where possible. -- Avoid bit operations on `u64` in the codec: write the 8-byte form as two - `u32` halves, so the generated TypeScript never needs 64-bit bitwise math. -- `Parameters` becomes Vec-backed, and decode paths stop branching on - `tracing::enabled!` (log after decoding instead). -- Benchmark the codec before and after (Criterion); it is on every message. - -Public API: breaks moq-net's `coding` module (Encode/Decode on primitives -go away), so this retargets to `dev`. Wire: none. diff --git a/rs/AGENTS.md b/rs/AGENTS.md index 06c123a21c..e4259e3414 100644 --- a/rs/AGENTS.md +++ b/rs/AGENTS.md @@ -50,7 +50,7 @@ Prefer poll. New logic is a `poll_*` with an `async` helper, not the other way a # Testing -- Tests are inline `#[cfg(test)] mod tests`. Time-dependent async tests call `tokio::time::pause()` first, unless they cross real networking that can't be mocked (sockets, smoke tests); those run on the wall clock and assert lower bounds. +- Tests are inline `#[cfg(test)] mod tests`. Time-dependent async tests call `tokio::time::pause()` first (moq-net uses `#[moq_net_sim::test]` instead), unless they cross real networking that can't be mocked (sockets, smoke tests); those run on the wall clock and assert lower bounds. - Run tests through `just` (nextest), not `cargo test`: nextest kills a wedged test as TIMEOUT, cargo hangs forever. A test flagged SLOW is a bug to fix, not a threshold to raise. - `just check` compiles default features only, like CI. `just rs features` (nightly) covers `--all-features` / `--no-default-features`. Keep a feature gate around the dependency, not the logic, so the logic's tests stay in the merge gate. - Local checks compile only the host platform; PR CI runs `just rs windows` / `macos` on those hosts, and `just rs wasm` covers `moq-wasm`. diff --git a/rs/hang/src/catalog/container.rs b/rs/hang/src/catalog/container.rs index 126f42601a..2466f3dbd2 100644 --- a/rs/hang/src/catalog/container.rs +++ b/rs/hang/src/catalog/container.rs @@ -15,7 +15,7 @@ use serde_with::{base64::Base64, serde_as}; /// rendition must be ignored by consumers. #[derive(Debug, Clone, PartialEq, Default)] pub enum Container { - /// A QUIC VarInt timestamp prefix followed by the raw codec payload. + /// A QUIC varint timestamp prefix followed by the raw codec payload. /// Timestamps are in microseconds. #[default] Legacy, diff --git a/rs/hang/src/container/frame.rs b/rs/hang/src/container/frame.rs index 77c7ae8955..384f2ed2a3 100644 --- a/rs/hang/src/container/frame.rs +++ b/rs/hang/src/container/frame.rs @@ -1,7 +1,6 @@ use super::MAX_AGE; use bytes::{Buf, BufMut, Bytes, BytesMut}; use derive_more::Debug; -use moq_net::VarInt; use crate::Error; @@ -9,7 +8,7 @@ pub use moq_net::{Timescale, Timestamp}; /// Canonical timescale for the hang legacy wire format: microseconds. /// -/// The legacy container's on-wire timestamp is a single VarInt with no scale tag, +/// The legacy container's on-wire timestamp is a single varint with no scale tag, /// so encoders normalize to this scale and decoders attach it. pub const TIMESCALE: Timescale = Timescale::MICRO; @@ -67,7 +66,7 @@ pub struct Frame { } impl Frame { - /// Encode the frame: VarInt timestamp prefix followed by the raw codec payload. + /// Encode the frame: varint timestamp prefix followed by the raw codec payload. /// /// The timestamp is normalized to [`TIMESCALE`] (microseconds) so peers using a /// different source scale (e.g. nanoseconds from MKV) can decode without knowing @@ -78,12 +77,12 @@ impl Frame { Ok(()) } - /// Decode a frame from raw bytes (VarInt timestamp prefix + payload). + /// Decode a frame from raw bytes (varint timestamp prefix + payload). /// /// Attaches [`TIMESCALE`] (microseconds) to the decoded timestamp, matching what /// [`Self::encode`] writes. Inverse of [`Self::encode`]. pub fn decode(mut buf: impl Buf) -> Result { - let value: u64 = VarInt::decode_quic(&mut buf).map_err(moq_net::Error::from)?.into(); + let value: u64 = moq_net::varint::decode_quic(&mut buf).map_err(moq_net::Error::from)?; let timestamp = Timestamp::new(value, TIMESCALE)?; let payload = buf.copy_to_bytes(buf.remaining()); @@ -117,11 +116,10 @@ impl Frame { Ok(()) } - /// Write the VarInt timestamp prefix, normalized to [`TIMESCALE`]. + /// Write the varint timestamp prefix, normalized to [`TIMESCALE`]. fn encode_header(&self, buf: &mut impl BufMut) -> Result<(), Error> { let timestamp = self.timestamp.convert(TIMESCALE)?; - let value = VarInt::try_from(timestamp.value()).map_err(moq_net::Error::from)?; - value.encode_quic(buf).map_err(moq_net::Error::from)?; + moq_net::varint::encode_quic(timestamp.value(), buf).map_err(moq_net::Error::from)?; Ok(()) } diff --git a/rs/moq-archive/src/segment.rs b/rs/moq-archive/src/segment.rs index c2fd68c5aa..9b1cac5a89 100644 --- a/rs/moq-archive/src/segment.rs +++ b/rs/moq-archive/src/segment.rs @@ -1,7 +1,7 @@ use std::ops::RangeInclusive; use bytes::{Buf, BufMut, Bytes, BytesMut}; -use moq_net::VarInt; +use moq_net::varint; use crate::path::{check_id, check_range}; use crate::{Error, Result, VERSION}; @@ -203,12 +203,11 @@ fn validate(groups: &[Group]) -> Result<()> { } fn write_varint(buf: &mut impl BufMut, value: u64) -> Result<()> { - let value = VarInt::try_from(value).map_err(|_| Error::Overflow)?; - value.encode_quic(buf).map_err(|_| Error::Overflow) + varint::encode_quic(value, buf).map_err(|_| Error::Overflow) } fn read_varint(buf: &mut impl Buf) -> Result { - Ok(VarInt::decode_quic(buf).map_err(|_| Error::Table)?.into_inner()) + varint::decode_quic(buf).map_err(|_| Error::Table) } fn read_count(buf: &mut impl Buf, min_entry: usize) -> Result { diff --git a/rs/moq-c/src/api.rs b/rs/moq-c/src/api.rs index 809a4c53ab..fea565c195 100644 --- a/rs/moq-c/src/api.rs +++ b/rs/moq-c/src/api.rs @@ -15,7 +15,7 @@ use tracing::Level; #[allow(non_camel_case_types)] #[derive(Clone, Copy, Debug)] pub enum moq_container_kind { - /// A QUIC VarInt timestamp prefix followed by the raw codec payload. + /// A QUIC varint timestamp prefix followed by the raw codec payload. /// Timestamps are in microseconds. MOQ_CONTAINER_KIND_LEGACY = 0, /// Fragmented MP4: each frame is a complete moof+mdat fragment, described by diff --git a/rs/moq-c/src/video.rs b/rs/moq-c/src/video.rs index f180d124d6..f37df0d75c 100644 --- a/rs/moq-c/src/video.rs +++ b/rs/moq-c/src/video.rs @@ -968,7 +968,7 @@ pub unsafe extern "C" fn moq_decode_video_frame(id: u32, dst: *mut moq_video_fra let frame = State::lock().video.frame(id)?; let pixels = frame.pixels()?; *dst = moq_video_frame { - // The decoded Timestamp is bounded by a QUIC VarInt, so its microseconds fit. + // The decoded Timestamp is bounded by a QUIC varint, so its microseconds fit. timestamp_us: frame.frame.timestamp.as_micros() as u64, width: pixels.width, height: pixels.height, diff --git a/rs/moq-ffi/src/video.rs b/rs/moq-ffi/src/video.rs index e17779f610..537a76e2ed 100644 --- a/rs/moq-ffi/src/video.rs +++ b/rs/moq-ffi/src/video.rs @@ -625,7 +625,7 @@ pub struct MoqVideoDecodedFrame { impl MoqVideoDecodedFrame { /// Presentation timestamp, in microseconds. pub fn timestamp_us(&self) -> u64 { - // A decoded Timestamp is bounded by a QUIC VarInt, so its microseconds fit. + // A decoded Timestamp is bounded by a QUIC varint, so its microseconds fit. self.frame.timestamp.as_micros() as u64 } diff --git a/rs/moq-loc/src/lib.rs b/rs/moq-loc/src/lib.rs index af241919a5..1d19e7ab37 100644 --- a/rs/moq-loc/src/lib.rs +++ b/rs/moq-loc/src/lib.rs @@ -23,10 +23,10 @@ //! encode. Public properties are not handled here. They belong in the MoQ //! object header and are stripped by the transport layer. //! -//! Varint encoding is QUIC-style throughout via [`moq_net::VarInt`]. +//! Varint encoding is QUIC-style throughout via [`moq_net::varint`]. use bytes::{Buf, Bytes, BytesMut}; -use moq_net::{BoundsExceeded, DecodeError, EncodeError, VarInt}; +use moq_net::{BoundsExceeded, DecodeError, EncodeError, varint}; /// Property IDs recognized by this implementation. const PROP_TIMESCALE: u64 = 0x08; @@ -104,7 +104,7 @@ impl From for Error { /// Consumes the properties_length prefix, walks the bounded property block, /// and returns the remainder as `payload`. pub fn decode(mut buf: Bytes) -> Result { - let properties_length: u64 = VarInt::decode_quic(&mut buf)?.into(); + let properties_length = varint::decode_quic(&mut buf)?; let properties_length: usize = properties_length.try_into().map_err(|_| Error::MalformedProperties)?; if properties_length > buf.remaining() { @@ -119,7 +119,7 @@ pub fn decode(mut buf: Bytes) -> Result { let mut first = true; while props.has_remaining() { - let delta: u64 = VarInt::decode_quic(&mut props)?.into(); + let delta = varint::decode_quic(&mut props)?; let abs = if first { first = false; delta @@ -129,7 +129,7 @@ pub fn decode(mut buf: Bytes) -> Result { prev_type = abs; if abs % 2 == 0 { - let value: u64 = VarInt::decode_quic(&mut props)?.into(); + let value = varint::decode_quic(&mut props)?; match abs { PROP_TIMESTAMP | PROP_TIMESTAMP_DRAFT03 => timestamp = Some(value), PROP_TIMESCALE => { @@ -141,7 +141,7 @@ pub fn decode(mut buf: Bytes) -> Result { _ => {} } } else { - let len: u64 = VarInt::decode_quic(&mut props)?.into(); + let len = varint::decode_quic(&mut props)?; let len: usize = len.try_into().map_err(|_| Error::MalformedProperties)?; if len > props.remaining() { return Err(Error::MalformedProperties); @@ -167,11 +167,11 @@ pub fn decode(mut buf: Bytes) -> Result { /// catalog timescale to interpret `timestamp`. pub fn encode(timestamp: u64, payload: &[u8]) -> Result { let mut props = BytesMut::with_capacity(16); - VarInt::try_from(PROP_TIMESTAMP)?.encode_quic(&mut props)?; - VarInt::try_from(timestamp)?.encode_quic(&mut props)?; + varint::encode_quic(PROP_TIMESTAMP, &mut props)?; + varint::encode_quic(timestamp, &mut props)?; let mut out = BytesMut::with_capacity(props.len() + payload.len() + 8); - VarInt::try_from(props.len() as u64)?.encode_quic(&mut out)?; + varint::encode_quic(props.len() as u64, &mut out)?; out.extend_from_slice(&props); out.extend_from_slice(payload); @@ -184,7 +184,7 @@ mod tests { /// Test helper: write a u64 as a QUIC varint into `buf`. fn write_varint(buf: &mut BytesMut, value: u64) { - VarInt::try_from(value).unwrap().encode_quic(buf).unwrap(); + varint::encode_quic(value, buf).unwrap(); } #[test] diff --git a/rs/moq-mux/src/catalog/hang/container.rs b/rs/moq-mux/src/catalog/hang/container.rs index be57cfc8e2..f780a415e3 100644 --- a/rs/moq-mux/src/catalog/hang/container.rs +++ b/rs/moq-mux/src/catalog/hang/container.rs @@ -6,7 +6,7 @@ use crate::container::{Container as ContainerTrait, Frame, Kind, fmp4, legacy, l /// /// Built from a track's audio or video configuration, including its container. pub enum Container { - /// VarInt timestamp + raw codec bitstream. The original hang wire format. + /// varint timestamp + raw codec bitstream. The original hang wire format. Legacy(Kind), /// ISO-BMFF moof+mdat fragments. The wrapped [`fmp4::Wire`] holds /// the track's `trak` box so per-frame writes and reads have the diff --git a/rs/moq-mux/src/catalog/msf/consumer.rs b/rs/moq-mux/src/catalog/msf/consumer.rs index c8e2a256a0..1b488bb5d7 100644 --- a/rs/moq-mux/src/catalog/msf/consumer.rs +++ b/rs/moq-mux/src/catalog/msf/consumer.rs @@ -166,7 +166,7 @@ pub(crate) fn from_msf(msf: &moq_msf::Catalog) -> Result Result> { match &track.packaging { - // Neither is ISO-BMFF boxed, but they frame differently: a LOC property block against a VarInt + // Neither is ISO-BMFF boxed, but they frame differently: a LOC property block against a varint // timestamp prefix. Reading one as the other misparses the head of every frame. moq_msf::Packaging::Loc => Ok(Some(Container::Loc)), moq_msf::Packaging::Legacy => Ok(Some(Container::Legacy)), diff --git a/rs/moq-mux/src/container/legacy/mod.rs b/rs/moq-mux/src/container/legacy/mod.rs index 2307148f3b..c63ead824d 100644 --- a/rs/moq-mux/src/container/legacy/mod.rs +++ b/rs/moq-mux/src/container/legacy/mod.rs @@ -1,6 +1,6 @@ //! The original hang wire format. //! -//! Each moq frame holds one media frame: a VarInt-encoded timestamp +//! Each moq frame holds one media frame: a varint-encoded timestamp //! followed by the raw codec bitstream. Simple but ad-hoc; new //! broadcasts should use [`crate::container::loc`] instead. diff --git a/rs/moq-net/AGENTS.md b/rs/moq-net/AGENTS.md index 1fbbd48ffa..cdba32d624 100644 --- a/rs/moq-net/AGENTS.md +++ b/rs/moq-net/AGENTS.md @@ -11,6 +11,7 @@ The wire layer. Generic over the transport and media-agnostic: the relay never i # Testing +- Tests never use tokio or the wall clock. Prefer a plain `#[test]` that polls with explicit instants; async bodies use `#[moq_net_sim::test]`, whose simulated time jumps only when every task is idle. tokio is for the benches. - `just rs fuzz ` (one per file in `fuzz/fuzz_targets/`) needs nightly. The target bodies live in `src/fuzz.rs`; `fuzz/regressions//` replays under `just check`, so commit every crash input there. - `just rs loom` model-checks the concurrent handoffs. A hang is a lost wakeup, not a flake. - `just test interop --all` runs the cross-language interop matrix after a wire change. diff --git a/rs/moq-net/Cargo.toml b/rs/moq-net/Cargo.toml index 967afb259d..ac6573dd43 100644 --- a/rs/moq-net/Cargo.toml +++ b/rs/moq-net/Cargo.toml @@ -37,10 +37,12 @@ web-transport-trait = { workspace = true } [dev-dependencies] criterion = { workspace = true } +# The tests' executor, with simulated time. +moq-net-sim = { path = "sim" } serde_json = { workspace = true } -# test-util (tokio::time::pause/advance) is test-only and is NOT supported on -# wasm, so it must not leak into the normal dependency feature set. -tokio = { workspace = true, features = ["macros", "io-util", "sync", "test-util", "time", "rt"] } +# The benches run on tokio. test-util (tokio::time::pause/advance) is test-only and +# is NOT supported on wasm, so it must not leak into the normal dependency feature set. +tokio = { workspace = true, features = ["rt", "test-util", "time"] } # Model checks for the kio-backed model layer (tests/loom.rs). `cfg(loom)` is only # set by `just rs loom`, and the whole test file compiles away without it. @@ -75,6 +77,12 @@ name = "announce" harness = false required-features = ["fuzz"] +# Reaches the private message codecs through the hidden `fuzz` module. +[[bench]] +name = "codec" +harness = false +required-features = ["fuzz"] + [[bench]] name = "group" harness = false diff --git a/rs/moq-net/benches/codec.rs b/rs/moq-net/benches/codec.rs new file mode 100644 index 0000000000..39f6dad934 --- /dev/null +++ b/rs/moq-net/benches/codec.rs @@ -0,0 +1,72 @@ +//! The wire codec on its own: a fixed mix of moq-lite and moq-transport messages, and raw +//! varints in each wire form. Every message and frame header pays this cost. +//! +//! Run with `cargo bench -p moq-net --features fuzz --bench codec`. + +use criterion::{BenchmarkId, Criterion, Throughput, criterion_group, criterion_main}; +use moq_net::fuzz::{Messages, decode_varints, encode_varints}; + +type Encode = fn(&Messages, &mut Vec); +type Decode = fn(&Messages, &[u8]); + +/// Varints spanning the length classes of both wire forms. +fn varints() -> Vec { + (0..1_024u64) + .map(|n| match n % 4 { + 0 => n % 60, + 1 => 1_000 + n, + 2 => 1_000_000 + n, + _ => (1 << 40) + n, + }) + .collect() +} + +fn bench(c: &mut Criterion) { + let messages = Messages::default(); + + let mut group = c.benchmark_group("codec_messages"); + let protocols: [(&str, Encode, Decode); 2] = [ + ("lite", Messages::encode_lite, Messages::decode_lite), + ("ietf", Messages::encode_ietf, Messages::decode_ietf), + ]; + for (name, encode, decode) in protocols { + let mut encoded = Vec::new(); + encode(&messages, &mut encoded); + group.throughput(Throughput::Bytes(encoded.len() as u64)); + + group.bench_function(BenchmarkId::new("encode", name), |b| { + let mut out = Vec::with_capacity(encoded.len()); + b.iter(|| { + out.clear(); + encode(&messages, &mut out); + }); + }); + group.bench_function(BenchmarkId::new("decode", name), |b| { + b.iter(|| decode(&messages, &encoded)) + }); + } + group.finish(); + + let values = varints(); + let mut group = c.benchmark_group("codec_varint"); + group.throughput(Throughput::Elements(values.len() as u64)); + for (name, ietf) in [("quic", false), ("leading_ones", true)] { + let mut encoded = Vec::new(); + encode_varints(&values, ietf, &mut encoded); + + group.bench_function(BenchmarkId::new("encode", name), |b| { + let mut out = Vec::with_capacity(encoded.len()); + b.iter(|| { + out.clear(); + encode_varints(&values, ietf, &mut out); + }); + }); + group.bench_function(BenchmarkId::new("decode", name), |b| { + b.iter(|| decode_varints(&encoded, ietf)) + }); + } + group.finish(); +} + +criterion_group!(benches, bench); +criterion_main!(benches); diff --git a/rs/moq-net/benches/session.rs b/rs/moq-net/benches/session.rs index 77df70de90..0441c047fa 100644 --- a/rs/moq-net/benches/session.rs +++ b/rs/moq-net/benches/session.rs @@ -255,7 +255,7 @@ impl Cluster { let mut config = origin::Config::new(Hop::new(self.next_hop).unwrap()); config.pool = cache::Pool::new(cache::Config::default().with_capacity(capacity)); let (producer, driver) = origin::Producer::new(config); - tokio::spawn(support::harness::run(driver)); + support::harness::spawn(driver); producer } diff --git a/rs/moq-net/benches/track.rs b/rs/moq-net/benches/track.rs index 5592698996..a1ace095df 100644 --- a/rs/moq-net/benches/track.rs +++ b/rs/moq-net/benches/track.rs @@ -7,6 +7,9 @@ //! `track_aborted_scan` isolates the other half of delivery: how much a cached //! prefix of aborted groups costs the scan that has to walk past it. //! +//! `track_gc` is the cache's clock-driven half: the per-poll [`cache::Pool::gc`] +//! call every origin driver makes, and the due expiry pass it runs over every track. +//! //! Run with `cargo bench -p moq-net --bench track`. use std::hint::black_box; @@ -33,6 +36,9 @@ const CACHE_CAPACITY: u64 = 64 * 1024; /// Concurrent publishers sharing one relay-style cache pool. const WRITERS: [usize; 4] = [1, 2, 4, 8]; +/// Tracks sharing one expiring pool: the table a due collection pass walks. +const SWEPT: [usize; 4] = [1, 64, 1_024, 16_384]; + /// Keeps the ownership chain alive around the track and its subscribers. struct Fanout { _broadcast: broadcast::Producer, @@ -123,6 +129,69 @@ impl AbortedScan { } } +/// Live tracks with one cached group each, sharing an expiring pool. +/// +/// Each track's only group is its live latest, which expiry never reclaims, so +/// every pass walks the same table instead of emptying it. +struct Swept { + _broadcast: broadcast::Producer, + _tracks: Vec, + pool: cache::Pool, + now: Instant, +} + +impl Swept { + fn new(tracks: usize) -> Self { + let mut info = broadcast::Info::default(); + info.pool = cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY)); + let pool = info.pool.clone(); + let broadcast = broadcast::Producer::new(info); + let payload = Bytes::from_static(&[0; PAYLOAD]); + let tracks = (0..tracks) + .map(|i| { + let track = broadcast.create_track(format!("bench{i}"), None).unwrap(); + let mut group = track.append_group().unwrap(); + group.write_frame(Timestamp::ZERO, payload.clone()).unwrap(); + group.finish().unwrap(); + track + }) + .collect(); + // The first call starts the pool's clock and runs its first pass. + let now = Instant::now(); + pool.gc(now); + Self { + _broadcast: broadcast, + _tracks: tracks, + pool, + now, + } + } +} + +fn bench_gc(c: &mut Criterion) { + let mut group = c.benchmark_group("track_gc"); + for tracks in SWEPT { + // Between passes, which is every origin poll: should stay flat as tracks grow. + group.throughput(Throughput::Elements(1)); + group.bench_with_input(BenchmarkId::new("idle", tracks), &tracks, |b, &tracks| { + let swept = Swept::new(tracks); + b.iter(|| black_box(swept.pool.gc(swept.now))); + }); + // A due pass, stepping one sweep interval per call. + group.throughput(Throughput::Elements(tracks as u64)); + group.bench_with_input(BenchmarkId::new("sweep", tracks), &tracks, |b, &tracks| { + let swept = Swept::new(tracks); + let interval = cache::DEFAULT_EXPIRY / 2; + let mut now = swept.now; + b.iter(|| { + now += interval; + black_box(swept.pool.gc(now)) + }); + }); + } + group.finish(); +} + fn bench_fanout(c: &mut Criterion) { let mut group = c.benchmark_group("track_fanout_group"); for subscribers in FANOUT { @@ -206,5 +275,11 @@ fn bench_aborted_scan(c: &mut Criterion) { group.finish(); } -criterion_group!(benches, bench_fanout, bench_parallel_write, bench_aborted_scan); +criterion_group!( + benches, + bench_fanout, + bench_parallel_write, + bench_aborted_scan, + bench_gc +); criterion_main!(benches); diff --git a/rs/moq-net/sim/Cargo.toml b/rs/moq-net/sim/Cargo.toml new file mode 100644 index 0000000000..702761fa73 --- /dev/null +++ b/rs/moq-net/sim/Cargo.toml @@ -0,0 +1,16 @@ +# moq-net's test executor: single-threaded, deterministic, with simulated time. +# +# A workspace member so it shares the root lockfile. Never published: moq-net takes +# it as a path-only dev-dependency, which `cargo publish` strips. +[package] +name = "moq-net-sim" +version = "0.0.0" +license = "MIT OR Apache-2.0" +publish = false +edition = "2024" +rust-version.workspace = true + +[dependencies] +futures = { workspace = true } +kio = { workspace = true } +moq-net-sim-macros = { path = "macros" } diff --git a/rs/moq-net/sim/macros/Cargo.toml b/rs/moq-net/sim/macros/Cargo.toml new file mode 100644 index 0000000000..a4a34fb005 --- /dev/null +++ b/rs/moq-net/sim/macros/Cargo.toml @@ -0,0 +1,12 @@ +# The `#[moq_net_sim::test]` attribute. A separate crate only because a proc macro +# must be; `moq-net-sim` re-exports it. +[package] +name = "moq-net-sim-macros" +version = "0.0.0" +license = "MIT OR Apache-2.0" +publish = false +edition = "2024" +rust-version.workspace = true + +[lib] +proc-macro = true diff --git a/rs/moq-net/sim/macros/src/lib.rs b/rs/moq-net/sim/macros/src/lib.rs new file mode 100644 index 0000000000..a4131d905d --- /dev/null +++ b/rs/moq-net/sim/macros/src/lib.rs @@ -0,0 +1,46 @@ +//! The `#[moq_net_sim::test]` attribute: an `async fn` test run on the simulated executor. + +use proc_macro::{Delimiter, Group, Ident, Span, TokenStream, TokenTree}; + +/// Run an `async fn` test to completion on [`moq_net_sim::run`](../moq_net_sim/fn.run.html). +#[proc_macro_attribute] +pub fn test(args: TokenStream, item: TokenStream) -> TokenStream { + if !args.is_empty() { + return error("#[moq_net_sim::test] takes no arguments"); + } + + let mut tokens: Vec = item.into_iter().collect(); + let is = |token: &TokenTree, name: &str| matches!(token, TokenTree::Ident(ident) if ident.to_string() == name); + let Some(at) = tokens + .windows(2) + .position(|pair| is(&pair[0], "async") && is(&pair[1], "fn")) + else { + return error("#[moq_net_sim::test] expects an `async fn`"); + }; + tokens.remove(at); + + let body = match tokens.pop() { + Some(TokenTree::Group(body)) if body.delimiter() == Delimiter::Brace => body, + _ => return error("#[moq_net_sim::test] expects a function body"), + }; + + // `{ ::moq_net_sim::run(async move ) }` + let mut call: TokenStream = "::moq_net_sim::run".parse().unwrap(); + let block: TokenStream = [ + TokenTree::Ident(Ident::new("async", Span::call_site())), + TokenTree::Ident(Ident::new("move", Span::call_site())), + TokenTree::Group(body), + ] + .into_iter() + .collect(); + call.extend([TokenTree::Group(Group::new(Delimiter::Parenthesis, block))]); + + let mut out: TokenStream = "#[::core::prelude::v1::test]".parse().unwrap(); + out.extend(tokens); + out.extend([TokenTree::Group(Group::new(Delimiter::Brace, call))]); + out +} + +fn error(message: &str) -> TokenStream { + format!("::core::compile_error!({message:?});").parse().unwrap() +} diff --git a/rs/moq-net/sim/src/lib.rs b/rs/moq-net/sim/src/lib.rs new file mode 100644 index 0000000000..91ae336ae1 --- /dev/null +++ b/rs/moq-net/sim/src/lib.rs @@ -0,0 +1,665 @@ +//! A single-threaded, deterministic executor with simulated time, for moq-net's tests. +//! +//! Time never moves on its own. Once every task has stalled, [`run`] jumps the clock +//! to the earliest armed timer and wakes it, like tokio's paused clock, so a test +//! never waits on the wall clock. A run whose tasks are all parked with nothing +//! armed can never finish, so it panics rather than hang. +//! +//! Protocol drivers take time from their caller: [`drive`] polls one with the +//! simulated instant, and [`attach`] keeps an externally owned clock in step. + +use std::{ + cell::RefCell, + collections::BTreeMap, + future::{Future, IntoFuture}, + pin::{Pin, pin}, + rc::Rc, + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, + task::{Context, Poll, Wake, Waker}, + time::{Duration, Instant}, +}; + +use futures::{FutureExt, StreamExt, stream::FuturesUnordered}; + +pub use moq_net_sim_macros::test; + +thread_local! { + static CURRENT: RefCell>> = const { RefCell::new(None) }; +} + +/// A source of deadlines kept in step with simulated time: called with the new +/// instant whenever time moves, it returns its earliest deadline still pending. +type Source = Box Option>; + +type Task = Pin>>; + +struct Runtime { + /// Spawned since the run loop last collected them. + spawned: RefCell>, + time: RefCell