From 026b5389c123356723fafd230e5d23b3c4362bf1 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 16:12:16 -0800 Subject: [PATCH 1/9] Start requiring semver compatible changes on main. --- .github/workflows/dev.yml | 31 +++++++++++++++++++++++ .github/workflows/{check.yml => main.yml} | 3 ++- CLAUDE.md | 7 +++++ flake.nix | 1 + justfile | 11 +++++++- 5 files changed, 51 insertions(+), 2 deletions(-) create mode 100644 .github/workflows/dev.yml rename .github/workflows/{check.yml => main.yml} (97%) diff --git a/.github/workflows/dev.yml b/.github/workflows/dev.yml new file mode 100644 index 0000000000..b49a3c79e1 --- /dev/null +++ b/.github/workflows/dev.yml @@ -0,0 +1,31 @@ +name: Check + +permissions: + id-token: write + contents: read + +on: + pull_request: + branches: ["dev"] + +jobs: + check: + name: Check + runs-on: ubuntu-latest + + steps: + - name: Checkout + uses: actions/checkout@v6 + + - uses: DeterminateSystems/nix-installer-action@main + - uses: DeterminateSystems/magic-nix-cache-action@main + - uses: DeterminateSystems/flake-checker-action@main + + # Cache Rust dependencies and build artifacts + - name: Rust Cache + uses: Swatinem/rust-cache@v2 + with: + cache-on-failure: true + + # Run checks with cached dependencies + - run: nix develop --command just check diff --git a/.github/workflows/check.yml b/.github/workflows/main.yml similarity index 97% rename from .github/workflows/check.yml rename to .github/workflows/main.yml index 22c2be639d..0519d5e029 100644 --- a/.github/workflows/check.yml +++ b/.github/workflows/main.yml @@ -1,4 +1,4 @@ -name: Check +name: Stable permissions: id-token: write @@ -29,3 +29,4 @@ jobs: # Run checks with cached dependencies - run: nix develop --command just ci + diff --git a/CLAUDE.md b/CLAUDE.md index a46c4594f8..81e944ab26 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -89,6 +89,13 @@ Key architectural rule: The CDN/relay does not know anything about media. Anythi - Run `just fix` to automatically fix formating and easy things. - Rust tests are integrated within source files +## Branching Strategy + +- **`main`**: Stable branch for patch releases. Only non-breaking fixes and additions. +- **`dev`**: Development branch for breaking API changes. PRs with major API changes should target `dev`. +- When ready for a new minor/major release, merge `dev` into `main`. +- `cargo-semver-checks` enforces this on PRs to `main`. + ## Workflow When making changes to the codebase: diff --git a/flake.nix b/flake.nix index 707124e556..9405e6efb5 100644 --- a/flake.nix +++ b/flake.nix @@ -62,6 +62,7 @@ cargo-edit cargo-hack cargo-sweep + cargo-semver-checks ]; # JavaScript dependencies diff --git a/justfile b/justfile index 74015ef89f..0ddf768d85 100644 --- a/justfile +++ b/justfile @@ -11,7 +11,7 @@ default: # Install any dependencies. install: bun install - cargo install --locked cargo-shear cargo-sort cargo-upgrades cargo-edit cargo-hack cargo-sweep + cargo install --locked cargo-shear cargo-sort cargo-upgrades cargo-edit cargo-hack cargo-sweep cargo-semver-checks release-plz # Alias for dev. all: dev @@ -368,6 +368,15 @@ ci: echo "Checking all feature combinations..." cargo hack check --workspace --each-feature --no-dev-deps + # Check semver compatibility against crates.io + # requires: cargo install cargo-semver-checks + echo "Checking semver compatibility..." + cargo semver-checks check-release --workspace --exclude libmoq + +# Update versions and changelogs via release-plz +bump: + release-plz update + # Run the unit tests test: #!/usr/bin/env bash From 00219f8bfaecfa5ef4b071781435e77eda120565 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 17:09:17 -0800 Subject: [PATCH 2/9] =?UTF-8?q?API=20cleanup:=20rename=20close=E2=86=92abo?= =?UTF-8?q?rt,=20finish=E2=86=92append=5Ffinish,=20write=5Fchunk=E2=86=92w?= =?UTF-8?q?rite?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reduce semver violations and improve API clarity: - Rename `close(err)` → `abort(err)` on all producer types to distinguish error termination from clean close - Replace `TrackProducer::finish()` with `append_finish()` and `insert_finish(sequence)` for explicit control over the final group - Rename `FrameProducer::write_chunk()` → `write()` for simplicity - Add doc comments to public producer/consumer APIs Co-Authored-By: Claude Opus 4.6 --- rs/hang/src/container/frame.rs | 4 +-- rs/moq-lite/src/ietf/subscriber.rs | 16 +++++----- rs/moq-lite/src/lite/subscriber.rs | 18 +++++------ rs/moq-lite/src/model/broadcast.rs | 42 ++++++++++++++++-------- rs/moq-lite/src/model/frame.rs | 23 ++++++++++---- rs/moq-lite/src/model/group.rs | 23 +++++++++----- rs/moq-lite/src/model/state.rs | 6 ++-- rs/moq-lite/src/model/track.rs | 51 ++++++++++++++++++++---------- rs/moq-mux/src/catalog.rs | 4 +-- rs/moq-mux/src/import/fmp4.rs | 4 +-- rs/moq-relay/src/cluster.rs | 2 +- 11 files changed, 123 insertions(+), 70 deletions(-) diff --git a/rs/hang/src/container/frame.rs b/rs/hang/src/container/frame.rs index 5ac97ea01f..e99594cec3 100644 --- a/rs/hang/src/container/frame.rs +++ b/rs/hang/src/container/frame.rs @@ -46,9 +46,9 @@ impl Frame { let size = header.len() + self.payload.remaining(); let mut chunked = group.create_frame(size.into())?; - chunked.write_chunk(header.freeze())?; + chunked.write(header.freeze())?; for chunk in &self.payload { - chunked.write_chunk(chunk.clone())?; + chunked.write(chunk.clone())?; } chunked.finish()?; diff --git a/rs/moq-lite/src/ietf/subscriber.rs b/rs/moq-lite/src/ietf/subscriber.rs index 30d7c8e52b..e20c20067d 100644 --- a/rs/moq-lite/src/ietf/subscriber.rs +++ b/rs/moq-lite/src/ietf/subscriber.rs @@ -201,7 +201,7 @@ impl Subscriber { let mut state = self.state.lock(); if let Some(mut track) = state.subscribes.remove(&msg.request_id) { - let _ = track.producer.close(Error::Cancel); + let _ = track.producer.abort(Error::Cancel); if let Some(alias) = track.alias { state.aliases.remove(&alias); } @@ -221,7 +221,7 @@ impl Subscriber { let mut state = self.state.lock(); if let Some(mut track) = state.subscribes.remove(&msg.request_id) { - let _ = track.producer.close(Error::Cancel); + let _ = track.producer.abort(Error::Cancel); if let Some(alias) = track.alias { state.aliases.remove(&alias); } @@ -234,7 +234,7 @@ impl Subscriber { let mut state = self.state.lock(); if let Some(mut track) = state.subscribes.remove(&msg.request_id) { - let _ = track.producer.finish(); + let _ = track.producer.append_finish(); if let Some(alias) = track.alias { state.aliases.remove(&alias); } @@ -344,7 +344,7 @@ impl Subscriber { let _ = track.unused().await; tracing::info!(id = %request_id, broadcast = %self.origin.as_ref().unwrap().absolute(&broadcast), track = %track.info.name, "subscribe cancelled"); - let _ = track.close(Error::Cancel); + let _ = track.abort(Error::Cancel); Ok(()) } @@ -383,11 +383,11 @@ impl Subscriber { match res { Err(Error::Cancel) => { tracing::trace!(group = %producer.info.sequence, "group cancelled"); - let _ = producer.close(Error::Cancel); + let _ = producer.abort(Error::Cancel); } Err(err) => { tracing::debug!(%err, group = %producer.info.sequence, "group error"); - let _ = producer.close(err); + let _ = producer.abort(err); } _ => { tracing::trace!(group = %producer.info.sequence, "group complete"); @@ -433,7 +433,7 @@ impl Subscriber { let mut frame = producer.create_frame(Frame { size })?; if let Err(err) = self.run_frame(stream, frame.clone()).await { - let _ = frame.close(err.clone()); + let _ = frame.abort(err.clone()); return Err(err); } @@ -456,7 +456,7 @@ impl Subscriber { while remain > 0 { let chunk = stream.read(remain as usize).await?.ok_or(Error::WrongSize)?; remain = remain.checked_sub(chunk.len() as u64).ok_or(Error::WrongSize)?; - frame.write_chunk(chunk)?; + frame.write(chunk)?; } tracing::trace!(size = %frame.info.size, "read frame"); diff --git a/rs/moq-lite/src/lite/subscriber.rs b/rs/moq-lite/src/lite/subscriber.rs index 617d0c1307..2258ff9f7e 100644 --- a/rs/moq-lite/src/lite/subscriber.rs +++ b/rs/moq-lite/src/lite/subscriber.rs @@ -108,9 +108,9 @@ impl Subscriber { lite::Announce::Ended { suffix: path, .. } => { tracing::debug!(broadcast = %self.log_path(&path), "unannounced"); - // Close the producer. + // Abort the producer. let mut producer = producers.remove(&path.into_owned()).ok_or(Error::NotFound)?; - producer.close(Error::Cancel).ok(); + producer.abort(Error::Cancel).ok(); } } } @@ -198,15 +198,15 @@ impl Subscriber { match res { Err(Error::Cancel) => { tracing::info!(id, broadcast = %self.log_path(&broadcast), track = %track.info.name, "subscribe cancelled"); - let _ = track.close(Error::Cancel); + let _ = track.abort(Error::Cancel); } Err(err) => { tracing::warn!(id, broadcast = %self.log_path(&broadcast), track = %track.info.name, %err, "subscribe error"); - let _ = track.close(err); + let _ = track.abort(err); } _ => { tracing::info!(id, broadcast = %self.log_path(&broadcast), track = %track.info.name, "subscribe complete"); - let _ = track.finish(); + let _ = track.append_finish(); } } } @@ -262,11 +262,11 @@ impl Subscriber { match res { Err(Error::Cancel) => { tracing::trace!(group = %group.info.sequence, "group cancelled"); - let _ = group.close(Error::Cancel); + let _ = group.abort(Error::Cancel); } Err(err) => { tracing::debug!(%err, group = %group.info.sequence, "group error"); - let _ = group.close(err); + let _ = group.abort(err); } _ => { tracing::trace!(group = %group.info.sequence, "group complete"); @@ -286,7 +286,7 @@ impl Subscriber { let mut frame = group.create_frame(Frame { size })?; if let Err(err) = self.run_frame(stream, &mut frame).await { - let _ = frame.close(err.clone()); + let _ = frame.abort(err.clone()); return Err(err); } @@ -312,7 +312,7 @@ impl Subscriber { .await? .ok_or(Error::WrongSize)?; remain = remain.checked_sub(chunk.len() as u64).ok_or(Error::WrongSize)?; - frame.write_chunk(chunk)?; + frame.write(chunk)?; } tracing::trace!(size = %frame.info.size, "read frame"); diff --git a/rs/moq-lite/src/model/broadcast.rs b/rs/moq-lite/src/model/broadcast.rs index 18b43a66b5..d6218436b7 100644 --- a/rs/moq-lite/src/model/broadcast.rs +++ b/rs/moq-lite/src/model/broadcast.rs @@ -41,7 +41,10 @@ struct State { dynamic: usize, } -/// Receive broadcast/track requests and return if we can fulfill them. +/// Manages tracks within a broadcast. +/// +/// Insert tracks statically with [Self::insert_track] / [Self::create_track], +/// or handle on-demand requests via [Self::dynamic]. #[derive(Clone)] pub struct BroadcastProducer { state: Producer, @@ -91,33 +94,37 @@ impl BroadcastProducer { Ok(track) } + /// Create a dynamic producer that handles on-demand track requests from consumers. pub fn dynamic(&self) -> BroadcastDynamic { BroadcastDynamic::new(self.state.clone()) } + /// Create a consumer that can subscribe to tracks in this broadcast. pub fn consume(&self) -> BroadcastConsumer { BroadcastConsumer { state: self.state.consume(), } } - pub fn close(&mut self, err: Error) -> Result<(), Error> { + /// Abort the broadcast and all child tracks with the given error. + pub fn abort(&mut self, err: Error) -> Result<(), Error> { let mut state = self.state.modify()?; - // Cascade close to all child tracks. + // Cascade abort to all child tracks. for weak in state.tracks.values() { - weak.close(err.clone()); + weak.abort(err.clone()); } - // Close any pending dynamic track requests. + // Abort any pending dynamic track requests. for mut request in state.requests.drain(..) { - request.close(err.clone()).ok(); + request.abort(err.clone()).ok(); } - state.close(err); + state.abort(err); Ok(()) } + /// Return true if this is the same broadcast instance. pub fn is_clone(&self, other: &Self) -> bool { self.state.is_clone(&other.state) } @@ -134,6 +141,11 @@ impl BroadcastProducer { } } +/// Handles on-demand track creation for a broadcast. +/// +/// When a consumer requests a track that doesn't exist, a [TrackProducer] is created +/// and queued for the dynamic producer to fulfill via [Self::requested_track]. +/// Dropped when no longer needed; pending requests are automatically aborted. #[derive(Clone)] pub struct BroadcastDynamic { state: Producer, @@ -158,21 +170,25 @@ impl BroadcastDynamic { }) } + /// Block until a consumer requests a track, returning its producer. pub async fn requested_track(&mut self) -> Result, Error> { waiter_fn(move |waiter| self.poll_requested_track(waiter)).await } + /// Create a consumer that can subscribe to tracks in this broadcast. pub fn consume(&self) -> BroadcastConsumer { BroadcastConsumer { state: self.state.consume(), } } - pub fn close(&mut self, err: Error) -> Result<(), Error> { - self.state.close(err)?; + /// Abort the broadcast with the given error. + pub fn abort(&mut self, err: Error) -> Result<(), Error> { + self.state.abort(err)?; Ok(()) } + /// Return true if this is the same broadcast instance. pub fn is_clone(&self, other: &Self) -> bool { self.state.is_clone(&other.state) } @@ -189,7 +205,7 @@ impl Drop for BroadcastDynamic { // Abort all pending requests since there's no dynamic producer to handle them. for mut request in state.requests.drain(..) { - request.close(Error::Cancel).ok(); + request.abort(Error::Cancel).ok(); } } } @@ -330,8 +346,8 @@ mod test { let track1c = consumer.assert_subscribe_track(&track1.info); let track2 = consumer.assert_subscribe_track(&Track::new("track2")); - // Explicitly closing the broadcast should cascade to child tracks. - producer.close(Error::Cancel).unwrap(); + // Explicitly aborting the broadcast should cascade to child tracks. + producer.abort(Error::Cancel).unwrap(); // The requested TrackProducer should have been aborted. track2.assert_error(); @@ -392,7 +408,7 @@ mod test { // Get the requested producer and close it (simulating publisher disconnect) let mut producer1 = broadcast.assert_request(); producer1.append_group().unwrap(); - producer1.finish().unwrap(); + producer1.append_finish().unwrap(); drop(producer1); // The consumer should see the track as closed diff --git a/rs/moq-lite/src/model/frame.rs b/rs/moq-lite/src/model/frame.rs index 7bac5651fe..a0dfa67384 100644 --- a/rs/moq-lite/src/model/frame.rs +++ b/rs/moq-lite/src/model/frame.rs @@ -104,9 +104,12 @@ impl FrameState { } } -/// Used to write a frame's worth of data in chunks. +/// Writes a frame's payload in one or more chunks. +/// +/// The total bytes written must exactly match [Frame::size]. +/// Call [Self::finish] after writing all bytes to verify correctness. pub struct FrameProducer { - // Immutable stream state. + /// The frame header containing the expected size. pub info: Frame, // Mutable stream state. @@ -114,6 +117,7 @@ pub struct FrameProducer { } impl FrameProducer { + /// Create a new frame producer for the given frame header. pub fn new(info: Frame) -> Self { let state = FrameState { chunks: Vec::new(), @@ -125,13 +129,18 @@ impl FrameProducer { } } - pub fn write_chunk>(&mut self, chunk: B) -> Result<()> { + /// Write a chunk of data to the frame. + /// + /// Returns [Error::WrongSize] if the total bytes written would exceed [Frame::size]. + pub fn write>(&mut self, chunk: B) -> Result<()> { let chunk = chunk.into(); let mut state = self.state.modify()?; state.write_chunk(chunk) } - /// Optional: mark the frame as finished when all bytes have been written. + /// Verify that all bytes have been written. + /// + /// Returns [Error::WrongSize] if the bytes written don't match [Frame::size]. pub fn finish(&mut self) -> Result<()> { let state = self.state.modify()?; if state.remaining != 0 { @@ -140,8 +149,9 @@ impl FrameProducer { Ok(()) } - pub fn close(&mut self, err: Error) -> Result<()> { - self.state.close(err) + /// Abort the frame with the given error. + pub fn abort(&mut self, err: Error) -> Result<()> { + self.state.abort(err) } /// Create a new consumer for the frame. @@ -153,6 +163,7 @@ impl FrameProducer { } } + /// Block until there are no active consumers. pub async fn unused(&self) -> Result<()> { self.state.unused().await } diff --git a/rs/moq-lite/src/model/group.rs b/rs/moq-lite/src/model/group.rs index 287fa80c21..737a54b9bd 100644 --- a/rs/moq-lite/src/model/group.rs +++ b/rs/moq-lite/src/model/group.rs @@ -84,16 +84,21 @@ impl GroupState { } } -/// Create a group, frame-by-frame. +/// Writes frames to a group in order. +/// +/// Each group is delivered independently over a QUIC stream. +/// Use [Self::write_frame] for simple single-buffer frames, +/// or [Self::create_frame] for multi-chunk streaming writes. pub struct GroupProducer { // Mutable stream state. state: Producer, - // Immutable stream state. + /// The group header containing the sequence number. pub info: Group, } impl GroupProducer { + /// Create a new group producer. pub fn new(info: Group) -> Self { Self { info, @@ -111,7 +116,7 @@ impl GroupProducer { size: data.len() as u64, }; let mut frame = self.create_frame(frame)?; - frame.write_chunk(data)?; + frame.write(data)?; frame.finish()?; Ok(()) } @@ -133,26 +138,26 @@ impl GroupProducer { Ok(()) } - /// Clean termination of the group. + /// Mark the group as complete; no more frames will be written. pub fn finish(&mut self) -> Result<()> { let mut state = self.state.modify()?; state.fin = true; Ok(()) } - /// Close the group with the given error. + /// Abort the group with the given error. /// /// No updates can be made after this point. - pub fn close(&mut self, err: Error) -> Result<()> { + pub fn abort(&mut self, err: Error) -> Result<()> { let mut state = self.state.modify()?; // Abort all frames still in progress. for frame in state.frames.iter_mut() { // Ignore errors, we don't care if the frame was already closed. - frame.close(err.clone()).ok(); + frame.abort(err.clone()).ok(); } - state.close(err); + state.abort(err); Ok(()) } @@ -165,10 +170,12 @@ impl GroupProducer { } } + /// Block until the group is closed or aborted. pub async fn closed(&self) -> Error { self.state.closed().await } + /// Block until there are no active consumers. pub async fn unused(&self) -> Result<()> { self.state.unused().await } diff --git a/rs/moq-lite/src/model/state.rs b/rs/moq-lite/src/model/state.rs index 938e9491ce..a1d55bb043 100644 --- a/rs/moq-lite/src/model/state.rs +++ b/rs/moq-lite/src/model/state.rs @@ -98,8 +98,8 @@ impl Producer { } } - pub fn close(&mut self, err: Error) -> Result<(), Error> { - self.modify()?.close(err); + pub fn abort(&mut self, err: Error) -> Result<(), Error> { + self.modify()?.abort(err); Ok(()) } @@ -324,7 +324,7 @@ impl<'a, T> ProducerMut<'a, T> { } /// NOTE: This takes self so it's impossible to be in a closed state. - pub fn close(mut self, err: Error) { + pub fn abort(mut self, err: Error) { let state = self.state.as_mut().unwrap(); // We don't need to check for state.closed because we checked when making ProducerMut state.closed = Err(err); diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index 87d9d450b3..dccd0ca11e 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -56,7 +56,7 @@ struct State { duplicates: HashSet, offset: usize, max_sequence: Option, - fin: bool, + fin: Option, } impl State { @@ -71,7 +71,11 @@ impl State { } } - if self.fin { Poll::Ready(None) } else { Poll::Pending } + if self.fin.is_some() { + Poll::Ready(None) + } else { + Poll::Pending + } } fn poll_get_group(&self, sequence: u64) -> Poll> { @@ -91,7 +95,7 @@ impl State { return Poll::Ready(None); } - if self.fin { + if self.fin.is_some() { return Poll::Ready(None); } @@ -147,7 +151,9 @@ impl TrackProducer { let group = info.produce(); let mut state = self.state.modify()?; - if state.fin && group.info.sequence >= state.max_sequence.unwrap_or(0) { + if let Some(fin) = state.fin + && group.info.sequence >= fin + { return Err(Error::Closed); } @@ -166,7 +172,7 @@ impl TrackProducer { /// Create a new group with the next sequence number. pub fn append_group(&mut self) -> Result { let mut state = self.state.modify()?; - if state.fin { + if state.fin.is_some() { return Err(Error::Closed); } @@ -190,26 +196,39 @@ impl TrackProducer { Ok(()) } - /// Mark the last group of the track. + /// Mark the track as finished after the last appended group. + /// + /// Sets the final sequence to the current max_sequence. + /// No new groups at or above this sequence can be appended. + /// NOTE: Old groups with lower sequence numbers can still arrive. + pub fn append_finish(&mut self) -> Result<()> { + let mut state = self.state.modify()?; + let sequence = state.max_sequence.unwrap_or(0); + state.fin = Some(sequence); + Ok(()) + } + + /// Mark a specific group sequence as the final group. /// - /// NOTE: The track is not closed yet; old groups can still arrive. - pub fn finish(&mut self) -> Result<()> { + /// No new groups at or above this sequence can be created. + /// NOTE: Old groups with lower sequence numbers can still arrive. + pub fn insert_finish(&mut self, sequence: u64) -> Result<()> { let mut state = self.state.modify()?; - state.fin = true; + state.fin = Some(sequence); Ok(()) } /// Abort the track with the given error. - pub fn close(&mut self, err: Error) -> Result<()> { + pub fn abort(&mut self, err: Error) -> Result<()> { let mut state = self.state.modify()?; // Abort all groups still in progress. for (group, _) in state.groups.iter_mut().flatten() { // Ignore errors, we don't care if the group was already closed. - group.close(err.clone()).ok(); + group.abort(err.clone()).ok(); } - state.close(err); + state.abort(err); Ok(()) } @@ -272,17 +291,17 @@ pub(crate) struct TrackWeak { } impl TrackWeak { - pub fn close(&self, err: Error) { + pub fn abort(&self, err: Error) { // Upgrade to a temporary Producer so we can modify the state. let Ok(producer) = self.state.produce() else { return }; let Ok(mut state) = producer.modify() else { return }; - // Cascade close to all groups. + // Cascade abort to all groups. for (group, _) in state.groups.iter_mut().flatten() { - group.close(err.clone()).ok(); + group.abort(err.clone()).ok(); } - state.close(err); + state.abort(err); } pub fn is_closed(&self) -> bool { diff --git a/rs/moq-mux/src/catalog.rs b/rs/moq-mux/src/catalog.rs index dcf6fecf46..3a0f95a56a 100644 --- a/rs/moq-mux/src/catalog.rs +++ b/rs/moq-mux/src/catalog.rs @@ -58,8 +58,8 @@ impl CatalogProducer { /// Finish publishing to this catalog. pub fn finish(&mut self) -> Result<(), moq_lite::Error> { - self.hang_track.finish()?; - self.msf_track.finish()?; + self.hang_track.append_finish()?; + self.msf_track.append_finish()?; Ok(()) } } diff --git a/rs/moq-mux/src/import/fmp4.rs b/rs/moq-mux/src/import/fmp4.rs index d9c430046e..b3639ddf6a 100644 --- a/rs/moq-mux/src/import/fmp4.rs +++ b/rs/moq-mux/src/import/fmp4.rs @@ -612,8 +612,8 @@ impl Fmp4 { size: moof_raw.len() as u64 + mdat_raw.len() as u64, })?; - frame.write_chunk(moof_raw.clone())?; - frame.write_chunk(Bytes::copy_from_slice(mdat_raw))?; + frame.write(moof_raw.clone())?; + frame.write(Bytes::copy_from_slice(mdat_raw))?; frame.finish()?; track.group = Some(group); diff --git a/rs/moq-relay/src/cluster.rs b/rs/moq-relay/src/cluster.rs index a65d5ae57c..9f89333018 100644 --- a/rs/moq-relay/src/cluster.rs +++ b/rs/moq-relay/src/cluster.rs @@ -319,6 +319,6 @@ impl ClusterRegistration { impl Drop for ClusterRegistration { fn drop(&mut self) { tracing::info!(%self.node, "unregistered cluster client"); - self.broadcast.close(moq_lite::Error::Cancel).ok(); + let _ = self.broadcast.abort(moq_lite::Error::Cancel); } } From 5366ceefd6dfea76df60746638e0902cc01bc6b0 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 17:12:54 -0800 Subject: [PATCH 3/9] Update main.yml --- .github/workflows/main.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 0519d5e029..3e53c35dbe 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -1,4 +1,4 @@ -name: Stable +name: Main permissions: id-token: write From 230bd8a88662a379996d1d7dc130a0943b4f4f3b Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 19:56:10 -0800 Subject: [PATCH 4/9] Fix track finish semantics and document libmoq semver exclusion --- justfile | 1 + rs/moq-lite/src/model/track.rs | 108 +++++++++++++++++++++++++++++---- 2 files changed, 96 insertions(+), 13 deletions(-) diff --git a/justfile b/justfile index 0ddf768d85..7f946584ad 100644 --- a/justfile +++ b/justfile @@ -370,6 +370,7 @@ ci: # Check semver compatibility against crates.io # requires: cargo install cargo-semver-checks + # libmoq is an internal C-ABI crate and is intentionally excluded from published-crate semver checks. echo "Checking semver compatibility..." cargo semver-checks check-release --workspace --exclude libmoq diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index dccd0ca11e..ee2cfe3834 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -80,17 +80,15 @@ impl State { fn poll_get_group(&self, sequence: u64) -> Poll> { // Search for the group with the matching sequence, skipping tombstones. - // NOTE: Returns Ready(None) if the group was evicted (tombstoned) since - // it won't be found and max_sequence >= sequence. for (group, _) in self.groups.iter().flatten() { if group.info.sequence == sequence { return Poll::Ready(Some(group.clone())); } } - // If we've already seen a newer sequence, the group is gone (or was evicted). - if let Some(max) = self.max_sequence - && max >= sequence + // Once a fin boundary is set, groups at-or-after fin can never exist. + if let Some(fin) = self.fin + && sequence >= fin { return Poll::Ready(None); } @@ -172,14 +170,16 @@ impl TrackProducer { /// Create a new group with the next sequence number. pub fn append_group(&mut self) -> Result { let mut state = self.state.modify()?; - if state.fin.is_some() { + let sequence = state.max_sequence.map_or(0, |s| s + 1); + if let Some(fin) = state.fin + && sequence >= fin + { return Err(Error::Closed); } - let now = tokio::time::Instant::now(); - let sequence = state.max_sequence.map_or(0, |s| s + 1); let group = Group { sequence }.produce(); + let now = tokio::time::Instant::now(); state.duplicates.insert(sequence); state.max_sequence = Some(sequence); state.groups.push_back(Some((group.clone(), now))); @@ -203,18 +203,26 @@ impl TrackProducer { /// NOTE: Old groups with lower sequence numbers can still arrive. pub fn append_finish(&mut self) -> Result<()> { let mut state = self.state.modify()?; - let sequence = state.max_sequence.unwrap_or(0); - state.fin = Some(sequence); + if state.fin.is_some() { + return Err(Error::Closed); + } + let max = state.max_sequence.ok_or(Error::Closed)?; + state.fin = Some(max); Ok(()) } - /// Mark a specific group sequence as the final group. + /// Mark the track as finished, validating the caller's expected final sequence. /// - /// No new groups at or above this sequence can be created. + /// Freezes the final boundary at the current max_sequence. + /// No new groups at or above that sequence can be created. /// NOTE: Old groups with lower sequence numbers can still arrive. pub fn insert_finish(&mut self, sequence: u64) -> Result<()> { let mut state = self.state.modify()?; - state.fin = Some(sequence); + let max = state.max_sequence.ok_or(Error::Closed)?; + if state.fin.is_some() || sequence < max { + return Err(Error::Closed); + } + state.fin = Some(max); Ok(()) } @@ -605,4 +613,78 @@ mod test { // consume() starts at the last slot (seq 2). assert_eq!(group.info.sequence, 2); } + + #[test] + fn append_finish_requires_max_and_cannot_be_rewritten() { + let mut producer = Track::new("test").produce(); + + assert!(producer.append_finish().is_err()); + + producer.append_group().unwrap(); + assert!(producer.append_finish().is_ok()); + assert!(producer.append_finish().is_err()); + assert!(producer.append_group().is_err()); + } + + #[test] + fn insert_finish_validates_sequence_and_freezes_to_max() { + let mut producer = Track::new("test").produce(); + producer.create_group(Group { sequence: 5 }).unwrap(); + + assert!(producer.insert_finish(4).is_err()); + assert!(producer.insert_finish(10).is_ok()); + + { + let state = producer.state.borrow(); + assert_eq!(state.fin, Some(5)); + } + + assert!(producer.insert_finish(10).is_err()); + assert!(producer.create_group(Group { sequence: 4 }).is_ok()); + assert!(producer.create_group(Group { sequence: 5 }).is_err()); + } + + #[tokio::test] + async fn next_group_finishes_without_waiting_for_gaps() { + let mut producer = Track::new("test").produce(); + producer.create_group(Group { sequence: 1 }).unwrap(); + producer.insert_finish(1).unwrap(); + + let mut consumer = producer.consume(); + assert_eq!(consumer.assert_group().info.sequence, 1); + + let done = consumer + .next_group() + .now_or_never() + .expect("should not block") + .expect("would have errored"); + assert!(done.is_none(), "track should finish without waiting for gaps"); + } + + #[tokio::test] + async fn get_group_finishes_without_waiting_for_gaps() { + let mut producer = Track::new("test").produce(); + producer.create_group(Group { sequence: 1 }).unwrap(); + producer.insert_finish(1).unwrap(); + + let consumer = producer.consume(); + assert!( + consumer + .get_group(0) + .now_or_never() + .expect("should not block") + .expect("would have errored") + .is_none(), + "sequence below fin should not block forever" + ); + assert!( + consumer + .get_group(2) + .now_or_never() + .expect("sequence at-or-after fin should resolve") + .expect("should not error") + .is_none(), + "sequence at-or-after fin should not exist" + ); + } } From 3f9658dd72d3190487834b5e98ee345284850ffa Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 20:11:43 -0800 Subject: [PATCH 5/9] Harden track finish validation and overflow handling --- rs/moq-lite/src/model/track.rs | 30 +++++++++++++++++++++++------- 1 file changed, 23 insertions(+), 7 deletions(-) diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index ee2cfe3834..40ac8cc4e2 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -170,7 +170,10 @@ impl TrackProducer { /// Create a new group with the next sequence number. pub fn append_group(&mut self) -> Result { let mut state = self.state.modify()?; - let sequence = state.max_sequence.map_or(0, |s| s + 1); + let sequence = match state.max_sequence { + Some(s) => s.checked_add(1).ok_or(Error::BoundsExceeded)?, + None => 0, + }; if let Some(fin) = state.fin && sequence >= fin { @@ -211,15 +214,16 @@ impl TrackProducer { Ok(()) } - /// Mark the track as finished, validating the caller's expected final sequence. + /// Mark the track as finished at an exact final sequence. /// + /// The caller must pass the current max_sequence exactly. /// Freezes the final boundary at the current max_sequence. /// No new groups at or above that sequence can be created. /// NOTE: Old groups with lower sequence numbers can still arrive. pub fn insert_finish(&mut self, sequence: u64) -> Result<()> { let mut state = self.state.modify()?; let max = state.max_sequence.ok_or(Error::Closed)?; - if state.fin.is_some() || sequence < max { + if state.fin.is_some() || sequence != max { return Err(Error::Closed); } state.fin = Some(max); @@ -632,19 +636,20 @@ mod test { producer.create_group(Group { sequence: 5 }).unwrap(); assert!(producer.insert_finish(4).is_err()); - assert!(producer.insert_finish(10).is_ok()); + assert!(producer.insert_finish(10).is_err()); + assert!(producer.insert_finish(5).is_ok()); { let state = producer.state.borrow(); assert_eq!(state.fin, Some(5)); } - assert!(producer.insert_finish(10).is_err()); + assert!(producer.insert_finish(5).is_err()); assert!(producer.create_group(Group { sequence: 4 }).is_ok()); assert!(producer.create_group(Group { sequence: 5 }).is_err()); } - #[tokio::test] + #[tokio::test(start_paused = true)] async fn next_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); @@ -661,7 +666,7 @@ mod test { assert!(done.is_none(), "track should finish without waiting for gaps"); } - #[tokio::test] + #[tokio::test(start_paused = true)] async fn get_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); @@ -687,4 +692,15 @@ mod test { "sequence at-or-after fin should not exist" ); } + + #[test] + fn append_group_returns_bounds_exceeded_on_sequence_overflow() { + let mut producer = Track::new("test").produce(); + { + let mut state = producer.state.modify().unwrap(); + state.max_sequence = Some(u64::MAX); + } + + assert!(matches!(producer.append_group(), Err(Error::BoundsExceeded))); + } } From 48fe39aad8723d0491e53b6b909213a2b215772d Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 20:44:33 -0800 Subject: [PATCH 6/9] moq-lite: rename track finish API and deprecate legacy writes --- rs/hang/src/container/consumer.rs | 2 +- rs/moq-lite/src/ietf/subscriber.rs | 2 +- rs/moq-lite/src/lite/subscriber.rs | 2 +- rs/moq-lite/src/model/broadcast.rs | 2 +- rs/moq-lite/src/model/frame.rs | 8 ++++++++ rs/moq-lite/src/model/track.rs | 31 +++++++++++++++++++----------- rs/moq-mux/src/catalog.rs | 4 ++-- 7 files changed, 34 insertions(+), 17 deletions(-) diff --git a/rs/hang/src/container/consumer.rs b/rs/hang/src/container/consumer.rs index 8727c2ee5d..df9337912c 100644 --- a/rs/hang/src/container/consumer.rs +++ b/rs/hang/src/container/consumer.rs @@ -570,7 +570,7 @@ mod tests { let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(500)); write_group(&mut track, 0, &[ts(0)]); - track.close(moq_lite::Error::Cancel).unwrap(); + track.abort(moq_lite::Error::Cancel).unwrap(); // Consumer should not hang; it should return frames or error gracefully let result = tokio::time::timeout(Duration::from_millis(500), async { diff --git a/rs/moq-lite/src/ietf/subscriber.rs b/rs/moq-lite/src/ietf/subscriber.rs index e20c20067d..8b0e152198 100644 --- a/rs/moq-lite/src/ietf/subscriber.rs +++ b/rs/moq-lite/src/ietf/subscriber.rs @@ -234,7 +234,7 @@ impl Subscriber { let mut state = self.state.lock(); if let Some(mut track) = state.subscribes.remove(&msg.request_id) { - let _ = track.producer.append_finish(); + let _ = track.producer.finish(); if let Some(alias) = track.alias { state.aliases.remove(&alias); } diff --git a/rs/moq-lite/src/lite/subscriber.rs b/rs/moq-lite/src/lite/subscriber.rs index 2258ff9f7e..8e8e04b11d 100644 --- a/rs/moq-lite/src/lite/subscriber.rs +++ b/rs/moq-lite/src/lite/subscriber.rs @@ -206,7 +206,7 @@ impl Subscriber { } _ => { tracing::info!(id, broadcast = %self.log_path(&broadcast), track = %track.info.name, "subscribe complete"); - let _ = track.append_finish(); + let _ = track.finish(); } } } diff --git a/rs/moq-lite/src/model/broadcast.rs b/rs/moq-lite/src/model/broadcast.rs index d6218436b7..4266579614 100644 --- a/rs/moq-lite/src/model/broadcast.rs +++ b/rs/moq-lite/src/model/broadcast.rs @@ -408,7 +408,7 @@ mod test { // Get the requested producer and close it (simulating publisher disconnect) let mut producer1 = broadcast.assert_request(); producer1.append_group().unwrap(); - producer1.append_finish().unwrap(); + producer1.finish().unwrap(); drop(producer1); // The consumer should see the track as closed diff --git a/rs/moq-lite/src/model/frame.rs b/rs/moq-lite/src/model/frame.rs index a0dfa67384..86e0b9b8e7 100644 --- a/rs/moq-lite/src/model/frame.rs +++ b/rs/moq-lite/src/model/frame.rs @@ -138,6 +138,14 @@ impl FrameProducer { state.write_chunk(chunk) } + /// Write a chunk of data to the frame. + /// + /// Deprecated: use [`Self::write`] instead. + #[deprecated(note = "use write(chunk) instead")] + pub fn write_chunk>(&mut self, chunk: B) -> Result<()> { + self.write(chunk) + } + /// Verify that all bytes have been written. /// /// Returns [Error::WrongSize] if the bytes written don't match [Frame::size]. diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index 40ac8cc4e2..1062c84828 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -204,7 +204,7 @@ impl TrackProducer { /// Sets the final sequence to the current max_sequence. /// No new groups at or above this sequence can be appended. /// NOTE: Old groups with lower sequence numbers can still arrive. - pub fn append_finish(&mut self) -> Result<()> { + pub fn finish(&mut self) -> Result<()> { let mut state = self.state.modify()?; if state.fin.is_some() { return Err(Error::Closed); @@ -214,13 +214,22 @@ impl TrackProducer { Ok(()) } + /// Mark the track as finished after the last appended group. + /// + /// Deprecated: use [`Self::finish`] for this behavior, or + /// [`Self::finish_at`] to set an explicit final sequence. + #[deprecated(note = "use finish() or finish_at(sequence) instead")] + pub fn close(&mut self) -> Result<()> { + self.finish() + } + /// Mark the track as finished at an exact final sequence. /// /// The caller must pass the current max_sequence exactly. /// Freezes the final boundary at the current max_sequence. /// No new groups at or above that sequence can be created. /// NOTE: Old groups with lower sequence numbers can still arrive. - pub fn insert_finish(&mut self, sequence: u64) -> Result<()> { + pub fn finish_at(&mut self, sequence: u64) -> Result<()> { let mut state = self.state.modify()?; let max = state.max_sequence.ok_or(Error::Closed)?; if state.fin.is_some() || sequence != max { @@ -622,11 +631,11 @@ mod test { fn append_finish_requires_max_and_cannot_be_rewritten() { let mut producer = Track::new("test").produce(); - assert!(producer.append_finish().is_err()); + assert!(producer.finish().is_err()); producer.append_group().unwrap(); - assert!(producer.append_finish().is_ok()); - assert!(producer.append_finish().is_err()); + assert!(producer.finish().is_ok()); + assert!(producer.finish().is_err()); assert!(producer.append_group().is_err()); } @@ -635,16 +644,16 @@ mod test { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 5 }).unwrap(); - assert!(producer.insert_finish(4).is_err()); - assert!(producer.insert_finish(10).is_err()); - assert!(producer.insert_finish(5).is_ok()); + assert!(producer.finish_at(4).is_err()); + assert!(producer.finish_at(10).is_err()); + assert!(producer.finish_at(5).is_ok()); { let state = producer.state.borrow(); assert_eq!(state.fin, Some(5)); } - assert!(producer.insert_finish(5).is_err()); + assert!(producer.finish_at(5).is_err()); assert!(producer.create_group(Group { sequence: 4 }).is_ok()); assert!(producer.create_group(Group { sequence: 5 }).is_err()); } @@ -653,7 +662,7 @@ mod test { async fn next_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); - producer.insert_finish(1).unwrap(); + producer.finish_at(1).unwrap(); let mut consumer = producer.consume(); assert_eq!(consumer.assert_group().info.sequence, 1); @@ -670,7 +679,7 @@ mod test { async fn get_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); - producer.insert_finish(1).unwrap(); + producer.finish_at(1).unwrap(); let consumer = producer.consume(); assert!( diff --git a/rs/moq-mux/src/catalog.rs b/rs/moq-mux/src/catalog.rs index 3a0f95a56a..dcf6fecf46 100644 --- a/rs/moq-mux/src/catalog.rs +++ b/rs/moq-mux/src/catalog.rs @@ -58,8 +58,8 @@ impl CatalogProducer { /// Finish publishing to this catalog. pub fn finish(&mut self) -> Result<(), moq_lite::Error> { - self.hang_track.append_finish()?; - self.msf_track.append_finish()?; + self.hang_track.finish()?; + self.msf_track.finish()?; Ok(()) } } From 3a3220b984b694b66a76eb7695f06976ac5330b5 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 21:58:48 -0800 Subject: [PATCH 7/9] moq-lite: Fix BroadcastDynamic::abort cascade and track finish semantics BroadcastDynamic::abort now cascades to child tracks and drains pending requests, matching BroadcastProducer::abort behavior. Previously pending TrackProducers were left alive and the Drop handler couldn't clean up since the state was already closed. Track finish now stores final_sequence as one past max_sequence (the total number of groups), allowing finish() on empty tracks. Renamed `fin` to `final_sequence` for clarity. Co-Authored-By: Claude Opus 4.6 --- rs/moq-lite/src/model/broadcast.rs | 14 +++++++++- rs/moq-lite/src/model/track.rs | 42 ++++++++++++++++++------------ 2 files changed, 39 insertions(+), 17 deletions(-) diff --git a/rs/moq-lite/src/model/broadcast.rs b/rs/moq-lite/src/model/broadcast.rs index 4266579614..17a001349d 100644 --- a/rs/moq-lite/src/model/broadcast.rs +++ b/rs/moq-lite/src/model/broadcast.rs @@ -184,7 +184,19 @@ impl BroadcastDynamic { /// Abort the broadcast with the given error. pub fn abort(&mut self, err: Error) -> Result<(), Error> { - self.state.abort(err)?; + let mut state = self.state.modify()?; + + // Cascade abort to all child tracks. + for weak in state.tracks.values() { + weak.abort(err.clone()); + } + + // Abort any pending dynamic track requests. + for mut request in state.requests.drain(..) { + request.abort(err.clone()).ok(); + } + + state.abort(err); Ok(()) } diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index 1062c84828..c5fdb5e2a6 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -56,7 +56,7 @@ struct State { duplicates: HashSet, offset: usize, max_sequence: Option, - fin: Option, + final_sequence: Option, } impl State { @@ -71,7 +71,7 @@ impl State { } } - if self.fin.is_some() { + if self.final_sequence.is_some() { Poll::Ready(None) } else { Poll::Pending @@ -86,14 +86,14 @@ impl State { } } - // Once a fin boundary is set, groups at-or-after fin can never exist. - if let Some(fin) = self.fin + // Once final_sequence is set, groups at or past it can never exist. + if let Some(fin) = self.final_sequence && sequence >= fin { return Poll::Ready(None); } - if self.fin.is_some() { + if self.final_sequence.is_some() { return Poll::Ready(None); } @@ -149,7 +149,7 @@ impl TrackProducer { let group = info.produce(); let mut state = self.state.modify()?; - if let Some(fin) = state.fin + if let Some(fin) = state.final_sequence && group.info.sequence >= fin { return Err(Error::Closed); @@ -174,7 +174,7 @@ impl TrackProducer { Some(s) => s.checked_add(1).ok_or(Error::BoundsExceeded)?, None => 0, }; - if let Some(fin) = state.fin + if let Some(fin) = state.final_sequence && sequence >= fin { return Err(Error::Closed); @@ -201,16 +201,18 @@ impl TrackProducer { /// Mark the track as finished after the last appended group. /// - /// Sets the final sequence to the current max_sequence. + /// Sets the final sequence to one past the current max_sequence. /// No new groups at or above this sequence can be appended. /// NOTE: Old groups with lower sequence numbers can still arrive. pub fn finish(&mut self) -> Result<()> { let mut state = self.state.modify()?; - if state.fin.is_some() { + if state.final_sequence.is_some() { return Err(Error::Closed); } - let max = state.max_sequence.ok_or(Error::Closed)?; - state.fin = Some(max); + state.final_sequence = Some(match state.max_sequence { + Some(max) => max.checked_add(1).ok_or(Error::BoundsExceeded)?, + None => 0, + }); Ok(()) } @@ -226,16 +228,16 @@ impl TrackProducer { /// Mark the track as finished at an exact final sequence. /// /// The caller must pass the current max_sequence exactly. - /// Freezes the final boundary at the current max_sequence. + /// Freezes the final boundary at one past the current max_sequence. /// No new groups at or above that sequence can be created. /// NOTE: Old groups with lower sequence numbers can still arrive. pub fn finish_at(&mut self, sequence: u64) -> Result<()> { let mut state = self.state.modify()?; let max = state.max_sequence.ok_or(Error::Closed)?; - if state.fin.is_some() || sequence != max { + if state.final_sequence.is_some() || sequence != max { return Err(Error::Closed); } - state.fin = Some(max); + state.final_sequence = Some(max.checked_add(1).ok_or(Error::BoundsExceeded)?); Ok(()) } @@ -628,10 +630,18 @@ mod test { } #[test] - fn append_finish_requires_max_and_cannot_be_rewritten() { + fn append_finish_cannot_be_rewritten() { let mut producer = Track::new("test").produce(); + // Finishing an empty track is valid (fin = 0, total groups = 0). + assert!(producer.finish().is_ok()); assert!(producer.finish().is_err()); + assert!(producer.append_group().is_err()); + } + + #[test] + fn finish_after_groups() { + let mut producer = Track::new("test").produce(); producer.append_group().unwrap(); assert!(producer.finish().is_ok()); @@ -650,7 +660,7 @@ mod test { { let state = producer.state.borrow(); - assert_eq!(state.fin, Some(5)); + assert_eq!(state.final_sequence, Some(6)); } assert!(producer.finish_at(5).is_err()); From a85ab8c0985bc523f901240eda153295042fb670 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Mon, 2 Mar 2026 22:14:55 -0800 Subject: [PATCH 8/9] Replace start_paused = true with tokio::time::pause() Tests that use sleep/timeout now call tokio::time::pause() explicitly. Tests that don't need time simulation use plain #[tokio::test]. Co-Authored-By: Claude Opus 4.6 --- CLAUDE.md | 2 +- rs/hang/src/container/consumer.rs | 49 ++++++++++++++++++------------- rs/moq-lite/src/model/track.rs | 4 +-- 3 files changed, 32 insertions(+), 23 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 7765eb22ee..0d17dd326e 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -88,7 +88,7 @@ Key architectural rule: The CDN/relay does not know anything about media. Anythi - Run `just check` to execute all tests and linting. - Run `just fix` to automatically fix formating and easy things. - Rust tests are integrated within source files -- Async tests should use `#[tokio::test(start_paused = true)]` to simulate time +- Async tests that sleep should call `tokio::time::pause()` at the start to simulate time instantly ## Branching Strategy diff --git a/rs/hang/src/container/consumer.rs b/rs/hang/src/container/consumer.rs index df9337912c..fd77b93645 100644 --- a/rs/hang/src/container/consumer.rs +++ b/rs/hang/src/container/consumer.rs @@ -272,7 +272,7 @@ mod tests { // ---- Basic Reading ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn read_single_group() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -290,7 +290,7 @@ mod tests { assert!(consumer.read().await.unwrap().is_none()); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn read_multiple_frames_single_group() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -311,7 +311,7 @@ mod tests { assert!(!frames[2].keyframe); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn read_multiple_groups_within_latency() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -334,8 +334,9 @@ mod tests { // Meanwhile, subsequent finished groups accumulate in the pending queue via // next_group, allowing buffer_until to trigger latency-based skipping. - #[tokio::test(start_paused = true)] + #[tokio::test] async fn latency_skip_delivers_recent_groups() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(100)); @@ -371,8 +372,9 @@ mod tests { finisher.await.expect("finisher task panicked"); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn zero_latency_skips_aggressively() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::ZERO); @@ -415,8 +417,9 @@ mod tests { finisher.await.expect("finisher task panicked"); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn latency_skip_correctness() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(100)); @@ -471,8 +474,9 @@ mod tests { // ---- Group Ordering ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn groups_delivered_in_sequence_order() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(500)); @@ -508,7 +512,7 @@ mod tests { finisher.await.expect("finisher task panicked"); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn adjacent_group_flushed_immediately() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -526,7 +530,7 @@ mod tests { // ---- B-frames ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn bframes_within_group() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -546,8 +550,9 @@ mod tests { // ---- Track Lifecycle ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn empty_track_returns_none() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(500)); @@ -563,8 +568,9 @@ mod tests { } } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn track_closed_with_error() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(500)); @@ -585,8 +591,9 @@ mod tests { assert!(result.is_ok(), "Consumer should not hang after track error"); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn closed_resolves_when_track_ends() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(500)); @@ -611,7 +618,7 @@ mod tests { // ---- Gap Recovery ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn gap_in_group_sequence_recovery() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -632,7 +639,7 @@ mod tests { assert!(frames.len() >= 4, "Expected >= 4 frames, got {}", frames.len()); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn gap_at_start_of_sequence() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -651,7 +658,7 @@ mod tests { // ---- Frame Decoding ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn frame_timestamp_and_keyframe_decoding() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -673,7 +680,7 @@ mod tests { assert!(!frames[2].keyframe); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn frame_payload_preserved() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -714,8 +721,9 @@ mod tests { /// causes the select! to restart, creating a new buffer_until for group 1 /// which now has buffered frames. With the fix, read_unbuffered returns None /// and blocks; with the bug, read() re-reads buffered frames infinitely. - #[tokio::test(start_paused = true)] + #[tokio::test] async fn no_infinite_loop_with_buffered_frames() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_secs(10)); @@ -761,7 +769,7 @@ mod tests { // ---- Edge Cases ---- - #[tokio::test(start_paused = true)] + #[tokio::test] async fn large_timestamps() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -778,7 +786,7 @@ mod tests { assert_eq!(frames[0].timestamp.as_micros(), one_hour as u128); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn set_max_latency_changes_behavior() { let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); @@ -807,8 +815,9 @@ mod tests { /// Setup: group 0 is unfinished with B-frames, group 1 at ts(100ms), latency = 40ms. /// Bug: cutoff = 33ms + 40ms = 73ms → group 1's buffer_until sees 100ms >= 73ms → skip /// Fix: cutoff = 66ms + 40ms = 106ms → 100ms < 106ms → no skip, all groups delivered - #[tokio::test(start_paused = true)] + #[tokio::test] async fn max_timestamp_tracks_through_bframes() { + tokio::time::pause(); let mut track = moq_lite::Track::new("test").produce(); let consumer_track = track.consume(); let mut consumer = OrderedConsumer::new(consumer_track, Duration::from_millis(40)); diff --git a/rs/moq-lite/src/model/track.rs b/rs/moq-lite/src/model/track.rs index c5fdb5e2a6..8bfb4593a8 100644 --- a/rs/moq-lite/src/model/track.rs +++ b/rs/moq-lite/src/model/track.rs @@ -668,7 +668,7 @@ mod test { assert!(producer.create_group(Group { sequence: 5 }).is_err()); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn next_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); @@ -685,7 +685,7 @@ mod test { assert!(done.is_none(), "track should finish without waiting for gaps"); } - #[tokio::test(start_paused = true)] + #[tokio::test] async fn get_group_finishes_without_waiting_for_gaps() { let mut producer = Track::new("test").produce(); producer.create_group(Group { sequence: 1 }).unwrap(); From 245498a4c7538da047ae7c947ebc3b1bdf3119f7 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Tue, 3 Mar 2026 09:17:10 -0800 Subject: [PATCH 9/9] Revert the workflow changes. --- .github/workflows/{main.yml => check.yml} | 3 +-- .github/workflows/{dev.yml => semver.yml} | 6 +++--- justfile | 8 ++++---- 3 files changed, 8 insertions(+), 9 deletions(-) rename .github/workflows/{main.yml => check.yml} (94%) rename .github/workflows/{dev.yml => semver.yml} (87%) diff --git a/.github/workflows/main.yml b/.github/workflows/check.yml similarity index 94% rename from .github/workflows/main.yml rename to .github/workflows/check.yml index 3e53c35dbe..87ae78e4e3 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/check.yml @@ -1,4 +1,4 @@ -name: Main +name: Check permissions: id-token: write @@ -6,7 +6,6 @@ permissions: on: pull_request: - branches: ["main"] jobs: check: diff --git a/.github/workflows/dev.yml b/.github/workflows/semver.yml similarity index 87% rename from .github/workflows/dev.yml rename to .github/workflows/semver.yml index b49a3c79e1..92657d40a5 100644 --- a/.github/workflows/dev.yml +++ b/.github/workflows/semver.yml @@ -1,4 +1,4 @@ -name: Check +name: Semver permissions: id-token: write @@ -6,7 +6,7 @@ permissions: on: pull_request: - branches: ["dev"] + branches: ["main"] jobs: check: @@ -28,4 +28,4 @@ jobs: cache-on-failure: true # Run checks with cached dependencies - - run: nix develop --command just check + - run: nix develop --command just semver diff --git a/justfile b/justfile index 2f05afc6d6..7b74fe4d51 100644 --- a/justfile +++ b/justfile @@ -368,10 +368,10 @@ ci: echo "Checking all feature combinations..." cargo hack check --workspace --each-feature --no-dev-deps - # Check semver compatibility against crates.io - # requires: cargo install cargo-semver-checks - # libmoq is an internal C-ABI crate and is intentionally excluded from published-crate semver checks. - echo "Checking semver compatibility..." +# Check semver compatibility against crates.io +# requires: cargo install cargo-semver-checks +# libmoq is an internal C-ABI crate and is intentionally excluded from published-crate semver checks. +semver: cargo semver-checks check-release --workspace --exclude libmoq # Update versions and changelogs via release-plz