Skip to main content

Connection

Struct Connection 

Source
pub struct Connection {
    transport: Transport,
    endpoint: Endpoint,
    control_send: Option<FramedSendStream>,
    control_recv: Option<FramedRecvStream>,
    observer: Option<Box<dyn ConnectionObserver>>,
    pending_events: Vec<ClientEvent>,
}
Expand description

A live draft-07 MoQT connection over QUIC or WebTransport.

Fields§

§transport: Transport§endpoint: Endpoint§control_send: Option<FramedSendStream>§control_recv: Option<FramedRecvStream>§observer: Option<Box<dyn ConnectionObserver>>§pending_events: Vec<ClientEvent>

Setup events buffered during connect() and replayed when an observer attaches via set_observer — without this, an observer attached after connect returns would never see the handshake.

Implementations§

Source§

impl Connection

Source

pub async fn connect( addr: &str, config: ClientConfig, ) -> Result<Self, ConnectionError>

Connect to a draft-07 MoQT server as a client.

Source

pub async fn adopt( transport: Transport, config: ClientConfig, ) -> Result<Self, ConnectionError>

Run the MoQT setup handshake over a transport somebody else established.

For choosing the draft from what the server selected: dial once through crate::transport::dial_quic offering every ALPN, then bring the connection to the module its answer names. Self::connect cannot do this — it derives its single ALPN from the draft it was given.

config.draft must match this module. The transport is adopted as given; nothing here re-checks the ALPN it was negotiated with.

Source

async fn connect_quic( addr: &str, config: &ClientConfig, ) -> Result<Transport, ConnectionError>

Establish a raw QUIC connection.

Offers this draft’s ALPN alone; crate::transport::dial_quic holds the TLS and endpoint setup.

Source

async fn connect_webtransport( _url: &str, _config: &ClientConfig, ) -> Result<Transport, ConnectionError>

Stub for when the webtransport feature is not enabled.

Source

pub fn set_observer(&mut self, observer: Box<dyn ConnectionObserver>)

Attach an observer. Buffered handshake events from connect() are flushed in arrival order before this returns.

Source

pub fn clear_observer(&mut self)

Remove the observer.

Source

fn emit(&self, event: ClientEvent)

Emit an event to the observer, if one is attached.

Source

pub async fn send_control( &mut self, msg: &ControlMessage, ) -> Result<(), ConnectionError>

Send a control message on the control stream.

Source

pub async fn recv_control(&mut self) -> Result<ControlMessage, ConnectionError>

Read the next control message from the control stream.

Source

pub async fn recv_and_dispatch( &mut self, ) -> Result<ControlMessage, ConnectionError>

Read and dispatch the next incoming control message through the endpoint state machine. Returns the decoded message for inspection.

Source

pub async fn subscribe( &mut self, track_alias: VarInt, track_namespace: TrackNamespace, track_name: Vec<u8>, subscriber_priority: u8, group_order: GroupOrder, filter_type: FilterType, ) -> Result<VarInt, ConnectionError>

Send a SUBSCRIBE and return the allocated subscribe ID.

Source

pub async fn subscribe_range( &mut self, track_alias: VarInt, track_namespace: TrackNamespace, track_name: Vec<u8>, subscriber_priority: u8, group_order: GroupOrder, start_location: Location, end_location: Option<Location>, ) -> Result<VarInt, ConnectionError>

Send a SUBSCRIBE for a range of the track and return the allocated ID.

The Filter Type comes from the arguments, so the message cannot name a filter whose fields it does not carry.

Source

pub async fn unsubscribe( &mut self, subscribe_id: VarInt, ) -> Result<(), ConnectionError>

Send an UNSUBSCRIBE for the given subscribe ID.

Source

pub async fn subscribe_ok( &mut self, subscribe_id: VarInt, expires: VarInt, group_order: GroupOrder, parameters: Vec<KeyValuePair>, ) -> Result<(), ConnectionError>

Accept a subscription the peer opened, sending SUBSCRIBE_OK.

Source

pub async fn subscribe_error( &mut self, subscribe_id: VarInt, error_code: VarInt, reason_phrase: Vec<u8>, track_alias: VarInt, ) -> Result<(), ConnectionError>

Reject a subscription the peer opened, sending SUBSCRIBE_ERROR.

