Skip to main content

PacketLoopState

Struct PacketLoopState 

Source
pub(super) struct PacketLoopState {
    pub(super) users: KeyedSlotStore<TransportSessionKey, RtcSessionState, SessionSlot>,
    pub(super) routes: RouteTable,
    pub(super) session_media: BTreeMap<TransportSessionKey, SessionMediaLookup>,
    pub(super) rid_readiness_scratch: RidReadinessScratch,
    pub(super) incoming_bitrate_counters: BTreeMap<TransportMediaId, Arc<MediaBitrateCounter>>,
    pub(super) remote_addr_demux: RemoteAddrDemux,
    pub(super) mid_registry: BTreeMap<TransportMediaId, RegisteredMedia>,
    pub(super) dirty_sessions: Vec<SlotHandle<SessionSlot>>,
    pub(super) timeout_queue: BinaryHeap<Reverse<(Instant, SlotHandle<SessionSlot>)>>,
    pub(super) next_media_id: u64,
}
Expand description

authoritative mutable state for one RTC packet-loop worker

the packet loop owns this value without a mutex control commands, UDP ingress, str0m polling and relay fanout all pass through one mutable borrow so the media indexes can be updated together

command-facing APIs keep stable transport ids while hot queues use generation-checked handles

Fields§

§users: KeyedSlotStore<TransportSessionKey, RtcSessionState, SessionSlot>

live worker-local RTC sessions

§routes: RouteTable

source-scoped packet routing, relay and recovery state

§session_media: BTreeMap<TransportSessionKey, SessionMediaLookup>

session-scoped media lookup vectors for packet source resolution

§rid_readiness_scratch: RidReadinessScratch

reusable selected-RID readiness scratch vectors

§incoming_bitrate_counters: BTreeMap<TransportMediaId, Arc<MediaBitrateCounter>>

packet-loop write handles for incoming media bitrate accounting

§remote_addr_demux: RemoteAddrDemux

worker-local UDP ingress demux hints

§mid_registry: BTreeMap<TransportMediaId, RegisteredMedia>

primary media handle table keyed by stable transport media id

§dirty_sessions: Vec<SlotHandle<SessionSlot>>

sessions that must be polled before the worker waits again

§timeout_queue: BinaryHeap<Reverse<(Instant, SlotHandle<SessionSlot>)>>

Immutable deadline snapshots validated against the current session generation and RtcSessionState::next_timeout.

§next_media_id: u64

next worker-local media id from the disjoint range assigned at boot

Implementations§

Source§

impl PacketLoopState

Source

