Skip to main content

Connection

Struct Connection 

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

A live MoQT connection over QUIC or WebTransport, combining the endpoint state machine with actual network I/O.

Fields§

§transport: Transport§endpoint: Endpoint§draft: DraftVersion§control_send: Option<Mutex<FramedSendStream>>

Behind a lock because a control message is written from two kinds of place. Most of them are the caller’s own request, made through &mut self. The messages that answer a Malformed Track are not: the conditions that make a track malformed are detected on the data plane, where this connection is reached through a shared reference. The lock also makes one message the unit of writing, so two of them cannot interleave on the stream.

§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 MoQT server as a client.

Establishes a QUIC or WebTransport connection (based on config.transport), opens a bidirectional control stream, performs the CLIENT_SETUP / SERVER_SETUP handshake, and returns a ready-to-use connection.

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( &self, msg: &ControlMessage, ) -> Result<(), ConnectionError>

Send a control message on the control stream.

Wraps the draft-15 message in AnyControlMessage::Draft15 for framing.

Source

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

Read the next control message from the control stream.

Returns the AnyControlMessage and also extracts the draft-15 ControlMessage for internal endpoint dispatch.

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_namespace: TrackNamespace, track_name: Vec<u8>, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send a SUBSCRIBE and return the allocated request ID.

Source

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

Send an UNSUBSCRIBE for the given request ID.

Source

