Skip to main content

o_sfu_core/engine/media_transport/rtc/
demux.rs

1//! Worker-local indexes for UDP ingress routing.
2//!
3//! Learned source addresses provide fast-path session pins. Local ICE ufrags
4//! and signaled candidate addresses narrow unknown-source recovery while
5//! `Rtc::accepts()` remains the ownership authority. Mutations update forward
6//! and reverse indexes together so session teardown cannot leave routing hints.
7
8use std::{
9    collections::{BTreeMap, HashMap},
10    net::SocketAddr,
11};
12
13use crate::engine::media_transport::TransportSessionKey;
14
15/// Bidirectional demux indexes for worker-local UDP ingress recovery.
16#[derive(Debug, Default)]
17pub struct RemoteAddrDemux {
18    /// learned UDP source tuple to session pin
19    remote_addr_index: HashMap<SocketAddr, TransportSessionKey>,
20    /// reverse lookup for learned UDP source tuple cleanup
21    remote_addrs_by_session: BTreeMap<TransportSessionKey, Vec<SocketAddr>>,
22    /// local ICE ufrag to session recovery hint
23    local_ice_ufrag_index: HashMap<String, TransportSessionKey>,
24    /// reverse lookup for replacing or removing a session local ICE ufrag
25    local_ice_ufrag_by_session: BTreeMap<TransportSessionKey, String>,
26    /// signaled remote candidate address to possible sessions
27    remote_candidate_addr_index: HashMap<SocketAddr, Vec<TransportSessionKey>>,
28    /// reverse lookup for candidate hint cleanup after renegotiation or teardown
29    remote_candidate_addrs_by_session: BTreeMap<TransportSessionKey, Vec<SocketAddr>>,
30}
31
32impl RemoteAddrDemux {
33    /// returns the session currently pinned to a UDP source tuple
34    ///
35    /// this is a hot-path hint for cached ingress routing
36    /// callers must still re-check the packet with `Rtc::accepts()` before
37    /// feeding it into a session because ICE state can move after the pin was
38    /// learned
39    #[must_use]
40    pub fn session_key_for_remote_addr(
41        &self,
42        source_addr: SocketAddr,
43    ) -> Option<&TransportSessionKey> {
44        self.remote_addr_index.get(&source_addr)
45    }
46
47    /// pins a UDP source tuple to the session that just accepted traffic
48    ///
49    /// returns `true` when the visible mapping changed
50    /// remapping a tuple also removes it from the previous session reverse
51    /// index so session teardown can later clean every learned tuple with one
52    /// key
53    #[must_use]
54    pub fn remember_remote_addr(
55        &mut self,
56        source_addr: SocketAddr,
57        session_key: &TransportSessionKey,
58    ) -> bool {
59        if self
60            .remote_addr_index
61            .get(&source_addr)
62            .is_some_and(|current_session| current_session == session_key)
63        {
64            return false;
65        }
66        let previous_session = self
67            .remote_addr_index
68            .insert(source_addr, session_key.clone());
69        if let Some(previous_session) = previous_session {
70            self.remove_remote_addr(&previous_session, source_addr);
71        }
72        let session_addrs = self
73            .remote_addrs_by_session
74            .entry(session_key.clone())
75            .or_default();
76        if !session_addrs.contains(&source_addr) {
77            session_addrs.push(source_addr);
78        }
79        true
80    }
81
82    /// returns the session advertised by a local ICE ufrag
83    ///
84    /// this index is used to narrow STUN recovery when the USERNAME attribute
85    /// names the local fragment
86    /// the returned session is still only a candidate for `Rtc::accepts()`
87    pub(super) fn session_for_local_ufrag(
88        &self,
89        local_ice_ufrag: &str,
90    ) -> Option<&TransportSessionKey> {
91        self.local_ice_ufrag_index.get(local_ice_ufrag)
92    }
93
94    /// replaces the local ICE ufrag registered for a session
95    ///
96    /// each session owns at most one local ufrag
97    /// each local ufrag maps to at most one session
98    /// returning `false` means the existing mapping already expressed that
99    /// contract
100    pub(super) fn remember_local_ice_ufrag(
101        &mut self,
102        local_ice_ufrag: &str,
103        session_key: &TransportSessionKey,
104    ) -> bool {
105        if self
106            .local_ice_ufrag_index
107            .get(local_ice_ufrag)
108            .is_some_and(|current_session| current_session == session_key)
109        {
110            return false;
111        }
112        let previous_ufrag = self
113            .local_ice_ufrag_by_session
114            .insert(session_key.clone(), local_ice_ufrag.to_owned());
115        if let Some(previous_ufrag) = previous_ufrag {
116            self.local_ice_ufrag_index.remove(&previous_ufrag);
117        }
118        let previous_session = self
119            .local_ice_ufrag_index
120            .insert(local_ice_ufrag.to_owned(), session_key.clone());
121        if let Some(previous_session) = previous_session {
122            self.local_ice_ufrag_by_session.remove(&previous_session);
123        }
124        true
125    }
126
127    /// Returns sessions whose signaled candidates match the observed source.
128    ///
129    /// Candidate addresses are weaker than learned source pins because every
130    /// session advertising the address remains a candidate. `Rtc::accepts()`
131    /// remains the ownership authority.
132    #[must_use]
133    pub fn candidates_for_src_addr(
134        &self,
135        source_addr: SocketAddr,
136    ) -> Option<&[TransportSessionKey]> {
137        self.remote_candidate_addr_index
138            .get(&source_addr)
139            .map(Vec::as_slice)
140    }
141
142    /// replaces all signaled remote candidate addresses for one session
143    ///
144    /// this is called after answer application when str0m has updated the
145    /// candidate set
146    /// old hints are removed first so recovery cannot keep probing a session
147    /// through candidates that no longer belong to it
148    pub fn replace_remote_candidates<I>(
149        &mut self,
150        session_key: &TransportSessionKey,
151        candidate_addrs: I,
152    ) where
153        I: IntoIterator<Item = SocketAddr>,
154    {
155        self.forget_user_remote_candidates(session_key);
156        let session_candidate_addrs = self
157            .remote_candidate_addrs_by_session
158            .entry(session_key.clone())
159            .or_default();
160        for candidate_addr in candidate_addrs {
161            if session_candidate_addrs.contains(&candidate_addr) {
162                continue;
163            }
164            session_candidate_addrs.push(candidate_addr);
165            self.remote_candidate_addr_index
166                .entry(candidate_addr)
167                .or_default()
168                .push(session_key.clone());
169        }
170        if session_candidate_addrs.is_empty() {
171            self.remote_candidate_addrs_by_session.remove(session_key);
172        }
173    }
174
175    /// removes one learned UDP source tuple pin after cached `Rtc::accepts()` rejection
176    pub(super) fn forget_remote_addr(&mut self, source_addr: SocketAddr) {
177        let Some(session_key) = self.remote_addr_index.remove(&source_addr) else {
178            return;
179        };
180        self.remove_remote_addr(&session_key, source_addr);
181    }
182
183    /// removes every learned UDP source tuple when its session is removed
184    pub(super) fn forget_user_remote_addrs(&mut self, session_key: &TransportSessionKey) {
185        let Some(session_addrs) = self.remote_addrs_by_session.remove(session_key) else {
186            return;
187        };
188        for source_addr in session_addrs {
189            self.remote_addr_index.remove(&source_addr);
190        }
191    }
192
193    /// removes the local ICE ufrag recovery hint for a session
194    pub(super) fn forget_user_local_ice_ufrag(&mut self, session_key: &TransportSessionKey) {
195        let Some(local_ice_ufrag) = self.local_ice_ufrag_by_session.remove(session_key) else {
196            return;
197        };
198        self.local_ice_ufrag_index.remove(&local_ice_ufrag);
199    }
200
201    /// removes all remote candidate recovery hints owned by a session
202    ///
203    /// candidate address indexes can contain several sessions for one address
204    /// cleanup removes only the target session from each fanout list and drops
205    /// empty address entries afterward
206    pub(super) fn forget_user_remote_candidates(&mut self, session_key: &TransportSessionKey) {
207        let Some(candidate_addrs) = self.remote_candidate_addrs_by_session.remove(session_key)
208        else {
209            return;
210        };
211        for candidate_addr in candidate_addrs {
212            let should_remove_index_entry = self
213                .remote_candidate_addr_index
214                .get_mut(&candidate_addr)
215                .is_some_and(|session_keys| {
216                    session_keys.retain(|candidate_session| candidate_session != session_key);
217                    session_keys.is_empty()
218                });
219            if should_remove_index_entry {
220                self.remote_candidate_addr_index.remove(&candidate_addr);
221            }
222        }
223    }
224
225    #[cfg(test)]
226    pub(super) fn session_addrs_for(
227        &self,
228        session_key: &TransportSessionKey,
229    ) -> Option<&[SocketAddr]> {
230        self.remote_addrs_by_session
231            .get(session_key)
232            .map(Vec::as_slice)
233    }
234
235    #[cfg(test)]
236    pub(super) fn local_ice_ufrag_for(&self, session_key: &TransportSessionKey) -> Option<&str> {
237        self.local_ice_ufrag_by_session
238            .get(session_key)
239            .map(String::as_str)
240    }
241
242    #[cfg(test)]
243    pub(super) fn remote_candidate_addrs_for(
244        &self,
245        session_key: &TransportSessionKey,
246    ) -> Option<&[SocketAddr]> {
247        self.remote_candidate_addrs_by_session
248            .get(session_key)
249            .map(Vec::as_slice)
250    }
251
252    #[cfg(test)]
253    pub(super) fn is_empty(&self) -> bool {
254        self.remote_addrs_by_session.is_empty()
255            && self.local_ice_ufrag_by_session.is_empty()
256            && self.remote_candidate_addrs_by_session.is_empty()
257    }
258
259    fn remove_remote_addr(&mut self, session_key: &TransportSessionKey, source_addr: SocketAddr) {
260        let should_remove_session_entry = self
261            .remote_addrs_by_session
262            .get_mut(session_key)
263            .is_some_and(|session_addrs| {
264                if let Some(position) = session_addrs.iter().position(|addr| *addr == source_addr) {
265                    session_addrs.swap_remove(position);
266                }
267                session_addrs.is_empty()
268            });
269        if should_remove_session_entry {
270            self.remote_addrs_by_session.remove(session_key);
271        }
272    }
273}
274
275#[cfg(test)]
276#[path = "TESTS/demux.rs"]
277mod tests;