Skip to main content

o_sfu_core/engine/room/source_policy/
action.rs

1use std::time::Instant;
2
3use super::super::media_graph::{PendingUpgrade, SubscriptionKey};
4use crate::{
5    Bitrate,
6    engine::{
7        ConnectionId, UserId,
8        media_transport::{
9            ConsumerActivity, ConsumerRouteControl, SourcePacketGate, TransportConsumerRoute,
10        },
11        source_model::{
12            ConsumerSourceSelection, PolicyPauseReason, PublishedSourceId,
13            ReceiverVideoBudgetDiagnostics, SourceSelector,
14        },
15    },
16};
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
19pub(super) enum VideoRouteTransition {
20    Degraded,
21    Paused { reason: PolicyPauseReason },
22    Resumed { cleared_reason: PolicyPauseReason },
23}
24
25#[derive(Debug, Clone, Copy)]
26pub(super) struct VideoRouteAllocationState {
27    pub(super) selector: SourceSelector,
28    pub(super) policy_pause_reason: Option<PolicyPauseReason>,
29    pub(super) selected_bitrate: Bitrate,
30}
31
32impl VideoRouteAllocationState {
33    pub(super) fn matches_selection(self, selection: ConsumerSourceSelection) -> bool {
34        self.selector == selection.selector()
35            && self.policy_pause_reason == selection.policy_pause_reason()
36    }
37}
38
39#[derive(Debug)]
40pub(super) struct VideoRouteAllocation {
41    pub(super) key: SubscriptionKey,
42    pub(super) source_id: PublishedSourceId,
43    pub(super) route: TransportConsumerRoute,
44    pub(super) interrupts_upgrade: bool,
45    pub(super) captured: VideoRouteAllocationState,
46    pub(super) planned: VideoRouteAllocationState,
47}
48
49#[derive(Debug)]
50pub(super) struct ReceiverVideoBudgetPlan {
51    pub(super) receiver: UserId,
52    pub(super) planned_budget: ReceiverVideoBudgetDiagnostics,
53    /// Ordered by [`SubscriptionKey`] for linear reconciliation after transport work.
54    pub(super) routes: Vec<VideoRouteAllocation>,
55}
56
57/// Receiver timing captured before transport work, including cancellation.
58#[derive(Debug)]
59pub(super) struct ReceiverPolicyTiming {
60    pub receiver: UserId,
61    pub connection_id: ConnectionId,
62    pub soft_pause_deadline: Option<Instant>,
63    pub next_deadline: Option<Instant>,
64    /// Existing eligible wakeups survive rejection of an unrelated control.
65    pub retained_deadline: Option<Instant>,
66}
67
68/// Only changed upgrade state requires validation and a topology write.
69#[derive(Debug, Clone, Copy, PartialEq, Eq)]
70pub(super) enum UpgradeChange {
71    Unchanged,
72    Set(Option<PendingUpgrade>),
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub(in crate::engine::room) struct ConsumerPacketSelectionUpdate {
77    pub(super) key: SubscriptionKey,
78    pub(super) source_id: PublishedSourceId,
79    pub(in crate::engine::room) route: TransportConsumerRoute,
80    pub(super) selector: SourceSelector,
81    pub(super) policy_pause_reason: Option<PolicyPauseReason>,
82    pub(super) planned_budget: ReceiverVideoBudgetDiagnostics,
83    pub(super) transition: Option<VideoRouteTransition>,
84    pub(super) selected_estimated_bitrate: Option<Bitrate>,
85    pub(super) upgrade: UpgradeChange,
86    pub(super) interrupts_upgrade: bool,
87    pub(super) packet_gate: Option<SourcePacketGate>,
88    pub(super) route_activity_changed: bool,
89    pub(super) request_keyframe: bool,
90}
91
92impl ConsumerPacketSelectionUpdate {
93    pub(in crate::engine::room) fn route_activity(
94        key: SubscriptionKey,
95        source_id: PublishedSourceId,
96        route: TransportConsumerRoute,
97        current_selection: ConsumerSourceSelection,
98        policy_pause_reason: Option<PolicyPauseReason>,
99    ) -> Option<Self> {
100        (policy_pause_reason != current_selection.policy_pause_reason()).then(|| Self {
101            key,
102            source_id,
103            route,
104            selector: current_selection.selector(),
105            policy_pause_reason,
106            planned_budget: current_selection.budget(),
107            transition: None,
108            selected_estimated_bitrate: None,
109            upgrade: UpgradeChange::Unchanged,
110            interrupts_upgrade: false,
111            packet_gate: None,
112            route_activity_changed: true,
113            request_keyframe: false,
114        })
115    }
116
117    pub(super) const fn requires_media_transport_effect(&self) -> bool {
118        self.packet_gate.is_some() || self.route_activity_changed || self.request_keyframe
119    }
120
121    pub(in crate::engine::room) fn route_control(&self) -> ConsumerRouteControl {
122        let mut control =
123            ConsumerRouteControl::new(self.route.clone()).request_keyframe(self.request_keyframe);
124        if self.route_activity_changed {
125            control = control.activity(ConsumerActivity::from_active(self.route_active()));
126        }
127        if let Some(packet_gate) = &self.packet_gate {
128            control = control.packet_gate(packet_gate.clone());
129        }
130        control
131    }
132
133    pub(in crate::engine::room) const fn route_active(&self) -> bool {
134        self.policy_pause_reason.is_none()
135    }
136}
137
138#[derive(Debug, Clone, PartialEq, Eq)]
139pub(super) struct FeaturedUserUpdate {
140    pub(super) user_id: UserId,
141    pub(super) connection_id: ConnectionId,
142    pub(super) featured: Option<bool>,
143}