Skip to main content

moq_net/model/
broadcast.rs

1//! A broadcast is a named collection of tracks, split into a [Producer] and [Consumer] handle.
2//!
3//! A [Producer] creates tracks on demand: a [Consumer] subscribes by name, and the
4//! producer either serves a track it already has or is handed a [`track::Request`] to
5//! fill. Both handles are refcounted clones of one broadcast, which ends on
6//! [`Producer::close`] or when the last producer drops.
7//!
8//! [Info] is the broadcast's static metadata, fixed for its lifetime.
9use crate::{cache, stats, track};
10use std::{
11	collections::{HashMap, VecDeque},
12	sync::Arc,
13	task::{Poll, ready},
14};
15
16use crate::Error;
17use crate::origin::Route;
18
19use super::origin_impl::Announcer;
20use super::{Requests, WeakCache};
21
22/// A collection of media tracks that can be published and subscribed to.
23///
24/// Create via [`Info::produce`] to obtain both [`Producer`] and [`Consumer`] pair.
25/// This is the broadcast's static identity, fixed for its lifetime.
26#[derive(Clone, Debug)]
27#[non_exhaustive]
28pub struct Info {
29	/// The cache pool this broadcast's tracks and groups inherit.
30	pub pool: cache::Pool,
31
32	/// Ceiling on each track's media-timestamp retention window.
33	pub cache_duration: std::time::Duration,
34
35	/// The path this broadcast is named by, which relative references in a catalog it
36	/// serves (hang's `broadcast` field) resolve against.
37	///
38	/// [`origin::Producer::create_broadcast`](super::origin::Producer::create_broadcast) stamps
39	/// the path the broadcast was created at, relative to the origin root (including through a
40	/// scoped producer). Every [`Consumer`] an origin hands out is then re-stamped with the path
41	/// *that handle* was requested or announced at, relative to its cursor's root, since the
42	/// same broadcast can be reached under more than one name: a dynamic handler may serve a
43	/// standalone broadcast at any path, and a rooted cursor names a broadcast more tightly
44	/// than the origin does.
45	///
46	/// Empty (the default) for a standalone broadcast with no origin, which is then its own
47	/// root: any `..` reference escapes.
48	pub path: crate::PathOwned,
49}
50
51impl Default for Info {
52	fn default() -> Self {
53		Self {
54			pool: cache::Pool::new(cache::Config::default().with_expiry(cache::DEFAULT_EXPIRY)),
55			cache_duration: std::time::Duration::MAX,
56			path: crate::PathOwned::default(),
57		}
58	}
59}
60
61impl Info {
62	/// Create a new broadcast with default metadata.
63	pub fn new() -> Self {
64		Self::default()
65	}
66
67	/// Consume this [Info] to create a producer that carries its metadata.
68	///
69	/// Keep the returned [`Producer`] alive for as long as the broadcast should stay
70	/// available, and end it with [`Producer::close`]. See the note on [`Producer`].
71	pub fn produce(self) -> Producer {
72		Producer::new(self)
73	}
74}
75
76#[derive(Default)]
77struct BroadcastState {
78	// Weak references for deduplication. Doesn't prevent track auto-close.
79	// Keyed by the track's shared `Arc<str>` name (the same Arc the handle holds).
80	// The cache reclaims closed entries incrementally on insert so a long-lived
81	// broadcast churning distinct track names stays bounded by the live count.
82	tracks: WeakCache<Arc<str>, track::TrackWeak>,
83
84	// Shared across suffixes and producer clones; cache eviction must not reset IDs.
85	unique: u64,
86
87	// Pending requests keyed by track name, coalescing concurrent `track()` calls
88	// and waiting for a dynamic handler to accept or deny them. A request leaves
89	// here once handed out (the handler caches it in `tracks`, so lookups keep
90	// coalescing onto it there).
91	requests: Requests<Arc<str>, track::Request>,
92
93	// Route-fed mode (a relay/origin "front"): tracks are spliced logical tracks
94	// joined across per-session tracks. `None` for an ordinary broadcast.
95	spliced: Option<SplicedState>,
96
97	// Set once the broadcast ends: `Producer::close()`, an abort, or the last
98	// producer-side handle dropping. Every lookup after it answers `Unroutable`.
99	closing: bool,
100
101	// Set only by the deprecated `Producer::finish()`, for `Consumer::is_finished`.
102	finished: bool,
103
104	// The error passed to `Producer::abort()`, reported by `Consumer::closed`.
105	// `None` for a finish or a dropped producer (reported as `Error::Dropped`).
106	abort: Option<Error>,
107}
108
109/// The spliced (route-fed) half of a broadcast: logical tracks that outlive any
110/// single session, plus the queue of tracks awaiting a serving route.
111#[derive(Default)]
112struct SplicedState {
113	// Logical tracks by name, owned strongly: they live as long as the broadcast
114	// (the origin's front), not as long as any consumer.
115	tracks: HashMap<Arc<str>, super::resume::Producer>,
116
117	// Names awaiting assignment to a route, in request order.
118	pending: VecDeque<Arc<str>>,
119}
120
121impl BroadcastState {
122	/// Insert a track weak handle into the lookup, returning an error if a live
123	/// track already holds the name. A closed entry under the name is reclaimed.
124	fn insert_track(&mut self, weak: track::TrackWeak) -> Result<(), Error> {
125		match self.tracks.insert(weak.name().clone(), weak) {
126			Some(_) => Err(Error::Duplicate),
127			None => Ok(()),
128		}
129	}
130
131	/// Resolve every name the broadcast never filled, so subscribers waiting on a
132	/// [`track::Info`] that can no longer arrive fail with `err` instead of parking.
133	///
134	/// Covers a reservation nobody accepted and a request still queued for a handler.
135	/// A request a [`Dynamic`] already took is left alone: it may be in flight to a peer
136	/// still serving it, and the handler answers or drops it. So is a track that carries
137	/// its info: an end there is that publisher's call, and its cache stays readable.
138	fn reject_unserved(&mut self, err: Error) {
139		for request in self.requests.drain_queued() {
140			request.reject(err.clone());
141		}
142		for track in self.tracks.iter() {
143			track.reject(err.clone());
144		}
145	}
146
147	/// Live demand: a subscribed spliced track (route-fed broadcast), or a
148	/// pending request / consumed track (ordinary broadcast). See [`Demand`].
149	fn is_used(&self) -> bool {
150		if let Some(spliced) = &self.spliced {
151			return spliced.tracks.values().any(|track| track.is_used());
152		}
153		!self.requests.is_empty() || self.tracks.iter().any(|track| track.is_used())
154	}
155
156	/// Park `waiter` on every per-track channel feeding [`Self::is_used`]: the
157	/// consumer counts live on those channels, and their flips don't write this
158	/// state, so a watcher registered here alone would miss the edge. `want`
159	/// picks the direction; each channel only arms while its side is unmet.
160	fn register_demand(&self, waiter: &kio::Waiter, want: bool) {
161		if let Some(spliced) = &self.spliced {
162			for track in spliced.tracks.values() {
163				let _ = match want {
164					true => track.poll_used(waiter),
165					false => track.poll_unused(waiter),
166				};
167			}
168			return;
169		}
170		for track in self.tracks.iter() {
171			match want {
172				true => track.poll_used(waiter),
173				false => track.poll_unused(waiter),
174			}
175		}
176	}
177}
178
179/// Manages tracks within a broadcast.
180///
181/// Create tracks up front with [Self::create_track], reserve a name to fill in
182/// later with [Self::reserve_track], or handle on-demand consumer requests via
183/// [Self::dynamic].
184///
185/// # Lifetime
186///
187/// A broadcast lives until [`Self::close`] or until the last [`Producer`] (or
188/// [`Dynamic`]) drops, whichever comes first; both end it the same way. Children
189/// do *not* keep it alive: cloning a [`Consumer`] or holding a [`track::Producer`]
190/// does nothing for the broadcast's lifetime.
191#[derive(Clone)]
192pub struct Producer {
193	// Held behind an Arc so each track born from this broadcast can inherit a shared
194	// handle (threaded down by [`Self::create_track`] / [`Self::reserve_track`]).
195	info: Arc<Info>,
196
197	// Broadcast liveness, shared with every `Dynamic`. Consumers watch it (read-only)
198	// for close; the guard ends the broadcast when the last of those handles drops.
199	alive: Arc<Alive>,
200
201	// Track registry plus the dynamic request queue, mutated by producers and
202	// consumers alike under one lock.
203	state: kio::Shared<BroadcastState>,
204
205	// Ingress stats scope, set by a tagged `origin::Producer` at
206	// `create_broadcast`. Inherited by the tracks this producer creates. Empty
207	// (no-op) for an untagged broadcast.
208	stats: stats::Scope,
209}
210
211impl Producer {
212	/// Create a producer for the given broadcast metadata. Prefer [`Info::produce`].
213	pub fn new(info: Info) -> Self {
214		let state = kio::Shared::<BroadcastState>::default();
215		Self {
216			info: Arc::new(info),
217			alive: Alive::new(state.clone()),
218			state,
219			stats: stats::Scope::default(),
220		}
221	}
222
223	/// Attach an ingress stats scope, inherited by the tracks created on this
224	/// broadcast. Set by a tagged `origin::Producer` at `create_broadcast`.
225	pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
226		self.stats = scope;
227		self
228	}
229
230	/// Attach the advertisement of this broadcast's exact path. Set by
231	/// `origin::Producer::create_broadcast`; a standalone broadcast has none.
232	pub(crate) fn with_announcer(self, announcer: Announcer) -> Self {
233		*self.alive.announcer.lock() = Some(announcer);
234		self
235	}
236
237	/// Advertise this broadcast's exact path as a route, or re-price the standing
238	/// advertisement in place.
239	///
240	/// Until this is called the broadcast exists for nobody: announce cursors do
241	/// not list it and requests for its path fail with [`Error::Unroutable`], for
242	/// local consumers and peers alike. Call it once the tracks a subscriber needs
243	/// first (a catalog) exist, so the advertisement lands with them in place:
244	/// consumers act on it immediately. The route retracts on [`unannounce`](Self::unannounce),
245	/// [`close`](Self::close), or the last producer dropping.
246	///
247	/// Fails with [`Error::Closed`] on a standalone broadcast (one not created
248	/// through an origin, so there is nothing to announce into), once the broadcast
249	/// has closed, or once the origin's driver has been dropped.
250	pub fn announce(&self, route: Route) -> Result<(), Error> {
251		let mut announcer = self.alive.announcer.lock();
252		let announcer = announcer.as_mut().ok_or(Error::Closed)?;
253		announcer.announce(route)
254	}
255
256	/// Retract this broadcast's advertisement, if any, from local consumers and
257	/// peers alike. New requests for the path fail with [`Error::Unroutable`] and
258	/// the broadcast the origin served from it ends, while tracks already in
259	/// flight carry on to their own end. [`announce`](Self::announce) brings it
260	/// back.
261	pub fn unannounce(&self) {
262		self.alive.unannounce();
263	}
264
265	/// Create a route-fed (spliced) broadcast: consumer track lookups mint logical
266	/// tracks that are spliced across per-session tracks, queued for a route to
267	/// serve. Used by the origin for broadcasts reached over the network.
268	pub(crate) fn new_spliced(info: Info) -> Self {
269		let state = kio::Shared::new(BroadcastState {
270			spliced: Some(SplicedState::default()),
271			..Default::default()
272		});
273		Self {
274			info: Arc::new(info),
275			alive: Alive::new(state.clone()),
276			state,
277			// The origin-owned spliced broadcast stays untagged: egress attribution is
278			// applied when a tagged `origin::Consumer` hands the consumer out.
279			stats: stats::Scope::default(),
280		}
281	}
282
283	/// The broadcast's static metadata, fixed when it was created.
284	pub fn info(&self) -> &Info {
285		&self.info
286	}
287
288	/// A watch-only handle to the broadcast's demand. See [`Demand`].
289	pub fn demand(&self) -> Demand {
290		Demand {
291			alive: self.alive.token.consume().weak(),
292			state: self.state.clone(),
293		}
294	}
295
296	/// Produce a new track and insert it into the broadcast.
297	///
298	/// Pass a name and an optional [`track::Info`], so a bare name works:
299	/// `create_track("video", None)`.
300	pub fn create_track(
301		&self,
302		name: impl Into<Arc<str>>,
303		info: impl Into<Option<track::Info>>,
304	) -> Result<track::Producer, Error> {
305		let name = name.into();
306		let info = info.into().unwrap_or_default();
307		let mut state = self.state.lock();
308
309		// A consumer may have requested this name before it existed (a live
310		// [`Dynamic`] queues such requests). Creating the track fulfills that
311		// request: its consumers resolve against this very producer. Without
312		// this they would be stranded, since the name is taken the moment the
313		// track exists, so no handler could ever serve their queue entry.
314		if let Some(request) = state.requests.take(name.as_ref()) {
315			let track = request.with_stats(self.stats.clone()).accept(info);
316			// Cache it like a served request so concurrent lookups coalesce; a
317			// live same-name entry cannot exist (its presence would have kept
318			// the request from queuing).
319			let _ = state.tracks.insert(name, track.weak());
320			return Ok(track);
321		}
322
323		let track = track::Producer::new(self.info.clone(), name, info).with_stats(self.stats.clone());
324		state.insert_track(track.weak())?;
325		Ok(track)
326	}
327
328	/// Reserve a track by name without finalizing its [`track::Info`].
329	///
330	/// Returns a [`track::Request`] already discoverable by consumers; call
331	/// [`track::Request::accept`] to set its info and start producing. Use this when
332	/// the producer can't pick the track's properties (e.g. timescale) until it has
333	/// inspected the media, the same shape as a consumer-driven
334	/// [`Dynamic::requested_track`].
335	///
336	/// Subscribers wait on the name until it is accepted, so a reservation the producer
337	/// ends up never filling has to be dropped or rejected. Ending the broadcast
338	/// resolves whatever is left.
339	pub fn reserve_track(&self, name: impl Into<Arc<str>>) -> Result<track::Request, Error> {
340		let request = track::Request::new(self.info.clone(), name).with_stats(self.stats.clone());
341		self.state.lock().insert_track(request.weak())?;
342		Ok(request)
343	}
344
345	/// Create a track with a unique name using the given suffix.
346	///
347	/// Uses [`Self::unique_name`]; minted names are never reused, even after closure.
348	pub fn unique_track(&self, suffix: &str, info: impl Into<Option<track::Info>>) -> Result<track::Producer, Error> {
349		let name = self.unique_name(suffix);
350		self.create_track(name, info)
351	}
352
353	/// Generate a unique track name from a suffix without creating the track.
354	///
355	/// Returns `{id}{suffix}` with an increasing ID shared across all suffixes and
356	/// producer clones in this broadcast, skipping names already in the lookup.
357	/// A digit-leading suffix gets a `-` separator so it cannot be confused with the ID.
358	/// Minted names are never reused, even if no track is created or it is closed.
359	/// Explicit calls to [`Self::create_track`] can still reuse names.
360	///
361	/// # Panics
362	///
363	/// Panics if the broadcast exhausts its `u64` IDs.
364	pub fn unique_name(&self, suffix: &str) -> String {
365		let mut state = self.state.lock();
366		let separator = if suffix.starts_with(|c: char| c.is_ascii_digit()) {
367			"-"
368		} else {
369			""
370		};
371		loop {
372			let id = state.unique;
373			state.unique = id.checked_add(1).expect("unique track IDs exhausted");
374			let name = format!("{id}{separator}{suffix}");
375			if !state.tracks.contains_key(name.as_str()) {
376				return name;
377			}
378		}
379	}
380
381	/// Create a dynamic producer that handles on-demand track requests from consumers.
382	pub fn dynamic(&self) -> Dynamic {
383		Dynamic::new(
384			self.info.clone(),
385			self.alive.clone(),
386			self.state.clone(),
387			self.stats.clone(),
388		)
389	}
390
391	/// Poll for the next spliced track awaiting a serving route, returning its name
392	/// and logical producer. Route-fed broadcasts only.
393	pub(crate) fn poll_spliced_assigned(&self, waiter: &kio::Waiter) -> Poll<(Arc<str>, super::resume::Producer)> {
394		let mut state = ready!(self.state.poll(waiter, |state| {
395			match &state.spliced {
396				Some(spliced) if !spliced.pending.is_empty() => Poll::Ready(()),
397				_ => Poll::Pending,
398			}
399		}));
400
401		let spliced = state.spliced.as_mut().expect("predicate guaranteed spliced");
402		let name = spliced.pending.pop_front().expect("predicate guaranteed a request");
403		let producer = spliced.tracks.get(&name).expect("pending name without a track").clone();
404		Poll::Ready((name, producer))
405	}
406
407	/// Let go of every spliced track, aborting with `err` the ones never handed
408	/// out by [`Self::poll_spliced_assigned`]. Called when the broadcast ends:
409	/// whoever took the others decides how they end.
410	pub(crate) fn release_spliced(&self, err: Error) {
411		let mut state = self.state.lock();
412		if let Some(spliced) = state.spliced.as_mut() {
413			for name in std::mem::take(&mut spliced.pending) {
414				if let Some(producer) = spliced.tracks.get_mut(&name) {
415					let _ = producer.abort(err.clone());
416				}
417			}
418			spliced.tracks.clear();
419		}
420	}
421
422	/// Remove the spliced track `producer` from under `name`, unless it has a reader.
423	///
424	/// Returns false only when a reader holds it: lookups hand out readers under the
425	/// same lock, so one arriving after the caller decided is never cut off. Anything
426	/// else (removed, or already replaced under the name) is gone as far as the caller
427	/// is concerned.
428	pub(crate) fn forget_spliced(&self, name: &str, producer: &super::resume::Producer) -> bool {
429		let mut state = self.state.lock();
430		let Some(spliced) = state.spliced.as_mut() else {
431			return true;
432		};
433		match spliced.tracks.get(name) {
434			Some(current) if current.is_clone(producer) => {
435				if current.is_used() {
436					return false;
437				}
438				spliced.tracks.remove(name);
439				true
440			}
441			_ => true,
442		}
443	}
444
445	/// Create a consumer of this one publisher's broadcast.
446	///
447	/// A view of this broadcast object, not of its path: a new publisher at the same
448	/// path is never spliced into it, so it ends when this broadcast does. Go through
449	/// an origin for a consumer that should not care which publisher serves the path.
450	pub fn consume(&self) -> Consumer {
451		Consumer {
452			info: self.info.clone(),
453			alive: self.alive.token.consume(),
454			state: self.state.clone(),
455			stats: stats::Scope::default(),
456		}
457	}
458
459	/// End the broadcast for good, whether or not other clones are still alive.
460	///
461	/// Retracts its announcement and local discovery; a later [`Self::announce`] fails
462	/// with [`Error::Closed`]. Tracks already handed out carry on and end with their
463	/// own finish or abort. Every later [`Consumer::track`], and every name reserved or
464	/// still queued for a handler, answers [`Error::Unroutable`]: the same answer an
465	/// origin gives for a path nobody publishes. A request a [`Dynamic`] already took is
466	/// left for that handler to answer.
467	///
468	/// Dropping the last producer does the same. Closing twice is a no-op.
469	pub fn close(&self) {
470		self.alive.close();
471	}
472
473	#[doc(hidden)]
474	#[deprecated(note = "use close(); a broadcast end carries no cause")]
475	pub fn finish(&self) {
476		self.alive.end(true);
477	}
478
479	#[doc(hidden)]
480	#[deprecated(note = "use close(); a broadcast end carries no cause")]
481	pub fn abort(self, err: Error) -> Result<(), Error> {
482		{
483			let mut state = self.state.lock();
484			if state.closing {
485				return Err(Error::Closed);
486			}
487			state.closing = true;
488			state.abort = Some(err.clone());
489			// Same as a finish: an unserved name is answerable now, with the reason the
490			// broadcast ended. Published tracks keep their cache (no cascade).
491			state.reject_unserved(err);
492		}
493		let _ = self.alive.token.close();
494		self.alive.retire();
495		Ok(())
496	}
497
498	/// Return true if this is the same broadcast instance.
499	pub fn is_clone(&self, other: &Self) -> bool {
500		self.state.same_channel(&other.state)
501	}
502}
503
504/// Ends the broadcast on [`Producer::close`] or when the last [`Producer`] or
505/// [`Dynamic`] drops, closing the liveness channel every [`Consumer`] watches.
506///
507/// A refcount rather than a "am I the last one?" check inside `Drop`: that answer is
508/// a snapshot, and acting on it is exactly what invalidates it.
509struct Alive {
510	token: kio::Producer<()>,
511	state: kio::Shared<BroadcastState>,
512	// The advertisement of the broadcast's exact path, owned here so it retracts
513	// with the broadcast: on close, abort, or the last producer-side handle
514	// dropping. `None` for a standalone broadcast.
515	announcer: kio::Lock<Option<Announcer>>,
516}
517
518impl Alive {
519	fn new(state: kio::Shared<BroadcastState>) -> Arc<Self> {
520		Arc::new(Self {
521			token: kio::Producer::default(),
522			state,
523			announcer: kio::Lock::new(None),
524		})
525	}
526
527	/// Withdraw the path's advertisement, if any; the broadcast stays alive.
528	fn unannounce(&self) {
529		if let Some(announcer) = self.announcer.lock().as_mut() {
530			announcer.withdraw();
531		}
532	}
533
534	/// End the broadcast. See [`Producer::close`].
535	fn close(&self) {
536		self.end(false);
537	}
538
539	/// End the broadcast, recording the deprecated `finished` flag in the same locked
540	/// transition that claims the end, so a racing `abort` can't win after it's set.
541	fn end(&self, finished: bool) {
542		{
543			let mut state = self.state.lock();
544			if std::mem::replace(&mut state.closing, true) {
545				return;
546			}
547			state.finished = finished;
548			// A name that was reserved or queued but never served can't arrive now,
549			// and `Consumer::track` answers `Unroutable` for one asked about after this
550			// point. Say the same to whoever asked earlier.
551			state.reject_unserved(Error::Unroutable);
552		}
553		let _ = self.token.close();
554		self.retire();
555	}
556
557	/// End the broadcast's advertising for good: retract the standing advertisement
558	/// and drop the announcer, so a later `announce` fails with `Closed`.
559	fn retire(&self) {
560		let announcer = self.announcer.lock().take();
561		// Dropped outside the announcer lock: the entry's removal re-syncs the
562		// origin's cursors under the origin's own lock.
563		drop(announcer);
564	}
565}
566
567impl Drop for Alive {
568	fn drop(&mut self) {
569		self.close();
570	}
571}
572
573#[cfg(test)]
574#[allow(missing_docs)] // test-only assertion helpers
575impl Producer {
576	pub fn assert_create_track(
577		&mut self,
578		name: impl Into<Arc<str>>,
579		info: impl Into<Option<track::Info>>,
580	) -> track::Producer {
581		self.create_track(name, info).expect("should not have errored")
582	}
583}
584
585/// A session-owned handle to a source broadcast created via
586/// [`crate::origin::Producer::create_broadcast`], closing it on drop even while the
587/// session's serve machines still hold its [`Dynamic`].
588///
589/// A peer's retraction and a dead session end the source the same way: a broadcast
590/// carries no end cause past this hop. Shared by the lite and IETF subscribers.
591pub(crate) struct SourceGuard(Producer);
592
593impl SourceGuard {
594	pub fn new(producer: Producer) -> Self {
595		Self(producer)
596	}
597}
598
599impl Drop for SourceGuard {
600	fn drop(&mut self) {
601		self.0.close();
602	}
603}
604
605/// Handles on-demand track creation for a broadcast.
606///
607/// When a consumer requests a track that doesn't exist, the dynamic producer
608/// picks up the request via [`Self::requested_track`] and either
609/// [`track::Request::accept`]s it with a concrete [`track::Info`] or
610/// [`track::Request::reject`]s it. Dropped when no longer needed; pending requests
611/// are automatically aborted.
612#[derive(Clone)]
613pub struct Dynamic {
614	info: Arc<Info>,
615	// Keeps the broadcast alive while a handler exists (mirrors a producer).
616	alive: Arc<Alive>,
617	state: kio::Shared<BroadcastState>,
618	// Ingress stats scope, applied to the tracks this handler serves. Empty (no-op)
619	// for an untagged broadcast.
620	stats: stats::Scope,
621	// Declared after `alive` so it drops second: when this was the broadcast's last
622	// handle, `Alive` has already ended it and answered every queued request
623	// `Unroutable`, so the handler's own `Dropped` rejection finds nothing left.
624	_handler: Handler,
625}
626
627/// Counts one live [`Dynamic`], rejecting the queued requests when the last one drops.
628struct Handler(kio::Shared<BroadcastState>);
629
630impl Handler {
631	fn new(state: kio::Shared<BroadcastState>) -> Self {
632		state.lock().requests.add_handler();
633		Self(state)
634	}
635}
636
637impl Clone for Handler {
638	fn clone(&self) -> Self {
639		// Count each live handle, or dropping a clone would flip the handler count to
640		// zero and future `track` calls would return `NotFound`.
641		Self::new(self.0.clone())
642	}
643}
644
645impl Drop for Handler {
646	fn drop(&mut self) {
647		// Decrement and reject under one lock, so a `track` call that saw a live
648		// handler through the same lock can't slip a request past the rejection.
649		let mut state = self.0.lock();
650		if state.requests.remove_handler() {
651			// No handlers left to fulfill pending requests; reject them so consumers
652			// don't block forever on tracks nobody will serve.
653			for request in state.requests.drain_queued() {
654				request.reject(Error::Dropped);
655			}
656		}
657	}
658}
659
660impl Dynamic {
661	fn new(info: Arc<Info>, alive: Arc<Alive>, state: kio::Shared<BroadcastState>, stats: stats::Scope) -> Self {
662		Self {
663			info,
664			alive,
665			_handler: Handler::new(state.clone()),
666			state,
667			stats,
668		}
669	}
670
671	/// The broadcast's static metadata, fixed when it was created.
672	pub fn info(&self) -> &Info {
673		&self.info
674	}
675
676	/// Poll for the next consumer-requested track, without blocking.
677	///
678	/// Returns [`Error::Closed`] once the broadcast has ended, so a serving loop
679	/// knows to stop and release its handle.
680	pub fn poll_requested_track(&mut self, waiter: &kio::Waiter) -> Poll<Result<track::Request, Error>> {
681		let mut state = ready!(self.state.poll(waiter, |state| {
682			if state.requests.has_queued() || state.closing {
683				Poll::Ready(())
684			} else {
685				Poll::Pending
686			}
687		}));
688
689		if state.closing && !state.requests.has_queued() {
690			return Poll::Ready(Err(Error::Closed));
691		}
692
693		let name = state.requests.pop().expect("predicate guaranteed a request");
694		let pending = state.requests.remove(&name).expect("popped key must be pending");
695		// Cache the served track so concurrent lookups coalesce onto it. If a live track already
696		// holds the name (a publish raced the request), `insert` keeps it rather than shadowing it.
697		let _ = state.tracks.insert(name, pending.weak());
698		// Attribute the served track to this broadcast's ingress scope (no-op untagged).
699		Poll::Ready(Ok(pending.claim().with_stats(self.stats.clone())))
700	}
701
702	/// Block until a consumer requests a track, returning a [`track::Request`] to serve.
703	pub async fn requested_track(&mut self) -> Result<track::Request, Error> {
704		kio::wait(|waiter| self.poll_requested_track(waiter)).await
705	}
706
707	/// Create a consumer that can subscribe to tracks in this broadcast.
708	pub fn consume(&self) -> Consumer {
709		Consumer {
710			info: self.info.clone(),
711			alive: self.alive.token.consume(),
712			state: self.state.clone(),
713			stats: stats::Scope::default(),
714		}
715	}
716
717	/// Block until the broadcast ends, by [`Producer::close`] or every producer dropping.
718	///
719	/// Returns [`Error::Dropped`], or the error passed to the deprecated `abort`.
720	pub async fn closed(&self) -> Error {
721		kio::wait(|waiter| self.poll_closed(waiter)).await
722	}
723
724	/// Poll-based variant of [`Self::closed`].
725	pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
726		ready!(self.alive.token.poll_closed(waiter));
727		Poll::Ready(self.state.read().abort.clone().unwrap_or(Error::Dropped))
728	}
729
730	/// Return true if this is the same broadcast instance.
731	pub fn is_clone(&self, other: &Self) -> bool {
732		self.state.same_channel(&other.state)
733	}
734}
735
736#[cfg(test)]
737use futures::FutureExt;
738
739#[cfg(test)]
740#[allow(missing_docs)] // test-only assertion helpers
741impl Dynamic {
742	pub fn assert_request(&mut self) -> track::Request {
743		self.requested_track()
744			.now_or_never()
745			.expect("should not have blocked")
746			.expect("should not have errored")
747	}
748
749	pub fn assert_no_request(&mut self) {
750		assert!(self.requested_track().now_or_never().is_none(), "should have blocked");
751	}
752}
753
754/// Subscribe to arbitrary broadcast/tracks.
755///
756/// Its close signal means this broadcast object ended, not that the path went
757/// offline: announcements say whether a path is live, and a new publisher may
758/// announce the same path again.
759pub struct Consumer {
760	info: Arc<Info>,
761	// Broadcast liveness (read-only): watched for close.
762	alive: kio::Consumer<()>,
763	// Track registry plus request queue; `track()` reads the registry and enqueues requests.
764	state: kio::Shared<BroadcastState>,
765	// Egress stats scope, set by a tagged `origin::Consumer` at the broadcast
766	// handoff. Inherited by the tracks subscribed through this handle. Empty (no-op)
767	// for an untagged broadcast.
768	stats: stats::Scope,
769}
770
771impl Clone for Consumer {
772	fn clone(&self) -> Self {
773		Self {
774			info: self.info.clone(),
775			alive: self.alive.clone(),
776			state: self.state.clone(),
777			stats: self.stats.clone(),
778		}
779	}
780}
781
782impl Consumer {
783	/// Attach an egress stats scope, inherited by the tracks subscribed through this
784	/// handle. Set by a tagged `origin::Consumer` at the broadcast handoff.
785	pub(crate) fn with_stats(mut self, scope: stats::Scope) -> Self {
786		self.stats = scope;
787		self
788	}
789
790	/// Stamp the path this handle was handed out at, overriding [`Info::path`].
791	///
792	/// The origin applies it to every broadcast it resolves, because the name belongs to
793	/// the (broadcast, cursor) pair rather than to the broadcast: what a catalog's relative
794	/// references resolve against is where the *reader* found the broadcast, not where its
795	/// producer happened to create it. Free when the two already agree, which is the case
796	/// for a broadcast created at the path an unrooted cursor asks for.
797	pub(crate) fn with_path(mut self, path: crate::PathOwned) -> Self {
798		if self.info.path != path {
799			let mut info = (*self.info).clone();
800			info.path = path;
801			self.info = Arc::new(info);
802		}
803		self
804	}
805
806	/// The broadcast's metadata, as reached through this handle.
807	pub fn info(&self) -> &Info {
808		&self.info
809	}
810
811	/// Get a handle to a track on this broadcast.
812	///
813	/// Fails with [`Error::Unroutable`] once the broadcast has ended.
814	pub fn track(&self, name: &str) -> Result<track::Consumer, Error> {
815		// Rebind the track to *this* handle's view of the broadcast, so a catalog track
816		// resolves its relative references against the path we were handed out at rather
817		// than the one the producer was created at, and tag it with this broadcast's egress
818		// scope so its subscriptions, fetches, and groups are attributed to the same broadcast.
819		self.track_inner(name)
820			.map(|track| track.with_broadcast(self.info.clone()).with_stats(self.stats.clone()))
821	}
822
823	fn track_inner(&self, name: &str) -> Result<track::Consumer, Error> {
824		let mut state = self.state.lock();
825
826		// An ended broadcast serves nothing new, not even a track it still has:
827		// the lookup answers what a fresh `request_broadcast` for the path would.
828		// Tracks already handed out are untouched.
829		if state.closing {
830			return Err(Error::Unroutable);
831		}
832
833		// A route-fed broadcast mints spliced logical tracks: they outlive any
834		// session, and a route is asked (via the pending queue) to start serving.
835		if let Some(spliced) = state.spliced.as_mut() {
836			// An aborted logical track is a verdict from the sources attached at
837			// the time, not a property of the name: a publisher that had not yet
838			// created the track may have it now. Drop it so this request reaches a
839			// source again, exactly as the plain lookup below reclaims a closed
840			// entry. A *finished* one stays, since its cache is still readable,
841			// until the front forgets it after going unread for its linger.
842			//
843			// So a name, once finished, is never spliced onto again: a publisher
844			// that finishes a track and publishes it again is serving new content,
845			// not resuming this one, and a subscriber has to re-read the catalog and
846			// re-initialize rather than be spliced onto it. Resuming the same
847			// content across routes is the transparent case, and that is what
848			// `resume::Producer` already does. Publish new content under a new
849			// name.
850			if spliced.tracks.get(name).is_some_and(|track| track.is_aborted()) {
851				spliced.tracks.remove(name);
852			}
853			if let Some(producer) = spliced.tracks.get(name) {
854				return Ok(track::Consumer::spliced(
855					name.into(),
856					self.info.clone(),
857					producer.consume(),
858				));
859			}
860			let name: Arc<str> = name.into();
861			let producer = super::resume::Producer::new();
862			let consumer = producer.consume();
863			spliced.tracks.insert(name.clone(), producer);
864			spliced.pending.push_back(name.clone());
865			return Ok(track::Consumer::spliced(name, self.info.clone(), consumer));
866		}
867
868		// Reuse a live producer if one is already publishing the track. `get` drops a
869		// closed entry and returns `None`, so we fall through to a fresh request.
870		if let Some(weak) = state.tracks.get(name) {
871			match weak.try_consume() {
872				Some(consumer) => return Ok(consumer),
873				// It closed between the liveness probe and the count bump (an idle
874				// teardown committing under us). Reclaim it and request the track
875				// again, rather than handing back a consumer of a dead track.
876				None => {
877					state.tracks.remove(name);
878				}
879			}
880		}
881
882		if let Some(pending) = state.requests.join(name) {
883			// Coalesce onto a queued request for the same name.
884			return Ok(pending.consume());
885		}
886
887		// Allocate the name once and share the same Arc across the request, the
888		// requests map, and the FIFO order. The request inherits the broadcast's
889		// cache pool through its `Arc<Info>`, same as a producer-created track.
890		let name: Arc<str> = name.into();
891		let request = track::Request::new(self.info.clone(), name.clone());
892		let consumer = request.consume();
893
894		// With no handler alive to serve it, the request is dropped: `NotFound` beats
895		// handing back a consumer that would only resolve `Dropped`.
896		if state.requests.insert(name, request).is_err() {
897			return Err(Error::NotFound);
898		}
899
900		Ok(consumer)
901	}
902
903	/// A watch-only handle to the broadcast's demand. See [`Demand`].
904	///
905	/// The consumer-side sibling of [`Producer::demand`], for a holder that has
906	/// only a read handle. Holding this handle, or the [`Consumer`] it came from,
907	/// is not itself demand.
908	///
909	/// Demand going away is [`Demand::unused`] resolving. The broadcast going
910	/// away is [`Error::Dropped`], which here means the upstream producer ended.
911	pub fn demand(&self) -> Demand {
912		Demand {
913			alive: self.alive.weak(),
914			state: self.state.clone(),
915		}
916	}
917
918	/// Block until the broadcast ends, by [`Producer::close`] or every producer dropping.
919	///
920	/// Returns [`Error::Dropped`], or the error passed to the deprecated `abort`.
921	pub async fn closed(&self) -> Error {
922		self.alive.closed().await;
923		self.state.read().abort.clone().unwrap_or(Error::Dropped)
924	}
925
926	/// Returns true once the broadcast has ended.
927	pub fn is_closed(&self) -> bool {
928		self.alive.is_closed()
929	}
930
931	/// Whether the broadcast has ended, observed under the state lock. The origin's
932	/// dispatcher treats a rejection from such a source as imminent detach rather
933	/// than a strike.
934	pub(crate) fn is_closing(&self) -> bool {
935		self.state.read().closing
936	}
937
938	#[doc(hidden)]
939	#[deprecated(note = "a broadcast end carries no cause")]
940	pub fn is_finished(&self) -> bool {
941		self.state.read().finished
942	}
943
944	/// Register a [`kio::Waiter`] that fires when the broadcast closes.
945	///
946	/// Returns [`Poll::Ready`] if already closed, otherwise [`Poll::Pending`] after
947	/// arming the waiter. Useful for composing close-detection into a larger poll
948	/// without spawning a task per broadcast.
949	pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
950		self.alive.poll_closed(waiter)
951	}
952
953	/// Check if this is the exact same instance of a broadcast.
954	pub fn is_clone(&self, other: &Self) -> bool {
955		self.state.same_channel(&other.state)
956	}
957
958	/// Create a weak reference that doesn't keep the broadcast alive.
959	///
960	/// Used to deduplicate dynamically-served broadcasts in the origin: a live weak yields
961	/// a shared clone, a closed one is discarded so the next request re-serves.
962	pub(crate) fn weak(&self) -> WeakConsumer {
963		WeakConsumer {
964			info: self.info.clone(),
965			alive: self.alive.weak(),
966			state: self.state.clone(),
967		}
968	}
969}
970
971/// A weak reference to a broadcast that doesn't prevent it from closing.
972///
973/// Mirrors [`track::TrackWeak`]: held by the origin's dynamic cache to share one
974/// dynamically-served broadcast across repeat requests without pinning it alive.
975/// Only the `alive` handle needs to be weak; a [`kio::Shared`] carries no liveness,
976/// so holding the state outright pins nothing.
977#[derive(Clone)]
978pub(crate) struct WeakConsumer {
979	info: Arc<Info>,
980	alive: kio::ConsumerWeak<()>,
981	state: kio::Shared<BroadcastState>,
982}
983
984impl WeakConsumer {
985	/// Upgrade to a full [`Consumer`] sharing the same broadcast state.
986	pub fn consume(&self) -> Consumer {
987		Consumer {
988			info: self.info.clone(),
989			alive: self.alive.consume(),
990			state: self.state.clone(),
991			stats: stats::Scope::default(),
992		}
993	}
994}
995
996impl super::WeakEntry for WeakConsumer {
997	fn is_closed(&self) -> bool {
998		self.alive.is_closed()
999	}
1000
1001	fn same_channel(&self, other: &Self) -> bool {
1002		self.state.same_channel(&other.state)
1003	}
1004}
1005
1006/// A cloneable, watch-only handle to a broadcast's subscriber demand.
1007///
1008/// Obtained from [`Producer::demand`] or [`Consumer::demand`]; the broadcast-level sibling of
1009/// [`track::Demand`](crate::track::Demand). Demand means live interest in the
1010/// broadcast's content: a subscribed spliced track on a route-fed broadcast, or
1011/// a pending track request / a consumed track on an ordinary one. A publisher
1012/// uses it to run expensive work only while someone is watching, and routing
1013/// uses it to advertise a warm copy at zero cost.
1014///
1015/// It's a weak handle: it neither keeps the broadcast alive nor counts as
1016/// demand itself. Once every producer is gone, [`used`](Self::used) /
1017/// [`unused`](Self::unused) return [`Error::Dropped`].
1018#[derive(Clone)]
1019pub struct Demand {
1020	alive: kio::ConsumerWeak<()>,
1021	state: kio::Shared<BroadcastState>,
1022}
1023
1024impl Demand {
1025	/// Whether the broadcast has live demand right now.
1026	///
1027	/// A point-in-time snapshot with no registration; use [`Self::used`] /
1028	/// [`Self::unused`] (or their `poll_*` forms) to wait for the edge.
1029	pub fn is_used(&self) -> bool {
1030		self.state.read().is_used()
1031	}
1032
1033	/// Block until the broadcast has demand. Resolves immediately if it already
1034	/// does; returns [`Error::Dropped`] once every producer is gone.
1035	pub async fn used(&self) -> Result<(), Error> {
1036		kio::wait(|waiter| self.poll_used(waiter)).await
1037	}
1038
1039	/// Block until the broadcast has no demand. Resolves immediately if it has
1040	/// none; returns [`Error::Dropped`] once every producer is gone.
1041	pub async fn unused(&self) -> Result<(), Error> {
1042		kio::wait(|waiter| self.poll_unused(waiter)).await
1043	}
1044
1045	/// Poll-based variant of [`Self::used`].
1046	pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1047		self.poll_demand(waiter, true)
1048	}
1049
1050	/// Poll-based variant of [`Self::unused`].
1051	pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<(), Error>> {
1052		self.poll_demand(waiter, false)
1053	}
1054
1055	fn poll_demand(&self, waiter: &kio::Waiter, want: bool) -> Poll<Result<(), Error>> {
1056		// Closure is checked first, matching `track::Demand`: a dead broadcast
1057		// reports Dropped rather than pretending to answer.
1058		if self.alive.poll_closed(waiter).is_ready() {
1059			return Poll::Ready(Err(Error::Dropped));
1060		}
1061		let ready = self.state.poll(waiter, |state| {
1062			// The consumer counts live on the per-track channels, whose flips
1063			// don't write this state: park on those channels too so the edge
1064			// wakes us, then recompute here.
1065			state.register_demand(waiter, want);
1066			match state.is_used() == want {
1067				true => Poll::Ready(()),
1068				false => Poll::Pending,
1069			}
1070		});
1071		match ready {
1072			Poll::Ready(_) => Poll::Ready(Ok(())),
1073			Poll::Pending => Poll::Pending,
1074		}
1075	}
1076}
1077
1078#[cfg(test)]
1079#[allow(missing_docs)] // test-only assertion helpers
1080impl Consumer {
1081	pub fn assert_not_closed(&self) {
1082		assert!(self.closed().now_or_never().is_none(), "should not be closed");
1083	}
1084
1085	pub fn assert_closed(&self) {
1086		assert!(self.closed().now_or_never().is_some(), "should be closed");
1087	}
1088}
1089
1090#[cfg(test)]
1091mod test {
1092	use super::*;
1093	use std::time::Duration;
1094
1095	#[test]
1096	fn unique_names_are_never_reused() {
1097		let producer = Info::new().produce();
1098		let name = producer.unique_name(".opus");
1099		assert_eq!(name, "0.opus");
1100		let track = producer.create_track(name.clone(), None).unwrap();
1101		assert_eq!(producer.unique_name(".opus"), "1.opus");
1102		drop(track);
1103	}
1104
1105	#[test]
1106	fn unique_names_survive_closed_track_pruning() {
1107		let producer = Info::new().produce();
1108		let consumer = producer.consume();
1109		let track = producer.unique_track(".opus", None).unwrap();
1110		assert_eq!(track.name(), "0.opus");
1111		drop(track);
1112		assert!(matches!(consumer.track_inner("0.opus"), Err(Error::NotFound)));
1113		assert_eq!(producer.unique_name(".opus"), "1.opus");
1114	}
1115
1116	#[test]
1117	fn unique_name_skips_a_live_collision() {
1118		let producer = Info::new().produce();
1119		let track = producer.create_track("0.opus", None).unwrap();
1120		assert_eq!(producer.unique_name(".opus"), "1.opus");
1121		drop(track);
1122		assert_eq!(producer.unique_name(".opus"), "2.opus");
1123	}
1124
1125	#[test]
1126	fn unique_names_share_a_counter() {
1127		let producer = Info::new().produce();
1128		assert_eq!(producer.unique_name("-video"), "0-video");
1129		assert_eq!(producer.clone().unique_name("-audio"), "1-audio");
1130		assert_eq!(producer.unique_name("-video"), "2-video");
1131	}
1132
1133	#[test]
1134	fn unique_names_separate_numeric_suffixes() {
1135		let producer = Info::new().produce();
1136		assert_eq!(producer.unique_name(""), "0");
1137		let name = producer.unique_name("2");
1138		assert_eq!(name, "1-2");
1139		for _ in 2..12 {
1140			producer.unique_name("");
1141		}
1142		assert_eq!(producer.unique_name(""), "12");
1143	}
1144
1145	/// Await with a timeout so a missed demand wake fails the test instead of
1146	/// hanging it (time is paused, so the timeout fires instantly when idle).
1147	async fn expect<T>(fut: impl Future<Output = T>) -> T {
1148		tokio::time::timeout(Duration::from_secs(1), fut)
1149			.await
1150			.expect("timed out waiting for a demand edge")
1151	}
1152
1153	/// Demand on an ordinary broadcast tracks subscriber interest, not
1154	/// production: a live track producer alone is unused, a consumed track is
1155	/// used, and both edges wake parked waiters.
1156	#[tokio::test]
1157	async fn demand_ordinary() {
1158		tokio::time::pause();
1159
1160		let producer = Info::new().produce();
1161		let consumer = producer.consume();
1162		let demand = producer.demand();
1163
1164		// No demand yet; `unused` resolves immediately.
1165		assert!(!demand.is_used());
1166		demand.unused().await.unwrap();
1167
1168		// Producing alone is not demand.
1169		let _track = producer.create_track("a", None).unwrap();
1170		assert!(!demand.is_used());
1171
1172		// A consumer appearing wakes a parked `used`.
1173		let (used, handle) = tokio::join!(expect(demand.used()), async { consumer.track("a").unwrap() });
1174		used.unwrap();
1175		assert!(demand.is_used());
1176
1177		// The last consumer dropping wakes a parked `unused`.
1178		let (unused, ()) = tokio::join!(expect(demand.unused()), async { drop(handle) });
1179		unused.unwrap();
1180		assert!(!demand.is_used());
1181
1182		// Every producer gone: both edges report the closure.
1183		producer.close();
1184		assert!(matches!(demand.used().await, Err(Error::Dropped)));
1185		assert!(matches!(demand.unused().await, Err(Error::Dropped)));
1186	}
1187
1188	/// Demand on a spliced (route-fed) broadcast follows the logical tracks'
1189	/// consumers, which is what flips a relay's advertised cost.
1190	#[tokio::test]
1191	async fn demand_spliced() {
1192		tokio::time::pause();
1193
1194		let producer = Producer::new_spliced(Info::new());
1195		let consumer = producer.consume();
1196		let demand = producer.demand();
1197		let watched = consumer.demand();
1198
1199		assert!(!demand.is_used());
1200		assert!(!watched.is_used());
1201		let track = consumer.track("video").unwrap();
1202		assert!(demand.is_used());
1203		assert!(watched.is_used());
1204
1205		// Dropping the only consumer wakes a parked `unused`, even though the
1206		// logical track itself stays cached in the broadcast.
1207		let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1208		unused.unwrap();
1209		assert!(!demand.is_used());
1210		assert!(!watched.is_used());
1211
1212		// A repeat consumer for the cached track counts again.
1213		let _track = consumer.track("video").unwrap();
1214		assert!(demand.is_used());
1215	}
1216
1217	/// A consumer demand handle distinguishes lost demand from a dropped producer.
1218	#[tokio::test]
1219	async fn consumer_demand_reports_dropped_producer() {
1220		let producer = Producer::new_spliced(Info::new());
1221		let consumer = producer.consume();
1222		let watched = consumer.demand();
1223
1224		let track = consumer.track("video").unwrap();
1225		assert!(watched.is_used());
1226
1227		let (unused, ()) = tokio::join!(expect(watched.unused()), async { drop(track) });
1228		unused.unwrap();
1229
1230		drop(producer);
1231		assert!(matches!(watched.used().await, Err(Error::Dropped)));
1232		assert!(matches!(watched.unused().await, Err(Error::Dropped)));
1233	}
1234
1235	/// Subscribe and assert the result hasn't resolved yet (it stays pending until
1236	/// a publisher accepts). Returns the pending subscription to resolve after accepting.
1237	macro_rules! subscribe_pending {
1238		($consumer:expr, $name:expr) => {{
1239			let pending = $consumer.track($name).unwrap().subscribe(None);
1240			assert!(
1241				pending.poll_ok(&kio::Waiter::noop()).is_pending(),
1242				"subscribe should stay pending until the request is accepted"
1243			);
1244			pending
1245		}};
1246	}
1247
1248	#[tokio::test]
1249	async fn insert() {
1250		let mut producer = Info::new().produce();
1251
1252		// Create the track before any consumer exists.
1253		let track1 = producer.assert_create_track("track1", None);
1254		track1.append_group().unwrap();
1255
1256		let consumer = producer.consume();
1257
1258		// The track already exists, so subscribe resolves immediately.
1259		let mut track1_sub = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1260		track1_sub.assert_group();
1261
1262		let track2 = producer.assert_create_track("track2", None);
1263
1264		let consumer2 = producer.consume();
1265		let mut track2_consumer = consumer2.track("track2").unwrap().subscribe(None).await.unwrap();
1266		track2_consumer.assert_no_group();
1267
1268		track2.append_group().unwrap();
1269
1270		track2_consumer.assert_group();
1271	}
1272
1273	#[tokio::test]
1274	async fn closed() {
1275		let mut producer = Info::new().produce();
1276		let dynamic = producer.dynamic();
1277
1278		let consumer = producer.consume();
1279		consumer.assert_not_closed();
1280
1281		// Create a new track and insert it into the broadcast (resolves immediately).
1282		let track1 = producer.assert_create_track("track1", None);
1283		let mut track1c = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1284
1285		// A track nobody publishes stays pending until accepted.
1286		let track2_fut = subscribe_pending!(consumer, "track2");
1287
1288		// Dropping the last dynamic handler rejects pending requests, but must NOT
1289		// cascade to externally-owned tracks.
1290		drop(dynamic);
1291
1292		// track2 was a pending dynamic request, so its subscribe surfaces the rejection.
1293		assert!(track2_fut.await.is_err());
1294
1295		// track1's producer is held outside the broadcast, so it survives.
1296		assert!(!track1.is_closed());
1297		track1c.assert_not_closed();
1298	}
1299
1300	/// `close()` ends the broadcast for every clone at once, and a second close is a no-op.
1301	#[tokio::test]
1302	async fn close_ends_every_clone() {
1303		let producer = Info::new().produce();
1304		let clone = producer.clone();
1305		let consumer = producer.consume();
1306
1307		producer.close();
1308		assert!(matches!(consumer.closed().await, Error::Dropped));
1309		assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
1310		assert!(matches!(clone.consume().track("video"), Err(Error::Unroutable)));
1311
1312		clone.close();
1313		producer.close();
1314	}
1315
1316	/// Dropping the last producer ends the broadcast exactly like `close()`.
1317	#[tokio::test]
1318	async fn drop_ends_like_close() {
1319		let producer = Info::new().produce();
1320		let consumer = producer.consume();
1321		drop(producer);
1322		assert!(matches!(consumer.closed().await, Error::Dropped));
1323		assert!(matches!(consumer.track("video"), Err(Error::Unroutable)));
1324	}
1325
1326	/// The deprecated end APIs keep their old causes until they are removed.
1327	#[tokio::test]
1328	#[allow(deprecated)]
1329	async fn deprecated_end_causes() {
1330		let producer = Info::new().produce();
1331		let consumer = producer.consume();
1332		producer.abort(Error::Timeout).unwrap();
1333		assert!(matches!(consumer.closed().await, Error::Timeout));
1334		assert!(!consumer.is_finished());
1335
1336		let producer = Info::new().produce();
1337		let consumer = producer.consume();
1338		producer.finish();
1339		assert!(matches!(consumer.closed().await, Error::Dropped));
1340		assert!(consumer.is_finished());
1341	}
1342
1343	#[tokio::test]
1344	async fn requests() {
1345		let mut producer = Info::new().produce().dynamic();
1346
1347		let consumer = producer.consume();
1348		let consumer2 = consumer.clone();
1349
1350		// Two subscribers to the same name coalesce into one request.
1351		let track1_fut = subscribe_pending!(consumer, "track1");
1352		let track2_fut = subscribe_pending!(consumer2, "track1");
1353
1354		// There should be exactly one request to serve.
1355		let request = producer.assert_request();
1356		producer.assert_no_request();
1357		assert_eq!(request.name(), "track1");
1358
1359		// Accept it, which resolves both waiting subscribers.
1360		let track3 = request.accept(None);
1361		let mut track1 = track1_fut.await.unwrap();
1362		let mut track2 = track2_fut.await.unwrap();
1363
1364		track1.assert_not_closed();
1365		track1.assert_is_clone(&track2);
1366		track3.subscribe(None).assert_is_clone(&track1);
1367
1368		// Append a group and make sure they all get it.
1369		track3.append_group().unwrap();
1370		track1.assert_group();
1371		track2.assert_group();
1372
1373		// A pending request is cancelled when the dynamic producer is dropped.
1374		let track4_fut = subscribe_pending!(consumer, "track2");
1375		drop(producer);
1376		assert!(track4_fut.await.is_err());
1377
1378		// With no dynamic producer left, requesting the handle fails outright.
1379		let track5 = consumer2.track("track3");
1380		assert!(track5.is_err(), "should have errored");
1381	}
1382
1383	#[tokio::test]
1384	async fn stale_producer() {
1385		let mut broadcast = Info::new().produce().dynamic();
1386		let consumer = broadcast.consume();
1387
1388		// Subscribe to a track and serve it.
1389		let track1_fut = subscribe_pending!(consumer, "track1");
1390		let producer1 = broadcast.assert_request().accept(None);
1391		let mut track1 = track1_fut.await.unwrap();
1392
1393		// Close the producer (simulating publisher disconnect).
1394		producer1.append_group().unwrap();
1395		producer1.finish().unwrap();
1396		drop(producer1);
1397
1398		// The consumer should see the track as closed.
1399		track1.assert_closed();
1400
1401		// Subscribe again to the same track: should get a NEW producer, not the stale one.
1402		let track2_fut = subscribe_pending!(consumer, "track1");
1403		let producer2 = broadcast.assert_request().accept(None);
1404		let mut track2 = track2_fut.await.unwrap();
1405		track2.assert_not_closed();
1406		track2.assert_not_clone(&track1);
1407
1408		// The new consumer should receive the new group.
1409		producer2.append_group().unwrap();
1410		track2.assert_group();
1411	}
1412
1413	#[tokio::test(start_paused = true)]
1414	async fn requested_unused() {
1415		let mut broadcast = Info::new().produce().dynamic();
1416		let bc = broadcast.consume();
1417
1418		// Subscribe to a track that doesn't exist yet, then serve it.
1419		let c1_fut = subscribe_pending!(bc, "unknown_track");
1420		let producer1 = broadcast.assert_request().accept(None);
1421		let consumer1 = c1_fut.await.unwrap();
1422
1423		// The producer should NOT be unused yet because there's a consumer.
1424		assert!(
1425			producer1.unused().now_or_never().is_none(),
1426			"track producer should be used"
1427		);
1428
1429		// A second subscriber reuses the live producer (fast path / dedup).
1430		let consumer2 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1431		consumer2.assert_is_clone(&consumer1);
1432
1433		drop(consumer1);
1434		assert!(
1435			producer1.unused().now_or_never().is_none(),
1436			"track producer should be used"
1437		);
1438
1439		drop(consumer2);
1440		assert!(
1441			producer1.unused().now_or_never().is_some(),
1442			"track producer should be unused after all consumers are dropped"
1443		);
1444
1445		// While the producer is still alive, re-subscribing to the same name reuses
1446		// it (no new request). This is what lets the relay linger upstream
1447		// subscriptions across transient consumer churn.
1448		let consumer3 = bc.track("unknown_track").unwrap().subscribe(None).await.unwrap();
1449		consumer3.assert_is_clone(&producer1.subscribe(None));
1450		broadcast.assert_no_request();
1451		drop(consumer3);
1452
1453		// Aborting the producer closes its lookup entry; the next subscribe sees the
1454		// stale weak, evicts it, and creates a fresh request.
1455		producer1.abort(Error::Cancel).unwrap();
1456
1457		let c4_fut = subscribe_pending!(bc, "unknown_track");
1458		let producer2 = broadcast.assert_request().accept(None);
1459		let consumer4 = c4_fut.await.unwrap();
1460		drop(consumer4);
1461		assert!(
1462			producer2.unused().now_or_never().is_some(),
1463			"new track producer should be unused after its consumer is dropped"
1464		);
1465	}
1466
1467	/// Creating a track a consumer already requested fulfills that request: the
1468	/// waiting subscriber resolves against the created producer, and no handler
1469	/// ever sees the (now-taken) name. Without this the requester is stranded:
1470	/// the name exists the moment the track does, so the queue entry could
1471	/// never be served under it.
1472	#[tokio::test]
1473	async fn create_track_fulfills_queued_request() {
1474		let producer = Info::new().produce();
1475		let mut dynamic = producer.dynamic();
1476		let bc = dynamic.consume();
1477
1478		// Queue a request for a track that doesn't exist yet.
1479		let subscribing = subscribe_pending!(bc, "video");
1480
1481		// The producer creates the track before any handler drains the queue.
1482		let track = producer.create_track("video", None).unwrap();
1483		let mut sub = subscribing.await.expect("fulfilled by create_track");
1484
1485		// The fulfilled subscription is live against this very producer.
1486		track.append_group().unwrap();
1487		sub.recv_group().await.expect("recv").expect("group");
1488
1489		// The handler never sees the request; a fresh subscribe reuses the track.
1490		dynamic.assert_no_request();
1491		let again = bc.track("video").unwrap().subscribe(None).await.unwrap();
1492		again.assert_is_clone(&track.subscribe(None));
1493	}
1494
1495	// Cloning a `Dynamic` and dropping the clone must not flip the handler
1496	// count to zero. The relay's lite subscriber clones the
1497	// dynamic per spawned subscribe; if Clone skipped the increment, the
1498	// first finished subscribe would tear down the broadcast and any
1499	// follow-up `track` would return `NotFound`.
1500	#[tokio::test]
1501	async fn dynamic_clone_keeps_alive() {
1502		let broadcast = Info::new().produce().dynamic();
1503		let consumer = broadcast.consume();
1504
1505		let clone = broadcast.clone();
1506		drop(clone);
1507
1508		// Original handle is still live, so the request registers (stays pending)
1509		// instead of failing with NotFound.
1510		let _fut = subscribe_pending!(consumer, "track1");
1511	}
1512
1513	/// A reserved name nobody accepts is the parking case a publisher has to be able to
1514	/// end. Ending the broadcast is where it does: `Consumer::track` answers
1515	/// `Unroutable` for a name asked about after this point, so whoever asked earlier gets
1516	/// the same answer instead of waiting on info that can never arrive.
1517	#[tokio::test]
1518	async fn close_resolves_a_reserved_name() {
1519		let producer = Info::new().produce();
1520		let consumer = producer.consume();
1521
1522		let _request = producer.reserve_track("track1").unwrap();
1523		let pending = subscribe_pending!(consumer, "track1");
1524
1525		producer.close();
1526		assert!(matches!(pending.await, Err(Error::Unroutable)));
1527	}
1528
1529	/// The deprecated abort says why the broadcast ended, and an unserved name resolves
1530	/// with that reason.
1531	#[tokio::test]
1532	#[allow(deprecated)]
1533	async fn abort_resolves_a_reserved_name_with_its_reason() {
1534		let producer = Info::new().produce();
1535		let consumer = producer.consume();
1536
1537		let request = producer.reserve_track("track1").unwrap();
1538		let pending = subscribe_pending!(consumer, "track1");
1539
1540		producer.abort(Error::Cancel).unwrap();
1541		assert!(matches!(pending.await, Err(Error::Cancel)));
1542
1543		let track = request.accept(None);
1544		let mut subscriber = track.subscribe(None);
1545		assert!(matches!(subscriber.recv_group().await, Err(Error::Cancel)));
1546	}
1547
1548	/// A request still queued for a handler is the same parking case reached from the
1549	/// consumer side, so it ends the same way.
1550	#[tokio::test]
1551	async fn close_resolves_a_queued_request() {
1552		let producer = Info::new().produce();
1553		let dynamic = producer.dynamic();
1554		let consumer = dynamic.consume();
1555
1556		let pending = subscribe_pending!(consumer, "track1");
1557
1558		producer.close();
1559		assert!(matches!(pending.await, Err(Error::Unroutable)));
1560		drop(dynamic);
1561	}
1562
1563	/// A queued request answers the same when the broadcast ends by its last handle
1564	/// dropping, even when that handle is the `Dynamic` that would have served it.
1565	#[tokio::test]
1566	async fn dropping_the_last_handle_resolves_a_queued_request() {
1567		let dynamic = Info::new().produce().dynamic();
1568		let consumer = dynamic.consume();
1569
1570		let pending = subscribe_pending!(consumer, "track1");
1571
1572		drop(dynamic);
1573		assert!(matches!(pending.await, Err(Error::Unroutable)));
1574	}
1575
1576	/// With a producer still alive, losing the last handler is not the broadcast ending:
1577	/// the queued request fails as `Dropped`.
1578	#[tokio::test]
1579	async fn dropping_the_last_handler_resolves_a_queued_request_dropped() {
1580		let producer = Info::new().produce();
1581		let dynamic = producer.dynamic();
1582		let consumer = dynamic.consume();
1583
1584		let pending = subscribe_pending!(consumer, "track1");
1585
1586		drop(dynamic);
1587		assert!(matches!(pending.await, Err(Error::Dropped)));
1588		producer.close();
1589	}
1590
1591	/// A request a handler already took is the handler's to answer: it may be in flight
1592	/// to a peer, and a retraction does not disturb subscriptions already in flight.
1593	/// Whatever the handler decides still reaches the consumer.
1594	#[tokio::test]
1595	async fn close_leaves_a_claimed_request_to_its_handler() {
1596		let producer = Info::new().produce();
1597		let mut dynamic = producer.dynamic();
1598		let consumer = dynamic.consume();
1599
1600		let accepted = subscribe_pending!(consumer, "track1");
1601		let request = dynamic.requested_track().await.unwrap();
1602		let dropped = subscribe_pending!(consumer, "track2");
1603		let abandoned = dynamic.requested_track().await.unwrap();
1604
1605		producer.close();
1606		assert!(
1607			accepted.poll_ok(&kio::Waiter::noop()).is_pending(),
1608			"close rejected a claimed request"
1609		);
1610		assert!(
1611			dropped.poll_ok(&kio::Waiter::noop()).is_pending(),
1612			"close rejected a claimed request"
1613		);
1614
1615		let _track = request.accept(None);
1616		assert!(accepted.await.is_ok(), "the handler's accept reaches the consumer");
1617		drop(abandoned);
1618		assert!(dropped.await.is_err(), "the handler dropping it rejects the consumer");
1619		drop(dynamic);
1620	}
1621
1622	/// A reverse fetch can install the track metadata before the live request is
1623	/// accepted, but it does not create a live publisher. Closing the broadcast
1624	/// must still reject that name so an arrival-order subscriber does not park on
1625	/// backfill that is deliberately absent from its queue.
1626	#[tokio::test]
1627	async fn close_resolves_an_unaccepted_track_with_fetched_info() {
1628		let producer = Info::new().produce();
1629		let consumer = producer.consume();
1630
1631		let request = producer.reserve_track("track1").unwrap();
1632		let dynamic = request.dynamic();
1633		let track = consumer.track("track1").unwrap();
1634		let pending_fetch = track.fetch_group(0, None);
1635		let fetch = dynamic.requested_group().await.unwrap();
1636		let group = fetch.accept(None).unwrap();
1637		group.finish().unwrap();
1638		pending_fetch.await.unwrap();
1639
1640		let mut subscriber = track.subscribe(None).await.unwrap();
1641		producer.close();
1642		assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
1643
1644		let stale = request.accept(None);
1645		assert!(stale.append_group().is_err());
1646	}
1647
1648	/// Ending the broadcast doesn't cascade into a track someone is publishing: it keeps
1649	/// its cache and its publisher decides when it ends. Only a new lookup is refused.
1650	#[tokio::test]
1651	async fn close_spares_a_served_track() {
1652		let producer = Info::new().produce();
1653		let consumer = producer.consume();
1654
1655		let track = producer.create_track("track1", None).unwrap();
1656		let mut subscriber = consumer.track("track1").unwrap().subscribe(None).await.unwrap();
1657
1658		producer.close();
1659		assert!(matches!(consumer.track("track1"), Err(Error::Unroutable)));
1660
1661		track.append_group().unwrap();
1662		subscriber.assert_group();
1663		track.finish().unwrap();
1664	}
1665
1666	/// The publisher may still be holding the `track::Request` for a name the broadcast
1667	/// just gave up on. Accepting it afterwards must not resurrect the track, or a
1668	/// subscriber that was told `Unroutable` could be contradicted by a later one.
1669	#[tokio::test]
1670	async fn close_leaves_a_stale_reservation_inert() {
1671		let producer = Info::new().produce();
1672		let consumer = producer.consume();
1673
1674		let request = producer.reserve_track("track1").unwrap();
1675		let pending = subscribe_pending!(consumer, "track1");
1676
1677		producer.close();
1678		assert!(matches!(pending.await, Err(Error::Unroutable)));
1679
1680		let track = request.accept(None);
1681		assert!(track.append_group().is_err());
1682		let mut subscriber = track.subscribe(None);
1683		assert!(matches!(subscriber.recv_group().await, Err(Error::Unroutable)));
1684		assert!(consumer.track("track1").is_err());
1685	}
1686
1687	/// Dropping a `track::Request` is not a verdict about the name, so it resolves as
1688	/// `Dropped` (a handler lost to a crashed publisher or a dead transport), never as
1689	/// `NotFound`. Only an explicit rejection may claim the track is absent.
1690	#[tokio::test]
1691	async fn dropping_a_reserved_request_resolves_dropped() {
1692		let producer = Info::new().produce();
1693		let consumer = producer.consume();
1694
1695		let request = producer.reserve_track("track1").unwrap();
1696		let pending = subscribe_pending!(consumer, "track1");
1697
1698		drop(request);
1699		assert!(matches!(pending.await, Err(Error::Dropped)));
1700		producer.close();
1701	}
1702
1703	/// `track::Request::reject` carries its reason the same way, which is what lets a
1704	/// subscriber tell "no such track" from "the publisher went away".
1705	#[tokio::test]
1706	async fn rejecting_a_reserved_request_carries_the_reason() {
1707		let producer = Info::new().produce();
1708		let consumer = producer.consume();
1709
1710		let request = producer.reserve_track("track1").unwrap();
1711		let pending = subscribe_pending!(consumer, "track1");
1712
1713		request.reject(Error::NotFound);
1714		assert!(matches!(pending.await, Err(Error::NotFound)));
1715		producer.close();
1716	}
1717
1718	/// The interleave every unused-driven teardown has to survive: a wire subscriber
1719	/// observes zero consumers, and a viewer looks the track up again before the
1720	/// teardown commits. The returning viewer keeps the track alive, and once it
1721	/// really does commit the cached handle is reclaimed rather than handed out
1722	/// cancelled.
1723	#[tokio::test]
1724	async fn an_idle_teardown_yields_to_a_returning_viewer() {
1725		let producer = Info::new().produce();
1726		let consumer = producer.consume();
1727		let track = producer.create_track("video", None).unwrap();
1728
1729		// The unused wake a teardown acts on.
1730		assert!(track.poll_unused(&kio::Waiter::noop()).is_ready());
1731
1732		// Demand returns in the gap before it commits.
1733		let viewer = consumer.track("video").unwrap();
1734		let track = track
1735			.abort_unused(Error::Cancel)
1736			.expect_err("viewer keeps the track alive");
1737
1738		// So the viewer holds a live track, not a cancelled one.
1739		assert!(!track.is_closed());
1740		let mut subscriber = viewer.subscribe(None).await.unwrap();
1741		subscriber.assert_no_group();
1742		track.append_group().unwrap();
1743		assert!(subscriber.recv_group().await.unwrap().is_some());
1744
1745		// Once the viewer really leaves, the same teardown commits, and the lookup
1746		// re-requests the track instead of resolving the closed one.
1747		drop(subscriber);
1748		drop(viewer);
1749		assert!(track.abort_unused(Error::Cancel).is_ok());
1750		assert!(matches!(consumer.track("video"), Err(Error::NotFound)));
1751
1752		producer.close();
1753	}
1754
1755	#[test]
1756	fn abort_unused_accepts_an_already_closed_track_with_consumers() {
1757		let producer = Info::new().produce();
1758		let consumer = producer.consume();
1759		let track = producer.create_track("video", None).unwrap();
1760		let _viewer = consumer.track("video").unwrap();
1761		assert!(track.is_used());
1762		track.clone().abort(Error::Cancel).unwrap();
1763		assert!(!track.is_used());
1764		assert!(track.abort_unused(Error::Cancel).is_ok());
1765		producer.close();
1766	}
1767}