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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion doc/lib/c/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ and `target/include/moq.h`.
- **Raw decode output.** `moq_video_decoder_output` selects the decoded CPU pixel format (`MOQ_VIDEO_PIXEL_FORMAT_I420` or `_RGBA`) and target size (`width`/`height`, both zero for native; otherwise even and non-zero). Unknown formats and invalid sizes fail `moq_decode_video` before subscribing; accepted requests deliver exactly that layout or fail on the terminal callback.
- **Encoded video metadata.** `moq_video_init.hint` is a zero-initialized `moq_video_hint` with `has_*` flags for coded dimensions, bitrate (bits per second), frame rate, and latency preference. Hints seed a video codec track's catalog; detected dimensions take precedence.
- **Client config.** A zeroed `moq_client_config` means the defaults for every knob, which is what lets a new one be appended without disturbing callers. Fields cover protocol (`versions`), TLS (`tls_fingerprints`, `tls_roots`, `tls_cert`/`_key`, `tls_host_name`), transport (`bind`, `connect_timeout_us`, the Happy Eyeballs delays, `websocket_enabled`), and tuning (reconnect backoff, `quic_*`). Every duration is in microseconds. A knob whose default isn't zero carries a `has_*` flag, so setting `backoff_timeout_us = 0` needs `has_backoff_timeout = true` to mean "retry forever" rather than "use the default". `moq_client_defaults()` reports what a NULL config dials with.
- **Demand.** A watcher on a published track (`moq_publish_track_demand`, `moq_publish_media_demand`, `moq_encode_video_demand`, `moq_encode_audio_demand`) calls `on_demand` with `MOQ_DEMAND_USED` or `MOQ_DEMAND_UNUSED` right away and again on every change, so an encoder on a battery-powered device runs only while someone is watching. The first call is the current state, so a track that went unused before the watcher existed still reports it. `moq_publish_demand_cancel` stops it; the terminal callback still fires. A container has no single demand and is refused. Demand follows the last real subscriber: an origin that served the track drops its source copy on the unused edge and keeps only the groups it already cached warm for 30 seconds, so the cache linger does not delay the unused edge.
- **Demand.** A watcher on a published track (`moq_publish_track_demand`, `moq_publish_media_demand`, `moq_encode_video_demand`, `moq_encode_audio_demand`) calls `on_demand` with `MOQ_DEMAND_USED` or `MOQ_DEMAND_UNUSED` right away and again on every change, so an encoder on a battery-powered device runs only while someone is watching. The first call is the current state, so a track that went unused before the watcher existed still reports it. `moq_publish_demand_cancel` stops it; the terminal callback still fires. A container has no single demand and is refused. Demand follows the last real subscriber: an origin that served the track drops its source copy on the unused edge and keeps only the finished groups it already cached warm for 30 seconds, so the cache linger does not delay the unused edge.
- **Requests.** `moq_publish_dynamic` serves subscriptions to tracks the broadcast never declared: each arrives as a request handle, read its name with `moq_track_request_name`, then `moq_track_request_accept` (a raw track handle), `moq_track_request_video` / `_audio` (the media handle `moq_publish_video` / `_audio` return), or `moq_track_request_abort` with an application code the subscriber sees. Without a live handler an unknown name is refused. `moq_publish_track_dynamic` does the same for fetches of groups a track no longer has cached, delivered as `moq_group_request_*` (`sequence`, `priority`, `frame_start`); `moq_group_request_accept` starts the producer at `frame_start` so written frames keep their group indices. Register it with `moq_track_request_dynamic` before accepting a track that was itself requested by a fetch, so that pending group survives the transition. Both handlers stop with `moq_publish_dynamic_cancel`.
- **Everything the bindings can do** ([list](/lib/#what-every-binding-can-do)): media publish and consume with the catalog managed for you, raw pixels and PCM with the codec inside (`moq_encode_video`, `moq_encode_audio`, and the `moq_decode_*` mirrors), raw tracks with timestamps and datagrams, JSON snapshot and stream tracks, group fetch, catalog sections, shared video properties, and stalled hints. The three advertising operations are `moq_origin_create_broadcast` (locally discoverable producer), `moq_publish_announce` / `moq_publish_unannounce` (exact-path advertisement), and `moq_origin_dynamic` (a claim over a path prefix and everything beneath it; `""` for everything). A route is a capability, not an inventory. `moq_origin_announced` takes a literal prefix and an optional relative pattern filter; `moq_announce_update.prefix` stays relative to the origin, while `captures` reports what each wildcard matched when `has_captures` is true.

Expand Down
5 changes: 5 additions & 0 deletions rs/moq-net/src/model/group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -803,6 +803,11 @@ impl Producer {
self.alive.aborted.load(Ordering::Acquire)
}

/// Whether the group was finished: it holds every frame it will ever have.
pub(crate) fn is_finished(&self) -> bool {
self.state.read().fin.is_some()
}

/// The index of the first frame this group still holds, or `None` once it has been
/// aborted. Non-zero when the group started later (see [`Self::start_at`]); a reader
/// positioned below it is [`Error::Lagged`].
Expand Down
9 changes: 7 additions & 2 deletions rs/moq-net/src/model/origin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1697,12 +1697,17 @@ impl Drop for WarmCopy {
}
}

/// Cache `source`'s groups on a new local track the origin owns.
/// Cache `source`'s finished groups on a new local track the origin owns.
fn warm_copy(source: &track::Consumer) -> Option<WarmCopy> {
let info = source.cached_info()?;
let mut track = track::Producer::new(Arc::new(source.broadcast().clone()), source.name(), info);
for (group, visible) in source.cached_groups() {
let _ = track.adopt_group(group, visible);
// An open group is left for the re-splice to deliver whole. Dropping the source
// copy resets it mid-transfer, and its dead head would anchor the next takeover
// mid-group, asking upstream for a tail no returning reader can use.
if group.is_finished() {
let _ = track.adopt_group(group, visible);
}
}
let dynamic = track.dynamic();
Some(WarmCopy {
Expand Down
99 changes: 99 additions & 0 deletions rs/moq-net/tests/rejoin.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,99 @@
//! A reader leaving and rejoining a track relayed over the in-memory mock transport.

mod support;

use std::time::Duration;

use moq_net::{Hop, Timestamp, Version, group, track};
use support::harness::{MockConnectOptions, connect_mock};

/// Maximum time any single test may run before being treated as a deadlock.
const TEST_TIMEOUT: Duration = Duration::from_secs(10);

/// Build an origin producer, spawning its driver on the ambient runtime.
fn produce_origin(hop: u64) -> moq_net::origin::Producer {
let (producer, driver) = moq_net::origin::Producer::new(moq_net::origin::Config::new(Hop::new(hop).unwrap()));
tokio::spawn(support::harness::run(driver));
producer
}

async fn read_all(group: &mut group::Consumer) -> Vec<Vec<u8>> {
let mut frames = Vec::new();
while let Some(frame) = group.read_frame().await.expect("group aborted") {
frames.push(frame.payload.to_vec());
}
frames
}

/// The group in flight when the last reader left comes back whole on rejoin.
///
/// The relay cancels its idle upstream subscription, resetting that group mid-transfer.
/// Resuming the rejoin at the frame where the reset landed would ask upstream for a tail
/// whose head is gone, so the group would never reach the returning reader.
#[tokio::test]
async fn rejoin_recovers_the_group_reset_on_leave() {
for version in ["moq-lite-05", "moq-lite-06"] {
tokio::time::timeout(TEST_TIMEOUT, async {
let publisher = produce_origin(1);
let relay = produce_origin(2);

let broadcast = publisher.create_broadcast("bench").unwrap();
let track = broadcast.create_track("video", None).unwrap();
broadcast.announce(Default::default()).unwrap();

let mut options = MockConnectOptions::new(version.parse::<Version>().unwrap());
options.server_publish = Some(publisher);
options.client_subscribe = Some(relay.clone());
let _pair = connect_mock(options).await;

let consumer = relay.consume();
consumer.routed("bench").await.unwrap();
let remote = consumer.request_broadcast("bench").await.unwrap();

let ts = |ms| Timestamp::from_millis(ms).unwrap();
let prefs = || track::Subscription::default().with_max_age(Duration::from_secs(10));

let mut group = track.append_group().unwrap();
group.write_frame(ts(0), b"a0".as_ref()).unwrap();
group.finish().unwrap();
let mut open = track.append_group().unwrap();
open.write_frame(ts(100), b"b0".as_ref()).unwrap();

let mut sub = remote.track("video").unwrap().subscribe(prefs()).await.unwrap();
loop {
let mut group = sub.recv_group().await.unwrap().unwrap();
group.read_frame().await.unwrap().unwrap();
if group.sequence == 1 {
break;
}
}

// Leave mid-group; the publisher cuts the open group once demand is gone.
drop(sub);
track.unused().await.unwrap();
open.write_frame(ts(133), b"b1".as_ref()).unwrap();
open.finish().unwrap();

let mut live = track.append_group().unwrap();
live.write_frame(ts(5000), b"c0".as_ref()).unwrap();
live.finish().unwrap();

let mut sub = remote.track("video").unwrap().subscribe(prefs()).await.unwrap();
let mut rejoined = Vec::new();
loop {
let mut group = sub.recv_group().await.unwrap().unwrap();
let sequence = group.sequence;
rejoined.push((sequence, read_all(&mut group).await));
if sequence == 2 {
break;
}
}
assert!(
rejoined.contains(&(1, vec![b"b0".to_vec(), b"b1".to_vec()])),
"{version}: the reset group never came back whole: {rejoined:?}"
);
})
.await
.unwrap_or_else(|_| panic!("{version} timed out"));
}
}
Loading