o_sfu_core/engine/media_transport/
policy_invalidation.rs1use std::{
4 collections::{BTreeMap, BTreeSet},
5 mem,
6 sync::{Arc, Mutex},
7 time::{Duration, Instant},
8};
9
10use tokio::{
11 sync::Notify,
12 time::{Instant as TokioInstant, sleep_until},
13};
14
15use super::MediaTransport;
16use crate::{RoomInstanceId, engine::sync::lock_unpoisoned};
17
18const OBSERVATION_RETRY_DELAY: Duration = Duration::from_secs(1);
19
20#[derive(Debug, Default)]
21struct PendingSourcePolicyUpdates {
22 rooms: BTreeSet<RoomInstanceId>,
23 scheduled_by_room: BTreeMap<RoomInstanceId, Instant>,
24 deadlines: BTreeSet<(Instant, RoomInstanceId)>,
25}
26
27impl PendingSourcePolicyUpdates {
28 fn next_deadline(&self) -> Option<Instant> {
29 self.deadlines.first().map(|(deadline, _)| *deadline)
30 }
31
32 fn take(&mut self, now: Instant) -> BTreeSet<RoomInstanceId> {
33 while let Some(&(deadline, room)) = self.deadlines.first() {
34 if deadline > now {
35 break;
36 }
37 self.deadlines.pop_first();
38 self.scheduled_by_room.remove(&room);
39 self.rooms.insert(room);
40 }
41 mem::take(&mut self.rooms)
42 }
43}
44
45#[derive(Debug, Default)]
46struct SourcePolicyUpdates {
47 pending: Mutex<PendingSourcePolicyUpdates>,
48 notify: Notify,
49}
50
51#[derive(Debug, Clone)]
56pub struct SourcePolicyUpdateSubscription(Arc<SourcePolicyUpdates>);
57
58impl SourcePolicyUpdateSubscription {
59 pub async fn wait_for_update(&self) -> BTreeSet<RoomInstanceId> {
61 loop {
62 let notified = self.0.notify.notified();
65 let deadline = {
66 let mut pending = lock_unpoisoned(&self.0.pending);
67 let rooms = pending.take(TokioInstant::now().into_std());
68 if !rooms.is_empty() {
69 return rooms;
70 }
71 pending.next_deadline()
72 };
73 if let Some(deadline) = deadline {
74 tokio::select! {
75 () = notified => {},
76 () = sleep_until(deadline.into()) => {},
77 }
78 } else {
79 notified.await;
80 }
81 }
82 }
83
84 #[must_use]
86 pub fn take_pending_updates(&self) -> BTreeSet<RoomInstanceId> {
87 lock_unpoisoned(&self.0.pending).take(TokioInstant::now().into_std())
88 }
89}
90
91#[derive(Debug, Clone, Default)]
93pub struct SourcePolicySignal(Arc<SourcePolicyUpdates>);
94
95impl SourcePolicySignal {
96 #[must_use]
98 pub fn subscribe(&self) -> SourcePolicyUpdateSubscription {
99 SourcePolicyUpdateSubscription(Arc::clone(&self.0))
100 }
101
102 pub fn mark_dirty(&self, room_instance_id: RoomInstanceId) {
104 self.mark_dirty_rooms([room_instance_id]);
105 }
106
107 pub fn mark_dirty_rooms(&self, room_instance_ids: impl IntoIterator<Item = RoomInstanceId>) {
110 let mut room_instance_ids = room_instance_ids.into_iter();
111 let Some(first) = room_instance_ids.next() else {
112 return;
113 };
114 let mut pending = lock_unpoisoned(&self.0.pending);
115 let notify = pending.rooms.is_empty();
116 pending.rooms.insert(first);
117 pending.rooms.extend(room_instance_ids);
118 drop(pending);
119 if notify {
120 self.0.notify.notify_one();
121 }
122 }
123
124 fn schedule_observation_retry(&self, room: RoomInstanceId) {
126 let retry_at = Instant::now() + OBSERVATION_RETRY_DELAY;
127 let mut pending = lock_unpoisoned(&self.0.pending);
128 if pending
129 .scheduled_by_room
130 .get(&room)
131 .is_some_and(|deadline| *deadline <= retry_at)
132 {
133 return;
134 }
135 let previous_earliest = pending.next_deadline();
136 if let Some(previous) = pending.scheduled_by_room.insert(room, retry_at) {
137 pending.deadlines.remove(&(previous, room));
138 }
139 pending.deadlines.insert((retry_at, room));
140 let notify = previous_earliest != pending.next_deadline();
141 drop(pending);
142 if notify {
143 self.0.notify.notify_one();
144 }
145 }
146
147 fn set_deadline(&self, room: RoomInstanceId, deadline: Option<Instant>) {
149 let mut pending = lock_unpoisoned(&self.0.pending);
150 if pending.scheduled_by_room.get(&room).copied() == deadline {
151 return;
152 }
153 let previous_earliest = pending.next_deadline();
154 if let Some(previous) = pending.scheduled_by_room.remove(&room) {
155 pending.deadlines.remove(&(previous, room));
156 }
157 if let Some(deadline) = deadline {
158 pending.scheduled_by_room.insert(room, deadline);
159 pending.deadlines.insert((deadline, room));
160 }
161 let notify = previous_earliest != pending.next_deadline();
162 drop(pending);
163 if notify {
164 self.0.notify.notify_one();
165 }
166 }
167}
168
169impl MediaTransport {
170 pub(in crate::engine) fn schedule_source_policy_retry(&self, room: RoomInstanceId) {
171 self.source_policy_signal.schedule_observation_retry(room);
172 }
173
174 pub(in crate::engine) fn set_source_policy_deadline(
176 &self,
177 room: RoomInstanceId,
178 deadline: Option<Instant>,
179 ) {
180 self.source_policy_signal.set_deadline(room, deadline);
181 }
182}
183
184#[cfg(test)]
185#[path = "TESTS/policy_invalidation.rs"]
186mod tests;