1use std::{
49 collections::{BTreeMap, BTreeSet, btree_map::Entry},
50 mem,
51};
52
53use o_sfu_router::topology::RoutedConsumerId;
54
55use super::{
56 ConsumerSourceSelection, SubscriptionKey, consumer_setup::ConsumerSetupTarget,
57 remove_from_index_set,
58};
59use crate::engine::{
60 ConnectionId, MediaWorkerId, UserId,
61 media_transport::{
62 RelayRouteActivity, TransportConsumerRoute, TransportMediaId, TransportRelayRouteAction,
63 TransportSessionKey, TransportSourceKey,
64 },
65 source_model::{PolicyPauseReason, PublishedSourceId, SourceSubscriptionIntent},
66};
67
68#[derive(Debug, Default)]
69pub(super) struct RouteGraph {
70 entries: BTreeMap<SubscriptionKey, Subscription>,
71 by_receiver: BTreeMap<UserId, BTreeSet<SubscriptionKey>>,
72 by_source: BTreeMap<PublishedSourceId, BTreeSet<SubscriptionKey>>,
73 relays: BTreeMap<RelayRouteKey, RelayOwners>,
74 next_reservation: RouteReservationId,
75}
76
77type RelayOwners = BTreeMap<SubscriptionKey, RelayRouteActivity>;
78
79#[derive(Debug, Default)]
80struct Subscription {
81 intent: SourceSubscriptionIntent,
82 current: Option<CurrentPublication>,
83}
84
85#[derive(Debug)]
86pub(super) struct CurrentPublication {
87 pub source_id: PublishedSourceId,
88 pub selection: ConsumerSourceSelection,
89 realization: ConsumerRealization,
90}
91
92#[derive(Debug, Default)]
93enum ConsumerRealization {
94 #[default]
95 Absent,
96 Pending(RouteReservationId, Option<RouteRelay>),
97 Committed(
98 TransportConsumerRoute,
99 String,
100 RoutedConsumerId,
101 Option<RouteRelay>,
102 ),
103}
104
105#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
106struct RouteReservationId(u64);
107
108#[derive(Debug)]
109pub struct ConsumerRouteReservation {
110 key: SubscriptionKey,
111 source_id: PublishedSourceId,
112 selection: ConsumerSourceSelection,
113 id: RouteReservationId,
114}
115
116#[derive(Debug, Clone, PartialEq, Eq)]
117struct RouteRelay {
118 route: RelayRouteKey,
119 activity: RelayRouteActivity,
120}
121
122struct TakenPending {
123 relay: Option<RouteRelay>,
124}
125
126#[derive(Debug, Clone, PartialEq, Eq)]
127pub struct RelayRouteEffect {
128 pub route: RelayRouteKey,
129 pub action: TransportRelayRouteAction,
130}
131
132#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
133pub struct RelayRouteKey {
134 pub source_user: UserId,
135 pub source_connection: ConnectionId,
136 pub source_media: TransportMediaId,
137 pub target_worker: MediaWorkerId,
138}
139
140#[derive(Debug, Default)]
141pub(super) struct RemovedRoutes {
142 pub routes: Vec<TransportConsumerRoute>,
143 pub consumers: Vec<RoutedConsumerId>,
144 pub relays: Vec<RelayRouteEffect>,
145}
146
147impl RemovedRoutes {
148 pub(super) fn extend(&mut self, mut other: Self) {
149 self.routes.append(&mut other.routes);
150 self.consumers.append(&mut other.consumers);
151 self.relays.append(&mut other.relays);
152 }
153}
154
155impl RouteGraph {
156 pub(super) fn subscription_count(&self) -> usize {
157 self.entries
158 .values()
159 .filter(|entry| entry.has_consumer_setup_or_route())
160 .count()
161 }
162
163 #[cfg(any(test, feature = "testing-transport"))]
164 pub(super) fn count(&self) -> usize {
165 self.attached()
166 .filter(|(_, current)| current.committed().is_some())
167 .count()
168 }
169
170 #[cfg(test)]
171 pub(super) fn record_count(&self) -> usize {
172 self.entries.len()
173 }
174
175 pub(super) fn merge_intent(&mut self, key: SubscriptionKey, update: SourceSubscriptionIntent) {
176 if update.is_empty() {
177 return;
178 }
179 let entry = self.entry(key);
180 entry.intent.merge(update);
181 if let (Some(active), Some(current)) = (update.active(), entry.current.as_mut()) {
182 current.selection.set_active(active);
183 }
184 }
185
186 pub(super) fn intent(&self, key: &SubscriptionKey) -> SourceSubscriptionIntent {
187 self.entries
188 .get(key)
189 .map_or_else(SourceSubscriptionIntent::default, |entry| entry.intent)
190 }
191
192 pub(super) fn attach_for_setup(
193 &mut self,
194 key: SubscriptionKey,
195 source_id: PublishedSourceId,
196 ) -> bool {
197 let entry = self.entry(key.clone());
198 if let Some(current) = &entry.current {
199 return current.source_id == source_id
200 && matches!(current.realization, ConsumerRealization::Absent);
201 }
202 entry.current = Some(CurrentPublication {
203 source_id,
204 selection: ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true)),
205 realization: ConsumerRealization::Absent,
206 });
207 self.by_source.entry(source_id).or_default().insert(key);
208 true
209 }
210
211 pub(super) fn set_activity(
212 &mut self,
213 key: &SubscriptionKey,
214 source_id: PublishedSourceId,
215 connection_id: ConnectionId,
216 active: bool,
217 policy_pause_reason: Option<PolicyPauseReason>,
218 ) -> Option<Vec<RelayRouteEffect>> {
219 let relay = {
220 let current = self
221 .entries
222 .get_mut(key)
223 .and_then(|entry| entry.current.as_mut())
224 .filter(|current| current.source_id == source_id)?;
225 if let ConsumerRealization::Committed(route, ..) = ¤t.realization
226 && route.consumer_session_key().connection_id() != connection_id
227 {
228 return None;
229 }
230 current.selection.set_active(active);
231 if let Some(reason) = policy_pause_reason {
232 current.selection.set_policy_pause_reason(Some(reason));
233 }
234 current
235 .realization
236 .set_relay_activity(RelayRouteActivity::from_active(active))
237 };
238 Some(relay.map_or_else(Vec::new, |relay| self.set_relay_owner(key, &relay, false)))
239 }
240
241 pub(super) fn reserve_consumer_setup(
242 &mut self,
243 key: SubscriptionKey,
244 source_id: PublishedSourceId,
245 selection: ConsumerSourceSelection,
246 ) -> Option<ConsumerRouteReservation> {
247 let current = self.entries.get(&key)?.current.as_ref()?;
248 if current.source_id != source_id
249 || !matches!(current.realization, ConsumerRealization::Absent)
250 {
251 return None;
252 }
253 let id = self.next_reservation();
254 let current = self.entries.get_mut(&key)?.current.as_mut()?;
255 current.selection = selection;
256 current.realization = ConsumerRealization::Pending(id, None);
257 Some(ConsumerRouteReservation {
258 key,
259 source_id,
260 selection,
261 id,
262 })
263 }
264
265 pub(super) fn release_consumer_setup(
266 &mut self,
267 reservation: ConsumerRouteReservation,
268 ) -> Vec<RelayRouteEffect> {
269 let Some(TakenPending { relay }) = self.take_pending(&reservation) else {
270 return Vec::new();
271 };
272 let key = reservation.key;
273 relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay))
274 }
275
276 pub(super) fn commit(
277 &mut self,
278 reservation: ConsumerRouteReservation,
279 route: TransportConsumerRoute,
280 mid: String,
281 selection: ConsumerSourceSelection,
282 accept: impl FnOnce() -> Option<RoutedConsumerId>,
283 ) -> Result<(), Vec<RelayRouteEffect>> {
284 let pending = self.take_pending(&reservation);
287 let ConsumerRouteReservation { key, source_id, .. } = reservation;
288 let Some(TakenPending { relay }) = pending else {
289 return Err(Vec::new());
290 };
291 let Some(current) = self.current_mut(&key, source_id) else {
292 return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
293 };
294 let Some(routed) = accept() else {
295 return Err(relay.map_or_else(Vec::new, |relay| self.release_relay(&key, &relay)));
296 };
297 current.selection = selection;
298 current.realization = ConsumerRealization::Committed(route, mid, routed, relay);
299 Ok(())
300 }
301
302 pub(super) fn update_selection(
303 &mut self,
304 key: &SubscriptionKey,
305 source_id: PublishedSourceId,
306 route: &TransportConsumerRoute,
307 update: impl FnOnce(&mut ConsumerSourceSelection),
308 ) -> bool {
309 let Some(current) = self
310 .entries
311 .get_mut(key)
312 .and_then(|entry| entry.current.as_mut())
313 else {
314 return false;
315 };
316 let ConsumerRealization::Committed(current_route, ..) = ¤t.realization else {
317 return false;
318 };
319 if current.source_id != source_id || current_route != route {
320 return false;
321 }
322 update(&mut current.selection);
323 true
324 }
325
326 pub(super) fn selection(
327 &self,
328 key: &SubscriptionKey,
329 source_id: PublishedSourceId,
330 ) -> Option<ConsumerSourceSelection> {
331 let current = self.entries.get(key)?.current.as_ref()?;
332 (current.source_id == source_id).then_some(current.selection)
333 }
334
335 pub(super) fn attached(&self) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
336 self.entries
337 .iter()
338 .filter_map(|(key, entry)| Some((key, entry.current.as_ref()?)))
339 }
340
341 pub(super) fn current(
342 &self,
343 key: &SubscriptionKey,
344 ) -> Option<(&SubscriptionKey, &CurrentPublication)> {
345 let (key, entry) = self.entries.get_key_value(key)?;
346 Some((key, entry.current.as_ref()?))
347 }
348
349 pub(super) fn attached_for_receiver(
350 &self,
351 receiver: &UserId,
352 ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
353 self.by_receiver
354 .get(receiver)
355 .into_iter()
356 .flat_map(BTreeSet::iter)
357 .filter_map(|key| self.current(key))
358 }
359
360 pub(super) fn attached_for_source(
361 &self,
362 source_id: PublishedSourceId,
363 ) -> impl Iterator<Item = (&SubscriptionKey, &CurrentPublication)> {
364 self.by_source
365 .get(&source_id)
366 .into_iter()
367 .flat_map(BTreeSet::iter)
368 .filter_map(|key| self.current(key))
369 }
370
371 pub(super) fn detach_source(&mut self, source_id: PublishedSourceId) -> RemovedRoutes {
372 let keys = self.by_source.remove(&source_id).unwrap_or_default();
373 let mut removed = RemovedRoutes::default();
374 for key in keys {
375 let (realization, prune) = {
376 let Some(entry) = self.entries.get_mut(&key) else {
377 continue;
378 };
379 let Some(current) = entry.current.take() else {
380 continue;
381 };
382 debug_assert_eq!(current.source_id, source_id);
383 (current.realization, entry.intent.is_empty())
384 };
385 self.collect_removed_realization(&key, realization, &mut removed);
386 if prune {
387 self.entries.remove(&key);
388 remove_from_index_set(&mut self.by_receiver, &key.receiver, &key);
389 }
390 }
391 removed
392 }
393
394 pub(super) fn reset_receiver_for_replacement(
395 &mut self,
396 receiver: &UserId,
397 ) -> Vec<RelayRouteEffect> {
398 let keys = self.by_receiver.get(receiver).cloned().unwrap_or_default();
399 let mut relays = Vec::new();
400 for key in keys {
401 let relay = {
402 let Some(entry) = self.entries.get_mut(&key) else {
403 continue;
404 };
405 let Some(current) = entry.current.as_mut() else {
406 continue;
407 };
408 current.selection =
409 ConsumerSourceSelection::open(entry.intent.active().unwrap_or(true));
410 match mem::take(&mut current.realization) {
411 ConsumerRealization::Absent => None,
412 ConsumerRealization::Pending(_, relay)
413 | ConsumerRealization::Committed(_, _, _, relay) => relay,
414 }
415 };
416 if let Some(relay) = relay {
417 relays.extend(self.release_relay(&key, &relay));
418 }
419 }
420 relays
421 }
422
423 pub(super) fn remove_receiver(&mut self, receiver: &UserId) -> RemovedRoutes {
424 let keys = self.by_receiver.remove(receiver).unwrap_or_default();
425 let mut removed = RemovedRoutes::default();
426 for key in keys {
427 let Some(entry) = self.entries.remove(&key) else {
428 continue;
429 };
430 let Some(current) = entry.current else {
431 continue;
432 };
433 remove_from_index_set(&mut self.by_source, ¤t.source_id, &key);
434 self.collect_removed_realization(&key, current.realization, &mut removed);
435 }
436 removed
437 }
438
439 pub(super) fn detach_declined_consumers(
440 &mut self,
441 session: &TransportSessionKey,
442 declined: &[TransportMediaId],
443 ) -> RemovedRoutes {
444 let keys = self
445 .by_receiver
446 .get(session.user_id())
447 .cloned()
448 .unwrap_or_default();
449 let mut removed = RemovedRoutes::default();
450 for key in keys {
451 let realization = {
452 let Some(current) = self
453 .entries
454 .get_mut(&key)
455 .and_then(|entry| entry.current.as_mut())
456 else {
457 continue;
458 };
459 let ConsumerRealization::Committed(route, ..) = ¤t.realization else {
460 continue;
461 };
462 if route.consumer_session_key() != session
463 || !declined.contains(&route.consumer_transport_media_id())
464 {
465 continue;
466 }
467 mem::take(&mut current.realization)
468 };
469 self.collect_removed_realization(&key, realization, &mut removed);
470 }
471 removed
472 }
473
474 pub(super) fn reserve_relay(
475 &mut self,
476 reservation: &ConsumerRouteReservation,
477 target: &ConsumerSetupTarget,
478 target_worker: MediaWorkerId,
479 active: bool,
480 ) -> Vec<RelayRouteEffect> {
481 let (previous, relay) = {
482 let Some(current) = self.current_mut_for(reservation) else {
483 return Vec::new();
484 };
485 let ConsumerRealization::Pending(id, relay) = &mut current.realization else {
486 return Vec::new();
487 };
488 if *id != reservation.id {
489 return Vec::new();
490 }
491 let next = RouteRelay {
492 route: target.relay_route_key(target_worker),
493 activity: RelayRouteActivity::from_active(active),
494 };
495 if relay.as_ref() == Some(&next) {
496 return Vec::new();
497 }
498 (relay.replace(next.clone()), next)
499 };
500 self.replace_relay(&reservation.key, previous, &relay)
501 }
502
503 pub(super) fn source_activity_target_workers<'a>(
504 &'a self,
505 source: &'a TransportSourceKey,
506 ) -> impl Iterator<Item = MediaWorkerId> + 'a {
507 let session = source.session_key();
508 self.relays
509 .keys()
510 .filter(move |route| {
511 route.source_user == *session.user_id()
512 && route.source_connection == session.connection_id()
513 && route.source_media == source.transport_media_id()
514 })
515 .map(|route| route.target_worker)
516 }
517
518 fn entry(&mut self, key: SubscriptionKey) -> &mut Subscription {
519 self.by_receiver
520 .entry(key.receiver.clone())
521 .or_default()
522 .insert(key.clone());
523 self.entries.entry(key).or_default()
524 }
525
526 fn next_reservation(&mut self) -> RouteReservationId {
527 self.next_reservation.0 += 1;
528 self.next_reservation
529 }
530
531 fn current_mut(
532 &mut self,
533 key: &SubscriptionKey,
534 source_id: PublishedSourceId,
535 ) -> Option<&mut CurrentPublication> {
536 let current = self.entries.get_mut(key)?.current.as_mut()?;
537 (current.source_id == source_id).then_some(current)
538 }
539
540 fn current_mut_for(
541 &mut self,
542 reservation: &ConsumerRouteReservation,
543 ) -> Option<&mut CurrentPublication> {
544 self.current_mut(&reservation.key, reservation.source_id)
545 }
546
547 fn take_pending(&mut self, reservation: &ConsumerRouteReservation) -> Option<TakenPending> {
548 let current = self.current_mut_for(reservation)?;
549 let pending = mem::take(&mut current.realization);
550 match pending {
551 ConsumerRealization::Pending(id, relay) if id == reservation.id => {
552 Some(TakenPending { relay })
553 }
554 other => {
555 current.realization = other;
556 None
557 }
558 }
559 }
560
561 fn collect_removed_realization(
562 &mut self,
563 key: &SubscriptionKey,
564 realization: ConsumerRealization,
565 removed: &mut RemovedRoutes,
566 ) {
567 let relay = match realization {
568 ConsumerRealization::Absent => None,
569 ConsumerRealization::Pending(_, relay) => relay,
570 ConsumerRealization::Committed(route, _, consumer, relay) => {
571 removed.routes.push(route);
572 removed.consumers.push(consumer);
573 relay
574 }
575 };
576 if let Some(relay) = relay {
577 removed.relays.extend(self.release_relay(key, &relay));
578 }
579 }
580
581 fn replace_relay(
582 &mut self,
583 key: &SubscriptionKey,
584 previous: Option<RouteRelay>,
585 relay: &RouteRelay,
586 ) -> Vec<RelayRouteEffect> {
587 match previous {
588 None => self.set_relay_owner(key, relay, true),
589 Some(previous) if previous.route == relay.route => {
590 self.set_relay_owner(key, relay, false)
591 }
592 Some(previous) => {
593 let mut effects = self.release_relay(key, &previous);
594 effects.extend(self.set_relay_owner(key, relay, true));
595 effects
596 }
597 }
598 }
599
600 fn set_relay_owner(
601 &mut self,
602 key: &SubscriptionKey,
603 relay: &RouteRelay,
604 insert_missing: bool,
605 ) -> Vec<RelayRouteEffect> {
606 let route = relay.route.clone();
607 let owners = match self.relays.entry(route.clone()) {
608 Entry::Occupied(entry) => entry.into_mut(),
609 Entry::Vacant(entry) if insert_missing => entry.insert(RelayOwners::default()),
610 Entry::Vacant(_) => return Vec::new(),
611 };
612 let before = relay_aggregate(owners);
613 owners.insert(key.clone(), relay.activity);
614 relay_effects_for(route, before, relay_aggregate(owners))
615 }
616
617 fn release_relay(
618 &mut self,
619 key: &SubscriptionKey,
620 relay: &RouteRelay,
621 ) -> Vec<RelayRouteEffect> {
622 let route = relay.route.clone();
623 let Some(owners) = self.relays.get_mut(&route) else {
624 return Vec::new();
625 };
626 let before = relay_aggregate(owners);
627 if owners.remove(key).is_none() {
628 return Vec::new();
629 }
630 let after = relay_aggregate(owners);
631 if after.is_none() {
632 self.relays.remove(&route);
633 }
634 relay_effects_for(route, before, after)
635 }
636}
637
638impl ConsumerRouteReservation {
639 pub const fn selection(&self) -> ConsumerSourceSelection {
640 self.selection
641 }
642}
643
644impl Subscription {
645 fn has_consumer_setup_or_route(&self) -> bool {
646 self.current
647 .as_ref()
648 .is_some_and(|current| !matches!(current.realization, ConsumerRealization::Absent))
649 }
650}
651
652impl CurrentPublication {
653 pub(super) const fn is_pending(&self) -> bool {
654 matches!(self.realization, ConsumerRealization::Pending(..))
655 }
656
657 pub(super) fn committed(&self) -> Option<(&TransportConsumerRoute, &str)> {
658 match &self.realization {
659 ConsumerRealization::Committed(route, mid, ..) => Some((route, mid)),
660 ConsumerRealization::Absent | ConsumerRealization::Pending(..) => None,
661 }
662 }
663}
664
665impl ConsumerRealization {
666 fn set_relay_activity(&mut self, activity: RelayRouteActivity) -> Option<RouteRelay> {
667 let relay = match self {
668 Self::Absent => return None,
669 Self::Pending(_, relay) | Self::Committed(_, _, _, relay) => relay.as_mut()?,
670 };
671 if relay.activity == activity {
672 return None;
673 }
674 relay.activity = activity;
675 Some(relay.clone())
676 }
677}
678
679fn relay_aggregate(owners: &RelayOwners) -> Option<RelayRouteActivity> {
680 (!owners.is_empty()).then(|| {
681 RelayRouteActivity::from_active(owners.values().copied().any(RelayRouteActivity::is_active))
682 })
683}
684
685fn relay_effects_for(
686 route: RelayRouteKey,
687 before: Option<RelayRouteActivity>,
688 after: Option<RelayRouteActivity>,
689) -> Vec<RelayRouteEffect> {
690 let Some(activity) = after else {
691 return vec![RelayRouteEffect {
692 route,
693 action: TransportRelayRouteAction::Release,
694 }];
695 };
696 let mut effects = Vec::new();
697 if before.is_none() {
698 effects.push(RelayRouteEffect {
699 route: route.clone(),
700 action: TransportRelayRouteAction::Install,
701 });
702 }
703 if before.unwrap_or(RelayRouteActivity::Inactive) != activity {
704 effects.push(RelayRouteEffect {
705 route,
706 action: TransportRelayRouteAction::SetActivity(activity),
707 });
708 }
709 effects
710}