Skip to main content

moq_net/model/
group.rs

1//! A group is a stream of frames, split into a [Producer] and [Consumer] handle.
2//!
3//! A [Producer] writes an ordered stream of frames.
4//! Frames can be written all at once ([Producer::write_frame]), or in chunks
5//! ([Producer::create_frame]).
6//!
7//! A [Consumer] reads an ordered stream of frames.
8//! The reader can be cloned, in which case each reader receives a copy of each frame. (fanout)
9//!
10//! Frames are numbered from 0 in write order. A group can be short at its front or its
11//! back but never in the middle: [Producer::start_at] starts it later, so a handle can
12//! carry the tail of a group whose leading frames came from somewhere else, and
13//! [Producer::finish] ends it wherever writing stopped. [Consumer::set_frames] bounds a reader to a sub-range the same way [`track::Subscriber`]
14//! bounds group sequences.
15//!
16//! The stream is closed with [Error] when all writers or readers are dropped.
17use crate::cache;
18use crate::frame::{self, Frame, FrameBuf};
19use crate::{Cap, Timescale, stats, track};
20use std::collections::VecDeque;
21use std::mem::MaybeUninit;
22use std::ops::{Bound, RangeBounds};
23use std::sync::Arc;
24use std::sync::atomic::{AtomicBool, Ordering};
25use std::task::{Poll, ready};
26
27use crate::{Error, IntoBytes, Result, Timestamp};
28
29/// Maximum total size of frames in a group.
30///
31/// A write that would exceed this aborts the group with [`Error::GroupTooLarge`].
32/// Doubles as the per-frame size cap: a larger declared size is [`Error::FrameTooLarge`]
33/// before allocating, so one maximum-size frame can fill a group.
34pub const MAX_CACHE_BYTES: u64 = 32 * 1024 * 1024; // 32 MB
35
36/// Maximum number of frames in a group.
37///
38/// 8192 is the largest legal group; the 8193rd write returns [`Error::GroupTooLarge`]
39/// and aborts the group.
40pub const MAX_GROUP_FRAMES: usize = 8192;
41
42/// Slots `VecDeque` rounds a group's first frame up to.
43///
44/// A `RawVec` detail rather than a knob, so it is asserted rather than trusted: std
45/// handing out more would silently undercharge every cached group.
46const FRAME_SLOTS: usize = 4;
47
48/// Heap one cached group costs beyond its frame payloads, excluding the track-side
49/// bookkeeping in [`track::CACHE_OVERHEAD`].
50///
51/// A group is one kio channel (allocated whether or not anything ever parks on it), the
52/// `Arc<Alive>` its producer clones share, and the frame slots the first write rounds up
53/// to. Half of [`cache::ENTRY_OVERHEAD`]; see it for why this is derived rather than
54/// measured.
55pub(crate) const CACHE_OVERHEAD: u64 = (kio::Producer::<GroupState>::HEAP
56	// `Alive` behind an `Arc`'s two reference counts, which it is pointer-aligned to sit
57	// straight after.
58	+ 2 * size_of::<usize>()
59	+ size_of::<Alive>()
60	+ FRAME_SLOTS * size_of::<Frame>()) as u64;
61
62/// A group contains a sequence number because they can arrive out of order.
63///
64/// You can use [track::Producer::append_group] if you just want to +1 the sequence number.
65#[derive(Clone, Copy, Debug, Hash, Eq, PartialEq, Ord, PartialOrd)]
66pub struct Info {
67	/// Per-track sequence number used to detect ordering and gaps. Higher numbers
68	/// supersede lower ones; consumers may skip late arrivals.
69	pub sequence: u64,
70}
71
72impl Info {
73	/// Create an untimed producer for this group.
74	///
75	/// Test-only: real groups are created via [`track::Producer`], which
76	/// supplies the parent track's [`track::Info`]. This helper exists for in-crate
77	/// tests that don't exercise timestamps.
78	#[cfg(test)]
79	pub(crate) fn produce(self) -> Producer {
80		Producer::new(self, track::Info::default(), Default::default())
81	}
82}
83
84impl From<usize> for Info {
85	fn from(sequence: usize) -> Self {
86		Self {
87			sequence: sequence as u64,
88		}
89	}
90}
91
92impl From<u64> for Info {
93	fn from(sequence: u64) -> Self {
94		Self { sequence }
95	}
96}
97
98impl From<u32> for Info {
99	fn from(sequence: u32) -> Self {
100		Self {
101			sequence: sequence as u64,
102		}
103	}
104}
105
106impl From<u16> for Info {
107	fn from(sequence: u16) -> Self {
108		Self {
109			sequence: sequence as u64,
110		}
111	}
112}
113
114/// The in-flight (tail) frame being written. At most one exists at a time, since a
115/// group is a single ordered stream.
116pub(crate) struct Partial {
117	timestamp: Timestamp,
118	buf: FrameBuf,
119}
120
121/// Shared group state. `pub(crate)` so [`frame`] handles can observe the abort flag
122/// while streaming a partial frame.
123#[derive(Default)]
124pub(crate) struct GroupState {
125	// Completed frames, each a contiguous payload. `offset` is the first frame this
126	// handle holds, raised by [`Producer::start_at`].
127	pub(crate) frames: VecDeque<Frame>,
128
129	// The single in-flight frame, if one is open.
130	pub(crate) partial: Option<Partial>,
131
132	// Index of the first frame this handle holds: any the group deliberately started
133	// past (see [`Producer::start_at`]). Reading below it is [`Error::Lagged`]; the
134	// frames are not here.
135	pub(crate) offset: usize,
136
137	// The index the next frame written will get. Tracked separately from `frames` so it
138	// survives the cache being released: a route taking the track over needs to know where
139	// production stopped, and an abort is exactly when it asks.
140	next_index: usize,
141
142	// One past the last frame that was fully written. Trails `next_index` while a chunked
143	// frame is in flight, which is the frame a replacement route has to redeliver: only
144	// its opener saw the payload, and only partly.
145	committed: usize,
146
147	// The total size (in bytes) of all cached frames plus any in-flight frame.
148	pub(crate) cache: u64,
149
150	// Mirrors `cache` into the track's shared cache pool, so the group's bytes count
151	// against the byte budget tracks evict toward.
152	charge: cache::Charge,
153
154	// The first frame's timestamp, recorded once and never revised: the group's
155	// presentation start. Kept here rather than read off `frames` so an abort
156	// doesn't erase where the group sat in time. `None` until the first frame is
157	// written, which is the only honest answer: an empty group has not presented
158	// anything yet.
159	timestamp: Option<Timestamp>,
160
161	// The newest frame's timestamp: the group's presentation end so far. A reader
162	// that has taken every frame sits here, which is what a drift budget measures it
163	// against. Kept alongside `timestamp` for the same reasons.
164	latest: Option<Timestamp>,
165
166	// Once finalized, the total number of frames the group will ever contain. Recorded
167	// at finish so the count outlives an abort that clears the cache.
168	pub(crate) fin: Option<usize>,
169
170	// The error that caused the group to be aborted, if any. Mirrored into
171	// `Alive::aborted`, so [`Producer::abort`] stays the only writer: anything else
172	// setting this would leave track scans reading a group as live.
173	pub(crate) abort: Option<Error>,
174}
175
176impl GroupState {
177	/// Content still available to a reader of this group.
178	fn content(&self) -> stats::Content {
179		stats::Content {
180			bytes: self.cache,
181			frames: self.next_index.saturating_sub(self.offset) as u64,
182			groups: 1,
183			datagrams: 0,
184		}
185	}
186
187	/// Content in the half-open frame range that is still cached here.
188	pub(crate) fn content_range(&self, start: usize, end: usize) -> stats::Content {
189		let start = start.max(self.offset);
190		let end = end.min(self.next_index);
191		if start >= end {
192			return stats::Content::default();
193		}
194
195		let local_start = start.saturating_sub(self.offset).min(self.frames.len());
196		let local_end = end.saturating_sub(self.offset).min(self.frames.len());
197		let mut bytes = self
198			.frames
199			.range(local_start..local_end)
200			.map(|frame| frame.payload.len() as u64)
201			.sum();
202		if start <= self.committed
203			&& self.committed < end
204			&& let Some(partial) = &self.partial
205		{
206			bytes += partial.buf.capacity() as u64;
207		}
208
209		stats::Content {
210			bytes,
211			frames: (end - start) as u64,
212			groups: 0,
213			datagrams: 0,
214		}
215	}
216
217	/// Resolve the source for the frame at `index`: a completed frame (whole) or the
218	/// in-flight tail (streamed). Used by [`Consumer::poll_next_frame`].
219	fn poll_frame_source(&self, index: usize) -> Poll<Result<Option<(frame::Info, frame::Source)>>> {
220		if index < self.offset {
221			return Poll::Ready(Err(Error::Lagged));
222		}
223		let local = index - self.offset;
224		if let Some(f) = self.frames.get(local) {
225			// A frame read is a cache access: stamp it so expiry and the eviction
226			// walk spare a group a consumer is actively draining.
227			self.charge.refresh();
228			let info = frame::Info {
229				size: f.payload.len() as u64,
230				timestamp: f.timestamp,
231			};
232			return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
233		}
234		if local == self.frames.len()
235			&& let Some(p) = &self.partial
236		{
237			self.charge.refresh();
238			let info = frame::Info {
239				size: p.buf.capacity() as u64,
240				timestamp: p.timestamp,
241			};
242			return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
243		}
244		ready!(self.poll_terminal(index))?;
245		Poll::Ready(Ok(None))
246	}
247
248	/// Resolve the group's terminal state for a reader positioned at `index`.
249	///
250	/// A finished group is still aborted once its frames are released to free memory
251	/// (aged out of the track's max age window, or evicted by the cache pool). A reader
252	/// that already consumed every frame is missing nothing, so it gets the clean end of
253	/// group; one that fell short sees the abort rather than a silently truncated stream.
254	fn poll_terminal(&self, index: usize) -> Poll<Result<()>> {
255		match (self.fin, &self.abort) {
256			(Some(total), Some(err)) if index < total => Poll::Ready(Err(err.clone())),
257			(Some(_), _) => Poll::Ready(Ok(())),
258			(None, Some(err)) => Poll::Ready(Err(err.clone())),
259			(None, None) => Poll::Pending,
260		}
261	}
262
263	/// Resolve whether a reader at `index` can still make progress, answering the same
264	/// question as a read without consuming anything.
265	fn poll_end(&self, index: usize) -> Poll<Result<()>> {
266		if index < self.offset {
267			return Poll::Ready(Err(Error::Lagged));
268		}
269		self.poll_terminal(index)
270	}
271
272	/// Record where the group starts and currently ends in presentation time.
273	/// `timestamp` keeps the first frame only; `latest` follows every frame.
274	fn stamp(&mut self, timestamp: Timestamp) {
275		self.timestamp.get_or_insert(timestamp);
276		self.latest = Some(timestamp);
277	}
278
279	/// Whether adding `extra_frames` totaling `extra_bytes` would exceed the group budget.
280	fn would_overflow(&self, extra_frames: usize, extra_bytes: u64) -> bool {
281		self.next_index.saturating_sub(self.offset).saturating_add(extra_frames) > MAX_GROUP_FRAMES
282			|| self.cache.saturating_add(extra_bytes) > MAX_CACHE_BYTES
283	}
284
285	/// Drop the cached frames (and any in-flight tail) and release their pool charge.
286	fn release(&mut self) {
287		self.frames.clear();
288		self.partial = None;
289		self.cache = 0;
290		self.charge.clear();
291	}
292}
293
294fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
295	state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
296}
297
298/// Writes frames to a group in order.
299///
300/// Each group is delivered independently over a QUIC stream.
301/// Use [Self::write_frame] for simple single-buffer frames,
302/// or [Self::create_frame] for multi-chunk streaming writes.
303pub struct Producer {
304	// Mutable stream state.
305	state: kio::Producer<GroupState>,
306
307	// The group header containing the sequence number. A small `Copy` value,
308	// inherited by each frame (see [`Self::create_frame`]).
309	info: Info,
310
311	// The parent track's properties, inherited rather than passed piecemeal. Its
312	// `timescale` is used by [`Self::create_frame`] to normalize every frame's
313	// timestamp into the track scale before it enters the stream. Threaded down by
314	// value from [`track::Producer::create_group`] / `append_group`.
315	track: track::Info,
316
317	// The parent track's account against the shared cache pool. Held here as well as
318	// in the group's `cache::Charge` so a frame write can settle the track's eviction
319	// debt with the group lock released.
320	cache: Arc<cache::Track>,
321
322	// Ingress payload meter, set by a tagged [`track::Producer`] via
323	// [`Self::with_meter`]. Empty (no-op) for an untagged group.
324	stats: stats::Meter,
325
326	// Shared by every clone: its `Drop` is the abrupt-teardown, running exactly once
327	// when the last of them goes.
328	alive: Arc<Alive>,
329}
330
331/// Ends the group when the last [`Producer`] clone drops, including the clone the
332/// parent track holds in its cache.
333///
334/// A refcount rather than a "am I the last one?" check inside `Drop`: that answer is
335/// a snapshot, and acting on it is exactly what can invalidate it. Holding a producer
336/// of its own also keeps the state writable until the teardown has run, whatever order
337/// the last owner's fields drop in.
338struct Alive {
339	info: Info,
340	state: kio::Producer<GroupState>,
341	// Monotone mirror of `GroupState::abort` for track scans that already hold the
342	// track lock. A stale false only hands out a group that is concurrently aborting;
343	// true is stored after the abort exists, so it can never hide a live group.
344	//
345	// Only ever read on its own. A decision that pairs the abort with something else
346	// out of `GroupState` has to read both under one guard, or the two halves can
347	// straddle the abort: see `Producer::live_first_frame`.
348	aborted: AtomicBool,
349	// The cache stamp `GroupState::charge` maintains, held here as well so the
350	// eviction and expiry walks can weigh a candidate without taking the group lock
351	// they already hold the track lock over.
352	access: Arc<cache::Access>,
353}
354
355impl Drop for Alive {
356	fn drop(&mut self) {
357		// See track::Alive: the last producer dropping without a clean finish releases
358		// the cached frames so a stale consumer can't pin their buffers forever. A
359		// finished group keeps its cache so consumers can drain.
360		//
361		// Check Ok and Err: Ok is unreachable after a deliberate close.
362		match self.state.write() {
363			Ok(mut state) => {
364				if state.fin.is_some() || state.abort.is_some() {
365					return;
366				}
367				tracing::warn!(
368					sequence = self.info.sequence,
369					"group::Producer dropped without finish() or abort()"
370				);
371				state.release();
372			}
373			Err(state) => {
374				if state.fin.is_some() || state.abort.is_some() {
375					return;
376				}
377				tracing::warn!(
378					sequence = self.info.sequence,
379					"group::Producer dropped without finish() or abort()"
380				);
381			}
382		}
383	}
384}
385
386impl std::ops::Deref for Producer {
387	type Target = Info;
388
389	fn deref(&self) -> &Self::Target {
390		&self.info
391	}
392}
393
394impl Producer {
395	/// Create a group producer bound to its parent track's [`track::Info`] and cache
396	/// account.
397	///
398	/// Crate-private: groups are only constructed via [`track::Producer`], which
399	/// threads both down so properties like the timescale are inherited rather than
400	/// passed in. Every frame added to this group is normalized to the track's
401	/// timescale by [`Self::create_frame`].
402	///
403	/// Charges the group into `cache`, so its cached bytes count against the budget the
404	/// track evicts toward under memory pressure.
405	pub(crate) fn new(info: Info, track: track::Info, cache: Arc<cache::Track>) -> Self {
406		let state = kio::Producer::<GroupState>::default();
407		let charge = cache.charge();
408		let access = charge.access();
409		state.write().ok().expect("a new group is open").charge = charge;
410		let alive = Arc::new(Alive {
411			info,
412			state: state.clone(),
413			aborted: AtomicBool::new(false),
414			access,
415		});
416		Self {
417			info,
418			state,
419			track,
420			cache,
421			stats: stats::Meter::default(),
422			alive,
423		}
424	}
425
426	/// Attach an ingress payload meter, counting this as one delivered group.
427	/// Called by a tagged [`track::Producer`] when it creates the group.
428	pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
429		meter.group();
430		self.stats = meter;
431		self
432	}
433
434	/// The group header.
435	pub(crate) fn info(&self) -> Info {
436		self.info
437	}
438
439	/// The parent track's timescale.
440	pub fn timescale(&self) -> Timescale {
441		self.track.timescale
442	}
443
444	/// Start the group at frame `index` rather than 0, so the first frame written lands
445	/// there.
446	///
447	/// A group can be short at its front or its back, never in the middle: this trims
448	/// the front, and simply stopping (then [`finish`](Self::finish)ing) trims the back.
449	/// The frames below `index` are not a gap this handle will ever fill, so a reader
450	/// positioned below it gets [`Error::Lagged`]. They belong to whoever produced the
451	/// head of the group, typically another route serving the same track (see
452	/// [`crate::track::Subscriber`]).
453	///
454	/// The counterpart of [`Consumer::set_frames`], which positions a *reader* the same
455	/// way. Where the group begins is part of its shape, so this must come before the
456	/// first frame; afterwards it returns [`Error::Closed`].
457	pub fn start_at(&mut self, index: u64) -> Result<()> {
458		let index = usize::try_from(index).map_err(|_| Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
459		if index == usize::MAX {
460			return Err(Error::BoundsExceeded(crate::coding::BoundsExceeded));
461		}
462
463		let mut state = modify(&self.state)?;
464		// Every write advances `next_index` past `offset`, so this is "nothing written
465		// yet".
466		if state.fin.is_some() || state.next_index != state.offset {
467			return Err(Error::Closed);
468		}
469		state.offset = index;
470		state.next_index = index;
471		state.committed = index;
472		Ok(())
473	}
474
475	/// A helper method to write a frame from a single byte buffer.
476	///
477	/// If you want to write multiple chunks, use [Self::create_frame] to get a frame producer.
478	/// But an upfront size is required.
479	///
480	/// `timestamp` is converted into the parent track's timescale. For data without
481	/// a presentation time, pass [`Timestamp::now`] explicitly.
482	pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
483		let timestamp = timestamp
484			.convert(self.track.timescale)
485			.map_err(|_| Error::TimestampMismatch)?;
486		let payload = data.into_bytes();
487		if payload.len() as u64 > MAX_CACHE_BYTES {
488			return Err(Error::FrameTooLarge);
489		}
490
491		let mut state = modify(&self.state)?;
492		if state.fin.is_some() {
493			return Err(Error::Closed);
494		}
495		if state.partial.is_some() {
496			return Err(Error::FrameOpen);
497		}
498		let next_index = state
499			.next_index
500			.checked_add(1)
501			.ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
502		debug_assert!(state.partial.is_none(), "a frame is already open");
503		let size = payload.len() as u64;
504		if state.would_overflow(1, size) {
505			return Err(self.abort_too_large(state));
506		}
507		state.cache += size;
508		let now = state.charge.add(size);
509		state.frames.push_back(Frame { timestamp, payload });
510		state.next_index = next_index;
511		state.committed = state.next_index;
512		state.stamp(timestamp);
513		drop(state);
514
515		// With the group lock released (lock order is track then group), settle
516		// eviction debt if enough has been written since the track last paid.
517		self.cache.settle(now);
518
519		// Ingress payload: one whole frame written.
520		self.stats.frames(1);
521		self.stats.bytes(size);
522		Ok(())
523	}
524
525	/// Write a whole batch of frames at once, draining `frames`.
526	///
527	/// One lock covers the batch, so an ingest with several frames in hand pays the
528	/// group mutex and the track's eviction settle once rather than per frame. Build
529	/// the batch with [`frame::Buffer::push`].
530	///
531	/// The batch is validated before anything is written, so a rejected frame leaves
532	/// both the group and the buffer exactly as they were, ready to retry or redirect.
533	/// Returns [`Error::FrameOpen`] if another handle is streaming a frame into this
534	/// group, since appending around it would reorder the group.
535	pub fn write_frames<const N: usize>(&mut self, frames: &mut frame::Buffer<N>) -> Result<()> {
536		for frame in frames.filled() {
537			frame
538				.timestamp
539				.convert(self.track.timescale)
540				.map_err(|_| Error::TimestampMismatch)?;
541			if frame.payload.len() as u64 > MAX_CACHE_BYTES {
542				return Err(Error::FrameTooLarge);
543			}
544		}
545
546		let count = frames.len();
547		let bytes: u64 = frames.filled().iter().map(|frame| frame.payload.len() as u64).sum();
548		let mut state = modify(&self.state)?;
549		if state.fin.is_some() {
550			return Err(Error::Closed);
551		}
552		if state.partial.is_some() {
553			return Err(Error::FrameOpen);
554		}
555		let next_index = state
556			.next_index
557			.checked_add(count)
558			.ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
559		if state.would_overflow(count, bytes) {
560			return Err(self.abort_too_large(state));
561		}
562
563		// The last frame's tick, reused below so settling does not re-read the clock.
564		let mut now = None;
565		for mut frame in frames.drain() {
566			frame.timestamp = frame
567				.timestamp
568				.convert(self.track.timescale)
569				.expect("timestamp scale checked above");
570			let size = frame.payload.len() as u64;
571			state.cache += size;
572			now = state.charge.add(size);
573			state.stamp(frame.timestamp);
574			state.frames.push_back(frame);
575		}
576		state.next_index = next_index;
577		state.committed = next_index;
578		drop(state);
579
580		self.cache.settle(now);
581		self.stats.frames(count as u64);
582		self.stats.bytes(bytes);
583		Ok(())
584	}
585
586	/// Create a frame with an upfront size and presentation timestamp, streamed in
587	/// chunks. Borrows the group exclusively until the returned [`frame::Producer`]
588	/// is finished or dropped, so only one frame is open at a time.
589	///
590	/// The `timestamp` is converted into the parent track's timescale, so the scale you
591	/// build it with doesn't have to match the track. Returns [`Error::FrameTooLarge`]
592	/// if the declared size exceeds the group's byte budget (refused before allocating)
593	/// or [`Error::TimestampMismatch`] if the timestamp can't be converted (overflow).
594	pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
595		let timestamp = frame
596			.timestamp
597			.convert(self.track.timescale)
598			.map_err(|_| Error::TimestampMismatch)?;
599		if frame.size > MAX_CACHE_BYTES {
600			return Err(Error::FrameTooLarge);
601		}
602		let buf = FrameBuf::new(frame.size as usize);
603
604		let mut state = modify(&self.state)?;
605		if state.fin.is_some() {
606			return Err(Error::Closed);
607		}
608		if state.partial.is_some() {
609			return Err(Error::FrameOpen);
610		}
611		let next_index = state
612			.next_index
613			.checked_add(1)
614			.ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
615		if state.would_overflow(1, frame.size) {
616			return Err(self.abort_too_large(state));
617		}
618		state.cache += frame.size;
619		let now = state.charge.add(frame.size);
620		state.partial = Some(Partial {
621			timestamp,
622			buf: buf.clone(),
623		});
624		state.next_index = next_index;
625		// Opening the frame is enough: the header carries the timestamp, so the group's
626		// place in time is known before a single payload byte streams in.
627		state.stamp(timestamp);
628		drop(state);
629
630		// With the group lock released (lock order is track then group), settle
631		// eviction debt if enough has been written since the track last paid.
632		self.cache.settle(now);
633
634		// Ingress payload: one frame opened; its bytes are counted per chunk as the
635		// frame::Producer writes them.
636		self.stats.frames(1);
637		let meter = self.stats.clone();
638
639		let info = frame::Info {
640			size: frame.size,
641			timestamp,
642		};
643		Ok(frame::Producer::new(self, buf, info).with_meter(meter))
644	}
645
646	/// The owned counterpart of [`Self::create_frame`], for the wire drivers that
647	/// stream a frame across polls and cannot hold the group borrowed inside their
648	/// state. The one-live-frame rule the borrow normally enforces becomes the
649	/// caller's promise; see [`frame::ProducerOwned`].
650	pub(crate) fn create_frame_owned(&mut self, frame: frame::Info) -> Result<frame::ProducerOwned> {
651		let timestamp = frame
652			.timestamp
653			.convert(self.track.timescale)
654			.map_err(|_| Error::TimestampMismatch)?;
655		if frame.size > MAX_CACHE_BYTES {
656			return Err(Error::FrameTooLarge);
657		}
658		let buf = FrameBuf::new(frame.size as usize);
659
660		let mut state = modify(&self.state)?;
661		if state.fin.is_some() {
662			return Err(Error::Closed);
663		}
664		if state.partial.is_some() {
665			return Err(Error::FrameOpen);
666		}
667		let next_index = state
668			.next_index
669			.checked_add(1)
670			.ok_or(Error::BoundsExceeded(crate::coding::BoundsExceeded))?;
671		if state.would_overflow(1, frame.size) {
672			return Err(self.abort_too_large(state));
673		}
674		state.cache += frame.size;
675		let now = state.charge.add(frame.size);
676		state.partial = Some(Partial {
677			timestamp,
678			buf: buf.clone(),
679		});
680		state.next_index = next_index;
681		// Opening the frame is enough: the header carries the timestamp, so the group's
682		// place in time is known before a single payload byte streams in.
683		state.stamp(timestamp);
684		drop(state);
685
686		// With the group lock released (lock order is track then group), settle
687		// eviction debt if enough has been written since the track last paid.
688		self.cache.settle(now);
689
690		// Ingress payload: one frame opened; its bytes are counted per chunk as the
691		// producer writes them.
692		self.stats.frames(1);
693		let meter = self.stats.clone();
694
695		let info = frame::Info {
696			size: frame.size,
697			timestamp,
698		};
699		Ok(frame::ProducerOwned::new(self.clone(), buf, info).with_meter(meter))
700	}
701
702	/// Wake consumers parked on the group channel (called after a partial write).
703	pub(crate) fn frame_notify(&self) {
704		// The chunk that was just written is a write access: restart the retention
705		// clock so a straggler group streaming a large frame isn't expired
706		// mid-write (its bytes were already charged when the frame was created).
707		// `record_write` takes `&mut`, which marks the guard modified: kio only
708		// notifies on a mutably-accessed guard's release, and that notify is what
709		// delivers the chunk to parked readers.
710		let now = self
711			.state
712			.write()
713			.ok()
714			.and_then(|mut state| state.charge.record_write());
715		// The payload was charged when the frame opened, but a long streamed frame
716		// still counts as track activity for the independent expiry time gate.
717		self.cache.settle(now);
718	}
719
720	/// Commit the in-flight frame as a completed frame (called by [`frame::Producer::finish`]).
721	pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
722		let mut state = modify(&self.state)?;
723		// Bytes were already counted against the cache (and the pool charge) when the
724		// frame was created; committing just moves the tail into the completed set.
725		state.partial = None;
726		state.frames.push_back(frame);
727		state.committed = state.next_index;
728		// Completing the frame is a write access like any chunk, and the only one the
729		// payload is guaranteed to get: the wire ingest defers its chunk notifications
730		// to the poll boundary, so a tail that arrives and completes in one turn never
731		// reaches [`Self::frame_notify`]. Without this, a group whose payload streamed
732		// in across an idle gap would expire the instant it finished.
733		let now = state.charge.record_write();
734		drop(state);
735
736		// With the group lock released (lock order is track then group), settle
737		// eviction debt and age idle content out, reusing the tick above.
738		self.cache.settle(now);
739		Ok(())
740	}
741
742	/// Fail the group because an in-flight frame couldn't complete (called by
743	/// [`frame::Producer::abort`] / its drop).
744	pub(crate) fn frame_abort(&mut self, err: Error) {
745		let _ = self.clone().abort(err);
746	}
747
748	/// One past the index of the last frame written (completed or in-flight), which is
749	/// also the index the next frame will get.
750	///
751	/// Counts any frames the group [started past](Self::start_at), so it's the group's
752	/// logical length rather than the number of frames this handle holds.
753	pub fn frame_count(&self) -> usize {
754		self.state.read().next_index
755	}
756
757	/// Mark the group as complete; no more frames will be written.
758	///
759	/// Borrows rather than consumes, so a later failure can still be reported through
760	/// [`abort`](Self::abort). The handle also keeps the cached frames readable.
761	pub fn finish(&self) -> Result<()> {
762		let mut state = modify(&self.state)?;
763		if state.partial.is_some() {
764			return Err(Error::FrameOpen);
765		}
766		state.fin = Some(state.next_index);
767		Ok(())
768	}
769
770	/// Abort the group with the given error.
771	///
772	/// Consumes the handle. Drops the cached frames so a stale [`Consumer`] can't pin
773	/// their buffers in memory forever; consumers that haven't drained yet surface the
774	/// abort error instead of the leftover cache.
775	pub fn abort(self, err: Error) -> Result<()> {
776		let mut guard = modify(&self.state)?;
777		guard.abort = Some(err);
778		self.alive.aborted.store(true, Ordering::Release);
779		guard.release();
780		guard.close();
781		Ok(())
782	}
783
784	/// Abort a write that would grow the group past its budget, holding the lock already
785	/// taken for that write so nothing else lands in between.
786	fn abort_too_large(&self, mut state: kio::Mut<'_, GroupState>) -> Error {
787		let err = Error::GroupTooLarge;
788		state.abort = Some(err.clone());
789		self.alive.aborted.store(true, Ordering::Release);
790		state.release();
791		state.close();
792		err
793	}
794
795	/// Whether the group has been aborted (including pool eviction). The track's
796	/// read paths treat an aborted cached group as absent.
797	///
798	/// Reads the mirror rather than the group's state, so a track scan holding the
799	/// track lock never takes the group's. Monotone, and only ever conservative: a
800	/// concurrent abort can still read as live for the length of [`Self::abort`],
801	/// which hands out a group whose consumer then surfaces the abort.
802	pub(crate) fn is_aborted(&self) -> bool {
803		self.alive.aborted.load(Ordering::Acquire)
804	}
805
806	/// Whether the group was finished: it holds every frame it will ever have.
807	pub(crate) fn is_finished(&self) -> bool {
808		self.state.read().fin.is_some()
809	}
810
811	/// The index of the first frame this group still holds, or `None` once it has been
812	/// aborted. Non-zero when the group started later (see [`Self::start_at`]); a reader
813	/// positioned below it is [`Error::Lagged`].
814	///
815	/// One guard for both halves, deliberately. The track asks this to decide whether a
816	/// cached slot can still answer a request, and reading the abort and the offset
817	/// separately lets the abort land between them: the slot reads live, then hands
818	/// back an offset it only has because it is dead. The mirror
819	/// ([`Self::is_aborted`]) is for scans that ask about the abort alone.
820	pub(crate) fn live_first_frame(&self) -> Option<usize> {
821		let state = self.state.read();
822		state.abort.is_none().then_some(state.offset)
823	}
824
825	/// One past the last frame committed to an unfinished group, when that is past its
826	/// first: where a replacement route resumes. `None` once the group is finished, or
827	/// while it holds nothing a replacement could splice onto.
828	///
829	/// The *committed* count, not the written one: a route dying midway through a
830	/// chunked frame leaves that frame unusable, so the replacement has to send it
831	/// again rather than start after it. Answered under one guard so the count can't be
832	/// weighed against an offset from a different moment.
833	///
834	/// An aborted group still answers: readers that already consumed its head want the
835	/// tail, and the count outlives the released cache.
836	pub(crate) fn resume_frame(&self) -> Option<usize> {
837		let state = self.state.read();
838		if state.fin.is_some() {
839			return None;
840		}
841		(state.committed > state.offset).then_some(state.committed)
842	}
843
844	/// Where the group starts in presentation time: its first frame's timestamp,
845	/// or `None` while no frame has been opened.
846	///
847	/// Stamped once, when the group's first frame arrives, so it measures the group's
848	/// place in the media timeline rather than when it happened to be delivered. That
849	/// is what lets the track tell a burst of old content apart from live content (see
850	/// [`track::Subscriber`]). On protocols whose wire can't carry a timestamp the
851	/// receiver stamps frames with [`Timestamp::now`], which makes this the local
852	/// receive time instead: an estimate that a burst compresses.
853	pub(crate) fn timestamp(&self) -> Option<Timestamp> {
854		self.state.read().timestamp
855	}
856
857	/// Where the group ends in presentation time: its newest frame's timestamp, or
858	/// `None` while no frame has been opened.
859	///
860	/// This is what a drift budget measures an untouched group against. A group is not
861	/// late because it *started* long ago: a two-second group whose tail is level with
862	/// the live edge still has content nobody has read. Only once its newest frame has
863	/// fallen behind is there nothing left worth delivering. Grows as the group does, so
864	/// a group still receiving frames stays fresh and a stalled one ages in place.
865	pub(crate) fn latest(&self) -> Option<Timestamp> {
866		self.state.read().latest
867	}
868
869	/// The group's full cached footprint (payload plus fixed overhead), used by the
870	/// track to size this group as an eviction victim.
871	pub(crate) fn cache_size(&self) -> u64 {
872		self.state.read().charge.size()
873	}
874
875	/// Tick of the group's last cache access, driving eviction protection and age
876	/// expiry (see [`cache::Pool::average`]).
877	pub(crate) fn cache_accessed(&self) -> u64 {
878		self.alive.access.get()
879	}
880
881	/// Coarse clock tick of the group's last cache access, used by age expiry.
882	pub(crate) fn cache_accessed_tick(&self, now: Option<u64>) -> Option<u64> {
883		self.alive.access.tick(now)
884	}
885
886	/// Enter the group into the evictable population: demoted from the live edge,
887	/// or inserted behind it. Idempotent; a no-op once the group is closed.
888	pub(crate) fn cache_demote(&self) {
889		if let Ok(mut state) = self.state.write() {
890			state.charge.demote();
891		}
892	}
893
894	/// Record a cache access (delivery to a subscriber, a FETCH hit, or a fetched
895	/// backfill's birth), protecting the group from eviction and restarting its
896	/// expiry clock. Stamps through a read guard, whose release never notifies, so
897	/// delivery can't wake every consumer parked on the group. Harmless on a
898	/// closed group: its charge is already cleared.
899	pub(crate) fn cache_refresh(&self) {
900		self.state.read().charge.refresh();
901	}
902
903	/// Create a new consumer for the group.
904	pub fn consume(&self) -> Consumer {
905		Consumer {
906			info: self.info,
907			track: self.track.clone(),
908			inner: ConsumerKind::Plain(Plain {
909				state: self.state.consume(),
910				index: 0,
911				end: None,
912				prefetch: Prefetch::default(),
913				cache: self.cache.clone(),
914				access: self.alive.access.clone(),
915				refreshed: self.cache.pool().now(),
916			}),
917			// Untagged: a tagged track attaches the egress meter via `with_meter`
918			// when it hands the consumer to a subscriber/fetch.
919			stats: stats::Meter::default(),
920			stale_stats: stats::Meter::default(),
921			expiry: None,
922			expired: false,
923			ended: false,
924			stale_counted: Arc::default(),
925		}
926	}
927
928	/// Register for the first-frame timestamp while the group is still unstamped.
929	pub(crate) fn poll_timestamp(&self, waiter: &kio::Waiter) -> Poll<()> {
930		match self.state.poll(waiter, |state| {
931			if state.timestamp.is_some() || state.fin.is_some() || state.abort.is_some() {
932				Poll::Ready(())
933			} else {
934				Poll::Pending
935			}
936		}) {
937			Poll::Ready(_) => Poll::Ready(()),
938			Poll::Pending => Poll::Pending,
939		}
940	}
941
942	/// Block until the group is closed or aborted.
943	pub async fn closed(&self) -> Error {
944		kio::wait(|waiter| self.poll_closed(waiter)).await
945	}
946
947	/// Poll until the group is closed or aborted; ready with the cause.
948	pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
949		self.state.poll_closed(waiter).map(|()| self.abort_reason())
950	}
951
952	/// Block until there is at least one active consumer.
953	pub async fn used(&self) -> Result<()> {
954		self.state.used().await.map_err(|_| self.abort_reason())
955	}
956
957	/// Block until there are no active consumers.
958	pub async fn unused(&self) -> Result<()> {
959		self.state.unused().await.map_err(|_| self.abort_reason())
960	}
961
962	/// The recorded abort reason, or [`Error::Dropped`] if the group closed without one.
963	fn abort_reason(&self) -> Error {
964		self.state.read().abort.clone().unwrap_or(Error::Dropped)
965	}
966}
967
968impl Clone for Producer {
969	fn clone(&self) -> Self {
970		Self {
971			info: self.info,
972			state: self.state.clone(),
973			track: self.track.clone(),
974			cache: self.cache.clone(),
975			stats: self.stats.clone(),
976			alive: self.alive.clone(),
977		}
978	}
979}
980
981/// A small inline batch of completed frames, drained from the shared group state
982/// under one lock and then handed out without re-locking.
983///
984/// Each [`Consumer::read_frame`] otherwise takes the group mutex and allocates a
985/// waker just to clone one `Bytes`; draining a batch amortizes both across `CAP`
986/// frames. Storage is inline and uninitialized (no heap), so a consumer that never
987/// reads whole frames, or drains through a higher-level buffer, pays nothing.
988struct Prefetch {
989	// Initialized, not-yet-taken frames are `frames[pos..len]`; the rest are uninitialized.
990	frames: [MaybeUninit<Frame>; Self::CAP],
991	pos: usize,
992	len: usize,
993}
994
995impl Prefetch {
996	const CAP: usize = 8;
997
998	/// Take the next buffered frame, or `None` if the batch is drained.
999	fn pop(&mut self) -> Option<Frame> {
1000		if self.pos == self.len {
1001			return None;
1002		}
1003		// SAFETY: `pos < len`, so this slot was written by `fill` and not yet taken.
1004		let frame = unsafe { self.frames[self.pos].assume_init_read() };
1005		self.pos += 1;
1006		Some(frame)
1007	}
1008
1009	/// Refill with up to `CAP` frames. Must be drained first (`pop` returned `None`).
1010	fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
1011		debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
1012		self.pos = 0;
1013		self.len = 0;
1014		for frame in frames.take(Self::CAP) {
1015			self.frames[self.len].write(frame);
1016			self.len += 1;
1017		}
1018	}
1019
1020	/// `(frame count, total payload bytes)` of the buffered, not-yet-taken frames.
1021	/// Read once per fill to bump the egress payload counters for the whole batch.
1022	fn buffered(&self) -> (u64, u64) {
1023		let mut bytes = 0u64;
1024		for slot in &self.frames[self.pos..self.len] {
1025			// SAFETY: slots in `pos..len` are initialized (written by `fill`, not yet popped).
1026			bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
1027		}
1028		((self.len - self.pos) as u64, bytes)
1029	}
1030}
1031
1032impl Default for Prefetch {
1033	fn default() -> Self {
1034		Self {
1035			frames: [const { MaybeUninit::uninit() }; Self::CAP],
1036			pos: 0,
1037			len: 0,
1038		}
1039	}
1040}
1041
1042impl Drop for Prefetch {
1043	fn drop(&mut self) {
1044		for slot in &mut self.frames[self.pos..self.len] {
1045			// SAFETY: slots in `pos..len` are initialized and were never taken.
1046			unsafe { slot.assume_init_drop() };
1047		}
1048	}
1049}
1050
1051/// Consume a group, frame-by-frame.
1052///
1053/// Usually a view of one [`Producer`], but a group served across a route change is
1054/// *spliced*: it reads each contributing route's copy in turn, joined at the frame the
1055/// takeover happened on, so the reader never sees the seam.
1056pub struct Consumer {
1057	inner: ConsumerKind,
1058
1059	// Immutable stream state.
1060	info: Info,
1061
1062	// The parent track's info, inherited from the producer. Its `timescale` lets the
1063	// wire publisher emit per-frame timestamps at the right scale for a fetched group.
1064	track: track::Info,
1065
1066	// Egress payload meter, set by a tagged track via [`Self::with_meter`]. Empty
1067	// (no-op) for an untagged group.
1068	stats: stats::Meter,
1069	// The meter that owns unread content discarded by expiry. Route-specific
1070	// cursors inside a spliced group inherit this without metering delivery twice.
1071	stale_stats: stats::Meter,
1072
1073	// Subscriber-specific drift policy. A group can become stale after the track
1074	// hands it out, while its reader is waiting for the first or next frame.
1075	expiry: Option<Arc<dyn Expiry>>,
1076	expired: bool,
1077	// Sticky: the budget gave up on a cursor that had already taken every frame, so
1078	// the group ends rather than fails. Recorded because `expired` alone would turn a
1079	// clean end into `Error::Old` on the next poll, and a caller is allowed to probe
1080	// again after the end.
1081	ended: bool,
1082	// Cloned cursors are parallel views of one handed-out delivery. Whichever
1083	// observes expiry first records its unread tail; the others must not repeat it.
1084	stale_counted: Arc<AtomicBool>,
1085}
1086
1087/// Subscriber-specific policy for expiring a group after it was handed out.
1088pub(crate) trait Expiry: Send + Sync {
1089	/// Return whether the group is stale, registering `waiter` for anything that
1090	/// could change the answer while the group remains live.
1091	fn is_expired(&self, waiter: &kio::Waiter) -> bool;
1092}
1093
1094// `Plain` is the hot path and carries an inline frame prefetch, so boxing it to even the
1095// variants out would cost an allocation per group to save a pointer chase on the rare one.
1096#[expect(clippy::large_enum_variant)]
1097enum ConsumerKind {
1098	Plain(Plain),
1099	// Boxed: the spliced cursor set dwarfs the plain one, and splicing is the rare case.
1100	Spliced(Box<super::resume::Group>),
1101}
1102
1103/// The cursor state for a group backed by a single [`Producer`].
1104struct Plain {
1105	// Shared state with the producer.
1106	state: kio::Consumer<GroupState>,
1107
1108	// The index of the next frame to read.
1109	// NOTE: Cloned readers inherit this offset, but then run in parallel.
1110	index: usize,
1111
1112	// Exclusive cap on `index`, set by [`Consumer::set_frames`]. Reads end cleanly at it.
1113	end: Option<usize>,
1114
1115	// A batch of completed frames drained ahead under one lock (whole-frame reads only).
1116	prefetch: Prefetch,
1117
1118	// Record prefetched reads without entering the group's state on every frame.
1119	cache: Arc<cache::Track>,
1120	access: Arc<cache::Access>,
1121	refreshed: u64,
1122}
1123
1124impl Clone for Plain {
1125	fn clone(&self) -> Self {
1126		// A clone shares the channel and inherits `index`, but starts with an empty
1127		// prefetch: it re-reads its batch from the shared state, in parallel.
1128		Self {
1129			state: self.state.clone(),
1130			index: self.index,
1131			end: self.end,
1132			prefetch: Prefetch::default(),
1133			cache: self.cache.clone(),
1134			access: self.access.clone(),
1135			refreshed: self.refreshed,
1136		}
1137	}
1138}
1139
1140impl Clone for Consumer {
1141	fn clone(&self) -> Self {
1142		Self {
1143			inner: match &self.inner {
1144				ConsumerKind::Plain(plain) => ConsumerKind::Plain(plain.clone()),
1145				ConsumerKind::Spliced(spliced) => ConsumerKind::Spliced(Box::new((**spliced).clone())),
1146			},
1147			info: self.info,
1148			track: self.track.clone(),
1149			// Inherit the meter without re-counting the group: the original already
1150			// counted it when the track handed it out.
1151			stats: self.stats.clone(),
1152			stale_stats: self.stale_stats.clone(),
1153			expiry: self.expiry.clone(),
1154			expired: self.expired,
1155			ended: self.ended,
1156			stale_counted: self.stale_counted.clone(),
1157		}
1158	}
1159}
1160
1161impl std::ops::Deref for Consumer {
1162	type Target = Info;
1163
1164	fn deref(&self) -> &Self::Target {
1165		&self.info
1166	}
1167}
1168
1169impl Consumer {
1170	/// Snapshot the content this cursor would discard if its group were skipped.
1171	pub(crate) fn content(&self) -> stats::Content {
1172		match &self.inner {
1173			ConsumerKind::Plain(plain) => plain.state.read().content(),
1174			// Drift is evaluated before a segment copy is wrapped as a spliced group.
1175			// Keep the group count honest if a future caller reaches this fallback.
1176			ConsumerKind::Spliced(_) => stats::Content {
1177				groups: 1,
1178				..Default::default()
1179			},
1180		}
1181	}
1182
1183	/// Content not already attributed as delivered by this handed-out cursor.
1184	fn unread_content(&self) -> stats::Content {
1185		match &self.inner {
1186			ConsumerKind::Plain(plain) => plain.unread_content(),
1187			// Each route-specific plain cursor enforces expiry inside a spliced group.
1188			ConsumerKind::Spliced(_) => stats::Content::default(),
1189		}
1190	}
1191
1192	/// Rebuild this consumer as the head of a group assembled across route changes,
1193	/// keeping the group's identity and its track's properties. See [`super::resume`].
1194	pub(crate) fn into_spliced(self, mut spliced: super::resume::Group) -> Self {
1195		spliced.set_stale_meter(self.stale_stats.clone());
1196		Self {
1197			inner: ConsumerKind::Spliced(Box::new(spliced)),
1198			info: self.info,
1199			track: self.track,
1200			stats: self.stats,
1201			stale_stats: self.stale_stats,
1202			// Each segment keeps its own route-specific expiry policy. Applying the
1203			// head segment's policy to the assembled group would use the wrong edge
1204			// after a takeover.
1205			expiry: None,
1206			expired: false,
1207			ended: false,
1208			stale_counted: self.stale_counted,
1209		}
1210	}
1211
1212	/// Attach an egress payload meter, counting this as one delivered group.
1213	/// Called by a tagged track when it hands the consumer to a subscriber or fetch.
1214	pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
1215		meter.group();
1216		self.stats = meter.clone();
1217		self.set_stale_meter(meter);
1218		self
1219	}
1220
1221	/// Attach only the meter that owns content discarded by expiry.
1222	pub(crate) fn set_stale_meter(&mut self, meter: stats::Meter) {
1223		if let ConsumerKind::Spliced(spliced) = &mut self.inner {
1224			spliced.set_stale_meter(meter.clone());
1225		}
1226		self.stale_stats = meter;
1227	}
1228
1229	/// Keep applying this subscription's drift budget while the group is read.
1230	pub(crate) fn with_expiry(mut self, expiry: Arc<dyn Expiry>) -> Self {
1231		self.expiry = Some(expiry);
1232		self
1233	}
1234
1235	/// Check the parent subscription while a wire publisher drains detached payload.
1236	pub(crate) fn poll_expired(&mut self, waiter: &kio::Waiter) -> bool {
1237		self.poll_expired_while_pending(waiter, false)
1238	}
1239
1240	/// Apply the drift budget to a read that found nothing and is about to park.
1241	///
1242	/// A group with frames in hand is always drained to its end: the budget bounds a
1243	/// group that has *stalled* while the live edge moved on, not one whose reader is
1244	/// merely slower than the wire. Judging every read instead would truncate the tail
1245	/// of every group, since the arrival of the next group is exactly what makes the
1246	/// current one no longer newest.
1247	///
1248	/// Keeping it off the ready path also keeps it off the hot path: evaluating the
1249	/// policy walks the track's group cache under its lock, which is shared by every
1250	/// subscriber of that track.
1251	///
1252	/// `Some(false)` ends the group cleanly and `Some(true)` fails it with
1253	/// [`Error::Old`]; see [`Self::expired_truncates`].
1254	fn poll_expired_if_blocked(&mut self, waiter: &kio::Waiter) -> Option<bool> {
1255		if self.ended {
1256			return Some(false);
1257		}
1258		if !self.poll_expired(waiter) {
1259			return None;
1260		}
1261		let truncates = self.expired_truncates();
1262		self.ended = !truncates;
1263		Some(truncates)
1264	}
1265
1266	/// Whether giving up on this group now loses the reader anything.
1267	///
1268	/// A frame-level read only parks once the cursor has taken every frame the group
1269	/// holds, so expiring there costs nothing: the reader got everything that exists,
1270	/// and the group ends rather than fails. What it was still waiting for was the
1271	/// producer's FIN, and a group abandoned at its own end is indistinguishable from
1272	/// one that ended. A cursor that still holds unread content (a wire publisher with
1273	/// buffered frames, a half-read payload) is genuinely truncated and reports it.
1274	fn expired_truncates(&self) -> bool {
1275		let unread = self.unread_content();
1276		unread.frames > 0 || unread.bytes > 0
1277	}
1278
1279	/// Keep checking expiry while a wire publisher still owns buffered group data.
1280	pub(crate) fn poll_expired_while_pending(&mut self, waiter: &kio::Waiter, pending: bool) -> bool {
1281		if !self.expired
1282			&& (pending || self.expiry_pending())
1283			&& self.expiry.as_ref().is_some_and(|expiry| expiry.is_expired(waiter))
1284		{
1285			self.expired = true;
1286			if !self.stale_counted.swap(true, Ordering::Relaxed) {
1287				self.stale_stats.stale(self.unread_content());
1288			}
1289		}
1290		self.expired
1291	}
1292
1293	/// Whether expiry can still discard content or unblock a group that may grow.
1294	fn expiry_pending(&self) -> bool {
1295		match &self.inner {
1296			ConsumerKind::Plain(plain) => plain.expiry_pending(),
1297			// The route-specific plain cursors own expiry for a spliced group.
1298			ConsumerKind::Spliced(_) => false,
1299		}
1300	}
1301
1302	/// Whether this cursor failed because its subscription max age budget expired.
1303	pub(crate) fn latency_expired(&self) -> bool {
1304		self.expired
1305	}
1306
1307	/// Whether the group has been aborted (including pool eviction); the abort
1308	/// dropped the cached frames, so a held consumer has nothing left to read.
1309	///
1310	/// A spliced group spans several routes, so no single abort empties it; only a
1311	/// plain cursor can answer.
1312	pub(crate) fn is_aborted(&self) -> bool {
1313		match &self.inner {
1314			ConsumerKind::Plain(plain) => plain.state.read().abort.is_some(),
1315			ConsumerKind::Spliced(_) => false,
1316		}
1317	}
1318
1319	/// Mark the group as still being read, so a slow batch drain does not expire it.
1320	pub fn keep_alive(&self) {
1321		if let ConsumerKind::Plain(plain) = &self.inner {
1322			plain.state.read().charge.refresh();
1323		}
1324	}
1325
1326	/// Record a cache access from the consumer side: a parked group re-offered to
1327	/// its subscriber. Same stamp as [`Producer::cache_refresh`].
1328	pub(crate) fn cache_refresh(&self) {
1329		self.keep_alive();
1330	}
1331
1332	/// Park `waiter` until the group closes (finish, abort, or eviction). Spliced
1333	/// subscribers register on parked groups so an eviction wakes them; a group
1334	/// that already closed cleanly can never abort, so no waiter is needed.
1335	///
1336	/// A spliced group reads as closed without registering anything: no single abort
1337	/// empties it, so [`Self::is_aborted`] can never turn true and there is nothing
1338	/// a wakeup would change.
1339	pub(crate) fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<()> {
1340		match &self.inner {
1341			ConsumerKind::Plain(plain) => plain.state.poll_closed(waiter),
1342			ConsumerKind::Spliced(_) => Poll::Ready(()),
1343		}
1344	}
1345
1346	/// The parent track's timescale.
1347	pub fn timescale(&self) -> Timescale {
1348		self.track.timescale
1349	}
1350
1351	/// The index of the next frame this consumer will return.
1352	///
1353	/// Starts at 0, or at the group's first available frame once [`Self::set_frames`] has
1354	/// clamped it, and advances by one per frame read.
1355	pub fn index(&self) -> u64 {
1356		match &self.inner {
1357			ConsumerKind::Plain(plain) => plain.index as u64,
1358			ConsumerKind::Spliced(spliced) => spliced.index(),
1359		}
1360	}
1361
1362	/// Limit subsequent reads to these frame indices without rewinding read progress.
1363	///
1364	/// `2..=5` includes frames 2 through 5; `2..5` excludes frame 5. An omitted
1365	/// start preserves read progress, and an omitted end removes the cap.
1366	/// Raising the cap makes unread cached frames available again.
1367	pub fn set_frames(&mut self, frames: impl RangeBounds<u64>) {
1368		let (start, end) = super::subscription::sequence_bounds(frames);
1369		self.start_at(start);
1370		self.end_at(end.map_or(Bound::Unbounded, Bound::Excluded));
1371	}
1372
1373	/// Skip ahead so the next frame returned is `index`, discarding anything buffered
1374	/// below it.
1375	///
1376	/// Clamped *up* to the group's first available frame: frames the group never held
1377	/// (see [`Producer::start_at`]) can't be returned, so asking for one just starts at
1378	/// the first that exists. Read [`Self::index`] back to learn where the cursor
1379	/// actually landed.
1380	/// Only moves forward; a lower `index` is ignored, since the frames behind the
1381	/// cursor may already have been handed out.
1382	pub(crate) fn start_at(&mut self, index: u64) {
1383		match &mut self.inner {
1384			ConsumerKind::Plain(plain) => plain.start_at(index),
1385			ConsumerKind::Spliced(spliced) => spliced.start_at(index),
1386		}
1387	}
1388
1389	/// Advance the read cursor to `index`, skipping every frame below it.
1390	///
1391	/// Unlike [`Self::set_frames`], this does not clamp past a requested frame the group
1392	/// never held. A [`Producer::start_at`] floor above `index` still surfaces as
1393	/// [`Error::Lagged`].
1394	pub fn skip_to(&mut self, index: u64) {
1395		match &mut self.inner {
1396			ConsumerKind::Plain(plain) => plain.skip_to(index),
1397			ConsumerKind::Spliced(spliced) => spliced.start_at(index),
1398		}
1399	}
1400
1401	/// Stop reading at `end`, or remove the cap with `..`.
1402	///
1403	/// `..=2` reads through frame 2, `..2` stops before it, and `..0` is the empty range:
1404	/// no frame is delivered. Reads past the cap end cleanly (`None`), as if the group
1405	/// finished there. The cap can move in either direction: raising it re-offers
1406	/// frames that are still cached.
1407	pub(crate) fn end_at(&mut self, end: impl Into<Cap>) {
1408		let end = end.into().exclusive();
1409		match &mut self.inner {
1410			ConsumerKind::Plain(plain) => {
1411				plain.end = end.map(|end| usize::try_from(end).unwrap_or(usize::MAX));
1412			}
1413			ConsumerKind::Spliced(spliced) => spliced.end_at(end),
1414		}
1415	}
1416
1417	/// The number of frames written so far (completed plus any in-flight), independent of
1418	/// how many this consumer has read. The final total once the group is finished.
1419	pub fn frame_count(&self) -> usize {
1420		match &self.inner {
1421			ConsumerKind::Plain(plain) => {
1422				let state = plain.state.read();
1423				state.fin.unwrap_or(state.next_index)
1424			}
1425			ConsumerKind::Spliced(spliced) => spliced.frame_count(),
1426		}
1427	}
1428
1429	/// Return a consumer for the next frame for chunked reading.
1430	pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
1431		kio::wait(|waiter| self.poll_next_frame(waiter)).await
1432	}
1433
1434	/// Poll for the next frame, without blocking.
1435	///
1436	/// Returns None if the group is finished and the index is out of range, or the cursor
1437	/// passed the [`Self::set_frames`] cap.
1438	pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
1439		if self.ended {
1440			return Poll::Ready(Ok(None));
1441		}
1442		if self.expired {
1443			return Poll::Ready(Err(Error::Old));
1444		}
1445		let stats = self.stats.clone();
1446		let expiry = self
1447			.expiry
1448			.as_ref()
1449			.map(|policy| frame::Expiry::new(policy.clone(), self.stale_stats.clone(), self.stale_counted.clone()));
1450		let res = match &mut self.inner {
1451			ConsumerKind::Plain(plain) => plain.poll_next_frame(waiter, &stats, expiry),
1452			ConsumerKind::Spliced(spliced) => {
1453				// The per-route copies underneath are untagged, so meter the spliced
1454				// stream here: it is the one the subscriber actually reads.
1455				let res = ready!(spliced.poll_next_frame(waiter))?;
1456				if res.is_some() {
1457					stats.frames(1);
1458				}
1459				Poll::Ready(Ok(res.map(|frame| frame.with_meter(stats))))
1460			}
1461		};
1462		match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1463			Some(true) => Poll::Ready(Err(Error::Old)),
1464			Some(false) => Poll::Ready(Ok(None)),
1465			None => res,
1466		}
1467	}
1468
1469	/// Read the next frame (timestamp and payload) all at once, without blocking.
1470	pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
1471		if self.ended {
1472			return Poll::Ready(Ok(None));
1473		}
1474		if self.expired {
1475			return Poll::Ready(Err(Error::Old));
1476		}
1477		let stats = self.stats.clone();
1478		let res = match &mut self.inner {
1479			ConsumerKind::Plain(plain) => plain.poll_read_frame(waiter, &stats),
1480			ConsumerKind::Spliced(spliced) => {
1481				let res = ready!(spliced.poll_read_frame(waiter))?;
1482				if let Some(frame) = &res {
1483					stats.frames(1);
1484					stats.bytes(frame.payload.len() as u64);
1485				}
1486				Poll::Ready(Ok(res))
1487			}
1488		};
1489		match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1490			Some(true) => Poll::Ready(Err(Error::Old)),
1491			Some(false) => Poll::Ready(Ok(None)),
1492			None => res,
1493		}
1494	}
1495
1496	/// Read the next frame (timestamp and payload) all at once.
1497	pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
1498		// A prefetched frame is already buffered, so the drift budget (which only judges
1499		// a read that would park) can never apply to it.
1500		if !self.expired
1501			&& let ConsumerKind::Plain(plain) = &mut self.inner
1502		{
1503			// Serve from the prefetched batch without building a future or allocating a waker.
1504			if !plain.capped()
1505				&& let Some(frame) = plain.prefetch.pop()
1506			{
1507				plain.refresh_if_stale();
1508				plain.index += 1;
1509				return Ok(Some(frame));
1510			}
1511		}
1512		kio::wait(|waiter| self.poll_read_frame(waiter)).await
1513	}
1514
1515	/// Fill `out` with every frame that is ready, up to its capacity, without blocking.
1516	///
1517	/// This is a short read: it returns as soon as anything is ready rather than
1518	/// waiting for `out` to fill. A zero count means the group ended when the buffer
1519	/// has non-zero capacity.
1520	pub fn poll_read_frames<const N: usize>(
1521		&mut self,
1522		waiter: &kio::Waiter,
1523		out: &mut frame::Buffer<N>,
1524	) -> Poll<Result<usize>> {
1525		out.clear();
1526		if out.capacity() == 0 {
1527			return Poll::Ready(Ok(0));
1528		}
1529
1530		while !out.is_full() {
1531			match self.poll_read_frame(waiter) {
1532				Poll::Ready(Ok(Some(frame))) => out.push(frame).expect("buffer capacity checked"),
1533				Poll::Ready(Ok(None)) => break,
1534				Poll::Ready(Err(err)) => {
1535					if out.is_empty() {
1536						return Poll::Ready(Err(err));
1537					}
1538					break;
1539				}
1540				Poll::Pending if !out.is_empty() => break,
1541				Poll::Pending => return Poll::Pending,
1542			}
1543		}
1544
1545		Poll::Ready(Ok(out.len()))
1546	}
1547
1548	/// Fill `out` with every frame that is ready, blocking until a frame arrives or
1549	/// the group ends. Returns the current batch, empty only at the end of the group.
1550	pub async fn read_frames<'a, const N: usize>(
1551		&mut self,
1552		out: &'a mut frame::Buffer<N>,
1553	) -> Result<&'a mut [frame::Frame]> {
1554		kio::wait(|waiter| self.poll_read_frames(waiter, out)).await?;
1555		Ok(out.filled_mut())
1556	}
1557
1558	/// Poll until the group terminates, returning this cursor's next frame index.
1559	pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
1560		if self.ended {
1561			return Poll::Ready(Ok(self.index()));
1562		}
1563		if self.expired {
1564			return Poll::Ready(Err(Error::Old));
1565		}
1566		let res = match &mut self.inner {
1567			ConsumerKind::Plain(plain) => {
1568				let index = plain.index;
1569				plain
1570					.poll(waiter, |state| state.poll_end(index))
1571					.map(|res| res.map(|()| index as u64))
1572			}
1573			ConsumerKind::Spliced(spliced) => spliced.poll_finished(waiter),
1574		};
1575		match res.is_pending().then(|| self.poll_expired_if_blocked(waiter)).flatten() {
1576			Some(true) => Poll::Ready(Err(Error::Old)),
1577			// The group ended where the cursor stands, so that is its frame count.
1578			Some(false) => Poll::Ready(Ok(self.index())),
1579			None => res,
1580		}
1581	}
1582
1583	/// Block until the group terminates, returning this cursor's next frame index.
1584	///
1585	/// This answers for the cursor, not the group: a reader that drained every frame gets the
1586	/// clean end even if the group was aborted afterwards to release its cache, while one that
1587	/// stopped short gets that abort. A prior [`Self::skip_to`] contributes to the index even
1588	/// though those frames were not read. Use [`Self::frame_count`] for the producer's total.
1589	pub async fn finished(&mut self) -> Result<u64> {
1590		kio::wait(|waiter| self.poll_finished(waiter)).await
1591	}
1592}
1593
1594impl Plain {
1595	/// Whether this cursor still has unread content or may receive another frame.
1596	fn expiry_pending(&self) -> bool {
1597		if self.capped() {
1598			return false;
1599		}
1600
1601		let state = self.state.read();
1602		state.abort.is_none() && state.fin.is_none_or(|fin| self.index < fin)
1603	}
1604
1605	/// Content this cursor has neither returned nor already counted in a prefetch batch.
1606	fn unread_content(&self) -> stats::Content {
1607		let prefetched = self.prefetch.buffered().0 as usize;
1608		let start = self.index.saturating_add(prefetched);
1609		let end = self.end.unwrap_or(usize::MAX);
1610		self.state.read().content_range(start, end)
1611	}
1612
1613	/// Record prefetched reads, updating the eviction rank once per sampled tick.
1614	fn refresh_if_stale(&mut self) {
1615		self.access.touch();
1616		let tick = self.cache.pool().now();
1617		if tick != self.refreshed {
1618			self.state.read().charge.refresh();
1619			self.refreshed = tick;
1620		}
1621	}
1622
1623	// A helper to automatically apply Dropped if the state is closed without an error.
1624	fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
1625	where
1626		F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
1627	{
1628		Poll::Ready(match ready!(self.state.poll(waiter, f)) {
1629			Ok(res) => res,
1630			// We try to clone abort just in case the function forgot to check for terminal state.
1631			Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
1632		})
1633	}
1634
1635	/// Whether the cursor has passed the `end_at` cap.
1636	fn capped(&self) -> bool {
1637		self.end.is_some_and(|end| self.index >= end)
1638	}
1639
1640	fn start_at(&mut self, index: u64) {
1641		let index = usize::try_from(index).unwrap_or(usize::MAX);
1642		let index = index.max(self.state.read().offset);
1643		if index <= self.index {
1644			return;
1645		}
1646		self.index = index;
1647		// The batch was drained from below the new cursor, so it can't be reused.
1648		self.prefetch = Prefetch::default();
1649	}
1650
1651	fn skip_to(&mut self, index: u64) {
1652		let index = usize::try_from(index).unwrap_or(usize::MAX);
1653		if index <= self.index {
1654			return;
1655		}
1656		self.index = index;
1657		self.prefetch = Prefetch::default();
1658	}
1659
1660	fn poll_next_frame(
1661		&mut self,
1662		waiter: &kio::Waiter,
1663		stats: &stats::Meter,
1664		expiry: Option<frame::Expiry>,
1665	) -> Poll<Result<Option<frame::Consumer>>> {
1666		if self.capped() {
1667			return Poll::Ready(Ok(None));
1668		}
1669		let end = self.end.unwrap_or(usize::MAX);
1670
1671		// Hand out any frames a prior read_frame prefetched before touching the tail.
1672		// Their bytes were already counted at the batch fill, so the frame::Consumer
1673		// carries no meter.
1674		if let Some(frame) = self.prefetch.pop() {
1675			self.refresh_if_stale();
1676			self.index += 1;
1677			let tail = self.index.saturating_add(self.prefetch.buffered().0 as usize)..end;
1678			let info = frame::Info {
1679				size: frame.payload.len() as u64,
1680				timestamp: frame.timestamp,
1681			};
1682			let source = frame::Source::Complete(frame.payload);
1683			let frame = frame::Consumer::new(self.state.clone(), info, source);
1684			return Poll::Ready(Ok(Some(match expiry {
1685				Some(expiry) => frame.with_expiry(expiry.for_frame(tail, false)),
1686				None => frame,
1687			})));
1688		}
1689
1690		let index = self.index;
1691		let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
1692			return Poll::Ready(Ok(None));
1693		};
1694
1695		self.index += 1;
1696		// A direct read (not prefetched): count the frame here; the frame::Consumer
1697		// counts its bytes per chunk as they're read out.
1698		stats.frames(1);
1699		let frame = frame::Consumer::new(self.state.clone(), info, source).with_meter(stats.clone());
1700		Poll::Ready(Ok(Some(match expiry {
1701			Some(expiry) => frame.with_expiry(expiry.for_frame(self.index..end, true)),
1702			None => frame,
1703		})))
1704	}
1705
1706	fn poll_read_frame(&mut self, waiter: &kio::Waiter, stats: &stats::Meter) -> Poll<Result<Option<frame::Frame>>> {
1707		if self.capped() {
1708			return Poll::Ready(Ok(None));
1709		}
1710
1711		// Fast path: serve from the prefetched batch without locking or allocating a waker.
1712		if let Some(frame) = self.prefetch.pop() {
1713			self.refresh_if_stale();
1714			self.index += 1;
1715			return Poll::Ready(Ok(Some(frame)));
1716		}
1717
1718		// The batch is drained: refill it under a single lock, registering the waiter if
1719		// nothing is ready. Borrow the two fields disjointly so the closure can fill.
1720		let index = self.index;
1721		// Never buffer past the cap: `end_at` can be raised later, and those frames must
1722		// come from the shared state then, not from a batch drained under the old cap.
1723		let budget = self.end.map_or(usize::MAX, |end| end.saturating_sub(index));
1724		let prefetch = &mut self.prefetch;
1725		let res = self.state.poll(waiter, |state| {
1726			if index < state.offset {
1727				return Poll::Ready(Err(Error::Lagged));
1728			}
1729			// `local` can run past the buffered count when frames were cleared out from
1730			// under us (abort, unfinished drop); clamp so `range` never panics on an
1731			// out-of-bounds start. `fill` always resets the batch, so an empty range
1732			// leaves `len == 0` and the terminal checks below resolve abort/fin/pending.
1733			let local = (index - state.offset).min(state.frames.len());
1734			prefetch.fill(state.frames.range(local..).take(budget).cloned());
1735			if prefetch.len > 0 {
1736				// One stamp covers the whole batch: frames popped from the prefetch
1737				// don't re-stamp until the next refill, which `CAP` bounds.
1738				state.charge.refresh();
1739				return Poll::Ready(Ok(()));
1740			}
1741			// Nothing completed at `index`: an in-flight tail waits, otherwise resolve
1742			// the terminal state (whole-frame reads never stream the partial).
1743			state.poll_terminal(index)
1744		});
1745
1746		match ready!(res) {
1747			Ok(Ok(())) => {}
1748			Ok(Err(err)) => return Poll::Ready(Err(err)),
1749			Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
1750		}
1751
1752		// The refill already updated the eviction rank under the group lock.
1753		self.refreshed = self.cache.pool().now();
1754
1755		// A fresh batch was just filled (empty only on a clean end). Count the whole
1756		// batch once here, under no lock, so the drained pops that follow stay free.
1757		let (frames, bytes) = self.prefetch.buffered();
1758		stats.frames(frames);
1759		stats.bytes(bytes);
1760
1761		Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
1762			self.index += 1;
1763		})))
1764	}
1765}
1766
1767/// Options for a one-shot [`track::Consumer::fetch_group`] of a past group.
1768#[derive(Clone, Debug, Default)]
1769#[non_exhaustive]
1770pub struct Fetch {
1771	/// Delivery priority for the fetched group's stream. Defaults to 0.
1772	pub priority: u8,
1773
1774	/// Index of the first frame to fetch within the group. Defaults to 0, the whole group.
1775	///
1776	/// Use this to fill a hole left by a route change: the group's head is already
1777	/// cached locally and only the tail is missing.
1778	///
1779	/// There is no matching end: a fetch always runs to the end of the group, and a
1780	/// caller wanting less caps the returned consumer with [`Consumer::set_frames`]. Stopping
1781	/// the *fetch* short would put a group in the cache that is indistinguishable from a
1782	/// complete one, so a later fetch of the whole group would resolve from it and come
1783	/// up short.
1784	pub frame_start: u64,
1785}
1786
1787impl Fetch {
1788	/// Set the delivery priority, returning `self` for chaining.
1789	pub fn with_priority(mut self, priority: u8) -> Self {
1790		self.priority = priority;
1791		self
1792	}
1793
1794	/// Set the first frame to fetch, returning `self` for chaining.
1795	pub fn with_frame_start(mut self, frame_start: u64) -> Self {
1796		self.frame_start = frame_start;
1797		self
1798	}
1799}
1800
1801/// A consumer's request for a single past group, handed to a handler via
1802/// [`track::Dynamic::requested_group`].
1803///
1804/// The handler fulfills it by calling [`Self::accept`], which inserts the group
1805/// into the track cache (resolving every [`track::Consumer::fetch_group`] that joined the
1806/// attempt) and returns a [`Producer`] to fill. A relay typically opens a wire
1807/// FETCH, reads FETCH_OK, then accepts. The request carries its own producer handle,
1808/// so it works the same whether or not the track has been accepted yet.
1809pub struct Request {
1810	pub(crate) state: kio::Producer<track::TrackState>,
1811	pub(crate) fetch: kio::Shared<track::FetchState>,
1812	pub(crate) sequence: u64,
1813	pub(crate) priority: u8,
1814	pub(crate) frame_start: u64,
1815	pub(crate) result: kio::Producer<track::FetchOutcome>,
1816	pub(crate) done: bool,
1817}
1818
1819#[cfg(test)]
1820mod test {
1821	use super::*;
1822	use crate::model::test_tracing::count_drop_warnings;
1823	use bytes::Bytes;
1824	use futures::FutureExt;
1825
1826	/// [`FRAME_SLOTS`] is std's rounding, not ours, so measure it: a larger real value
1827	/// would undercharge every cached group without touching a line of this crate.
1828	#[test]
1829	fn one_frame_fits_the_charged_slots() {
1830		let mut frames: VecDeque<Frame> = VecDeque::new();
1831		frames.push_back(Frame {
1832			timestamp: Timestamp::ZERO,
1833			payload: Bytes::new(),
1834		});
1835		let capacity = frames.capacity();
1836		assert!(
1837			capacity <= FRAME_SLOTS,
1838			"a one-frame deque now allocates {capacity} slots"
1839		);
1840	}
1841
1842	#[test]
1843	fn basic_frame_reading() {
1844		let mut producer = Info { sequence: 0 }.produce();
1845		producer
1846			.write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
1847			.unwrap();
1848		producer
1849			.write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
1850			.unwrap();
1851		producer.finish().unwrap();
1852
1853		let mut consumer = producer.consume();
1854		let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1855		assert_eq!(f0.size, 6);
1856		let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1857		assert_eq!(f1.size, 6);
1858		let end = consumer.next_frame().now_or_never().unwrap().unwrap();
1859		assert!(end.is_none());
1860	}
1861
1862	#[test]
1863	fn read_frame_all_at_once() {
1864		let mut producer = Info { sequence: 0 }.produce();
1865		producer
1866			.write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
1867			.unwrap();
1868		producer.finish().unwrap();
1869
1870		let mut consumer = producer.consume();
1871		let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1872		assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1873	}
1874
1875	#[test]
1876	fn read_frame_preserves_timestamp() {
1877		let mut producer = Info { sequence: 0 }.produce();
1878		let timestamp = Timestamp::from_micros(20_000).unwrap();
1879		producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
1880		producer.finish().unwrap();
1881
1882		let mut consumer = producer.consume();
1883		let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1884		assert_eq!(frame.timestamp.as_micros(), 20_000);
1885		assert_eq!(frame.payload, Bytes::from_static(b"hello"));
1886	}
1887
1888	#[test]
1889	fn chunked_frame_reads_whole() {
1890		let mut producer = Info { sequence: 0 }.produce();
1891		{
1892			let mut frame = producer
1893				.create_frame(frame::Info {
1894					size: 10,
1895					timestamp: Timestamp::ZERO,
1896				})
1897				.unwrap();
1898			frame.write(Bytes::from_static(b"hello")).unwrap();
1899			frame.write(Bytes::from_static(b"world")).unwrap();
1900			frame.finish().unwrap();
1901		}
1902		producer.finish().unwrap();
1903
1904		// Frame data is held in a single per-frame buffer; a whole-frame read returns
1905		// the full contents in one slice.
1906		let mut consumer = producer.consume();
1907		let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
1908		assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
1909	}
1910
1911	#[test]
1912	fn chunked_frame_streams_partial() {
1913		let mut producer = Info { sequence: 0 }.produce();
1914		let mut consumer = producer.consume();
1915
1916		let mut frame = producer
1917			.create_frame(frame::Info {
1918				size: 6,
1919				timestamp: Timestamp::ZERO,
1920			})
1921			.unwrap();
1922		frame.write(Bytes::from_static(b"foo")).unwrap();
1923
1924		// A consumer can stream the in-flight tail before it's finished.
1925		let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
1926		let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
1927		assert_eq!(c1, Some(Bytes::from_static(b"foo")));
1928		assert!(f.read_chunk().now_or_never().is_none());
1929
1930		frame.write(Bytes::from_static(b"bar")).unwrap();
1931		frame.finish().unwrap();
1932
1933		let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
1934		assert_eq!(c2, Some(Bytes::from_static(b"bar")));
1935		let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
1936		assert_eq!(c3, None);
1937	}
1938
1939	#[test]
1940	fn group_finish_returns_none() {
1941		let producer = Info { sequence: 0 }.produce();
1942		producer.finish().unwrap();
1943
1944		let mut consumer = producer.consume();
1945		let end = consumer.next_frame().now_or_never().unwrap().unwrap();
1946		assert!(end.is_none());
1947	}
1948
1949	#[test]
1950	fn abort_propagates() {
1951		let producer = Info { sequence: 0 }.produce();
1952		let mut consumer = producer.consume();
1953		producer.abort(crate::Error::Cancel).unwrap();
1954
1955		let result = consumer.next_frame().now_or_never().unwrap();
1956		assert!(matches!(result, Err(crate::Error::Cancel)));
1957	}
1958
1959	#[test]
1960	fn abort_clears_cached_frames() {
1961		let mut producer = Info { sequence: 0 }.produce();
1962		producer
1963			.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1964			.unwrap();
1965
1966		// A stale consumer that never reads must not pin the cached frames.
1967		let _consumer = producer.consume();
1968		assert_eq!(producer.state.read().frames.len(), 1);
1969
1970		producer.clone().abort(crate::Error::Cancel).unwrap();
1971
1972		let state = producer.state.read();
1973		assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
1974		assert_eq!(state.cache, 0);
1975	}
1976
1977	#[test]
1978	fn drop_unfinished_clears_cached_frames() {
1979		let producer = Info { sequence: 0 }.produce();
1980		let mut writer = producer.clone();
1981		writer
1982			.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
1983			.unwrap();
1984
1985		// A stale consumer keeps the channel (and thus the cache) alive.
1986		let mut consumer = producer.consume();
1987		assert_eq!(producer.state.read().frames.len(), 1);
1988
1989		// Drop every producer without finishing: the cache is released.
1990		drop(writer);
1991		drop(producer);
1992
1993		let result = consumer.next_frame().now_or_never().unwrap();
1994		assert!(matches!(result, Err(crate::Error::Dropped)));
1995	}
1996
1997	#[test]
1998	fn drop_after_abort_does_not_warn() {
1999		let warns = count_drop_warnings("group::Producer dropped without finish", || {
2000			let producer = Info { sequence: 0 }.produce();
2001			let keep = producer.clone();
2002			let mut writer = producer.clone();
2003			writer
2004				.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2005				.unwrap();
2006			let _consumer = producer.consume();
2007			writer.abort(crate::Error::Cancel).unwrap();
2008			drop(keep);
2009		});
2010		assert_eq!(warns, 0, "abort-then-drop must not emit unfinished-producer WARN");
2011	}
2012
2013	#[test]
2014	fn drop_unfinished_warns() {
2015		let warns = count_drop_warnings("group::Producer dropped without finish", || {
2016			let producer = Info { sequence: 0 }.produce();
2017			let mut writer = producer.clone();
2018			writer
2019				.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2020				.unwrap();
2021			let _consumer = producer.consume();
2022			drop(writer);
2023			drop(producer);
2024		});
2025		assert!(warns >= 1, "unfinished drop must emit unfinished-producer WARN");
2026	}
2027
2028	#[test]
2029	fn drop_finished_keeps_cached_frames() {
2030		let mut producer = Info { sequence: 0 }.produce();
2031		producer
2032			.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2033			.unwrap();
2034		producer.finish().unwrap();
2035
2036		let mut consumer = producer.consume();
2037		drop(producer);
2038
2039		// A cleanly finished group keeps its cache so the consumer can still drain.
2040		let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2041		assert_eq!(frame.payload, Bytes::from_static(b"data"));
2042	}
2043
2044	#[tokio::test]
2045	async fn pending_then_ready() {
2046		let mut producer = Info { sequence: 0 }.produce();
2047		let mut consumer = producer.consume();
2048
2049		// Consumer blocks because no frames yet.
2050		assert!(consumer.next_frame().now_or_never().is_none());
2051
2052		producer
2053			.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
2054			.unwrap();
2055		producer.finish().unwrap();
2056
2057		let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
2058		assert_eq!(frame.size, 4);
2059	}
2060
2061	#[test]
2062	fn overflow_aborts_the_group() {
2063		let mut producer = Info { sequence: 0 }.produce();
2064		let mut consumer = producer.consume();
2065
2066		let big = Bytes::from(vec![0u8; MAX_CACHE_BYTES as usize]);
2067		producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
2068		assert!(matches!(
2069			producer.write_frame(Timestamp::ZERO, big),
2070			Err(Error::GroupTooLarge)
2071		));
2072
2073		{
2074			let state = producer.state.read();
2075			assert!(matches!(state.abort, Some(Error::GroupTooLarge)));
2076			assert!(state.frames.is_empty());
2077			assert_eq!(state.offset, 0);
2078		}
2079
2080		let result = consumer.next_frame().now_or_never().unwrap();
2081		assert!(matches!(result, Err(Error::GroupTooLarge)));
2082	}
2083
2084	#[test]
2085	fn no_overflow_under_budget() {
2086		let mut producer = Info { sequence: 0 }.produce();
2087		// 8192 one-byte frames is the largest legal group; they all stay cached.
2088		for _ in 0..MAX_GROUP_FRAMES {
2089			producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2090		}
2091		producer.finish().unwrap();
2092
2093		let state = producer.state.read();
2094		assert_eq!(state.offset, 0);
2095		assert_eq!(state.frames.len(), MAX_GROUP_FRAMES);
2096		assert!(state.abort.is_none());
2097	}
2098
2099	#[test]
2100	fn writer_sees_group_too_large_on_the_8193rd_frame() {
2101		let mut producer = Info { sequence: 0 }.produce();
2102		for _ in 0..MAX_GROUP_FRAMES {
2103			producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2104		}
2105		assert!(matches!(
2106			producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")),
2107			Err(Error::GroupTooLarge)
2108		));
2109		assert!(matches!(producer.state.read().abort, Some(Error::GroupTooLarge)));
2110	}
2111
2112	#[test]
2113	fn clone_consumer_independent() {
2114		let mut producer = Info { sequence: 0 }.produce();
2115		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2116
2117		let mut c1 = producer.consume();
2118		// Read one frame from c1
2119		let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
2120
2121		// Clone c1, inheriting its index (past first frame).
2122		let mut c2 = c1.clone();
2123
2124		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2125		producer.finish().unwrap();
2126
2127		// c2 should get the second frame (inherited index)
2128		let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
2129		assert_eq!(f.size, 1); // "b"
2130
2131		let end = c2.next_frame().now_or_never().unwrap().unwrap();
2132		assert!(end.is_none());
2133	}
2134
2135	fn prefetched_consumer(pool: &cache::Pool, max_age: std::time::Duration) -> (Producer, Consumer) {
2136		let cache = cache::Track::new(pool.clone(), kio::Weak::new());
2137		let track = track::Info::default().with_max_age(max_age);
2138		let mut producer = Producer::new(Info { sequence: 0 }, track, cache);
2139		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2140		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2141		producer.finish().unwrap();
2142
2143		let mut consumer = producer.consume();
2144		consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2145		(producer, consumer)
2146	}
2147
2148	#[test]
2149	fn prefetch_refresh_honors_pool_expiry() {
2150		let config = cache::Config::default().with_expiry(std::time::Duration::from_secs(1));
2151		let pool = cache::Pool::new(config);
2152		let (producer, mut consumer) = prefetched_consumer(&pool, std::time::Duration::MAX);
2153		let before = producer.cache_accessed();
2154
2155		crate::model::clock::advance(std::time::Duration::from_millis(600));
2156		consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2157
2158		assert!(producer.cache_accessed() > before, "the pool cadence is used");
2159	}
2160
2161	#[test]
2162	fn prefetch_refresh_honors_track_max_age() {
2163		let config = cache::Config::default().with_expiry(std::time::Duration::from_secs(30));
2164		let pool = cache::Pool::new(config);
2165		let (producer, mut consumer) = prefetched_consumer(&pool, std::time::Duration::from_secs(1));
2166		let before = producer.cache_accessed();
2167
2168		crate::model::clock::advance(std::time::Duration::from_millis(600));
2169		consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2170
2171		assert!(producer.cache_accessed() > before, "the track cadence remains in force");
2172	}
2173
2174	/// Reading more than one prefetch batch drains every frame in order across the
2175	/// batch boundary (the refill starts exactly where the previous batch ended).
2176	#[test]
2177	fn read_frame_crosses_prefetch_batches() {
2178		let n = Prefetch::CAP * 3 + 5;
2179		let mut producer = Info { sequence: 0 }.produce();
2180		for i in 0..n {
2181			producer
2182				.write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
2183				.unwrap();
2184		}
2185		producer.finish().unwrap();
2186
2187		let mut consumer = producer.consume();
2188		for i in 0..n {
2189			let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2190			assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
2191		}
2192		assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2193	}
2194
2195	/// A finished group is still aborted once its frames are released to free memory (the
2196	/// track's max age window, or the cache pool). A reader that already drained every frame
2197	/// is missing nothing, so it must see the clean end of group rather than the abort.
2198	#[test]
2199	fn abort_after_finish_keeps_the_clean_end_for_a_drained_reader() {
2200		let mut producer = Info { sequence: 0 }.produce();
2201		producer
2202			.write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
2203			.unwrap();
2204		producer.finish().unwrap();
2205
2206		let mut drained = producer.consume();
2207		let mut behind = producer.consume();
2208		let frame = drained.read_frame().now_or_never().unwrap().unwrap().unwrap();
2209		assert_eq!(frame.payload, Bytes::from_static(b"hello"));
2210
2211		producer.abort(Error::Old).unwrap();
2212
2213		// Drained everything before the abort: nothing is missing.
2214		assert!(drained.read_frame().now_or_never().unwrap().unwrap().is_none());
2215		assert!(drained.next_frame().now_or_never().unwrap().unwrap().is_none());
2216
2217		// Never read the frame, and its bytes are gone: a truncated stream, not a clean end.
2218		assert!(matches!(behind.read_frame().now_or_never().unwrap(), Err(Error::Old)));
2219	}
2220
2221	/// `finished` answers for the cursor: a drained reader gets the clean end even after the
2222	/// abort that released the cache, and one that stopped short gets that abort. The
2223	/// producer's total stays available on `frame_count`.
2224	#[test]
2225	fn finished_answers_for_the_cursor() {
2226		let mut producer = Info { sequence: 0 }.produce();
2227		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2228		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2229		producer.finish().unwrap();
2230
2231		let mut drained = producer.consume();
2232		let mut behind = producer.consume();
2233		while drained.read_frame().now_or_never().unwrap().unwrap().is_some() {}
2234		behind.read_frame().now_or_never().unwrap().unwrap().unwrap();
2235
2236		producer.abort(Error::Old).unwrap();
2237
2238		assert_eq!(drained.finished().now_or_never().unwrap().unwrap(), 2);
2239		assert!(matches!(behind.finished().now_or_never().unwrap(), Err(Error::Old)));
2240		assert_eq!(behind.frame_count(), 2);
2241	}
2242
2243	/// A cursor on a group aborted for overflowing its budget can never reach the end,
2244	/// so `finished` reports that abort instead of parking forever.
2245	#[test]
2246	fn finished_reports_a_group_too_large() {
2247		let mut producer = Info { sequence: 0 }.produce();
2248		let mut consumer = producer.consume();
2249
2250		let big = Bytes::from(vec![0u8; MAX_CACHE_BYTES as usize]);
2251		producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
2252		assert!(matches!(
2253			producer.write_frame(Timestamp::ZERO, big),
2254			Err(Error::GroupTooLarge)
2255		));
2256
2257		assert!(matches!(
2258			consumer.finished().now_or_never().unwrap(),
2259			Err(Error::GroupTooLarge)
2260		));
2261	}
2262
2263	/// `next_frame` drains frames a prior `read_frame` prefetched, preserving order.
2264	#[test]
2265	fn interleave_read_and_next_frame() {
2266		let mut producer = Info { sequence: 0 }.produce();
2267		for i in 0..5u8 {
2268			producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
2269		}
2270		producer.finish().unwrap();
2271
2272		let mut consumer = producer.consume();
2273		// The first whole-frame read prefetches all five frames into the batch.
2274		let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2275		assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
2276
2277		// next_frame must continue from the batch, not skip ahead or repeat.
2278		for i in 1..5u8 {
2279			let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
2280			let data = f.read_all().now_or_never().unwrap().unwrap();
2281			assert_eq!(data, Bytes::from(vec![i; 1]));
2282		}
2283		assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
2284	}
2285
2286	/// A `read_frame` whose index sits past the buffered frames (cleared by an abort)
2287	/// must surface the error, not panic on an out-of-range `range(local..)`.
2288	#[test]
2289	fn read_frame_past_cleared_frames_does_not_panic() {
2290		let mut producer = Info { sequence: 0 }.produce();
2291		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2292		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
2293
2294		let mut consumer = producer.consume();
2295		consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2296		consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2297
2298		// Abort clears the cached frames but leaves the consumer's index (2) past them, so the
2299		// refill's `local` (2) exceeds `frames.len()` (0).
2300		producer.abort(Error::Cancel).unwrap();
2301
2302		let result = consumer.read_frame().now_or_never().unwrap();
2303		assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
2304	}
2305
2306	/// Dropping a consumer mid-batch must drop the buffered-but-untaken frames
2307	/// (exercises the `MaybeUninit` Drop path; run under miri to catch leaks/UB).
2308	#[test]
2309	fn drop_with_partial_batch() {
2310		let mut producer = Info { sequence: 0 }.produce();
2311		for _ in 0..Prefetch::CAP {
2312			producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2313		}
2314		producer.finish().unwrap();
2315
2316		let mut consumer = producer.consume();
2317		// Take one frame so the batch is filled but only partially drained.
2318		let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
2319		drop(consumer);
2320	}
2321
2322	/// A parked chunk reader is woken by each chunk write. kio only notifies when
2323	/// a write guard was mutably accessed, so `frame_notify` must mark the guard
2324	/// modified; a guard dropped untouched wakes nobody and the reader would
2325	/// stall until the frame completed.
2326	#[tokio::test]
2327	async fn chunk_write_wakes_parked_reader() {
2328		let mut producer = Info { sequence: 0 }.produce();
2329		let mut consumer = producer.consume();
2330		let mut frame = producer
2331			.create_frame(frame::Info {
2332				size: 6,
2333				timestamp: Timestamp::ZERO,
2334			})
2335			.unwrap();
2336		let mut f = consumer.next_frame().await.unwrap().unwrap();
2337		let handle = tokio::spawn(async move { f.read_chunk().await });
2338		// Let the reader park on the empty partial before the chunk lands.
2339		tokio::time::sleep(std::time::Duration::from_millis(50)).await;
2340		frame.write(Bytes::from_static(b"foo")).unwrap();
2341		let chunk = tokio::time::timeout(std::time::Duration::from_secs(2), handle)
2342			.await
2343			.expect("parked chunk reader was never woken by the chunk write")
2344			.unwrap()
2345			.unwrap();
2346		assert_eq!(chunk, Some(Bytes::from_static(b"foo")));
2347	}
2348
2349	/// A frame whose timestamp is at a different scale is converted to the group's
2350	/// scale by `create_frame`.
2351	#[test]
2352	fn create_frame_converts_mismatched_scale() {
2353		use crate::{Timescale, Timestamp};
2354
2355		let mut producer = Producer::new(
2356			Info { sequence: 0 },
2357			track::Info::default().with_timescale(Timescale::MICRO),
2358			Default::default(),
2359		);
2360		let frame = frame::Info {
2361			size: 3,
2362			timestamp: Timestamp::from_millis(1).unwrap(), // 1ms -> 1000µs
2363		};
2364		let writer = producer.create_frame(frame).unwrap();
2365		assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
2366		assert_eq!(writer.timestamp.value(), 1000);
2367	}
2368
2369	/// An explicit current timestamp is converted to the group's scale.
2370	#[tokio::test]
2371	async fn create_frame_converts_current_timestamp() {
2372		use crate::Timescale;
2373
2374		let mut producer = Producer::new(
2375			Info { sequence: 0 },
2376			track::Info::default().with_timescale(Timescale::MICRO),
2377			Default::default(),
2378		);
2379		let writer = producer
2380			.create_frame(frame::Info {
2381				size: 3,
2382				timestamp: Timestamp::now(),
2383			})
2384			.unwrap();
2385		assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
2386		assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
2387	}
2388
2389	/// A group can start partway in, so a route can serve the tail of a group whose
2390	/// head came from somewhere else.
2391	#[test]
2392	fn start_at_starts_the_group_later() {
2393		let mut producer = Info { sequence: 0 }.produce();
2394		producer.start_at(3).unwrap();
2395		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"d")).unwrap();
2396		producer.finish().unwrap();
2397
2398		// The frame landed at index 3, so the group's length counts the missing head.
2399		assert_eq!(producer.frame_count(), 4);
2400
2401		let mut consumer = producer.consume();
2402		assert_eq!(consumer.frame_count(), 4);
2403
2404		// A reader positioned at the start is missing the head, and `finished` answers
2405		// for that cursor.
2406		assert!(matches!(
2407			consumer.finished().now_or_never().unwrap(),
2408			Err(Error::Lagged)
2409		));
2410		assert!(matches!(
2411			consumer.read_frame().now_or_never().unwrap(),
2412			Err(Error::Lagged)
2413		));
2414	}
2415
2416	/// Seeking to the group's first available frame is how a spliced reader picks up
2417	/// the tail; a lower index clamps up rather than failing.
2418	#[test]
2419	fn start_at_clamps_up_to_the_first_frame() {
2420		let mut producer = Info { sequence: 0 }.produce();
2421		producer.start_at(3).unwrap();
2422		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"d")).unwrap();
2423		producer.finish().unwrap();
2424
2425		let mut consumer = producer.consume();
2426		consumer.start_at(1);
2427		assert_eq!(consumer.index(), 3, "clamped up to the first frame that exists");
2428		assert_eq!(
2429			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload,
2430			Bytes::from_static(b"d")
2431		);
2432	}
2433
2434	/// `end_at` ends the read cleanly at the cap, and raising it re-offers the frames
2435	/// still cached behind it.
2436	#[test]
2437	fn end_at_caps_and_reopens() {
2438		let mut producer = Info { sequence: 0 }.produce();
2439		for i in 0..4u8 {
2440			producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i])).unwrap();
2441		}
2442		producer.finish().unwrap();
2443
2444		let mut consumer = producer.consume();
2445		consumer.set_frames(..2);
2446		assert_eq!(
2447			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2448			0
2449		);
2450		assert_eq!(
2451			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2452			1
2453		);
2454		assert!(
2455			consumer.read_frame().now_or_never().unwrap().unwrap().is_none(),
2456			"capped reads end cleanly"
2457		);
2458
2459		consumer.set_frames(..);
2460		assert_eq!(
2461			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2462			2
2463		);
2464	}
2465
2466	#[test]
2467	fn frame_ranges_preserve_progress_and_make_inclusion_explicit() {
2468		let mut producer = Info { sequence: 0 }.produce();
2469		for i in 0..4u8 {
2470			producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i])).unwrap();
2471		}
2472		producer.finish().unwrap();
2473		let mut consumer = producer.consume();
2474		consumer.set_frames(1..=1);
2475		assert_eq!(
2476			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2477			1
2478		);
2479		assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2480		consumer.set_frames(..3);
2481		assert_eq!(
2482			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2483			2
2484		);
2485		assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
2486		consumer.set_frames(0..=3);
2487		assert_eq!(
2488			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload[0],
2489			3
2490		);
2491	}
2492
2493	/// An exclusive cap at 0 is the empty range: no frame is delivered, and raising
2494	/// it re-offers the held frames.
2495	#[test]
2496	fn end_at_zero_is_empty() {
2497		let mut producer = Info { sequence: 0 }.produce();
2498		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
2499		producer.finish().unwrap();
2500
2501		let mut consumer = producer.consume();
2502		consumer.set_frames(..0);
2503		assert!(
2504			consumer.read_frame().now_or_never().unwrap().unwrap().is_none(),
2505			"empty cap delivers nothing"
2506		);
2507
2508		consumer.set_frames(..1);
2509		assert_eq!(
2510			consumer.read_frame().now_or_never().unwrap().unwrap().unwrap().payload,
2511			Bytes::from_static(b"x")
2512		);
2513	}
2514
2515	/// Where the group begins is part of its shape, so it can't move once frames exist.
2516	#[test]
2517	fn start_at_rejected_after_a_frame() {
2518		let mut producer = Info { sequence: 0 }.produce();
2519		// Re-declaring before the first frame is fine; the shape isn't committed yet.
2520		producer.start_at(2).unwrap();
2521		producer.start_at(3).unwrap();
2522
2523		producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
2524		assert!(matches!(producer.start_at(4), Err(Error::Closed)));
2525		assert_eq!(producer.frame_count(), 4, "the frame landed at index 3");
2526
2527		// Finishing likewise settles the shape.
2528		let mut producer = Info { sequence: 1 }.produce();
2529		producer.finish().unwrap();
2530		assert!(matches!(producer.start_at(1), Err(Error::Closed)));
2531	}
2532
2533	/// The start must leave room for at least one frame index.
2534	#[test]
2535	fn start_at_rejects_the_largest_index() {
2536		let mut producer = Info { sequence: 0 }.produce();
2537		assert!(matches!(
2538			producer.start_at(usize::MAX as u64),
2539			Err(Error::BoundsExceeded(_))
2540		));
2541	}
2542
2543	/// The per-frame size cap (the group byte budget) is enforced before allocating.
2544	#[test]
2545	fn create_frame_rejects_oversized() {
2546		let mut producer = Info { sequence: 0 }.produce();
2547		let result = producer.create_frame(frame::Info {
2548			size: MAX_CACHE_BYTES + 1,
2549			timestamp: Timestamp::ZERO,
2550		});
2551		assert!(matches!(result, Err(Error::FrameTooLarge)));
2552	}
2553}