pub fn register_consumer_route( &mut self, registration: ConsumerRouteRegistration<'_>, ) -> TransportMediaId

Registers a consumer identity, destination stream and its MID route index together.

Source

pub fn set_consumer_active( &mut self, route: &TransportConsumerRoute, active: bool, ) -> TransportResult<()>

Applies consumer activity and invalidates repair state before returning.

§Errors

Returns TransportAdapterError::InvalidInput for mismatched source or consumer ownership or an incompatible media handle. Returns TransportAdapterError::TransportUnavailable when source registration, consumer media or its indexed destination is missing or stale.

Source

pub fn set_consumer_packet_gates( &mut self, source: &TransportSourceKey, updates: impl ExactSizeIterator<Item = (TransportConsumerRoute, PacketLayerGate)>, ) -> Vec<TransportResult<()>> ⓘ

Applies ordered consumer gates with repair invalidation and one source-gate refresh.

Successful entries remain applied when another entry fails. Results retain input order. The aggregate gate is refreshed once if any route changed.

§Errors

Each entry can return the errors from Self::set_consumer_active. A route for another source returns TransportAdapterError::InvalidInput.

Source

pub fn remove_consumer_route( &mut self, consumer_key: &TransportSessionKey, consumer_media: TransportMediaId, src_media: TransportMediaId, )

Removes a consumer route, repairs displaced MID indexes and retires its stream.

Source

pub fn remove_source_route(&mut self, src_media: TransportMediaId)

Removes every destination and receiver stream retained by one source route.

Source

fn update_consumer_route( &mut self, route: &TransportConsumerRoute, mutation: ConsumerRouteMutation, ) -> TransportResult<bool>

Source

fn release_destination_stream(&mut self, destination: &MediaRouteDestination)

Source§

impl PacketLoopState

Source

pub fn unregister_media_handle( &mut self, bitrate_registry: &Mutex<BitrateRegistry>, transport_media_id: TransportMediaId, ) -> Result<(), TransportAdapterError>

Removes media identity and every dependent route, stream and accounting entry.

Callers stage the final MID’s negotiated removal before committing teardown. Receive or transmit streams are retired before their NACK totals are cleared.

§Errors

Returns TransportAdapterError::InvalidInput when the media is unregistered.

Source§

impl PacketLoopState

Source

pub fn ensure_route_src_registered( &mut self, route_owner_session_key: &TransportSessionKey, source: &TransportSourceKey, remote_source_control: Option<RemoteSourceControl>, ) -> Result<RouteSourceKind, TransportAdapterError>

validates source ownership for a route that is about to be created

local sources must still be producer handles owned by the declared source session remote sources require remote_source_control because later route refreshes need a command path back to the producer worker

§Errors

Returns TransportAdapterError::InvalidInput when a source id belongs to another owner, when the media id names a consumer or when a remote source is missing its control path Returns TransportAdapterError::TransportUnavailable when the local source media id does not exist

Source

pub fn ensure_existing_route_src( &self, route_owner_session_key: &TransportSessionKey, source: &TransportSourceKey, ) -> Result<RouteSourceKind, TransportAdapterError>

Validates source ownership without creating remote-source state.

Existing-route commands can lag teardown. They fail after registration disappears instead of recreating it from stale control input.

§Errors

Returns TransportAdapterError::InvalidInput when the media id exists but belongs to a different source owner or is not a producer source for a local route Returns TransportAdapterError::TransportUnavailable when the expected local producer or remote source registration is gone

Source

pub fn ensure_local_producer_mid( &self, src_key: &TransportSessionKey, src_media: TransportMediaId, ) -> Result<Mid, TransportAdapterError>

returns the MID for a local producer after enforcing source ownership

this is the strict form used by command paths that must reject stale or misaddressed producer effects

§Errors

Returns TransportAdapterError::InvalidInput when the media id exists without a live producer owned by src_key Returns TransportAdapterError::TransportUnavailable when the media id is not registered

Source§

impl PacketLoopState

Source

pub fn apply_producer_activity( &mut self, source: &TransportSourceKey, update: SourceActivityUpdate, ) -> Result<bool, TransportAdapterError>

Applies producer activity with NACK policy and consumer repair invalidation.

Accepted revisions that preserve activity keep existing receiver repair.

§Errors

Returns TransportAdapterError::InvalidInput when the media belongs to another owner or is not a producer. Returns TransportAdapterError::TransportUnavailable when the producer or its route state is missing.

Source

pub fn set_remote_source_activity( &mut self, source: &TransportSourceKey, update: SourceActivityUpdate, ) -> Result<(), TransportAdapterError>

Applies registered remote-source activity before retiring previous repair.

Missing registrations and obsolete revisions are unchanged.

§Errors

Returns TransportAdapterError::InvalidInput when the source identity differs from the current remote registration.

Source

pub fn update_decoder_readiness( &mut self, src_media: TransportMediaId, incoming_rid: Option<Rid>, is_keyframe: bool, now: Instant, scratch: &mut RidReadinessScratch, ) -> RidReadinessRouteUpdate

Commits decoder gates and invalidates affected repair before feedback dispatch.

Source packet liveness must already include the current observation. Scratch vectors retain capacity and collect the RIDs requiring feedback.

Source

fn invalidate_source_repair(&mut self, src_media: TransportMediaId)

Source§

impl PacketLoopState

Source

pub(in engine::media_transport::rtc) fn mark_session_dirty( &mut self, session_key: &TransportSessionKey, )

schedule a live session for the next packet-loop poll

missing sessions are ignored because teardown may race with already queued wakeups each live session can appear at most once until Self::collect_ready_sessions clears its dirty bit

Source

pub(in engine::media_transport::rtc) fn has_dirty_sessions( &self, ) -> bool

report whether the worker has session work that is due immediately

Source

pub(in engine::media_transport::rtc) fn collect_ready_sessions( &mut self, now: Instant, ready_sessions: &mut Vec<SlotHandle<SessionSlot>>, )

drain dirty sessions and due str0m timeouts into caller-owned scratch

this method is the session scheduler for the packet loop it clears dirty bits for live sessions, skips removed sessions and lazily discards timeout heap entries whose deadline no longer matches super::RtcSessionState::next_timeout

stale handles are skipped before replacement sessions can be polled

the output is sorted and deduplicated so a session that is both dirty and timed out is polled once in the current turn

Source

pub(in engine::media_transport::rtc) fn update_session_timeout_by_handle( &mut self, session_handle: SlotHandle<SessionSlot>, next_timeout: Option<Instant>, )

replace the next str0m timeout deadline by worker-local handle

stale handles are ignored because the session has already left this worker or the slot now belongs to a later generation

Source

pub(in engine::media_transport::rtc) fn next_timeout_deadline( &mut self, ) -> Option<Instant>

return the earliest live str0m timeout deadline

stale heap entries are removed while searching this includes entries for handles whose slot generation no longer names a live session callers may invoke this before awaiting because it does not borrow any session state after returning

Source

pub(in engine::media_transport::rtc) fn clear_session_schedule( &mut self, session_key: &TransportSessionKey, )

remove all explicit scheduler state for a session being torn down

stale timeout heap entries can remain because the session no longer validates their deadline stale handles are also rejected by generation checks before polling

Source§

impl PacketLoopState

Source

pub fn remove_session( &mut self, session_key: &TransportSessionKey, ) -> Option<RtcSessionState>

Retires a session and repairs every index retained by other sessions.

Missing sessions still clear stale scheduling, demux and media state. Returns the removed session for lifetime accounting after route cleanup.

Source§

impl PacketLoopState

Source

pub(in engine::media_transport::rtc) fn register_incoming_bitrate_counter( &mut self, transport_media_id: TransportMediaId, counter: Arc<MediaBitrateCounter>, )

Source

pub(in engine::media_transport::rtc) fn remove_incoming_bitrate_counter( &mut self, transport_media_id: TransportMediaId, )

Source

pub(in engine::media_transport::rtc) fn record_incoming_bitrate( &self, transport_media_id: TransportMediaId, now: Instant, payload_bytes: usize, ) -> Option<IncomingBitrateObservation>

Source§

impl PacketLoopState

Source

pub(in engine::media_transport::rtc) fn register_media_handle( &mut self, handle: RegisteredMediaHandle, ) -> TransportMediaId

Source

pub(in engine::media_transport::rtc) fn resolve_mid( &self, transport_media_id: TransportMediaId, ) -> Option<Mid>

Source

pub(in engine::media_transport::rtc) fn media_handle( &self, transport_media_id: TransportMediaId, ) -> Option<&RegisteredMediaHandle>

Source

pub(in engine::media_transport::rtc) fn producer_media_snapshot( &self, session_key: &TransportSessionKey, ) -> Vec<(TransportMediaId, Mid)>

Source

pub(in engine::media_transport::rtc) fn consumer_media_snapshot( &self, session_key: &TransportSessionKey, ) -> Vec<(TransportMediaId, Mid, TransportMediaId)>

Source

pub(in engine::media_transport::rtc) fn mid_is_shared( &self, session_key: &TransportSessionKey, mid: Mid, excluded_transport_media_id: TransportMediaId, ) -> bool

Source

pub(super) fn remove_media_handle( &mut self, transport_media_id: TransportMediaId, ) -> Option<RegisteredMediaHandle>

remove one media handle and every dependent reverse index owned by it

producer removal clears source packet policy, incoming bitrate counters, decoder-refresh metadata, live RID state and SSRC lookups Consumer removal clears its MID lookup. Complete route and stream teardown belongs to Self::unregister_media_handle.

Source

pub(in engine::media_transport::rtc) fn session_has_registered_media( &self, session_key: &TransportSessionKey, ) -> bool

Source

fn prune_empty_session_media(&mut self, session_key: &TransportSessionKey)

Source

pub(in engine::media_transport::rtc) fn take_expired_speaker_rooms( &mut self, now: Instant, ) -> BTreeSet<RoomInstanceId>

Source

pub(in engine::media_transport::rtc) fn src_media_for_mid( &self, src_key: &TransportSessionKey, source_mid: Mid, ) -> Option<TransportMediaId>

Source

pub(in engine::media_transport::rtc) fn src_media_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<TransportMediaId>

Source

pub fn producer_binding_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<(TransportMediaId, Option<Rid>)>

Source

pub(in engine::media_transport::rtc) fn source_rid_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<Rid>

Source

pub(in engine::media_transport::rtc) fn bind_producer_packet( &mut self, session_handle: SlotHandle<SessionSlot>, cached_media: Option<TransportMediaId>, mid: Option<Mid>, binding: &mut ProducerStreamBinding, ) -> Option<(TransportMediaId, RoomInstanceId, ProducerSsrcUpdate)>

Resolves and admits one authenticated local packet in the session’s media index.

Cached media precedes MID and SSRC. An indexed binding for that media supplies the current RID, including RID-less after renegotiation.

Source

pub fn bind_producer_stream( &mut self, source: &ForwardedPacketSource, transport_media_id: TransportMediaId, binding: ProducerStreamBinding, ) -> ProducerSsrcUpdate

Commits the current primary and repair identities for one negotiated encoding.

Both indexes retain at most two SSRCs per encoding. The previous primary is only a bounded rejection tombstone, never a demultiplexing binding. A collision, unknown RID or preceding primary returns Rejected without changing either index. Callers must not forward rejected packets.

Source

pub(in engine::media_transport::rtc) fn active_consumer_kf_target( &self, consumer_key: &TransportSessionKey, consumer_mid: Mid, feedback_rid: Option<Rid>, ) -> Option<ConsumerKeyframeTarget>

resolve consumer RTCP feedback to the currently active producer target

this is the packet-loop feedback path missing indexes, removed routes and inactive routes are treated as stale feedback because those can race with teardown after str0m emits a request

selected destination gates override the feedback RID so RID-less browser PLI stays scoped to the routed simulcast layer

Source

pub(in engine::media_transport::rtc) fn set_consumer_dst_idx( &mut self, consumer_key: &TransportSessionKey, consumer_mid: Mid, consumer_media: TransportMediaId, src_media: TransportMediaId, dst_idx: Option<usize>, )

updates the cached destination slot for one consumer MID binding

callers pass both sides of the route identity so a late repair from an old route cannot relink a MID to a different source

Source

pub(in engine::media_transport::rtc) fn consumer_dst_idx( &self, consumer_key: &TransportSessionKey, consumer_mid: Mid, consumer_media: TransportMediaId, src_media: TransportMediaId, ) -> Option<usize>

Source

fn source_room_instance_id( &self, src_media: TransportMediaId, ) -> Option<RoomInstanceId>

Source

pub(in engine::media_transport::rtc) fn remove_session_media_handles( &mut self, session_key: &TransportSessionKey, ) -> Vec<TransportMediaId>

Source

pub(in engine::media_transport::rtc) fn refresh_answer_producer_ssrcs( &mut self, session_key: &TransportSessionKey, producer_mids: &[Mid], refreshed_parameters: &[(Mid, MediaStream)], )

Source

fn refresh_producer_ssrcs_with_previous( &mut self, session_key: &TransportSessionKey, transport_media_id: TransportMediaId, parameters: &MediaStream, previous_encodings: &[ProducerEncoding], )

Source

pub(in engine::media_transport::rtc) fn apply_producer_nack_policy( &mut self, session_key: &TransportSessionKey, transport_media_id: TransportMediaId, )

Reapplies NACK suppression because SDP answers can recreate StreamRx.

Source

pub(in engine::media_transport::rtc) fn clear_producer_ssrcs_for_mid( &mut self, session_key: &TransportSessionKey, mid: Mid, )

Source

fn clear_producer_ssrcs( &mut self, session_key: &TransportSessionKey, transport_media_id: TransportMediaId, )

Trait Implementations§

Source§

impl Default for PacketLoopState

Source§

fn default() -> PacketLoopState

Returns the “default value” for a type. Read more

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> FutureExt for T

§

fn with_context(self, otel_cx: Context) -> WithContext<Self>

Attaches the provided Context to this type, returning a WithContext wrapper. Read more
§

fn with_current_context(self) -> WithContext<Self>

Attaches the current Context to this type, returning a WithContext wrapper. Read more
§

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> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
§

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

§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] only if self and other return Action::Follow. Read more
§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns [Action::Follow] if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

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

Source§

type Error = Infallible

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

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

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<V, T> VZip<V> for T
where V: MultiLane<T>,

§

fn vzip(self) -> V

§

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
§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,