Skip to main content

FramedRecvStream

Struct FramedRecvStream 

Source
pub struct FramedRecvStream {
    inner: RecvStream,
    buf: BytesMut,
    draft: DraftVersion,
    subgroup_io: Option<SubgroupObjectReader>,
    fetch_io: Option<FetchObjectReader>,
    tracking: Option<(TrackObjects, u64)>,
}
Expand description

A framed reader for a recv stream. Handles MoQT varint-length decoding.

Fields§

§inner: RecvStream§buf: BytesMut§draft: DraftVersion§subgroup_io: Option<SubgroupObjectReader>

Stateful subgroup object reader.

§fetch_io: Option<FetchObjectReader>

Stateful fetch object reader, seeded by FramedRecvStream::read_fetch_header. Draft-17 needs nothing the fetch header does not carry, so unlike draft-18 there is no separate call to start it.

§tracking: Option<(TrackObjects, u64)>

The record this stream’s objects are measured against, and the Group ID its header named.

One group for the whole stream: a subgroup header names it once and no object header repeats it. None on a stream that was never given one - a stream for an alias no live binding names, and every stream built outside Connection::accept_subgroup_stream.

Implementations§

Source§

impl FramedRecvStream

Source

pub fn new(inner: RecvStream, draft: DraftVersion) -> Self

Create a new framed receive stream for the given draft version.

Source

pub fn stream_id(&self) -> u64

Get the transport-level stream ID.

Source

fn measure_objects_against(&mut self, objects: TrackObjects, group: u64)

Measure this stream’s objects against objects, all of them in group.

Source

fn note_subgroup_object( &self, object: u64, status: Option<u64>, ) -> Result<(), ConnectionError>

Judge one object this stream carried against where its track ended.

Source

async fn fill(&mut self) -> Result<bool, ConnectionError>

Read more data from the stream into the internal buffer.

Source

async fn ensure(&mut self, n: usize) -> Result<(), ConnectionError>

Ensure at least n bytes are available in the buffer.

Source

pub async fn peek_stream_type(&mut self) -> Result<u64, ConnectionError>

Read this stream’s leading variable-length integer without consuming it, and return its value.

Every unidirectional MoQT stream on draft-17 opens with a varint naming what it is (Section 3.4): 0x05 for FETCH_HEADER, 0x10-0x1D for SUBGROUP_HEADER, and CONTROL_STREAM_TYPE for SETUP. Telling the peer’s control stream apart from a data stream means reading that varint, and taking it off the transport would destroy it: the control stream’s type varint is the SETUP message’s type field, so a stream whose type had been stripped would no longer decode as a SETUP.

Nothing is stripped. The bytes land in this reader’s own buffer, and every other method here — read_control, read_subgroup_header, read_fetch_header — decodes out of that buffer and advances it only on a successful decode. A stream this was called on is indistinguishable from one it was not, which is what makes it safe to peek a stream and then hand it to whichever reader the type turned out to call for.

It reads, so it can block: a peer that opens a stream and then writes nothing leaves this pending until a byte arrives or the stream ends.

§Errors

Whatever did arrive stays in the buffer in every case.

Source

pub async fn read_control( &mut self, capture_raw: bool, ) -> Result<(AnyControlMessage, Option<Vec<u8>>), ConnectionError>

Read a control message from the stream.

When capture_raw is true, the returned tuple includes a clone of the framed wire bytes (for observer emission). When false, the second element is None and the payload clone is skipped.

Source

pub async fn read_subgroup_header( &mut self, ) -> Result<AnySubgroupHeader, ConnectionError>

Read a subgroup stream header. Also initializes the internal delta-decoding state.

Source

pub async fn read_fetch_header( &mut self, ) -> Result<AnyFetchHeader, ConnectionError>

Read a fetch response header.

Source

pub async fn read_subgroup_object( &mut self, ) -> Result<SubgroupObject, ConnectionError>

Read the next draft-17 subgroup object from this stream using the stateful reader seeded by FramedRecvStream::read_subgroup_header.

Errors with ConnectionError::PropertiesOnNonNormalStatus on an Object that carries properties on a status other than Normal, which draft-17 Section 10.2.1.2 answers with a session close. The Object is consumed from the stream before the check, so the reader stays in step with the wire and a caller that reports the violation and reads on sees the following Object rather than a re-parse of this one.

Source

pub async fn read_fetch_stream_header( &mut self, ) -> Result<FetchHeader, ConnectionError>

Read the next draft-17 fetch header from this stream.

The typed twin of read_fetch_header, and it has to do the same two things that one does.

It seeds the object reader, because the two consume the same bytes: a version that left fetch_io unset would put the stream in a state no caller can leave, with the next read_fetch_object returning its fetch header not read yet refusal about a header this method has just read, and the bytes it would need already spent.

It fills before decoding, and treats a varint that ran out of buffer as a short read rather than a malformed header. FetchHeader::decode reports that as CodecError::VarInt(VarIntError::UnexpectedEnd) where AnyFetchHeader reports a bare CodecError::UnexpectedEnd, so a loop matching only the latter never reaches its own fill — and since the buffer starts empty, that is every first call on a fresh stream.

Source

pub fn stop(&mut self, code: u64) -> Result<(), ConnectionError>

Stop accepting data on this stream with code as the STOP_SENDING application error code, discarding anything unread.

Dropping a receive stream also stops it, but with a hard-coded 0. See RecvStream::stop.

Source

pub async fn received_reset(&mut self) -> Result<Option<u64>, ConnectionError>

Wait for the peer to reset this stream, consuming nothing.

See RecvStream::received_reset for what Ok(None) means and why a caller must not re-poll after it.

Source

pub async fn read_fetch_object( &mut self, ) -> Result<(FetchObject, Vec<u8>), ConnectionError>

Read the next draft-17 fetch object and its payload.

The mirror of FramedSendStream::write_fetch_object. What comes back is a FetchObject rather than a header: draft-17’s reader resolves the fields the Serialization Flags left off the wire, so the Group ID, Subgroup ID, Object ID and Priority it carries are the frame’s own values, not the flags that say where to find them. object.header is the frame exactly as it arrived, for a caller forwarding the bytes on.

The payload comes back with it for the reason the codec leaves it on the wire: payload_length says how many bytes follow, and a reader that takes the wrong number of them desynchronises every later object on the stream. Doing it here is the only place that count and the buffer are both in hand.

§Errors

ConnectionError::DataStreamState when FramedRecvStream::read_fetch_header has not been read first, since that is what seeds the reader; ConnectionError::UnexpectedEnd when the stream ends inside the header or inside the payload it declared; and ConnectionError::Codec on a Serialization Flags value the draft does not define, and on every rule Section 10.4.4.1 answers with a session close — a first Object that inherits a field no Object before it established, or an inherited Object ID or Subgroup ID one past the end of the space.

Source

async fn read_object_payload( &mut self, length: &VarInt, ) -> Result<Vec<u8>, ConnectionError>

Take the length payload bytes that follow a fetch object’s header.

Separate from the header read because the header is decoded from a probe cursor that may have to be retried after a fill, and the payload is a flat byte count that never is.

Source

pub fn draft(&self) -> DraftVersion

Returns the draft version this stream is framed for.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

§

impl<T> Instrument for T

§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided [Span], returning an Instrumented wrapper. Read more
§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
§

impl<T> WithSubscriber for T

§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a [WithDispatch] wrapper. Read more
§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a [WithDispatch] wrapper. Read more