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/concept/moq-lite.md
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,7 @@ Patterns never travel as announcements. `origin.dynamic(prefix, route)`
advertises a prefix: the call claims that `prefix` and every path beneath it
*can* be served, not that any exist. A route is a capability, not an inventory:
a subscriber must not treat a prefix as a concrete broadcast name. Use
`create_broadcast(path)` and `announce(route)` when the path is known; use
`publish(path, route)` when the path is known; use
`dynamic` when the set of paths is not, and refuse the requests you will not
serve. A pattern lives in two places: the token, which scopes what a session
may publish and subscribe to, and a local filter a consumer applies to the
Expand Down
9 changes: 4 additions & 5 deletions doc/lib/rs/index.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,12 +40,12 @@ The reference implementation. Every crate is on
```rust
// The Origin is the local hub: the session fills it with remote broadcasts
// and serves your local broadcasts out of it.
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();

let client = moq_tokio::connect::Config::default().init(Default::default())?;
let url = url::Url::parse("https://cdn.moq.dev/anon")?;
// Reconnects on its own; `closed()` resolves when it gives up.
let session = client.with_subscriber(origin.clone()).with_publisher(&origin).connect(url);
let session = client.with_origin(origin.clone()).connect(url);

// Subscribe: wait for a route, resolve the broadcast at its path, read the catalog.
let consumer = origin.consume();
Expand All @@ -62,10 +62,9 @@ while let Some(update) = announced.next().await {
```

```rust
// Publish: create a broadcast on the origin, fill it, then announce its path.
let mut broadcast = origin.create_broadcast("my-stream.hang")?;
// Publish: create and announce a broadcast on the origin, then fill it.
let mut broadcast = origin.publish("my-stream.hang", Default::default())?;
// moq-mux (from a container) or moq-video / moq-audio (from a device) fill it.
broadcast.announce(Default::default())?;
// The route retracts on `unannounce()` or when the broadcast ends. To serve a whole
// subtree on demand instead, `origin.dynamic("room", Default::default())?` yields
// each requested path for the application to accept or reject.
Expand Down
1 change: 1 addition & 0 deletions doc/lib/rs/moq-net.md
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ assert_eq!(

Three operations, on an origin:

- `origin.publish(path, route)` creates and advertises a broadcast in one call.
- `origin.create_broadcast(path)` returns a producer. The broadcast is
reachable by exact path immediately and invisible to discovery until
advertised.
Expand Down
4 changes: 2 additions & 2 deletions doc/lib/rs/moq-room.md
Original file line number Diff line number Diff line change
Expand Up @@ -23,11 +23,11 @@ cargo add moq-room
```

