Skip to main content

o_sfu_core/engine/room/source_policy/video/
solver.rs

1//! receiver-video policy turn
2//!
3//! [`SourcePolicyTransaction`]: filter -> budget -> plan -> admit -> fit -> dwell -> projection
4//! packet-gate changes stay behind [`projection`] so planning never builds transport gates directly
5
6use std::{cmp::Reverse, time::Instant};
7
8use itertools::Itertools;
9
10use super::{
11    super::{
12        action::{
13            ReceiverPolicyTiming, ReceiverVideoBudgetPlan, VideoRouteAllocation,
14            VideoRouteAllocationState,
15        },
16        input::SourcePolicySnapshot,
17        turn::SourcePolicyTransaction,
18    },
19    input::{ReceiverVideoRouteInput, receiver_video_routes},
20    projection,
21};
22use crate::{
23    Bitrate, VideoAdaptationTuning,
24    engine::{
25        media_transport::ReceiverBweTargetUpdate,
26        room::{media_graph::PendingUpgrade, state::RoomState},
27        source_model::{
28            PolicyPauseReason, PublishedSourceDescriptor, PublishedSourceId,
29            ReceiverVideoBudgetDiagnostics, SourceAdaptationPolicy, SourceEncodingDescriptor,
30            SourceRoomPolicySelector, SourceRoutePriority, SourceSelector,
31        },
32    },
33};
34
35/// ordering key shared by policy refresh and pre-setup receiver admission
36#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
37pub struct VideoAdmissionRank {
38    /// lower values admit more important layout roles first
39    priority: u8,
40    /// active-speaker order, or `usize::MAX` for non-speakers
41    active_speaker_rank: usize,
42    /// deterministic tie-breaker after priority and speaker rank
43    source_id: u64,
44}
45
46impl VideoAdmissionRank {
47    pub const fn new(
48        priority: SourceRoutePriority,
49        active_speaker_rank: Option<usize>,
50        source_id: PublishedSourceId,
51    ) -> Self {
52        Self {
53            priority: match priority {
54                SourceRoutePriority::PinnedOrFeatured => 0,
55                SourceRoutePriority::ReadableDetail => 1,
56                SourceRoutePriority::ActiveSpeaker => 2,
57                SourceRoutePriority::VisibleThumbnail => 3,
58                SourceRoutePriority::HiddenOrOverflow => 4,
59            },
60            active_speaker_rank: match active_speaker_rank {
61                Some(rank) => rank,
62                None => usize::MAX,
63            },
64            source_id: source_id.as_u64(),
65        }
66    }
67}
68
69/// candidate selector and policy pause state before projection
70#[derive(Debug, Clone, Copy, PartialEq, Eq)]
71pub(super) struct ReceiverRouteSelection {
72    /// packet selector to commit if the route stays deliverable
73    pub(super) selector: SourceSelector,
74    /// policy reason that keeps receiver intent while blocking delivery
75    pub(super) policy_pause_reason: Option<PolicyPauseReason>,
76    pub(super) pending_upgrade: Option<PendingUpgrade>,
77    /// decoder refresh requested after selector changes or delivery resumes
78    pub(super) request_keyframe: bool,
79}
80
81impl ReceiverRouteSelection {
82    pub(super) fn interrupts_upgrade(self, current: Option<&PendingUpgrade>) -> bool {
83        current.is_some_and(|current| {
84            self.pending_upgrade.map_or_else(
85                || self.policy_pause_reason.is_some() || current.selector != self.selector,
86                |pending| current.selector != pending.selector,
87            )
88        })
89    }
90
91    const fn send(selector: SourceSelector, request_keyframe: bool) -> Self {
92        Self {
93            selector,
94            policy_pause_reason: None,
95            pending_upgrade: None,
96            request_keyframe,
97        }
98    }
99
100    const fn pause(selector: SourceSelector, reason: PolicyPauseReason) -> Self {
101        Self {
102            selector,
103            policy_pause_reason: Some(reason),
104            pending_upgrade: None,
105            request_keyframe: false,
106        }
107    }
108}
109
110/// planned receiver route passed to [`projection`] after admission and budget pressure
111///
112/// `selected_bitrate` and `selection` must change together through
113/// [`Self::apply_resolved_selection`], [`Self::send`] or [`Self::pause`]
114#[derive(Debug, Clone, Copy)]
115pub(super) struct PlannedReceiverRoute<'a> {
116    /// immutable route facts captured before policy mutation
117    pub(super) input: &'a ReceiverVideoRouteInput<'a>,
118    /// bitrate counted against receiver budget after this policy step
119    pub(super) selected_bitrate: Bitrate,
120    /// candidate selector, pause state and exact upgrade target
121    pub(super) selection: ReceiverRouteSelection,
122}
123
124impl<'a> PlannedReceiverRoute<'a> {
125    fn new(input: &'a ReceiverVideoRouteInput<'a>, selection: ReceiverRouteSelection) -> Self {
126        let selected_bitrate = if selection.policy_pause_reason.is_some() {
127            Bitrate::zero()
128        } else {
129            selector_bitrate(input, selection.selector).unwrap_or_default()
130        };
131        Self {
132            input,
133            selected_bitrate,
134            selection,
135        }
136    }
137
138    fn apply_resolved_selection(&mut self, selection: ReceiverRouteSelection) {
139        *self = Self::new(self.input, selection);
140    }
141
142    fn send(&mut self, selector: SourceSelector, selected_bitrate: Bitrate) {
143        self.selected_bitrate = selected_bitrate;
144        self.selection = ReceiverRouteSelection::send(selector, self.selection.request_keyframe);
145    }
146
147    fn pause(&mut self, reason: PolicyPauseReason) {
148        self.selected_bitrate = Bitrate::zero();
149        self.selection = ReceiverRouteSelection::pause(self.selection.selector, reason);
150    }
151}
152
153pub(in crate::engine::room::source_policy) fn append_receiver_video_policy(
154    tx: &mut SourcePolicyTransaction,
155    state: &RoomState,
156    input: &SourcePolicySnapshot<'_>,
157    now: Instant,
158) {
159    let routes = receiver_video_routes(state, input);
160    let max_video_updates = routes.len();
161    let max_video_downloads_per_receiver = input.media_limits.max_video_downloads_per_receiver();
162    let tuning = input.video_adaptation_tuning;
163    // `committed_consumer_routes` is ordered by `SubscriptionKey` with
164    // `receiver` first. Preserve that order so `chunk_by` sees each complete
165    // receiver allocation.
166    let mut receiver_routes = routes
167        .chunk_by(|left, right| left.key.receiver == right.key.receiver)
168        .peekable();
169    let mut receiver_bwe_targets = Vec::with_capacity(state.users.len());
170    for (receiver, user) in &state.users {
171        let mut target = input
172            .audio_reserve_by_connection
173            .get(&user.connection_id)
174            .copied()
175            .unwrap_or_else(Bitrate::zero);
176        if let Some(routes) = receiver_routes.next_if(|routes| {
177            routes
178                .first()
179                .is_some_and(|route| &route.key.receiver == receiver)
180        }) {
181            let video_demand = append_receiver_policy_updates(
182                tx,
183                routes,
184                max_video_updates,
185                max_video_downloads_per_receiver,
186                tuning,
187                user.video_soft_pause_deadline,
188                now,
189            );
190            // The target includes audio reserve. Add eventual admitted video
191            // demand so str0m can probe while overload pauses routes.
192            target = target.saturating_add(video_demand);
193        } else if user.video_soft_pause_deadline.is_some() {
194            tx.push_receiver_timing(ReceiverPolicyTiming {
195                receiver: receiver.clone(),
196                connection_id: user.connection_id,
197                soft_pause_deadline: None,
198                next_deadline: None,
199                retained_deadline: None,
200            });
201        }
202        // Include receivers without selected media to clear previous BWE demand.
203        receiver_bwe_targets.push(ReceiverBweTargetUpdate::new(
204            state.transport_user_key(receiver, user.connection_id),
205            target,
206        ));
207    }
208    tx.set_receiver_bwe_targets(receiver_bwe_targets);
209}
210
211/// Fits hard-admitted targets before applying time-based continuity holds.
212///
213/// Eligible downsteps commit immediately. Soft overload may exceed the receiver
214/// budget during its dwell. Admitted BWE demand remains independent of that hold.
215fn append_receiver_policy_updates<'a>(
216    tx: &mut SourcePolicyTransaction,
217    receiver_routes: &'a [ReceiverVideoRouteInput<'a>],
218    max_video_updates: usize,
219    max_video_downloads_per_receiver: usize,
220    tuning: VideoAdaptationTuning,
221    current_soft_pause_deadline: Option<Instant>,
222    now: Instant,
223) -> Bitrate {
224    let Some(first_route) = receiver_routes.first() else {
225        return Bitrate::zero();
226    };
227    let receiver_bandwidth = receiver_routes
228        .iter()
229        .find_map(|route| route.receiver_bandwidth);
230    let audio_reserve = first_route.audio_budget_reserve;
231    let video_budget = receiver_bandwidth
232        .map(|bandwidth| effective_video_budget(bandwidth, tuning, audio_reserve));
233    let mut planned_routes = receiver_routes
234        .iter()
235        .map(|route| PlannedReceiverRoute::new(route, route_plan(route, tuning)))
236        .collect::<Vec<_>>();
237    // Apply the hard route count before sharing bandwidth. Rejected routes must
238    // not consume receiver budget or force admitted routes down.
239    apply_video_download_limit(&mut planned_routes, max_video_downloads_per_receiver);
240    // str0m caps probes from the desired target. Use each hard-admitted route's
241    // highest allowed layer so current BWE cannot bound future discovery.
242    let eventual_admitted_video_bitrate = admitted_video_bwe_demand(&planned_routes);
243    if let Some(video_budget) = video_budget {
244        apply_overload_policy(&mut planned_routes, video_budget);
245    }
246    append_receiver_dwell(
247        tx,
248        &mut planned_routes,
249        current_soft_pause_deadline,
250        tuning,
251        now,
252    );
253    let planned_budget =
254        receiver_video_budget_diagnostics(&planned_routes, receiver_bandwidth, video_budget);
255    if planned_routes.iter().any(route_controls_changed) {
256        tx.push_receiver_video_budget_plan(ReceiverVideoBudgetPlan {
257            receiver: first_route.key.receiver.clone(),
258            planned_budget,
259            routes: planned_routes.iter().map(video_route_allocation).collect(),
260        });
261    }
262    for planned_route in planned_routes {
263        let Some(update) =
264            projection::consumer_packet_selection_update(&planned_route, planned_budget)
265        else {
266            if route_controls_changed(&planned_route) {
267                tx.expect_route_update(
268                    planned_route
269                        .input
270                        .route
271                        .consumer_session_key()
272                        .connection_id(),
273                );
274            }
275            continue;
276        };
277        if update.requires_media_transport_effect() {
278            tx.push_route_update(update);
279        } else {
280            tx.push_state_update(update, max_video_updates);
281        }
282    }
283    eventual_admitted_video_bitrate.max(planned_budget.selected_video_bitrate())
284}
285
286fn route_controls_changed(route: &PlannedReceiverRoute<'_>) -> bool {
287    let current = route.input.current_selection;
288    current.selector() != route.selection.selector
289        || current.policy_pause_reason() != route.selection.policy_pause_reason
290}
291
292/// Resolves post-fit continuity and carries new plus retained receiver wakeups.
293fn append_receiver_dwell(
294    tx: &mut SourcePolicyTransaction,
295    planned_routes: &mut [PlannedReceiverRoute<'_>],
296    current_soft_pause_deadline: Option<Instant>,
297    tuning: VideoAdaptationTuning,
298    now: Instant,
299) {
300    let soft_pause_deadline = planned_routes
301        .iter()
302        .any(|route| {
303            route
304                .selection
305                .policy_pause_reason
306                .is_some_and(is_soft_pause)
307        })
308        .then(|| current_soft_pause_deadline.unwrap_or_else(|| now + tuning.soft_pause_dwell));
309    let mut next_deadline = soft_pause_deadline.filter(|deadline| *deadline > now);
310    let mut retained_deadline =
311        next_deadline.filter(|deadline| Some(*deadline) == current_soft_pause_deadline);
312    for route in planned_routes.iter_mut() {
313        let selection = resolve_dwell(route, tuning, soft_pause_deadline, now);
314        if let Some(upgrade) = selection
315            .pending_upgrade
316            .filter(|upgrade| upgrade.deadline > now)
317        {
318            next_deadline = Some(
319                next_deadline.map_or(upgrade.deadline, |deadline| deadline.min(upgrade.deadline)),
320            );
321            if route.input.pending_upgrade == Some(&upgrade) {
322                retained_deadline = Some(
323                    retained_deadline
324                        .map_or(upgrade.deadline, |deadline| deadline.min(upgrade.deadline)),
325                );
326            }
327        }
328        route.apply_resolved_selection(selection);
329    }
330    if current_soft_pause_deadline.is_some()
331        || soft_pause_deadline.is_some()
332        || next_deadline.is_some()
333    {
334        let Some(first_route) = planned_routes.first().map(|route| route.input) else {
335            return;
336        };
337        tx.push_receiver_timing(ReceiverPolicyTiming {
338            receiver: first_route.key.receiver.clone(),
339            connection_id: first_route.route.consumer_session_key().connection_id(),
340            soft_pause_deadline,
341            next_deadline,
342            retained_deadline,
343        });
344    }
345}
346
347fn video_route_allocation(route: &PlannedReceiverRoute<'_>) -> VideoRouteAllocation {
348    let input = route.input;
349    let captured = input.current_selection;
350    let captured_selected_bitrate = if captured.policy_pause_reason().is_some() {
351        Bitrate::zero()
352    } else {
353        selector_bitrate(input, captured.selector()).unwrap_or_default()
354    };
355    VideoRouteAllocation {
356        key: input.key.clone(),
357        source_id: input.source.source_id(),
358        route: input.route.clone(),
359        interrupts_upgrade: route.selection.interrupts_upgrade(input.pending_upgrade),
360        captured: VideoRouteAllocationState {
361            selector: captured.selector(),
362            policy_pause_reason: captured.policy_pause_reason(),
363            selected_bitrate: captured_selected_bitrate,
364        },
365        planned: VideoRouteAllocationState {
366            selector: route.selection.selector,
367            policy_pause_reason: route.selection.policy_pause_reason,
368            selected_bitrate: route.selected_bitrate,
369        },
370    }
371}
372
373fn route_plan(
374    route: &ReceiverVideoRouteInput<'_>,
375    tuning: VideoAdaptationTuning,
376) -> ReceiverRouteSelection {
377    let policy = route.source.policy().adaptation();
378    let source_cap = route.source.policy().video_bitrate_cap();
379    if source_cap.is_some_and(|cap| source_exceeds_bitrate_cap(route, cap)) {
380        return ReceiverRouteSelection::pause(
381            route.current_selection.selector(),
382            PolicyPauseReason::SourceBitrateLimit,
383        );
384    }
385    let selection = match policy {
386        SourceAdaptationPolicy::ReadableDetail
387            if route.source.selectable_encoding_count() < 2 && source_cap.is_none() =>
388        {
389            None
390        }
391        SourceAdaptationPolicy::ReadableDetail | SourceAdaptationPolicy::None => {
392            highest_allowed_plan(route)
393        }
394        SourceAdaptationPolicy::ScalableVideo => scalable_plan(route, tuning),
395    }
396    .unwrap_or_else(|| {
397        if source_cap.is_some() {
398            ReceiverRouteSelection::pause(
399                route.current_selection.selector(),
400                PolicyPauseReason::SourceBitrateLimit,
401            )
402        } else {
403            ReceiverRouteSelection::send(route.current_selection.selector(), false)
404        }
405    });
406    if selection.policy_pause_reason.is_none()
407        && selector_bitrate(route, selection.selector).is_none()
408    {
409        return ReceiverRouteSelection::pause(
410            route.current_selection.selector(),
411            PolicyPauseReason::MissingUsableLayer,
412        );
413    }
414    selection
415}
416
417fn source_exceeds_bitrate_cap(route: &ReceiverVideoRouteInput<'_>, cap: Bitrate) -> bool {
418    if route.source.selectable_encoding_count() > 1 {
419        return false;
420    }
421    // Missing data cannot prove recovery from a cap violation. Keep the route
422    // paused until a measured source rate is at or below the cap.
423    route.source_bitrate.map_or_else(
424        || {
425            route.current_selection.policy_pause_reason()
426                == Some(PolicyPauseReason::SourceBitrateLimit)
427        },
428        |observed| observed > cap,
429    )
430}
431
432pub(super) fn selector_bitrate(
433    route: &ReceiverVideoRouteInput<'_>,
434    selector: SourceSelector,
435) -> Option<Bitrate> {
436    let source = route.source;
437    let declared = selector.selected_encoding().map_or_else(
438        || {
439            source
440                .selectable_encodings()
441                .filter_map(SourceEncodingDescriptor::max_bitrate)
442                .max()
443        },
444        |encoding_id| {
445            source
446                .selectable_encodings()
447                .find(|encoding| encoding.encoding_id() == encoding_id)
448                .and_then(SourceEncodingDescriptor::max_bitrate)
449        },
450    );
451    // A per-media observation aggregates every RID. Use it only when policy has
452    // at most one selectable encoding and needs no per-RID estimate.
453    let observed = route
454        .source_bitrate
455        .filter(|_| source.selectable_encoding_count() <= 1);
456    declared.max(observed)
457}
458
459fn highest_allowed_plan(route: &ReceiverVideoRouteInput<'_>) -> Option<ReceiverRouteSelection> {
460    let target_selector = highest_allowed_selector(route)?;
461    Some(ReceiverRouteSelection::send(
462        target_selector,
463        target_selector != route.current_selection.selector(),
464    ))
465}
466
467fn highest_allowed_selector(route: &ReceiverVideoRouteInput<'_>) -> Option<SourceSelector> {
468    let source_cap = route.source.policy().video_bitrate_cap();
469    if route.source.selectable_encoding_count() == 0 {
470        return source_cap
471            .is_none_or(|cap| route.source_bitrate.is_some_and(|rate| rate <= cap))
472            .then_some(SourceSelector::Open);
473    }
474    let target_index = allowed_encoding_indices(route.source, source_cap)
475        .next_back()
476        .or_else(|| {
477            (route.source.selectable_encoding_count() == 1
478                && source_cap
479                    .is_some_and(|cap| route.source_bitrate.is_some_and(|rate| rate <= cap)))
480            .then_some(0)
481        })?;
482    Some(SourceSelector::Encoding(
483        route
484            .source
485            .selectable_encoding_by_rank(target_index)?
486            .encoding_id(),
487    ))
488}
489
490/// Separates configured downlink headroom from audio consumption so receivers
491/// reserve audio only for admitted routes they can hear.
492fn effective_video_budget(
493    receiver_bandwidth: Bitrate,
494    tuning: VideoAdaptationTuning,
495    audio_reserve: Bitrate,
496) -> Bitrate {
497    let usable_percent = 100u64.saturating_sub(u64::from(tuning.receiver_budget_headroom_percent));
498    let after_overhead =
499        Bitrate::from_bps(receiver_bandwidth.as_bps().saturating_mul(usable_percent) / 100);
500    after_overhead.saturating_sub(audio_reserve)
501}
502
503/// Returns the next cheaper allowed encoding and its bitrate.
504///
505/// Routes using featured quality retain an intermediate rung during aggregate fitting.
506/// Their per-route target may still select the floor when the receiver cannot
507/// afford that intermediate quality. Equal bitrate ceilings do not create an
508/// intermediate operating point.
509fn step_down_selector(route: &PlannedReceiverRoute<'_>) -> Option<(SourceSelector, Bitrate)> {
510    let source = route.input.source;
511    let source_cap = source.policy().video_bitrate_cap();
512    let current_index = selector_index(source, route.selection.selector);
513    let mut cheaper_encodings = (0..current_index).rev().filter_map(|index| {
514        let encoding = source.selectable_encoding_by_rank(index)?;
515        let bitrate = encoding.max_bitrate()?;
516        (bitrate < route.selected_bitrate && source_cap.is_none_or(|cap| bitrate <= cap))
517            .then_some((SourceSelector::Encoding(encoding.encoding_id()), bitrate))
518    });
519    let (selector, bitrate) = cheaper_encodings.next()?;
520    if route.input.layout_role.uses_featured_quality()
521        && !cheaper_encodings.any(|(_selector, lower_bitrate)| lower_bitrate < bitrate)
522    {
523        return None;
524    }
525    Some((selector, bitrate))
526}
527
528fn scalable_plan(
529    route: &ReceiverVideoRouteInput<'_>,
530    tuning: VideoAdaptationTuning,
531) -> Option<ReceiverRouteSelection> {
532    let source_cap = route.source.policy().video_bitrate_cap();
533    if route.source.selectable_encoding_count() < 2 && source_cap.is_none() {
534        return None;
535    }
536    let current = route.current_selection;
537    let target_index = desired_encoding_index(route, source_cap, tuning)?;
538    let target_selector = SourceSelector::Encoding(
539        route
540            .source
541            .selectable_encoding_by_rank(target_index)?
542            .encoding_id(),
543    );
544    Some(ReceiverRouteSelection::send(
545        target_selector,
546        target_selector != current.selector(),
547    ))
548}
549
550fn desired_encoding_index(
551    route: &ReceiverVideoRouteInput<'_>,
552    source_cap: Option<Bitrate>,
553    tuning: VideoAdaptationTuning,
554) -> Option<usize> {
555    if route.user_count < tuning.multiparty_scalable_video_threshold {
556        return allowed_encoding_indices(route.source, source_cap).next_back();
557    }
558    let uses_featured_quality = route.layout_role.uses_featured_quality();
559    let Some(receiver_bandwidth) = route.receiver_bandwidth else {
560        return if uses_featured_quality {
561            allowed_encoding_indices(route.source, source_cap).next_back()
562        } else {
563            allowed_encoding_indices(route.source, source_cap).next()
564        };
565    };
566    let receiver_bandwidth =
567        effective_video_budget(receiver_bandwidth, tuning, route.audio_budget_reserve);
568    let budget = if uses_featured_quality || route.visible_scalable_route_count <= 1 {
569        receiver_bandwidth
570    } else {
571        let divisor = u64::try_from(route.visible_scalable_route_count)
572            .unwrap_or(u64::MAX)
573            // Bias secondary routes toward thumbnail quality before aggregate overload handling.
574            .saturating_mul(tuning.thumbnail_budget_divisor);
575        receiver_bandwidth.divided_by(divisor)
576    };
577    highest_affordable_encoding_index(route.source, budget, uses_featured_quality, source_cap)
578}
579
580fn highest_affordable_encoding_index(
581    source: &PublishedSourceDescriptor,
582    budget: Bitrate,
583    uses_featured_quality: bool,
584    source_cap: Option<Bitrate>,
585) -> Option<usize> {
586    if source
587        .selectable_encodings()
588        .all(|encoding| encoding.max_bitrate().is_none())
589    {
590        return source_cap.is_none().then_some(if uses_featured_quality {
591            source.selectable_encoding_count().saturating_sub(1)
592        } else {
593            0
594        });
595    }
596    allowed_encoding_indices(source, source_cap)
597        .rev()
598        .find_or_last(|index| {
599            source
600                .selectable_encoding_by_rank(*index)
601                .and_then(SourceEncodingDescriptor::max_bitrate)
602                .is_some_and(|bitrate| bitrate <= budget)
603        })
604}
605
606fn selector_index(source: &PublishedSourceDescriptor, selector: SourceSelector) -> usize {
607    selector
608        .selected_encoding()
609        .and_then(|encoding_id| {
610            source
611                .selectable_encodings()
612                .position(|encoding| encoding.encoding_id() == encoding_id)
613        })
614        .unwrap_or_else(|| source.selectable_encoding_count().saturating_sub(1))
615}
616
617fn allowed_encoding_indices(
618    source: &PublishedSourceDescriptor,
619    source_cap: Option<Bitrate>,
620) -> impl DoubleEndedIterator<Item = usize> + '_ {
621    (0..source.selectable_encoding_count()).filter(move |index| {
622        source_cap.is_none_or(|source_cap| {
623            source
624                .selectable_encoding_by_rank(*index)
625                .and_then(SourceEncodingDescriptor::max_bitrate)
626                .is_some_and(|bitrate| bitrate <= source_cap)
627        })
628    })
629}
630
631/// enforces receiver download limits by pausing the lowest-ranked routes using top-k selection
632///
633/// ```text
634/// all active routes (N)
635///   [ R0: rank 2, R1: rank 5, R2: rank 1, R3: rank 8, R4: rank 4 ]
636///   limit = 3  ==>  routes_to_pause (K) = 5 - 3 = 2
637///                              |
638///                              v  .k_largest_by_key(K = 2)
639///   +-------------------------------------------------------------+
640///   | min-heap of size K=2:                                       |
641///   | retains only the 2 highest ranks: [ R1: rank 5, R3: rank 8 ]|
642///   +-------------------------------------------------------------+
643///                              |
644///                              v  .rev()
645///   route.pause(VideoDownloadLimit) applied only to [ R3, R1 ]
646/// ```
647fn apply_video_download_limit(
648    routes: &mut [PlannedReceiverRoute<'_>],
649    max_video_downloads_per_receiver: usize,
650) {
651    let routes_to_pause =
652        active_route_count(routes).saturating_sub(max_video_downloads_per_receiver);
653    if routes_to_pause == 0 {
654        return;
655    }
656    for (_rank_key, route) in routes
657        .iter_mut()
658        .filter(|route| route.selection.policy_pause_reason.is_none())
659        .enumerate()
660        .map(|(position, route)| ((video_download_rank(route), position), route))
661        .k_largest_by_key(routes_to_pause, |(rank_key, _route)| *rank_key)
662        .rev()
663    {
664        route.pause(PolicyPauseReason::VideoDownloadLimit);
665    }
666}
667
668fn active_route_count(routes: &[PlannedReceiverRoute<'_>]) -> usize {
669    routes
670        .iter()
671        .filter(|route| route.selection.policy_pause_reason.is_none())
672        .count()
673}
674
675fn video_download_rank(route: &PlannedReceiverRoute<'_>) -> VideoAdmissionRank {
676    let input = route.input;
677    VideoAdmissionRank::new(
678        input.layout_role.priority(),
679        input.active_speaker_rank,
680        input.source.source_id(),
681    )
682}
683
684fn apply_overload_policy(routes: &mut [PlannedReceiverRoute<'_>], video_budget: Bitrate) {
685    let mut total_bitrate = selected_active_video_bitrate(routes);
686    if total_bitrate <= video_budget {
687        return;
688    }
689    // Missing-layer pauses stay limited to secondary routes. Widening encoding
690    // downsteps must not bypass the soft-pause dwell for important fixed video.
691    let mut missing_ladders = routes
692        .iter_mut()
693        .filter(|route| route_needs_missing_layer_pause(route))
694        .map(|route| (video_download_rank(route), route))
695        .collect::<Vec<_>>();
696    missing_ladders.sort_by_key(|(rank, _route)| Reverse(*rank));
697    for (_rank, route) in missing_ladders {
698        if total_bitrate <= video_budget {
699            break;
700        }
701        let selected_bitrate = route.selected_bitrate;
702        route.pause(PolicyPauseReason::MissingUsableLayer);
703        total_bitrate = total_bitrate.saturating_sub(selected_bitrate);
704    }
705    // Step the least important downgradable route down one layer at a time.
706    // Preserve intermediate layers and stop as soon as aggregate demand fits.
707    while total_bitrate > video_budget {
708        let Some((route, selector, bitrate)) = routes
709            .iter_mut()
710            .filter(|route| route_can_downgrade(route))
711            .filter_map(|route| {
712                step_down_selector(route).map(|(selector, bitrate)| (route, selector, bitrate))
713            })
714            .max_by_key(|(route, _selector, _bitrate)| {
715                (video_download_rank(route), route.selected_bitrate)
716            })
717        else {
718            break;
719        };
720        let selected_bitrate = route.selected_bitrate;
721        total_bitrate = total_bitrate
722            .saturating_sub(selected_bitrate)
723            .saturating_add(bitrate);
724        route.send(selector, bitrate);
725    }
726    if total_bitrate <= video_budget {
727        return;
728    }
729    let mut pause_order = Vec::with_capacity(routes.len());
730    for route in routes.iter_mut() {
731        if route.selection.policy_pause_reason.is_some() {
732            continue;
733        }
734        pause_order.push((video_download_rank(route), route));
735    }
736    pause_order.sort_by_key(|(rank, _)| Reverse(*rank));
737    for (_rank, route) in pause_order {
738        if total_bitrate <= video_budget {
739            break;
740        }
741        let selected_bitrate = route.selected_bitrate;
742        let pause_reason = pause_reason_for_route(route);
743        route.pause(pause_reason);
744        total_bitrate = total_bitrate.saturating_sub(selected_bitrate);
745    }
746}
747
748fn receiver_video_budget_diagnostics(
749    routes: &[PlannedReceiverRoute<'_>],
750    receiver_bandwidth: Option<Bitrate>,
751    video_budget: Option<Bitrate>,
752) -> ReceiverVideoBudgetDiagnostics {
753    let selected_video_bitrate = selected_active_video_bitrate(routes);
754    ReceiverVideoBudgetDiagnostics::new(
755        receiver_bandwidth,
756        video_budget,
757        active_route_count(routes),
758        selected_video_bitrate,
759    )
760}
761
762fn selected_active_video_bitrate(routes: &[PlannedReceiverRoute<'_>]) -> Bitrate {
763    routes
764        .iter()
765        .filter(|route| route.selection.policy_pause_reason.is_none())
766        .fold(Bitrate::zero(), |total, route| {
767            total.saturating_add(route.selected_bitrate)
768        })
769}
770
771fn admitted_video_bwe_demand(routes: &[PlannedReceiverRoute<'_>]) -> Bitrate {
772    routes
773        .iter()
774        .filter(|route| route.selection.policy_pause_reason.is_none())
775        .fold(Bitrate::zero(), |total, route| {
776            let route_demand = highest_allowed_selector(route.input)
777                .and_then(|selector| selector_bitrate(route.input, selector))
778                .unwrap_or(route.selected_bitrate);
779            total.saturating_add(route_demand.max(route.selected_bitrate))
780        })
781}
782
783fn route_can_downgrade(route: &PlannedReceiverRoute<'_>) -> bool {
784    route.selection.policy_pause_reason.is_none()
785        && route.input.source.policy().adaptation() == SourceAdaptationPolicy::ScalableVideo
786}
787
788fn route_needs_missing_layer_pause(route: &PlannedReceiverRoute<'_>) -> bool {
789    let input = route.input;
790    let source_cap = input.source.policy().video_bitrate_cap();
791    route_can_downgrade(route)
792        && !input.layout_role.uses_featured_quality()
793        && !input.source.selectable_encodings().any(|encoding| {
794            encoding
795                .max_bitrate()
796                .is_some_and(|bitrate| source_cap.is_none_or(|cap| bitrate <= cap))
797        })
798}
799
800fn pause_reason_for_route(route: &PlannedReceiverRoute<'_>) -> PolicyPauseReason {
801    let role = route.input.layout_role;
802    match role.priority() {
803        SourceRoutePriority::HiddenOrOverflow => match role {
804            SourceRoomPolicySelector::Hidden => PolicyPauseReason::HiddenTile,
805            SourceRoomPolicySelector::Overflow => PolicyPauseReason::OverflowTile,
806            _ => PolicyPauseReason::BudgetPressure,
807        },
808        _ => PolicyPauseReason::BudgetPressure,
809    }
810}
811
812const fn is_soft_pause(reason: PolicyPauseReason) -> bool {
813    matches!(
814        reason,
815        PolicyPauseReason::BudgetPressure
816            | PolicyPauseReason::HiddenTile
817            | PolicyPauseReason::OverflowTile
818    )
819}
820
821fn resolve_dwell(
822    route: &PlannedReceiverRoute<'_>,
823    tuning: VideoAdaptationTuning,
824    soft_pause_deadline: Option<Instant>,
825    now: Instant,
826) -> ReceiverRouteSelection {
827    let current = route.input.current_selection;
828    let mut selection = route.selection;
829    if let Some(reason) = selection.policy_pause_reason {
830        if !is_soft_pause(reason)
831            || current.policy_pause_reason().is_some()
832            || soft_pause_deadline.is_some_and(|deadline| deadline <= now)
833        {
834            return selection;
835        }
836        // Preserve a fitted downstep during grace. Rebuilding from `current`
837        // used to retain the expensive RID until the whole route paused.
838        if selection.selector != current.selector()
839            && selector_index(route.input.source, selection.selector)
840                > selector_index(route.input.source, current.selector())
841        {
842            selection.selector = current.selector();
843        }
844        selection.policy_pause_reason = None;
845        return selection;
846    }
847    let immediate = current.policy_pause_reason() == Some(PolicyPauseReason::VideoDownloadLimit)
848        || (current.policy_pause_reason().is_none()
849            && (selection.selector == current.selector()
850                || current.selector() == SourceSelector::Open
851                || selector_index(route.input.source, selection.selector)
852                    <= selector_index(route.input.source, current.selector())));
853    if immediate {
854        return selection;
855    }
856    // Resolve one exact target after aggregate fitting. A pre-fit upgrade that
857    // fitting reverses must not keep restarting its own eligibility deadline.
858    let pending = route
859        .input
860        .pending_upgrade
861        .filter(|pending| pending.selector == selection.selector)
862        .copied()
863        .unwrap_or_else(|| PendingUpgrade {
864            selector: selection.selector,
865            deadline: now + tuning.upgrade_dwell,
866        });
867    if pending.deadline <= now {
868        selection.request_keyframe = true;
869        return selection;
870    }
871    ReceiverRouteSelection {
872        selector: current.selector(),
873        policy_pause_reason: current.policy_pause_reason(),
874        pending_upgrade: Some(pending),
875        request_keyframe: false,
876    }
877}
878
879#[cfg(test)]
880#[path = "TESTS/solver.rs"]
881mod tests;