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: RouteTablesource-scoped packet routing, relay and recovery state
session_media: BTreeMap<TransportSessionKey, SessionMediaLookup>session-scoped media lookup vectors for packet source resolution
rid_readiness_scratch: RidReadinessScratchreusable selected-RID readiness scratch vectors
incoming_bitrate_counters: BTreeMap<TransportMediaId, Arc<MediaBitrateCounter>>packet-loop write handles for incoming media bitrate accounting
remote_addr_demux: RemoteAddrDemuxworker-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: u64next worker-local media id from the disjoint range assigned at boot
Implementations§
Source§impl PacketLoopState
impl PacketLoopState
Sourcepub fn register_consumer_route(
&mut self,
registration: ConsumerRouteRegistration<'_>,
) -> TransportMediaId
pub fn register_consumer_route( &mut self, registration: ConsumerRouteRegistration<'_>, ) -> TransportMediaId
Registers a consumer identity, destination stream and its MID route index together.
Sourcepub fn set_consumer_active(
&mut self,
route: &TransportConsumerRoute,
active: bool,
) -> TransportResult<()>
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.
Sourcepub fn set_consumer_packet_gates(
&mut self,
source: &TransportSourceKey,
updates: impl ExactSizeIterator<Item = (TransportConsumerRoute, PacketLayerGate)>,
) -> Vec<TransportResult<()>> ⓘ
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.
Sourcepub fn remove_consumer_route(
&mut self,
consumer_key: &TransportSessionKey,
consumer_media: TransportMediaId,
src_media: TransportMediaId,
)
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.
Sourcepub fn remove_source_route(&mut self, src_media: TransportMediaId)
pub fn remove_source_route(&mut self, src_media: TransportMediaId)
Removes every destination and receiver stream retained by one source route.
fn update_consumer_route( &mut self, route: &TransportConsumerRoute, mutation: ConsumerRouteMutation, ) -> TransportResult<bool>
fn release_destination_stream(&mut self, destination: &MediaRouteDestination)
Source§impl PacketLoopState
impl PacketLoopState
Sourcepub fn unregister_media_handle(
&mut self,
bitrate_registry: &Mutex<BitrateRegistry>,
transport_media_id: TransportMediaId,
) -> Result<(), TransportAdapterError>
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
impl PacketLoopState
Sourcepub fn ensure_route_src_registered(
&mut self,
route_owner_session_key: &TransportSessionKey,
source: &TransportSourceKey,
remote_source_control: Option<RemoteSourceControl>,
) -> Result<RouteSourceKind, TransportAdapterError>
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
Sourcepub fn ensure_existing_route_src(
&self,
route_owner_session_key: &TransportSessionKey,
source: &TransportSourceKey,
) -> Result<RouteSourceKind, TransportAdapterError>
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
Sourcepub fn ensure_local_producer_mid(
&self,
src_key: &TransportSessionKey,
src_media: TransportMediaId,
) -> Result<Mid, TransportAdapterError>
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
impl PacketLoopState
Sourcepub fn apply_producer_activity(
&mut self,
source: &TransportSourceKey,
update: SourceActivityUpdate,
) -> Result<bool, TransportAdapterError>
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.
Sourcepub fn set_remote_source_activity(
&mut self,
source: &TransportSourceKey,
update: SourceActivityUpdate,
) -> Result<(), TransportAdapterError>
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.
Sourcepub fn update_decoder_readiness(
&mut self,
src_media: TransportMediaId,
incoming_rid: Option<Rid>,
is_keyframe: bool,
now: Instant,
scratch: &mut RidReadinessScratch,
) -> RidReadinessRouteUpdate
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.
fn invalidate_source_repair(&mut self, src_media: TransportMediaId)
Source§impl PacketLoopState
impl PacketLoopState
Sourcepub(in engine::media_transport::rtc) fn mark_session_dirty(
&mut self,
session_key: &TransportSessionKey,
)
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
Sourcepub(in engine::media_transport::rtc) fn has_dirty_sessions(
&self,
) -> bool
pub(in engine::media_transport::rtc) fn has_dirty_sessions( &self, ) -> bool
report whether the worker has session work that is due immediately
Sourcepub(in engine::media_transport::rtc) fn collect_ready_sessions(
&mut self,
now: Instant,
ready_sessions: &mut Vec<SlotHandle<SessionSlot>>,
)
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
Sourcepub(in engine::media_transport::rtc) fn update_session_timeout_by_handle(
&mut self,
session_handle: SlotHandle<SessionSlot>,
next_timeout: Option<Instant>,
)
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
Sourcepub(in engine::media_transport::rtc) fn next_timeout_deadline(
&mut self,
) -> Option<Instant>
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
Sourcepub(in engine::media_transport::rtc) fn clear_session_schedule(
&mut self,
session_key: &TransportSessionKey,
)
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
impl PacketLoopState
Sourcepub fn remove_session(
&mut self,
session_key: &TransportSessionKey,
) -> Option<RtcSessionState>
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
impl PacketLoopState
pub(in engine::media_transport::rtc) fn register_incoming_bitrate_counter( &mut self, transport_media_id: TransportMediaId, counter: Arc<MediaBitrateCounter>, )
pub(in engine::media_transport::rtc) fn remove_incoming_bitrate_counter( &mut self, transport_media_id: TransportMediaId, )
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
impl PacketLoopState
pub(in engine::media_transport::rtc) fn register_media_handle( &mut self, handle: RegisteredMediaHandle, ) -> TransportMediaId
pub(in engine::media_transport::rtc) fn resolve_mid( &self, transport_media_id: TransportMediaId, ) -> Option<Mid>
pub(in engine::media_transport::rtc) fn media_handle( &self, transport_media_id: TransportMediaId, ) -> Option<&RegisteredMediaHandle>
pub(in engine::media_transport::rtc) fn producer_media_snapshot( &self, session_key: &TransportSessionKey, ) -> Vec<(TransportMediaId, Mid)>
pub(in engine::media_transport::rtc) fn consumer_media_snapshot( &self, session_key: &TransportSessionKey, ) -> Vec<(TransportMediaId, Mid, TransportMediaId)>
Sourcepub(super) fn remove_media_handle(
&mut self,
transport_media_id: TransportMediaId,
) -> Option<RegisteredMediaHandle>
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.
pub(in engine::media_transport::rtc) fn session_has_registered_media( &self, session_key: &TransportSessionKey, ) -> bool
fn prune_empty_session_media(&mut self, session_key: &TransportSessionKey)
pub(in engine::media_transport::rtc) fn take_expired_speaker_rooms( &mut self, now: Instant, ) -> BTreeSet<RoomInstanceId>
pub(in engine::media_transport::rtc) fn src_media_for_mid( &self, src_key: &TransportSessionKey, source_mid: Mid, ) -> Option<TransportMediaId>
pub(in engine::media_transport::rtc) fn src_media_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<TransportMediaId>
pub fn producer_binding_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<(TransportMediaId, Option<Rid>)>
pub(in engine::media_transport::rtc) fn source_rid_for_ssrc( &self, src_key: &TransportSessionKey, source_ssrc: Ssrc, ) -> Option<Rid>
Sourcepub(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)>
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.
Sourcepub fn bind_producer_stream(
&mut self,
source: &ForwardedPacketSource,
transport_media_id: TransportMediaId,
binding: ProducerStreamBinding,
) -> ProducerSsrcUpdate
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.
Sourcepub(in engine::media_transport::rtc) fn active_consumer_kf_target(
&self,
consumer_key: &TransportSessionKey,
consumer_mid: Mid,
feedback_rid: Option<Rid>,
) -> Option<ConsumerKeyframeTarget>
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
Sourcepub(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>,
)
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
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>
fn source_room_instance_id( &self, src_media: TransportMediaId, ) -> Option<RoomInstanceId>
pub(in engine::media_transport::rtc) fn remove_session_media_handles( &mut self, session_key: &TransportSessionKey, ) -> Vec<TransportMediaId>
pub(in engine::media_transport::rtc) fn refresh_answer_producer_ssrcs( &mut self, session_key: &TransportSessionKey, producer_mids: &[Mid], refreshed_parameters: &[(Mid, MediaStream)], )
fn refresh_producer_ssrcs_with_previous( &mut self, session_key: &TransportSessionKey, transport_media_id: TransportMediaId, parameters: &MediaStream, previous_encodings: &[ProducerEncoding], )
Sourcepub(in engine::media_transport::rtc) fn apply_producer_nack_policy(
&mut self,
session_key: &TransportSessionKey,
transport_media_id: TransportMediaId,
)
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.
pub(in engine::media_transport::rtc) fn clear_producer_ssrcs_for_mid( &mut self, session_key: &TransportSessionKey, mid: Mid, )
fn clear_producer_ssrcs( &mut self, session_key: &TransportSessionKey, transport_media_id: TransportMediaId, )
Trait Implementations§
Source§impl Default for PacketLoopState
impl Default for PacketLoopState
Source§fn default() -> PacketLoopState
fn default() -> PacketLoopState
Auto Trait Implementations§
impl Freeze for PacketLoopState
impl !RefUnwindSafe for PacketLoopState
impl Send for PacketLoopState
impl Sync for PacketLoopState
impl Unpin for PacketLoopState
impl UnsafeUnpin for PacketLoopState
impl UnwindSafe for PacketLoopState
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
§impl<T> FutureExt for T
impl<T> FutureExt for T
§fn with_context(self, otel_cx: Context) -> WithContext<Self>
fn with_context(self, otel_cx: Context) -> WithContext<Self>
§fn with_current_context(self) -> WithContext<Self>
fn with_current_context(self) -> WithContext<Self>
§impl<T> Instrument for T
impl<T> Instrument for T
§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self>
fn into_either(self, into_left: bool) -> Either<Self, Self>
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
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