Skip to main content

o_sfu_core/engine/media_transport/
policy_invalidation.rs

1//! Coalesces transport observations and absolute policy deadlines into room wakeups.
2
3use 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/// Shared drain for coalesced room source-policy invalidations.
52///
53/// Clones share one drain. The runtime must assign all clones to a single
54/// consumer task. Dropping a wait leaves room work and deadlines in the drain.
55#[derive(Debug, Clone)]
56pub struct SourcePolicyUpdateSubscription(Arc<SourcePolicyUpdates>);
57
58impl SourcePolicyUpdateSubscription {
59    /// Waits until immediate work or an absolute deadline requires a policy pass.
60    pub async fn wait_for_update(&self) -> BTreeSet<RoomInstanceId> {
61        loop {
62            // Retain the notification across the timer race and drain again
63            // before waiting. A timer winner must not consume an unseen dirty wake.
64            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    /// Drains immediate updates and expired deadlines after the previous wait.
85    #[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/// Sender for coalesced room source-policy updates.
92#[derive(Debug, Clone, Default)]
93pub struct SourcePolicySignal(Arc<SourcePolicyUpdates>);
94
95impl SourcePolicySignal {
96    /// Creates the runtime's single-consumer subscription.
97    #[must_use]
98    pub fn subscribe(&self) -> SourcePolicyUpdateSubscription {
99        SourcePolicyUpdateSubscription(Arc::clone(&self.0))
100    }
101
102    /// Marks one room as needing a source-policy pass.
103    pub fn mark_dirty(&self, room_instance_id: RoomInstanceId) {
104        self.mark_dirty_rooms([room_instance_id]);
105    }
106
107    /// Marks rooms dirty without changing their independently scheduled deadlines.
108    /// Empty input from packet-pump turns does not acquire the shared lock.
109    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    /// Preserves an earlier room deadline while retrying unavailable observations.
125    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    /// Replaces or cancels a room's earliest deadline. Equal deadlines are a no-op.
148    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    /// Replaces or cancels the room deadline recomputed by its ordered policy turn.
175    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;