Skip to main content

o_sfu_core/engine/
packet_sink_registry.rs

1use std::{
2    collections::HashMap,
3    fmt,
4    sync::{
5        Arc, RwLock,
6        atomic::{AtomicBool, AtomicU64, Ordering},
7    },
8    time::Instant,
9};
10
11use super::{
12    RoomInstanceId,
13    media_transport::{TransportMediaId, TransportSessionKey},
14    metrics::RtpForwardDestinationKind,
15    sync::{read_unpoisoned, write_unpoisoned},
16};
17
18#[cfg(test)]
19#[expect(non_snake_case, reason = "test modules map to local TESTS directories")]
20#[path = "packet_sink_registry/TESTS/mod.rs"]
21mod TESTS;
22
23/// Observes origin-side RTP payloads routed to a room packet sink.
24///
25/// Packet gates do not filter this source-side stream. Relayed packets are not
26/// observed again.
27pub trait PacketSink: Send + Sync {
28    /// Records one source RTP payload.
29    ///
30    /// `session_key` and `transport_media_id` identify the source transport
31    /// media. The RTC engine invokes this synchronously on a packet worker.
32    /// Different workers may invoke the same sink concurrently, so
33    /// implementations must return promptly and must not perform blocking I/O.
34    fn record_packet(
35        &self,
36        session_key: &TransportSessionKey,
37        transport_media_id: TransportMediaId,
38        received_at: Instant,
39        payload: &[u8],
40    );
41}
42
43#[derive(Clone)]
44pub struct RegisteredPacketSink {
45    sink: Arc<dyn PacketSink>,
46    forward_destination_kind: RtpForwardDestinationKind,
47}
48
49impl RegisteredPacketSink {
50    pub fn new(
51        sink: Arc<dyn PacketSink>,
52        forward_destination_kind: RtpForwardDestinationKind,
53    ) -> Self {
54        Self {
55            sink,
56            forward_destination_kind,
57        }
58    }
59
60    pub fn record_packet(
61        &self,
62        session_key: &TransportSessionKey,
63        transport_media_id: TransportMediaId,
64        received_at: Instant,
65        payload: &[u8],
66    ) {
67        self.sink
68            .record_packet(session_key, transport_media_id, received_at, payload);
69    }
70
71    #[must_use]
72    pub const fn forward_destination_kind(&self) -> RtpForwardDestinationKind {
73        self.forward_destination_kind
74    }
75}
76
77impl fmt::Debug for RegisteredPacketSink {
78    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
79        formatter
80            .debug_struct("RegisteredPacketSink")
81            .field("forward_destination_kind", &self.forward_destination_kind)
82            .finish_non_exhaustive()
83    }
84}
85
86pub struct RoomPacketSinkRegistry {
87    any_active: AtomicBool,
88    generation: AtomicU64,
89    active_rooms: RwLock<HashMap<RoomInstanceId, RegisteredPacketSink>>,
90}
91
92impl Default for RoomPacketSinkRegistry {
93    fn default() -> Self {
94        Self {
95            any_active: AtomicBool::new(false),
96            generation: AtomicU64::new(0),
97            active_rooms: RwLock::new(HashMap::new()),
98        }
99    }
100}
101
102#[derive(Default)]
103pub struct PacketSinkRouteCache {
104    generation: u64,
105    active_rooms: HashMap<RoomInstanceId, RegisteredPacketSink>,
106}
107
108impl PacketSinkRouteCache {
109    /// Refreshes the cached room routes to one registry generation.
110    ///
111    /// Later registry changes remain invisible through [`Self::sink_for_room`]
112    /// until `refresh_from` runs again.
113    pub fn refresh_from(&mut self, registry: &RoomPacketSinkRegistry) {
114        let generation = registry.generation();
115        if self.generation == generation {
116            return;
117        }
118        let snapshot = registry.snapshot();
119        self.generation = snapshot.generation;
120        self.active_rooms = snapshot.active_rooms;
121    }
122
123    #[inline]
124    pub fn sink_for_room(&self, room_instance_id: RoomInstanceId) -> Option<RegisteredPacketSink> {
125        self.active_rooms.get(&room_instance_id).cloned()
126    }
127}
128
129struct PacketSinkRegistrySnapshot {
130    generation: u64,
131    active_rooms: HashMap<RoomInstanceId, RegisteredPacketSink>,
132}
133
134impl RoomPacketSinkRegistry {
135    #[inline]
136    pub fn sink_for_room(&self, room_instance_id: RoomInstanceId) -> Option<RegisteredPacketSink> {
137        if !self.any_active.load(Ordering::Acquire) {
138            return None;
139        }
140        read_unpoisoned(&self.active_rooms)
141            .get(&room_instance_id)
142            .cloned()
143    }
144
145    fn generation(&self) -> u64 {
146        self.generation.load(Ordering::Acquire)
147    }
148
149    fn snapshot(&self) -> PacketSinkRegistrySnapshot {
150        // Both writers advance the generation before releasing this map's guard.
151        let active_rooms = read_unpoisoned(&self.active_rooms);
152        PacketSinkRegistrySnapshot {
153            generation: self.generation(),
154            active_rooms: active_rooms.clone(),
155        }
156    }
157
158    /// Registers `sink` for `room_instance_id`, replacing the current entry.
159    ///
160    /// Previously cloned [`RegisteredPacketSink`] handles are not revoked.
161    pub fn register_room(
162        &self,
163        room_instance_id: RoomInstanceId,
164        sink: Arc<dyn PacketSink>,
165        forward_destination_kind: RtpForwardDestinationKind,
166    ) {
167        let mut active_rooms = write_unpoisoned(&self.active_rooms);
168        active_rooms.insert(
169            room_instance_id,
170            RegisteredPacketSink::new(sink, forward_destination_kind),
171        );
172        self.any_active.store(true, Ordering::Release);
173        self.generation.fetch_add(1, Ordering::AcqRel);
174        drop(active_rooms);
175    }
176
177    /// Removes the current entry for `room_instance_id`.
178    ///
179    /// Cached or cloned sink handles may still receive packets after
180    /// `unregister_room` returns.
181    pub fn unregister_room(&self, room_instance_id: RoomInstanceId) {
182        let mut active_rooms = write_unpoisoned(&self.active_rooms);
183        active_rooms.remove(&room_instance_id);
184        self.any_active
185            .store(!active_rooms.is_empty(), Ordering::Release);
186        self.generation.fetch_add(1, Ordering::AcqRel);
187        drop(active_rooms);
188    }
189
190    fn active_room_count(&self) -> usize {
191        read_unpoisoned(&self.active_rooms).len()
192    }
193}
194
195impl fmt::Debug for RoomPacketSinkRegistry {
196    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
197        formatter
198            .debug_struct("RoomPacketSinkRegistry")
199            .field("any_active", &self.any_active.load(Ordering::Relaxed))
200            .field("active_room_count", &self.active_room_count())
201            .finish_non_exhaustive()
202    }
203}