The Track Alias travels back with the refusal: under the ‘Retry Track Alias’ code it is the alias the peer should try again with, and under any other code it is ignored.

Source

pub async fn subscribe_done( &mut self, subscribe_id: VarInt, status_code: VarInt, reason_phrase: Vec<u8>, ) -> Result<(), ConnectionError>

End a subscription this endpoint accepted, sending SUBSCRIBE_DONE.

Source

pub async fn fetch( &mut self, track_namespace: TrackNamespace, track_name: Vec<u8>, subscriber_priority: u8, group_order: GroupOrder, start_group: VarInt, start_object: VarInt, end_group: VarInt, end_object: VarInt, ) -> Result<VarInt, ConnectionError>

Send a FETCH and return the allocated subscribe ID.

Source

pub async fn fetch_cancel( &mut self, subscribe_id: VarInt, ) -> Result<(), ConnectionError>

Send a FETCH_CANCEL for the given subscribe ID.

Source

pub async fn fetch_ok( &mut self, subscribe_id: VarInt, group_order: GroupOrder, end_of_track: u8, largest_group_id: Option<VarInt>, largest_object_id: Option<VarInt>, parameters: Vec<KeyValuePair>, ) -> Result<(), ConnectionError>

Accept a fetch the peer opened, sending FETCH_OK.

The endpoint refuses a Joining Fetch naming a subscription this session cannot join and refuses a second answer to one FETCH, so nothing is written on the wire when it does either.

Source

pub async fn fetch_error( &mut self, subscribe_id: VarInt, error_code: VarInt, reason_phrase: Vec<u8>, ) -> Result<(), ConnectionError>

Refuse a fetch the peer opened, sending FETCH_ERROR.

The endpoint refuses a second answer to one FETCH, so nothing is written on the wire when it does.

Source

pub async fn subscribe_announces( &mut self, track_namespace_prefix: TrackNamespace, ) -> Result<(), ConnectionError>

Send a SUBSCRIBE_ANNOUNCES.

Source

pub async fn subscribe_announces_ok( &mut self, track_namespace_prefix: TrackNamespace, ) -> Result<(), ConnectionError>

Accept a namespace subscription the peer made, sending SUBSCRIBE_ANNOUNCES_OK.

The endpoint refuses a second answer to one SUBSCRIBE_ANNOUNCES, so nothing is written on the wire when it does.

Source

pub async fn subscribe_announces_error( &mut self, track_namespace_prefix: TrackNamespace, error_code: VarInt, reason_phrase: Vec<u8>, ) -> Result<(), ConnectionError>

Refuse a namespace subscription the peer made, sending SUBSCRIBE_ANNOUNCES_ERROR.

The other half of the same sentence: one answer, and this is the other one it can be.

Source

pub async fn announce( &mut self, track_namespace: TrackNamespace, ) -> Result<(), ConnectionError>

Send an ANNOUNCE.

Source

pub async fn unannounce( &mut self, track_namespace: TrackNamespace, ) -> Result<(), ConnectionError>

Send an UNANNOUNCE.

Source

pub async fn announce_ok( &mut self, track_namespace: TrackNamespace, ) -> Result<(), ConnectionError>

Accept an announcement the peer made, sending ANNOUNCE_OK.

The endpoint refuses a second answer to one ANNOUNCE, so nothing is written on the wire when it does.

Source

pub async fn announce_error( &mut self, track_namespace: TrackNamespace, error_code: VarInt, reason_phrase: Vec<u8>, ) -> Result<(), ConnectionError>

Refuse an announcement the peer made, sending ANNOUNCE_ERROR.

The other half of the same sentence: one answer, and this is the other one it can be.

Source

pub async fn announce_cancel( &mut self, track_namespace: TrackNamespace, error_code: VarInt, reason_phrase: Vec<u8>, ) -> Result<(), ConnectionError>

Revoke an acceptance, sending ANNOUNCE_CANCEL.

The endpoint refuses one for an announcement it never accepted, so nothing is written on the wire when it does.

Source

pub async fn track_status_request( &mut self, track_namespace: TrackNamespace, track_name: Vec<u8>, ) -> Result<(), ConnectionError>

Send a TRACK_STATUS_REQUEST.