pub async fn subscribe_ok( &mut self, request_id: VarInt, track_alias: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<(), ConnectionError>

Accept a subscription the peer opened, sending SUBSCRIBE_OK and giving its track a Track Alias.

The endpoint refuses an alias a live track of its own already holds and refuses a second answer to one SUBSCRIBE, so nothing is written on the wire when it does either.

Source

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

Refuse a request the peer opened, sending REQUEST_ERROR.

One message refuses a SUBSCRIBE, a FETCH, an announcement, a track status or a namespace subscription, and the endpoint finds which by the identifier. It refuses a second answer to any of them, and refuses a Joining Fetch’s refusal under any code but the one the draft names for it, so nothing is written on the wire when it does.

Source

pub async fn subscribe_update( &mut self, subscription_request_id: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Narrow a subscription this endpoint opened, sending SUBSCRIBE_UPDATE, and return the Request ID the update itself spent.

Source

pub async fn publish_ok( &mut self, request_id: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<(), ConnectionError>

Accept a PUBLISH the peer sent, which establishes the subscription it opened.

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

Source

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

Reject a PUBLISH the peer sent, which ends the subscription it opened before it was established.

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

Source

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

Send a standalone FETCH and return the allocated request ID.

Source

pub async fn joining_fetch( &mut self, joining_request_id: VarInt, joining_start: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send a Relative Joining Fetch and return the allocated request ID.

joining_start counts groups back from the subscription’s largest group. To name the starting group outright, use absolute_joining_fetch.

Source

pub async fn absolute_joining_fetch( &mut self, joining_request_id: VarInt, joining_start: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send an Absolute Joining Fetch and return the allocated request ID.

Here joining_start is the group to begin at rather than an offset, which is what an application that knows the group it wants has: draft-15 Section 9.16.2.1 has the publisher set the Start Location to {Joining Start, 0}.

Source

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

Send a FETCH_CANCEL for the given request ID.

Source

pub async fn fetch_ok( &mut self, request_id: VarInt, end_of_track: u8, end_group: VarInt, end_object: 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 subscribe_namespace( &mut self, namespace_prefix: TrackNamespace, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send a SUBSCRIBE_NAMESPACE and return the request ID.

Source

pub async fn publish_namespace( &mut self, track_namespace: TrackNamespace, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send a PUBLISH_NAMESPACE and return the request ID.

Source

pub async fn request_ok( &mut self, request_id: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<(), ConnectionError>

Accept a request the peer opened, sending REQUEST_OK.

The endpoint refuses a second answer to one request, so nothing is written on the wire when it does. On this draft an announcement, a track status and a namespace subscription are the requests REQUEST_OK accepts; a subscription, a publication and a fetch each have an acceptance of their own that carries more than this one can.

Source

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

Revoke an acceptance, sending PUBLISH_NAMESPACE_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 publish_namespace_done( &mut self, track_namespace: TrackNamespace, ) -> Result<(), ConnectionError>

Withdraw an announcement this endpoint made, sending PUBLISH_NAMESPACE_DONE.

The mirror of Self::publish_namespace, and the counterpart of Self::publish_namespace_cancel: this one ends an announcement of this endpoint’s, that one revokes the acceptance of one the peer made.

Source

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

Send a TRACK_STATUS and return the allocated request ID.

Source

pub async fn publish( &mut self, track_namespace: TrackNamespace, track_name: Vec<u8>, track_alias: VarInt, parameters: Vec<KeyValuePair>, ) -> Result<VarInt, ConnectionError>

Send a PUBLISH and return the allocated request ID.

Source

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

Send a PUBLISH_DONE for the given request ID.

Source

async fn withdraw_malformed_track( &self, alias: u64, condition: MalformedTrackCondition, )

Send what Section 2.4.2 asks for when this endpoint finds a track malformed.

“it MUST UNSUBSCRIBE any subscription and FETCH_CANCEL any fetch for that Track from that publisher” — one message per request, in Request ID order, and the endpoint decides which message each request takes.

A write that fails is not reported. The caller is on its way to returning an error that says what went wrong with the track, and a control stream that will not take an UNSUBSCRIBE is a session on its way out for a reason of its own; replacing the condition’s report with a transport error would lose the only account of why the track was withdrawn. The rest of the withdrawal is abandoned, because a stream that refused one message will refuse the next.

Source

async fn note_received_framing( &self, alias: u64, seen: ObjectForwardingPreference, ) -> Result<(), ConnectionError>

Record the framing an arriving object was sent with, and withdraw from the track when it is the second framing that track has been sent.

The receiving half of a pair. The two writing paths call the endpoint directly and answer a mixed track by refusing to write it, because Section 2.4.2’s sentence is a subscriber’s: an endpoint about to send an object is that object’s Original Publisher, and a publisher has no subscription of its own to withdraw and no fetch of its own to cancel.

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

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 async fn withdraw_for_data_stream(&self, err: &ConnectionError) -> bool

Withdraw from a track a data stream found malformed, reporting whether it did.

The Malformed Track twin of Connection::close_for_data_stream, and separate from it for the same reason and one more. The same one: a FramedRecvStream holds no connection, so the reader that finds the fault is not the object that can send an UNSUBSCRIBE. The one more: the two answers are opposites - that call ends the session, this one gives up a track and leaves it running - and a single entry point would have to decide between them from the error alone, which is exactly the decision a caller reproducing a capture wants to make itself.

The datagram path needs none of this. It is read through the connection, so Connection::recv_datagram answers the condition where it finds it, and this is only for the objects that arrive on a stream the caller holds.

Source

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

Close the session when a failure raised while reading a data stream is one draft-15 answers with a close. Reports whether it closed.

accept_subgroup_stream hands the caller a FramedRecvStream, which holds no connection and so cannot close one, and the read that raises this failure happens there. The caller is the only party holding both halves, which is what this is for.

Splitting it this way rather than closing inside the reader keeps a caller that is deliberately permissive — a tool reproducing a capture, say — able to read a violating stream and report it without tearing the session down. The rule is stated at endpoints, and this is where an endpoint decides it is one.

Answers the extension-header rule of Section 10.2.1.2, and any decode failure codec_session_error_code recognises, 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 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 draft(&self) -> DraftVersion

Returns the draft version this connection is using.

Source

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

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

Every variant listed here comes from a sentence in this draft that names the consequence, and the list is deliberately shorter than draft-17’s: the bounds are per draft, and answering one this draft does not state would close a session over traffic a conforming peer may send.

  • Reason Phrase, maximum 1024 bytes: “If an endpoint receives a length exceeding the maximum, it MUST close the session with a PROTOCOL_VIOLATION.”

  • KVP value, maximum 2^16-1 bytes, with the same sentence.

  • Track Namespace field count: “If an endpoint receives a Track Namespace consisting of 0 or greater than 32 Track Namespace Fields, it MUST close the session with a PROTOCOL_VIOLATION.” Note the lower bound — an empty tuple is refused here, where drafts 17 and later permit it.

  • Full Track Name, maximum 4,096 bytes. Draft-15 states this of the Full Track Name alone; draft-16 widened it to a Track Namespace on its own.

  • Duplicate parameters, a SHOULD rather than a MUST: “Receivers SHOULD check that there are no unauthorized duplicate parameters and close the session as a PROTOCOL_VIOLATION”

  • GOAWAY New Session URI, maximum 8,192 bytes: “If an endpoint receives a length exceeding the maximum, it MUST close the session with a PROTOCOL_VIOLATION.” Every draft from 11 to 19 states it; 07 through 10 state no maximum for the field at all.

  • Unknown control message type: “An endpoint that receives an unknown message type MUST close the session.” All thirteen drafts state it, in the same words, and the sentence names no code, so Protocol Violation is what carries it.

Not the unknown Message Parameter rule. Drafts 16 through 19 require a close for a Message Parameter whose type the negotiated version does not define. This draft states the opposite and states it about the same parameters: “Receivers MUST allow duplicates of unknown parameters”, which presumes an unknown parameter arrives and is carried. Refusing one here would close a session over an extension this draft leaves room for.

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 treating it as fatal would close sessions the draft does not ask to be closed. Splitting it is the way to bring the rest of those rules 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-15 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 with a PROTOCOL_VIOLATION” is a statement about the wire.

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