Skip to main content

moq_net/model/
frame.rs

1//! Frames are the leaf of the model: a sized, timestamped payload within a group.
2//!
3//! A group is a single ordered stream, so at most one frame is ever in flight.
4//! Completed frames are plain data ([`Frame`]); the in-flight frame is written
5//! through [`Producer`], which borrows its parent [`group::Producer`] exclusively so
6//! the borrow checker enforces that only one frame is open at a time. A [`Consumer`]
7//! reads one frame, sharing the group's channel rather than a per-frame one.
8use std::ops::Range;
9use std::sync::Arc;
10use std::sync::OnceLock;
11use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
12use std::task::{Poll, ready};
13
14use arrayvec::ArrayVec;
15use bytes::Bytes;
16
17use crate::group::{self, GroupState};
18use crate::{Error, IntoBytes, Result, Timestamp, stats};
19
20/// A chunk of data with an upfront size and a presentation timestamp.
21///
22/// This is just the header; the payload is carried separately (as a completed
23/// [`Frame`] or streamed via [`Producer`] / [`Consumer`]).
24#[derive(Clone, Copy, Debug)]
25pub struct Info {
26	/// Total payload size in bytes. Declared up front so consumers can preallocate.
27	pub size: u64,
28	/// Presentation timestamp.
29	///
30	/// [`group::Producer::create_frame`] converts it into the parent track's
31	/// timescale, so the scale you build it with doesn't have to match the track.
32	/// For data without a presentation time, pass [`Timestamp::now`] explicitly.
33	pub timestamp: Timestamp,
34}
35
36/// A completed frame: a timestamp and its full, contiguous payload.
37///
38/// This is the stored form of every finished frame in a group. The payload is a
39/// single [`Bytes`], so a consumer gets it with one zero-copy slice.
40#[derive(Clone, Debug)]
41pub struct Frame {
42	/// Presentation timestamp, at the parent track's timescale.
43	pub timestamp: Timestamp,
44	/// The full frame payload.
45	pub payload: Bytes,
46}
47
48/// A reusable batch of frames, filled by [`group::Consumer::read_frames`] and drained
49/// by [`group::Producer::write_frames`].
50///
51/// A fixed-capacity inline buffer: `N` frames of stack storage, never a heap
52/// allocation and never a spill. Allocate one per task and reuse it for the life of
53/// the group rather than one per read.
54///
55/// `N` defaults to 8. Most reads are not big batches: at the live edge a frame arrives
56/// at a time, so the capacity past the first frame or two is idle stack, and a
57/// publisher holds one buffer per in-flight group. 8 costs 384 bytes and still reads
58/// ~5x faster than a frame at a time (`benches/group.rs`). Ask for a larger `N` when
59/// you know you are draining a backlog: 32 is ~8x, for 1.5 KB.
60///
61/// A default on a const parameter only applies in type position, so name the type to
62/// get it: `let buf: frame::Buffer = Buffer::new()`.
63///
64/// A fill stamps the group's cache access once for the whole batch, which bounds
65/// frames rather than elapsed time. A reader that may take longer than the track's
66/// `max_age` to work through one batch calls
67/// [`group::Consumer::keep_alive`] between frames, or the rest of the group is
68/// expired out from under it.
69#[derive(Debug, Default)]
70pub struct Buffer<const N: usize = 8>(ArrayVec<Frame, N>);
71
72impl<const N: usize> Buffer<N> {
73	/// An empty buffer with room for `N` frames.
74	pub fn new() -> Self {
75		Self(ArrayVec::new())
76	}
77
78	/// How many frames a single fill can hold.
79	pub const fn capacity(&self) -> usize {
80		N
81	}
82
83	/// How many frames the buffer currently holds.
84	pub fn len(&self) -> usize {
85		self.0.len()
86	}
87
88	/// Whether the buffer holds no frames.
89	pub fn is_empty(&self) -> bool {
90		self.0.is_empty()
91	}
92
93	/// Whether the buffer is at capacity, so [`Self::push`] would refuse.
94	pub fn is_full(&self) -> bool {
95		self.0.is_full()
96	}
97
98	/// The frames from the most recent fill, in order.
99	pub fn filled(&self) -> &[Frame] {
100		&self.0
101	}
102
103	/// The frames from the most recent fill, mutably (to take payloads out, say).
104	pub fn filled_mut(&mut self) -> &mut [Frame] {
105		&mut self.0
106	}
107
108	/// Append a frame, handing it back if the buffer is already full.
109	///
110	/// Fill a buffer this way to hand a whole batch to
111	/// [`group::Producer::write_frames`].
112	pub fn push(&mut self, frame: Frame) -> std::result::Result<(), Frame> {
113		self.0.try_push(frame).map_err(|err| err.element())
114	}
115
116	/// Move every frame out, leaving the buffer empty.
117	///
118	/// Frames left in the iterator when it drops are dropped with it, so the buffer
119	/// ends up empty either way.
120	pub fn drain(&mut self) -> impl ExactSizeIterator<Item = Frame> + '_ {
121		self.0.drain(..)
122	}
123
124	/// Drop the current batch, leaving the buffer empty.
125	pub fn clear(&mut self) {
126		self.0.clear();
127	}
128}
129
130/// Payload storage for the single in-flight frame, shared between the writing
131/// [`Producer`] and any streaming [`Consumer`]s.
132///
133/// A whole-frame [`Bytes`] write is stored directly. Chunked writes fall back to one
134/// mutable heap allocation sized to the declared frame. The producer writes through
135/// the raw pointer (sole writer, guaranteed by the exclusive borrow of the parent
136/// group); `written` provides happens-before for cross-thread reads. Implements
137/// [AsRef]<[u8]> so it can back a [`Bytes::from_owner`].
138#[derive(Clone)]
139pub(crate) struct FrameBuf(Arc<FrameBufInner>);
140
141struct FrameBufInner {
142	capacity: usize,
143	written: AtomicUsize,
144	storage: OnceLock<FrameStorage>,
145}
146
147enum FrameStorage {
148	Shared(Bytes),
149	Mutable(MutableFrameBuf),
150}
151
152struct MutableFrameBuf {
153	// Owned heap allocation of `capacity` bytes (zero-initialized).
154	data: *mut u8,
155	capacity: usize,
156}
157
158// Safety: `data` is owned (Box-allocated, freed in Drop). The producer is the sole
159// writer and consumers only read bytes `< written`.
160unsafe impl Send for MutableFrameBuf {}
161unsafe impl Sync for MutableFrameBuf {}
162
163impl Drop for MutableFrameBuf {
164	fn drop(&mut self) {
165		// Safety: data was obtained from `Box::into_raw` of a `Box<[u8]>` of length
166		// `capacity` and is not aliased at drop (Arc refcount hit 0).
167		unsafe {
168			let slice = std::ptr::slice_from_raw_parts_mut(self.data, self.capacity);
169			drop(Box::from_raw(slice));
170		}
171	}
172}
173
174impl MutableFrameBuf {
175	fn new(size: usize) -> Self {
176		let boxed: Box<[u8]> = vec![0u8; size].into_boxed_slice();
177		let capacity = boxed.len();
178		let data = Box::into_raw(boxed) as *mut u8;
179		Self { data, capacity }
180	}
181}
182
183impl FrameBuf {
184	/// Allocate a buffer for a frame of `size` bytes.
185	///
186	/// The oversized-frame guard lives in [`group::Producer`], which rejects a declared
187	/// size larger than the group's byte budget before calling this.
188	pub(crate) fn new(size: usize) -> Self {
189		Self(Arc::new(FrameBufInner {
190			capacity: size,
191			written: AtomicUsize::new(0),
192			storage: OnceLock::new(),
193		}))
194	}
195
196	pub(crate) fn capacity(&self) -> usize {
197		self.0.capacity
198	}
199
200	pub(crate) fn written(&self, ord: Ordering) -> usize {
201		self.0.written.load(ord)
202	}
203
204	fn try_set_bytes(&self, bytes: Bytes) -> std::result::Result<(), Bytes> {
205		if bytes.len() != self.capacity() || self.written(Ordering::Acquire) != 0 {
206			return Err(bytes);
207		}
208		self.0
209			.storage
210			.set(FrameStorage::Shared(bytes))
211			.map_err(|storage| match storage {
212				FrameStorage::Shared(bytes) => bytes,
213				FrameStorage::Mutable(_) => unreachable!("try_set_bytes only installs shared storage"),
214			})
215	}
216
217	/// The mutable buffer for multi-chunk writes, lazily allocated.
218	///
219	/// Returns `None` once a whole-frame write has installed shared storage.
220	fn mutable(&self) -> Option<&MutableFrameBuf> {
221		match self
222			.0
223			.storage
224			.get_or_init(|| FrameStorage::Mutable(MutableFrameBuf::new(self.capacity())))
225		{
226			FrameStorage::Shared(_) => None,
227			FrameStorage::Mutable(buf) => Some(buf),
228		}
229	}
230
231	/// Safety: caller must be the sole producer and `new_written` must be `<= capacity`.
232	unsafe fn store_written(&self, new_written: usize) {
233		// Release pairs with consumers' Acquire load to publish prior writes.
234		self.0.written.store(new_written, Ordering::Release);
235	}
236
237	/// Append `src` at the current write offset and publish it.
238	///
239	/// Safety relies on the single-producer invariant: only one [`Producer`] exists for
240	/// a frame (it holds the exclusive borrow of the parent group), so this is the sole
241	/// writer even though it takes `&self`.
242	fn append(&self, src: &[u8]) {
243		if src.is_empty() {
244			return;
245		}
246		let prev = self.written(Ordering::Relaxed);
247		let Some(buf) = self.mutable() else {
248			// Only reachable if the frame is already complete via shared storage, which
249			// `Producer::write` rejects for a non-empty chunk. Nothing to copy.
250			return;
251		};
252		// Safety: sole writer; the caller bounds-checked `src` against the remaining
253		// capacity, and consumers only read `[..written]`.
254		unsafe {
255			std::ptr::copy_nonoverlapping(src.as_ptr(), buf.data.add(prev), src.len());
256			self.store_written(prev + src.len());
257		}
258	}
259
260	/// Freeze the buffer into the completed payload (`size` bytes).
261	///
262	/// Returns the shared [`Bytes`] directly for a whole-frame write (zero-copy), or
263	/// wraps the mutable allocation otherwise.
264	fn freeze(&self, size: usize) -> Bytes {
265		match self.0.storage.get() {
266			Some(FrameStorage::Shared(bytes)) => bytes.clone(),
267			_ => self.slice(0, size),
268		}
269	}
270
271	/// A zero-copy slice of the initialized region `[start..end]`.
272	fn slice(&self, start: usize, end: usize) -> Bytes {
273		Bytes::from_owner(self.clone()).slice(start..end)
274	}
275}
276
277impl AsRef<[u8]> for FrameBuf {
278	fn as_ref(&self) -> &[u8] {
279		// Snapshot the initialized region (bytes the producer has written so far).
280		// Acquire pairs with the producer's Release on `written`.
281		let written = self.0.written.load(Ordering::Acquire);
282		match self.0.storage.get() {
283			Some(FrameStorage::Shared(bytes)) => &bytes[..written],
284			Some(FrameStorage::Mutable(buf)) => {
285				// Safety: data..data+written is initialized (zero-init at alloc + producer
286				// writes up to `written`). The Arc keeps the allocation alive while any
287				// reference to the slice lives.
288				unsafe { std::slice::from_raw_parts(buf.data, written) }
289			}
290			None => &[],
291		}
292	}
293}
294
295/// The writer behind [`Producer`] and [`ProducerOwned`], generic over how it
296/// reaches the parent group (an exclusive borrow, or an owned clone).
297struct Raw<G: std::borrow::BorrowMut<group::Producer>> {
298	group: G,
299	buf: FrameBuf,
300	info: Info,
301	// Set once the frame is committed (finished) or aborted, so Drop is a no-op.
302	done: bool,
303	// Ingress payload meter, inherited from the parent group. Counts each written
304	// chunk's bytes. Empty (no-op) for an untagged group.
305	stats: stats::Meter,
306}
307
308impl<G: std::borrow::BorrowMut<group::Producer>> Raw<G> {
309	fn remaining(&self) -> usize {
310		self.buf.capacity() - self.buf.written(Ordering::Acquire)
311	}
312
313	fn write<B: IntoBytes>(&mut self, chunk: B) -> Result<()> {
314		let len = chunk.as_ref().len();
315		if len > self.remaining() {
316			return Err(Error::WrongSize);
317		}
318		// Ingress payload: count the chunk's bytes as they're written.
319		self.stats.bytes(len as u64);
320		// Fast path: a single whole-frame write keeps the caller's allocation.
321		if len == self.buf.capacity() && self.buf.written(Ordering::Acquire) == 0 {
322			match self.buf.try_set_bytes(chunk.into_bytes()) {
323				Ok(()) => {
324					let cap = self.buf.capacity();
325					// Safety: `try_set_bytes` checked the buffer exactly matches the declared
326					// size, so publishing all bytes is in bounds.
327					unsafe { self.buf.store_written(cap) };
328				}
329				Err(chunk) => self.buf.append(&chunk),
330			}
331		} else {
332			self.buf.append(chunk.as_ref());
333		}
334		Ok(())
335	}
336
337	fn finish(&mut self) -> Result<()> {
338		if self.buf.written(Ordering::Acquire) != self.buf.capacity() {
339			return Err(Error::WrongSize);
340		}
341		let payload = self.buf.freeze(self.buf.capacity());
342		self.group.borrow_mut().frame_commit(Frame {
343			timestamp: self.info.timestamp,
344			payload,
345		})?;
346		self.done = true;
347		Ok(())
348	}
349
350	fn abort(&mut self, err: Error) -> Result<()> {
351		self.group.borrow_mut().frame_abort(err);
352		self.done = true;
353		Ok(())
354	}
355}
356
357impl<G: std::borrow::BorrowMut<group::Producer>> Drop for Raw<G> {
358	fn drop(&mut self) {
359		if !self.done {
360			// An unfinished frame leaves the group stream broken; fail the group so
361			// consumers surface an error instead of hanging on the partial forever.
362			tracing::warn!(
363				group = self.group.borrow_mut().info().sequence,
364				"frame::Producer dropped before writing all bytes"
365			);
366			self.group.borrow_mut().frame_abort(Error::Dropped);
367		}
368	}
369}
370
371/// Writes the payload of the single in-flight frame in one or more chunks.
372///
373/// Borrows the parent [`group::Producer`] exclusively, so no other frame can be
374/// opened while this one is live. The total bytes written must exactly match
375/// [`Info::size`]; call [`Self::finish`] to commit the frame (or [`Self::abort`] to
376/// fail it). Dropping without either aborts the group, since an unfinished frame
377/// leaves the group's stream broken.
378///
379/// A single whole-frame [`write`](Self::write) keeps the caller's allocation
380/// (zero-copy); chunked writes copy into one buffer sized to the declared frame.
381pub struct Producer<'a>(Raw<&'a mut group::Producer>);
382
383impl std::ops::Deref for Producer<'_> {
384	type Target = Info;
385
386	fn deref(&self) -> &Self::Target {
387		&self.0.info
388	}
389}
390
391impl<'a> Producer<'a> {
392	pub(crate) fn new(group: &'a mut group::Producer, buf: FrameBuf, info: Info) -> Self {
393		Self(Raw {
394			group,
395			buf,
396			info,
397			done: false,
398			stats: stats::Meter::default(),
399		})
400	}
401
402	/// Attach the parent group's ingress meter, so written chunks bump `bytes`.
403	pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
404		self.0.stats = meter;
405		self
406	}
407
408	/// The parent group this frame belongs to.
409	pub fn group(&self) -> group::Info {
410		self.0.group.info()
411	}
412
413	/// Bytes still needed to complete the frame.
414	pub fn remaining(&self) -> usize {
415		self.0.remaining()
416	}
417
418	/// Write a chunk of data to the frame.
419	///
420	/// Returns [`Error::WrongSize`] if the chunk would exceed the remaining bytes.
421	pub fn write<B: IntoBytes>(&mut self, chunk: B) -> Result<()> {
422		self.0.write(chunk)?;
423		self.0.group.frame_notify();
424		Ok(())
425	}
426
427	/// Commit the frame, verifying that all bytes were written.
428	///
429	/// Returns [`Error::WrongSize`] if the bytes written don't match [`Info::size`].
430	pub fn finish(mut self) -> Result<()> {
431		self.0.finish()
432	}
433
434	/// Abort the frame (and its group) with the given error.
435	pub fn abort(mut self, err: Error) -> Result<()> {
436		self.0.abort(err)
437	}
438}
439
440/// The owned counterpart of [`Producer`], for the wire drivers that stream a
441/// frame across polls and cannot hold the group borrowed inside their state.
442///
443/// Crate-private on purpose: the exclusivity the public borrow enforces (one
444/// live frame per group) becomes the holder's promise here. Do not open another
445/// frame on the group until this one is finished or aborted.
446pub(crate) struct ProducerOwned(Raw<group::Producer>);
447
448impl std::ops::Deref for ProducerOwned {
449	type Target = Info;
450
451	fn deref(&self) -> &Self::Target {
452		&self.0.info
453	}
454}
455
456impl ProducerOwned {
457	pub(crate) fn new(group: group::Producer, buf: FrameBuf, info: Info) -> Self {
458		Self(Raw {
459			group,
460			buf,
461			info,
462			done: false,
463			stats: stats::Meter::default(),
464		})
465	}
466
467	/// Attach the parent group's ingress meter, so written chunks bump `bytes`.
468	pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
469		self.0.stats = meter;
470		self
471	}
472
473	/// Bytes still needed to complete the frame.
474	pub fn remaining(&self) -> usize {
475		self.0.remaining()
476	}
477
478	/// Write a chunk of payload *without* waking consumers; pair it with [`Self::notify`].
479	///
480	/// The wake is split out because the wire ingest drains every chunk the transport
481	/// has already buffered in one poll turn, and a consumer parked on the group cannot
482	/// run until that turn yields. Waking per chunk pays a group lock and a clock read
483	/// to publish bytes nobody can observe yet.
484	///
485	/// `coding::Reader::poll_read_frame` owns the pairing and is the only caller.
486	pub(crate) fn write<B: IntoBytes>(&mut self, chunk: B) -> Result<()> {
487		self.0.write(chunk)
488	}
489
490	/// Publish what has been written so far, waking consumers parked on the group.
491	pub(crate) fn notify(&self) {
492		self.0.group.frame_notify();
493	}
494
495	/// Commit the frame, verifying that all bytes were written.
496	pub fn finish(mut self) -> Result<()> {
497		self.0.finish()
498	}
499
500	/// Abort the frame (and its group) with the given error.
501	pub fn abort(mut self, err: Error) -> Result<()> {
502		self.0.abort(err)
503	}
504}
505
506/// The source of a [`Consumer`]'s payload: a finished frame (whole) or the in-flight
507/// tail (streamed).
508#[derive(Clone)]
509pub(crate) enum Source {
510	Complete(Bytes),
511	Partial(FrameBuf),
512}
513
514/// Subscriber expiry state carried across the group-to-frame handoff.
515#[derive(Clone)]
516pub(crate) struct Expiry {
517	policy: Arc<dyn group::Expiry>,
518	stale_stats: stats::Meter,
519	stale_counted: Arc<AtomicBool>,
520	tail: Range<usize>,
521	count_payload: bool,
522}
523
524impl Expiry {
525	pub(crate) fn new(
526		policy: Arc<dyn group::Expiry>,
527		stale_stats: stats::Meter,
528		stale_counted: Arc<AtomicBool>,
529	) -> Self {
530		Self {
531			policy,
532			stale_stats,
533			stale_counted,
534			tail: 0..0,
535			count_payload: false,
536		}
537	}
538
539	pub(crate) fn for_frame(mut self, tail: Range<usize>, count_payload: bool) -> Self {
540		self.tail = tail;
541		self.count_payload = count_payload;
542		self
543	}
544}
545
546/// Reads one frame's payload, streaming as bytes arrive for the in-flight tail.
547///
548/// Owns a handle to the parent group's channel (not a per-frame one), so a group with
549/// many frames doesn't allocate a channel per frame. Cloning yields an independent
550/// reader of the same frame.
551#[derive(Clone)]
552pub struct Consumer {
553	// The group's channel, used to park while a partial frame fills.
554	state: kio::Consumer<GroupState>,
555	info: Info,
556	source: Source,
557	// Byte offset consumed so far.
558	read_idx: usize,
559	// Egress payload meter, so chunks bump `bytes` exactly once as they're read out.
560	// Empty (no-op) for an untagged group.
561	stats: stats::Meter,
562	// The parent subscription can expire after this frame handle is returned.
563	expiry: Option<Expiry>,
564	expired: bool,
565}
566
567impl std::ops::Deref for Consumer {
568	type Target = Info;
569
570	fn deref(&self) -> &Self::Target {
571		&self.info
572	}
573}
574
575impl Consumer {
576	pub(crate) fn new(state: kio::Consumer<GroupState>, info: Info, source: Source) -> Self {
577		Self {
578			state,
579			info,
580			source,
581			read_idx: 0,
582			stats: stats::Meter::default(),
583			expiry: None,
584			expired: false,
585		}
586	}
587
588	/// Attach an egress meter so read-out chunks bump `bytes`. Used only for frames
589	/// read directly from the group (whose bytes weren't counted at a batch fill).
590	pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
591		self.stats = meter;
592		self
593	}
594
595	pub(crate) fn with_expiry(mut self, expiry: Expiry) -> Self {
596		self.expiry = Some(expiry);
597		self
598	}
599
600	fn size(&self) -> usize {
601		match &self.source {
602			Source::Complete(bytes) => bytes.len(),
603			Source::Partial(_) => self.info.size as usize,
604		}
605	}
606
607	/// Evaluate the parent subscription's drift budget.
608	///
609	/// Only called once a read has nothing buffered to return: the budget bounds a
610	/// *stalled* payload, so bytes already in hand are always drained rather than
611	/// truncated. `self.expired` is sticky, so the answer is only ever computed once.
612	fn poll_expired(&mut self, waiter: &kio::Waiter) -> bool {
613		if self.expired || self.read_idx >= self.size() {
614			return self.expired;
615		}
616		let Some(expiry) = &self.expiry else {
617			return false;
618		};
619		if !expiry.policy.is_expired(waiter) {
620			return false;
621		}
622
623		self.expired = true;
624		if !expiry.stale_counted.swap(true, Ordering::Relaxed) {
625			let mut stale = self.state.read().content_range(expiry.tail.start, expiry.tail.end);
626			if expiry.count_payload {
627				stale.bytes += self.size().saturating_sub(self.read_idx) as u64;
628			}
629			expiry.stale_stats.stale(stale);
630		}
631		true
632	}
633
634	/// Poll for the next chunk of bytes since the last read.
635	///
636	/// Returns `None` once the frame is finished and all bytes have been consumed.
637	pub fn poll_read_chunk(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Bytes>>> {
638		if self.expired {
639			return Poll::Ready(Err(Error::Old));
640		}
641		let buf = match &self.source {
642			Source::Complete(bytes) => {
643				if self.read_idx >= bytes.len() {
644					return Poll::Ready(Ok(None));
645				}
646				let out = bytes.slice(self.read_idx..);
647				self.read_idx = bytes.len();
648				self.stats.bytes(out.len() as u64);
649				return Poll::Ready(Ok(Some(out)));
650			}
651			Source::Partial(buf) => buf.clone(),
652		};
653
654		let size = self.info.size as usize;
655		loop {
656			let written = buf.written(Ordering::Acquire);
657			if written > self.read_idx {
658				let out = buf.slice(self.read_idx, written);
659				self.read_idx = written;
660				self.stats.bytes(out.len() as u64);
661				return Poll::Ready(Ok(Some(out)));
662			}
663			if written >= size {
664				return Poll::Ready(Ok(None));
665			}
666			// Nothing buffered and the frame isn't finished: this park is the stall the
667			// drift budget exists to bound, and the only place it can apply.
668			if self.poll_expired(waiter) {
669				return Poll::Ready(Err(Error::Old));
670			}
671			let read_idx = self.read_idx;
672			// Park on the group's channel; the producer notifies it on each write and
673			// on abort. Re-check the atomic on wake.
674			ready!(poll_state(&self.state, waiter, |state| {
675				if let Some(err) = &state.abort {
676					return Poll::Ready(Err(err.clone()));
677				}
678				let w = buf.written(Ordering::Acquire);
679				if w > read_idx || w >= size {
680					Poll::Ready(Ok(()))
681				} else {
682					Poll::Pending
683				}
684			})?);
685		}
686	}
687
688	/// Return the next chunk of bytes since the last read.
689	pub async fn read_chunk(&mut self) -> Result<Option<Bytes>> {
690		kio::wait(|waiter| self.poll_read_chunk(waiter)).await
691	}
692
693	/// Poll for all remaining bytes, resolving once the frame is finished.
694	pub fn poll_read_all(&mut self, waiter: &kio::Waiter) -> Poll<Result<Bytes>> {
695		if self.expired {
696			return Poll::Ready(Err(Error::Old));
697		}
698		let buf = match &self.source {
699			Source::Complete(bytes) => {
700				let out = bytes.slice(self.read_idx..);
701				self.read_idx = bytes.len();
702				self.stats.bytes(out.len() as u64);
703				return Poll::Ready(Ok(out));
704			}
705			Source::Partial(buf) => buf.clone(),
706		};
707
708		let size = self.info.size as usize;
709		let read_idx = self.read_idx;
710		// Waiting on the rest of the payload is a stall; see `poll_read_chunk`.
711		if buf.written(Ordering::Acquire) < size && self.poll_expired(waiter) {
712			return Poll::Ready(Err(Error::Old));
713		}
714		ready!(poll_state(&self.state, waiter, |state| {
715			if let Some(err) = &state.abort {
716				return Poll::Ready(Err(err.clone()));
717			}
718			if buf.written(Ordering::Acquire) >= size {
719				Poll::Ready(Ok(()))
720			} else {
721				Poll::Pending
722			}
723		})?);
724		let out = buf.slice(read_idx, size);
725		self.read_idx = size;
726		self.stats.bytes(out.len() as u64);
727		Poll::Ready(Ok(out))
728	}
729
730	/// Return all remaining bytes, blocking until the frame is finished.
731	pub async fn read_all(&mut self) -> Result<Bytes> {
732		kio::wait(|waiter| self.poll_read_all(waiter)).await
733	}
734}
735
736/// Poll the group channel, mapping a terminal close without an error to
737/// [`Error::Dropped`]. Mirrors [`group::Consumer`]'s internal helper.
738fn poll_state<F, R>(state: &kio::Consumer<GroupState>, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
739where
740	F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
741{
742	Poll::Ready(match ready!(state.poll(waiter, f)) {
743		Ok(res) => res,
744		Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
745	})
746}