Source

pub async fn track_status( &mut self, track_namespace: TrackNamespace, track_name: Vec<u8>, status_code: VarInt, last_group_id: VarInt, last_object_id: VarInt, ) -> Result<(), ConnectionError>

Answer a TRACK_STATUS_REQUEST the peer sent, sending TRACK_STATUS.

The endpoint refuses a second answer to one request, so nothing is written on the wire when it does.

Source

pub async fn open_subgroup_stream( &self, header: &AnySubgroupHeader, ) -> Result<FramedSendStream, ConnectionError>

Open a new unidirectional stream for sending subgroup data.

Source

pub async fn open_fetch_stream( &self, header: &AnyFetchHeader, ) -> Result<FramedSendStream, ConnectionError>

Open a new unidirectional stream for sending a FETCH’s objects.

The objects answering a FETCH do not go on the request’s own stream: they go on a unidirectional stream of their own, which opens with a FETCH_HEADER naming the request they belong to. This writes that header and hands back the stream, the same way open_subgroup_stream does for a subgroup.

The caller owns the stream that comes back. Nothing here remembers which request it belongs to, so an endpoint serving several fetches at once keeps its own map from Request ID to stream.

Source

pub async fn accept_fetch_stream( &self, ) -> Result<(AnyFetchHeader, FramedRecvStream), ConnectionError>

Accept the next unidirectional stream and read its fetch header.

accept_subgroup_stream’s twin. The two are separate because the header decides how every object after it is framed, so a caller has to know which it is expecting before the first byte is read.

Objects come off the returned stream with FramedRecvStream::read_fetch_object.

Source

pub async fn accept_subgroup_stream( &self, ) -> Result<(AnySubgroupHeader, FramedRecvStream), ConnectionError>

Accept an incoming unidirectional data stream and read its subgroup header.

Source

pub fn send_datagram( &self, header: &AnyDatagramHeader, payload: &[u8], ) -> Result<(), ConnectionError>

Send an object via datagram.

The header goes through AnyDatagramHeader::encode, which refuses a header whose Object Status the framing it names cannot carry. Such a header errors here and nothing is sent, rather than going out as an ordinary payload datagram with the status quietly dropped.

Source

pub async fn recv_datagram( &self, ) -> Result<(AnyDatagramHeader, Bytes), ConnectionError>

Receive a datagram and decode its header.

Source

pub fn endpoint(&self) -> &Endpoint

Access the underlying endpoint state machine.

Source

pub fn endpoint_mut(&mut self) -> &mut Endpoint

Mutable access to the endpoint state machine.

Source

pub fn negotiated_version(&self) -> Option<VarInt>

Get the negotiated MoQT version.

Source

fn codec_session_error_code(err: &CodecError) -> Option<SessionErrorCode>

The code to close the session with when a message could not be decoded because the peer broke a rule draft-07 answers with a close.

Every variant listed here comes from a sentence in this draft that names the consequence, and the list is per draft: answering a bound this draft does not state would close a session over traffic a conforming peer may send. This draft’s list is the shortest of the thirteen and shares only two entries with draft-08’s, which sits next to it.

  • Duplicate parameters, a SHOULD rather than a MUST: “Receivers SHOULD check that there are no duplicate parameters and close the session as a ‘Protocol Violation’ if found.” Unqualified here, as on drafts 08 through 10: there is no carve-out for repeats a message authorizes and none for duplicates of unknown parameters, both of which arrive at draft-11 and make the rule asymmetric there.
  • Unknown control message type: “An endpoint that receives an unknown message type MUST close the session.” The sentence names no code, so Protocol Violation is what carries it, as on every other draft.
  • A parameter whose value does not match the length its type implies — the one rule here not answered with a Protocol Violation: “If a receiver understands a parameter type, and the parameter length implied by that type does not match the Parameter Length field, the receiver MUST terminate the session with error code ‘Parameter Length Mismatch’.” Drafts 08, 09 and 10 carry the same sentence; drafts 11 and later drop it along with the Parameter framing it describes.