```rust
use moq_net::{Hop, Path};
use moq_net::Path;
use moq_room::{Kind, Room, claims};

let token = key.sign(&claims("meet/demo", "alice")?, None)?;
let origin = moq_tokio::origin::spawn(Hop::random());
let origin = moq_tokio::origin::spawn();
let mut room = Room::new(&origin.consume(), Some(Path::new("alice").to_owned()));
while let Some(event) = room.next().await {
if event.kind == Kind::Camera {
Expand Down
1 change: 0 additions & 1 deletion quest/m1/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,6 @@ the transport line in m2 assumes a single stack.

- [One QUIC backend](/quest/m1/quic-one-backend.md) - quinn and quiche are deleted; noq (and iroh on it) is the only QUIC stack, with the qmux fallbacks untouched
- [Announce event](/quest/m1/api-net-announce.md) - publishers announce prefixes on every wire, consumers scoped by a pattern read the covered path already trimmed, with no `as_prefix().expect()` at 89 call sites
- [Origin scoping](/quest/m1/api-net-origin.md) - `scope(root, patterns)` is one fallible call, a fresh origin has a random hop, and the handles stop derefing to `Hop`
- [Bindings announce match](/quest/m1/api-origin-scopes.md) - every binding takes a pattern scope and reports the announce match with its captures
- [PathPrefixes](/quest/m1/api-path-prefixes.md) - the unused moq_net::PathPrefixes type is deleted before the release
- [Rendition ownership](/quest/m1/api-mux-rendition.md) - one handle publishes a media track and reports its estimate, instead of five
Expand Down
53 changes: 0 additions & 53 deletions quest/m1/api-net-origin.md

This file was deleted.

1 change: 0 additions & 1 deletion quest/m1/api-review-gate.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@ file is deleted on completion) or deletes the quest with a note in
quest is deleted too. No code.

The list: [Announce event](/quest/m1/api-net-announce.md),
[Origin scoping](/quest/m1/api-net-origin.md),
[Rendition ownership](/quest/m1/api-mux-rendition.md),
[Gateway types](/quest/m1/api-gateways.md).

Expand Down
2 changes: 1 addition & 1 deletion rs/hang/examples/subscribe.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ async fn main() -> anyhow::Result<()> {
moq_tokio::Log::new(tracing::Level::DEBUG).init()?;

// Create an origin that the session can publish incoming broadcasts to.
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let consumer = origin.consume();

// Run the subscription and the session in parallel.
Expand Down
2 changes: 1 addition & 1 deletion rs/hang/examples/video.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ async fn main() -> anyhow::Result<()> {
moq_tokio::Log::new(tracing::Level::DEBUG).init()?;

// Create an origin that we can publish to and the session can consume from.
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();

// Run the broadcast production and the session in parallel.
// This is a simple example of how you can concurrently run multiple tasks.
Expand Down
2 changes: 1 addition & 1 deletion rs/libmoq/src/origin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,7 @@ impl Origin {
pub fn create(&mut self) -> Result<Id, Error> {
// Every FFI entry point runs inside `RUNTIME.enter()`, so the driver
// lands on the dedicated libmoq runtime.
self.active.insert(moq_tokio::origin::spawn(moq_net::Hop::random()))
self.active.insert(moq_tokio::origin::spawn())
}

pub fn get(&self, id: Id) -> Result<&moq_net::origin::Producer, Error> {
Expand Down
12 changes: 6 additions & 6 deletions rs/moq-bench/src/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, SystemTime, UNIX_EPOCH};

use moq_tokio::Status;
use moq_tokio::moq_net::{self, Hop, bytes::Bytes};
use moq_tokio::moq_net::{self, bytes::Bytes};
use moq_tokio::moq_net::{broadcast, group, track};
use rand::RngExt;
use serde::{Deserialize, Serialize};
Expand Down Expand Up @@ -96,9 +96,9 @@ pub async fn run(ctx: Connection) {
let url = config.client.url.clone().expect("url required");

// Publish side: an origin we fill with our broadcasts and hand to the session.
let publish = moq_tokio::origin::spawn(Hop::random());
let publish = moq_tokio::origin::spawn();
// Consume side: the session fills this with peer announcements.
let consume = moq_tokio::origin::spawn(Hop::random());
let consume = moq_tokio::origin::spawn();

let namespace = format!("{}/{run_id:08x}", config.name());
let discovery = if config.publishes() {
Expand Down Expand Up @@ -286,7 +286,7 @@ async fn produce(
fn discover(consume: &moq_net::origin::Producer, name: &str) -> moq_net::origin::Consumer {
consume
.consume()
.with_root(name)
.scope(name, &moq_net::Patterns::from(moq_net::Pattern::all()))
.expect("origin must permit the bench namespace")
}

Expand Down Expand Up @@ -737,7 +737,7 @@ mod tests {
tokio::time::pause();

let stats = Arc::new(Stats::default());
let origin = moq_tokio::origin::spawn(Hop::random());
let origin = moq_tokio::origin::spawn();

// The relay-internal broadcast: announced, but with no bench data track.
let _internal = origin.create_broadcast(".stats/node/host").unwrap();
Expand Down Expand Up @@ -776,7 +776,7 @@ mod tests {
#[tokio::test]
async fn named_subscription_waits_for_the_exact_broadcast() {
let stats = Arc::new(Stats::default());
let origin = moq_tokio::origin::spawn(Hop::random());
let origin = moq_tokio::origin::spawn();
let consume = origin.consume();
let task = tokio::spawn(subscribe_named(consume, "bench/run/chat".into(), stats.clone()));

Expand Down
6 changes: 3 additions & 3 deletions rs/moq-boy/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -225,7 +225,7 @@ async fn run(config: &Config) -> Result<()> {
let client = config.client.clone().init(config.quic.clone())?;

// Publish origin: the game session broadcast.
let publish_origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let publish_origin = moq_tokio::origin::spawn();
let default_game_prefix = format!("{}/game", config.prefix);
let default_viewer_prefix = format!("{}/viewer", config.prefix);
let game_prefix = config.prefix_game.as_deref().unwrap_or(&default_game_prefix);
Expand All @@ -243,9 +243,9 @@ async fn run(config: &Config) -> Result<()> {
// Consume origin: viewer broadcasts under the viewer prefix.
// JS publishes viewer feedback at "{viewer_prefix}/{name}/{viewerId}"
let viewer_path = format!("{viewer_prefix}/{name}");
let consume_origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let consume_origin = moq_tokio::origin::spawn();
let viewer_consumer = consume_origin
.with_root(&viewer_path)
.scope(&viewer_path, &moq_net::Patterns::from(moq_net::Pattern::all()))
.expect("viewer prefix should be valid")
.consume();

Expand Down
10 changes: 5 additions & 5 deletions rs/moq-cli/src/complete.rs
Original file line number Diff line number Diff line change
Expand Up @@ -596,7 +596,7 @@ async fn catalog(
/// what tells a relay two sessions carry the same content.
async fn dial(side: &MoqSide, deadline: Instant) -> Option<(moq_net::origin::Producer, moq_tokio::Connection)> {
let url = side.client.url.clone()?;
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();

// Building the client reads the TLS material off disk synchronously, so it goes on
// the blocking pool and under the deadline like everything else: a `--connect-tls-root`
Expand Down Expand Up @@ -775,7 +775,7 @@ mod tests {
/// because the shell exports a relay for the publishing it usually does.
#[tokio::test]
async fn the_environment_cannot_ask_for_a_moq_side() {
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let _alpha = origin.create_broadcast("alpha").expect("alpha");
_alpha.announce(Default::default()).expect("alpha");
let connect = relay(&origin);
Expand Down Expand Up @@ -867,7 +867,7 @@ mod tests {
#[tokio::test]
async fn a_relay_on_the_line_answers_broadcast() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let _alpha = origin.create_broadcast("alpha").expect("alpha");
_alpha.announce(Default::default()).expect("alpha");
let _nested = origin.create_broadcast("room/beta").expect("beta");
Expand All @@ -892,7 +892,7 @@ mod tests {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
use hang::catalog::{AudioCodec, AudioConfig, H264, VideoConfig};

let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();

// Two broadcasts with different renditions, so a completer reading the wrong
// one fails loudly instead of matching by luck.
Expand Down Expand Up @@ -939,7 +939,7 @@ mod tests {
#[tokio::test]
async fn the_catalog_format_on_the_line_is_honored() {
let _env = EnvGuard::clear(&["MOQ_CONNECT"]);
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let broadcast = origin.create_broadcast("room").expect("broadcast");
broadcast.announce(Default::default()).expect("broadcast");

Expand Down
2 changes: 1 addition & 1 deletion rs/moq-cli/src/hls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ pub async fn export(origin: moq_net::origin::Consumer, args: ExportArgs, name: S
moq_net::Pattern::subtree(&name).with_context(|| format!("invalid broadcast name `{name}`"))?,
);
let scoped = origin
.scope(&scope)
.scope("", &scope)
.with_context(|| format!("failed to scope origin to broadcast `{name}`"))?;

let mut config = moq_hls::export::Config::default();
Expand Down
5 changes: 2 additions & 3 deletions rs/moq-cli/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -225,14 +225,13 @@ async fn serve_client(

// What the grant allows, as origin handles rooted where the session dialed.
let token = lease.token();
let rooted = origin.with_root(&token.root);
let publish = directions
.publish
.then(|| rooted.as_ref().and_then(|o| o.scope(&token.subscribe)))
.then(|| origin.scope(&token.root, &token.subscribe).ok())
.flatten();
let subscribe = directions
.consume
.then(|| rooted.as_ref().and_then(|o| o.scope(&token.publish)))
.then(|| origin.scope(&token.root, &token.publish).ok())
.flatten();
if publish.is_none() && subscribe.is_none() {
request.reject(moq_tokio::server::Reject::Forbidden).await.ok();
Expand Down
2 changes: 1 addition & 1 deletion rs/moq-cli/src/play/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,7 @@ mod tests {
async fn subscribe_waits_for_the_announcement() {
tokio::time::pause();

let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let consumer = origin.consume();

// Direct resolution has no route before the announcement.
Expand Down
4 changes: 2 additions & 2 deletions rs/moq-cli/src/publish.rs
Original file line number Diff line number Diff line change
Expand Up @@ -558,7 +558,7 @@ mod tests {

async fn manufacture_input() -> Vec<u8> {
// Create the broadcast on a throwaway origin so the exporter can resolve it by path.
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let mut broadcast = origin.create_broadcast("cli").unwrap();
broadcast.announce(Default::default()).unwrap();
settle().await;
Expand Down Expand Up @@ -668,7 +668,7 @@ mod tests {
// Publish side: `Publish::new(Ts)` builds a `ts::Import<Ext>`, so the verbatim
// streams land in the broadcast instead of being dropped by the media-only path.
// The broadcast is created on a throwaway origin so the exporter can resolve it by path.
let origin = moq_tokio::origin::spawn(moq_net::Hop::random());
let origin = moq_tokio::origin::spawn();
let broadcast = origin.create_broadcast("cli").unwrap();
settle().await;
let mut publish = Publish::new(broadcast, &PublishFormat::Ts, Default::default()).unwrap();
Expand Down
6 changes: 3 additions & 3 deletions rs/moq-cli/src/rtc.rs
Original file line number Diff line number Diff line change
Expand Up @@ -71,11 +71,11 @@ pub async fn listen_export(origin: moq_net::origin::Consumer, name: String, list
moq_net::Pattern::subtree(&name).with_context(|| format!("invalid broadcast name `{name}`"))?,
);
let subscriber = origin
.scope(&scope)
.scope("", &scope)
.with_context(|| format!("failed to scope origin to broadcast `{name}`"))?;
// A WHEP server only reads; it still needs a publisher handle for the shared
// glue, so hand it an unused, empty Origin producer.
let publisher = moq_tokio::origin::spawn(moq_net::Hop::random());
let publisher = moq_tokio::origin::spawn();
let server = moq_rtc::Server::new(server_config(&listen), publisher, subscriber);
serve(server.subscribe_router(), "WHEP", listen).await
}
Expand All @@ -86,7 +86,7 @@ fn scope_producer(origin: &moq_net::origin::Producer, name: &str) -> anyhow::Res
moq_net::Pattern::subtree(name).with_context(|| format!("invalid broadcast name `{name}`"))?,
);
origin
.scope(&scope)
.scope("", &scope)
.with_context(|| format!("failed to scope origin to broadcast `{name}`"))
}

Expand Down
4 changes: 2 additions & 2 deletions rs/moq-cli/src/transcode.rs
Original file line number Diff line number Diff line change
Expand Up @@ -128,11 +128,11 @@ pub async fn run(moq: MoqSide, args: Args, net: Net) -> anyhow::Result<()> {
.url
.clone()
.context("`transcode` requires a relay: pass --connect <url>")?;
let publish = moq_tokio::origin::spawn(moq_net::Hop::random());
let publish = moq_tokio::origin::spawn();
// A session drop closes the source broadcast and ends the run: the outage is
// surfaced rather than transcoded over. The reconnect loop covers the dial;
// restarting after a mid-run drop is the caller's call.
let remote = moq_tokio::origin::spawn(moq_net::Hop::random());
let remote = moq_tokio::origin::spawn();
let session = net
.client(moq.client.clone())?
.with_publisher(&publish)
Expand Down
Loading
Loading