o_sfu_core/engine/
packet_sink_registry.rs1use 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
23pub trait PacketSink: Send + Sync {
28 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 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 let active_rooms = read_unpoisoned(&self.active_rooms);
152 PacketSinkRegistrySnapshot {
153 generation: self.generation(),
154 active_rooms: active_rooms.clone(),
155 }
156 }
157
158 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 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}