Not the Track Namespace tuple size, and this is where draft-07 parts company with every draft above it. Section 2.3 states the range — “an ordered N-tuple of bytes where N can be between 1 and 32” — and stops there. It is draft-08 Section 2.4.1 that adds the consequence: “If an endpoint receives a Track Namespace tuple with an N of 0 or more than 32, it MUST close the session with a Protocol Violation.” A tuple of 40 fields is malformed on this draft and the session survives it, so the codec does not raise the variant here and this table would have nothing to answer if it did.

Not the Reason Phrase maximum, the GOAWAY New Session URI maximum, the Full Track Name maximum or the parameter value maximum. Those enter the specification at draft-11 and this draft states none of them.

Not the end-of-track Object ID rule either, which drafts 08, 09 and 10 do state: this draft assigns no Object Status 0x5 for it to be about.

Not CodecError::UnexpectedEnd, which reports no rule at all: the reader raises it whenever a message is still arriving, and read_control loops on it. Closing over it would end a session on an ordinary short read.

Not the key-value pair serialization rule. Drafts 11 and later require a close with KEY_VALUE_FORMATTING_ERROR when a value does not match the serialization its Type defines; this draft has no Key-Value Pair at all. What it states instead is the Parameter Length Mismatch rule above, over the Parameter framing it has in its place.

Not the unknown Message Parameter rule, which enters at draft-16 and requires a close for a Message Parameter type the negotiated version does not define. This draft states nothing of the kind, and its parameters are not Key-Value-Pairs at all.

None for everything else, including CodecError::InvalidField. That variant is shared by a dozen unrelated malformations, only some of which the draft answers with a close, so a session cannot be ended on it without ending sessions the draft does not ask to be ended. Two of this draft’s own rules are stuck behind it — an unknown data stream type, and a ContentExists field holding anything but 0 or 1, both of which this draft calls a protocol error. Splitting it is the way to bring them under this function; widening the match is not.

Source

fn close_for_codec(&self, err: ConnectionError) -> ConnectionError

Close the session on the wire when a decode failure is one draft-07 answers with a close, and hand the error back unchanged. Without it every bound the decoder enforces would stop at this endpoint refused the frame while the peer, which is the one that broke the rule, saw a session that was still open and went on sending. “MUST close the session” is a statement about the wire.

Source

fn close_for(&self, err: &EndpointError)

Close the session on the wire when the endpoint says a violation is fatal to it, and hand the error back unchanged.

EndpointError::session_error_code answers Some for exactly the errors this draft ends the session over, and the endpoint has already moved its own state machine to Closed by the time this runs. Without this step that move is purely internal: the local endpoint refuses to start anything new while the peer, which is the one that broke the rule, sees a session that is still open and goes on sending. A rule that names a session termination code is a statement about the wire, so it takes a CONNECTION_CLOSE to satisfy it.

The reason phrase is the error’s own Display text, which names the rule rather than repeating the numeric code the close already carries.

Errors that answer None are recoverable and nothing is sent.

Source

fn close_if_session_fatal(&self, err: EndpointError) -> ConnectionError

close_for, then the error unchanged, for the common case where the endpoint’s error is also what the caller returns.

Source

pub fn close_for_data_stream(&self, err: &ConnectionError) -> bool

Close the session over a rule broken on a data stream, reporting whether it did.

A data stream cannot close for itself the way recv_control does: Connection::accept_subgroup_stream hands the caller a FramedRecvStream holding no connection, so the reader that finds the violation is not the object that can act on it. Keeping it a separate call is deliberate as well — a permissive caller, one reproducing a capture, can read a violating stream and report it without tearing the session down.

The rule this draft answers here is the unknown stream type, Section 7, which arrives as the very first varint on a unidirectional stream and nowhere else. Every draft from 07 to 19 states it, in one of two phrasings — this one names streams alone because draft-07 numbers its datagrams in the same table, and drafts 08 through 16, draft-08 Section 8 among them, say “an unknown stream or datagram type” for the two tables they split it into. It shares codec_session_error_code with the control path, so a rule is answered with one code whichever stream carried it.

Not every rule that reaches here is the decoder’s. A track whose objects mix forwarding preferences is the endpoint’s to notice — it takes the alias table to know which track an object belongs to — and it arrives on exactly these streams. Both kinds are asked for a code the same way, and a rule with no code is declined rather than guessed at.

Source

pub fn close(&self, code: u32, reason: &[u8])

Close the connection.

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