1use std::task::Poll;
11
12use crate::{Error, Result, track};
13
14#[derive(Copy, Clone, Debug, Default, PartialEq, Eq, PartialOrd, Ord, Hash)]
25pub struct Rate(u64);
26
27impl Rate {
28 pub const ZERO: Self = Self(0);
30
31 pub const fn from_bps(bps: u64) -> Self {
33 Self(bps)
34 }
35
36 pub const fn from_kbps(kbps: u64) -> Self {
38 Self(kbps.saturating_mul(1_000))
39 }
40
41 pub const fn from_mbps(mbps: u64) -> Self {
43 Self(mbps.saturating_mul(1_000_000))
44 }
45
46 pub const fn as_bps(self) -> u64 {
48 self.0
49 }
50
51 pub fn scaled(self, factor: f64) -> Self {
56 let scaled = self.0 as f64 * factor.max(0.0);
57 if scaled >= u64::MAX as f64 {
58 Self(u64::MAX)
59 } else {
60 Self(scaled as u64)
61 }
62 }
63
64 pub const fn abs_diff(self, other: Self) -> Self {
66 Self(self.0.abs_diff(other.0))
67 }
68}
69
70#[derive(Default)]
71struct State {
72 bitrate: Option<Rate>,
73 abort: Option<Error>,
74}
75
76#[derive(Clone)]
78pub struct Producer {
79 state: kio::Producer<State>,
80}
81
82impl Producer {
83 pub fn new() -> Self {
85 Self {
86 state: kio::Producer::default(),
87 }
88 }
89
90 pub fn set(&self, bitrate: Option<Rate>) -> Result<()> {
92 let mut state = self.modify()?;
93 if state.bitrate != bitrate {
94 state.bitrate = bitrate;
95 }
96 Ok(())
97 }
98
99 pub fn consume(&self) -> Consumer {
101 Consumer {
102 inner: Inner::Whole(self.state.consume()),
103 last: None,
104 }
105 }
106
107 pub fn abort(&self, err: Error) -> Result<()> {
109 let mut state = self.modify()?;
110 state.abort = Some(err);
111 state.close();
112 Ok(())
113 }
114
115 pub async fn closed(&self) -> Error {
117 self.state.closed().await;
118 self.close_error()
119 }
120
121 pub async fn unused(&self) -> Result<()> {
123 kio::wait(|waiter| self.poll_unused(waiter)).await
124 }
125
126 pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
128 self.state.poll_unused(waiter).map(|used| match used {
129 Some(()) => Ok(()),
130 None => Err(self.close_error()),
131 })
132 }
133
134 pub fn is_used(&self) -> bool {
136 self.state.is_used()
137 }
138
139 pub async fn used(&self) -> Result<()> {
141 kio::wait(|waiter| self.poll_used(waiter)).await
142 }
143
144 pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
146 self.state.poll_used(waiter).map(|used| match used {
147 Some(()) => Ok(()),
148 None => Err(self.close_error()),
149 })
150 }
151
152 fn modify(&self) -> Result<kio::Mut<'_, State>> {
153 self.state
154 .write()
155 .map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
156 }
157
158 fn close_error(&self) -> Error {
160 self.state.read().abort.clone().unwrap_or(Error::Dropped)
161 }
162}
163
164impl Default for Producer {
165 fn default() -> Self {
166 Self::new()
167 }
168}
169
170#[derive(Clone)]
184pub struct Allocator {
185 estimate: Consumer,
186 registry: kio::Producer<Registry>,
187}
188
189impl Allocator {
190 pub fn new(estimate: Consumer) -> Self {
193 Self {
194 estimate,
195 registry: kio::Producer::default(),
196 }
197 }
198
199 pub fn unlimited() -> Self {
206 Self {
207 estimate: Consumer {
208 inner: Inner::Unavailable,
209 last: None,
210 },
211 registry: kio::Producer::default(),
212 }
213 }
214
215 pub fn reserve(&self, track: &track::Demand, max: Rate) -> Reservation {
249 let priority = track.priority();
252 let demand = track.clone();
253
254 let id = {
255 let Ok(mut registry) = self.registry.write() else {
258 unreachable!("the allocator holds its own registry producer")
259 };
260 registry.entries.retain(|entry| !entry.demand.is_closed());
263
264 let id = registry.next_id;
265 registry.next_id += 1;
266 registry.entries.push(Entry {
267 id,
268 demand,
269 priority,
270 max,
271 });
272 id
273 };
274
275 Reservation {
276 share: Share {
277 estimate: self.estimate.clone(),
278 registry: self.registry.consume(),
279 id,
280 },
281 registry: self.registry.downgrade(),
282 }
283 }
284}
285
286#[must_use = "a dropped Reservation is released, so the sender claims nothing and its siblings take the room"]
293pub struct Reservation {
294 share: Share,
295 registry: kio::Weak<Registry>,
298}
299
300impl Reservation {
301 pub fn peek(&self) -> Option<Rate> {
306 self.share.grant()
307 }
308
309 pub fn consumer(&self) -> Consumer {
314 Consumer {
315 inner: Inner::Share(Box::new(self.share.clone())),
316 last: None,
317 }
318 }
319
320 pub fn update(&self, max: Rate) {
331 let Some(registry) = self.registry.upgrade() else {
332 return;
333 };
334 let Ok(mut registry) = registry.write() else {
335 return;
336 };
337 if let Some(entry) = registry.entries.iter_mut().find(|entry| entry.id == self.share.id) {
338 entry.max = max;
339 }
340 }
341}
342
343impl Drop for Reservation {
344 fn drop(&mut self) {
345 let Some(registry) = self.registry.upgrade() else {
346 return;
347 };
348 let Ok(mut registry) = registry.write() else {
349 return;
350 };
351 registry.entries.retain(|entry| entry.id != self.share.id);
352 }
353}
354
355impl Default for Allocator {
358 fn default() -> Self {
359 Self::unlimited()
360 }
361}
362
363impl std::fmt::Debug for Allocator {
366 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
367 f.debug_struct("Allocator")
368 .field("registered", &self.registry.read().entries.len())
369 .finish()
370 }
371}
372
373#[derive(Default)]
375struct Registry {
376 entries: Vec<Entry>,
377 next_id: u64,
378}
379
380struct Entry {
382 id: u64,
383 demand: track::Demand,
384 priority: u8,
385 max: Rate,
386}
387
388#[derive(Clone)]
390struct Share {
391 estimate: Consumer,
393 registry: kio::Consumer<Registry>,
394 id: u64,
395}
396
397#[derive(Copy, Clone, Debug, Eq, PartialEq)]
400struct Want {
401 id: u64,
402 priority: u8,
403 max: Rate,
404}
405
406impl Share {
407 fn grant(&self) -> Option<Rate> {
409 let estimate = self.estimate.peek()?;
410 let wants: Vec<Want> = self
411 .claims()
412 .into_iter()
413 .filter(|(_, demand)| demand.is_used())
414 .map(|(want, _)| want)
415 .collect();
416 allocate(estimate, &wants, self.id)
417 }
418
419 fn poll_grant(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Rate>>> {
422 loop {
429 match self.estimate.poll_changed(waiter) {
430 Poll::Ready(Ok(_)) => continue,
433 Poll::Ready(Err(err)) => return Poll::Ready(Err(err)),
434 Poll::Pending => break,
435 }
436 }
437
438 if let Poll::Ready(Err(_)) = self.registry.poll(waiter, |_| Poll::<()>::Pending) {
441 return Poll::Ready(Err(Error::Dropped));
442 }
443
444 let wants: Vec<Want> = self
445 .claims()
446 .into_iter()
447 .filter(|(_, demand)| match demand.poll_state(waiter) {
448 track::DemandState::Active => true,
449 track::DemandState::Idle => false,
452 track::DemandState::Closed => false,
454 })
455 .map(|(want, _)| want)
456 .collect();
457
458 let grant = self
459 .estimate
460 .peek()
461 .and_then(|estimate| allocate(estimate, &wants, self.id));
462 Poll::Ready(Ok(grant))
463 }
464
465 fn claims(&self) -> Vec<(Want, track::Demand)> {
467 self.registry
468 .read()
469 .entries
470 .iter()
471 .map(|entry| {
472 (
473 Want {
474 id: entry.id,
475 priority: entry.priority,
476 max: entry.max,
477 },
478 entry.demand.clone(),
479 )
480 })
481 .collect()
482 }
483}
484
485fn allocate(estimate: Rate, wants: &[Want], id: u64) -> Option<Rate> {
498 let mut budget = estimate.as_bps();
501 let mut tier = wants.iter().map(|want| want.priority).max();
502
503 while let Some(priority) = tier {
504 let mut members: Vec<&Want> = wants.iter().filter(|want| want.priority == priority).collect();
507 members.sort_by_key(|want| want.max);
508
509 let mut remaining = members.len() as u64;
510 for want in members {
511 let even = budget / remaining;
512 let grant = want.max.as_bps().min(even);
513 if want.id == id {
514 return Some(Rate::from_bps(grant));
515 }
516 budget -= grant;
517 remaining -= 1;
518 }
519
520 tier = wants
521 .iter()
522 .map(|want| want.priority)
523 .filter(|other| *other < priority)
524 .max();
525 }
526
527 None
528}
529
530#[derive(Clone)]
532pub struct Consumer {
533 inner: Inner,
534 last: Option<Rate>,
535}
536
537#[derive(Clone)]
539enum Inner {
540 Whole(kio::Consumer<State>),
541 Share(Box<Share>),
542 Unavailable,
544}
545
546impl Consumer {
547 pub fn peek(&self) -> Option<Rate> {
549 match &self.inner {
550 Inner::Whole(state) => state.read().bitrate,
551 Inner::Share(share) => share.grant(),
552 Inner::Unavailable => None,
553 }
554 }
555
556 pub fn poll_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<Rate>>> {
567 let last = self.last;
568
569 let bitrate = match &mut self.inner {
570 Inner::Whole(state) => match state.poll(waiter, |state| {
571 if state.bitrate != last {
572 Poll::Ready(state.bitrate)
573 } else {
574 Poll::Pending
575 }
576 }) {
577 Poll::Ready(Ok(bitrate)) => bitrate,
578 Poll::Ready(Err(state)) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
583 Poll::Pending => return Poll::Pending,
584 },
585 Inner::Unavailable => return Poll::Pending,
591 Inner::Share(share) => match share.poll_grant(waiter) {
592 Poll::Ready(Ok(grant)) if grant == last => return Poll::Pending,
593 Poll::Ready(Ok(grant)) => grant,
594 Poll::Ready(Err(err)) => return Poll::Ready(Err(err)),
595 Poll::Pending => return Poll::Pending,
596 },
597 };
598
599 self.last = bitrate;
600 Poll::Ready(Ok(bitrate))
601 }
602
603 pub async fn changed(&mut self) -> Result<Option<Rate>> {
611 kio::wait(|waiter| self.poll_changed(waiter)).await
612 }
613}
614
615#[cfg(test)]
616mod tests {
617 use super::*;
618 use crate::broadcast;
619
620 const AUDIO: u8 = 80;
622 const VIDEO: u8 = 60;
623
624 fn bps(bps: u64) -> Rate {
626 Rate::from_bps(bps)
627 }
628
629 fn want(id: u64, priority: u8, max: u64) -> Want {
630 Want {
631 id,
632 priority,
633 max: bps(max),
634 }
635 }
636
637 fn track(priority: u8) -> (broadcast::Producer, track::Producer) {
639 let broadcast = broadcast::Info::default().produce();
640 let track = broadcast
641 .create_track("t", track::Info::default().with_priority(priority))
642 .unwrap();
643 (broadcast, track)
644 }
645
646 #[test]
647 fn strict_priority_fills_the_top_tier_first() {
648 let wants = [want(0, AUDIO, 128_000), want(1, VIDEO, 4_000_000)];
649
650 assert_eq!(allocate(bps(2_000_000), &wants, 0), Some(bps(128_000)));
653 assert_eq!(allocate(bps(2_000_000), &wants, 1), Some(bps(1_872_000)));
654 }
655
656 #[test]
657 fn a_starved_tier_gets_nothing() {
658 let wants = [want(0, AUDIO, 2_000_000), want(1, VIDEO, 4_000_000)];
659
660 assert_eq!(allocate(bps(1_000_000), &wants, 0), Some(bps(1_000_000)));
661 assert_eq!(allocate(bps(1_000_000), &wants, 1), Some(bps(0)));
664 }
665
666 #[test]
672 fn one_tier_still_serves_audio_before_video() {
673 let flat = [want(0, 0, 128_000), want(1, 0, 4_000_000)];
674 let tiered = [want(0, AUDIO, 128_000), want(1, VIDEO, 4_000_000)];
675
676 for wants in [flat, tiered] {
677 assert_eq!(allocate(bps(2_000_000), &wants, 0), Some(bps(128_000)));
678 assert_eq!(allocate(bps(2_000_000), &wants, 1), Some(bps(1_872_000)));
679 }
680 }
681
682 #[test]
683 fn an_even_tier_splits_evenly() {
684 let wants = [want(0, VIDEO, 4_000_000), want(1, VIDEO, 4_000_000)];
685
686 assert_eq!(allocate(bps(6_000_000), &wants, 0), Some(bps(3_000_000)));
687 assert_eq!(allocate(bps(6_000_000), &wants, 1), Some(bps(3_000_000)));
688 }
689
690 #[test]
693 fn a_small_share_frees_what_it_does_not_want() {
694 let wants = [want(0, VIDEO, 1_000_000), want(1, VIDEO, 8_000_000)];
695
696 assert_eq!(allocate(bps(6_000_000), &wants, 0), Some(bps(1_000_000)));
697 assert_eq!(allocate(bps(6_000_000), &wants, 1), Some(bps(5_000_000)));
698 }
699
700 #[test]
704 fn surplus_is_left_unclaimed() {
705 assert_eq!(
706 allocate(bps(10_000_000), &[want(0, VIDEO, 4_000_000)], 0),
707 Some(bps(4_000_000))
708 );
709 }
710
711 #[test]
712 fn an_unregistered_share_has_no_grant() {
713 assert_eq!(allocate(bps(1_000_000), &[], 0), None);
714 assert_eq!(allocate(bps(1_000_000), &[want(0, VIDEO, 1_000)], 7), None);
715 }
716
717 #[tokio::test]
720 async fn concurrent_tracks_split_the_estimate() {
721 let estimate = Producer::new();
722 let allocator = Allocator::new(estimate.consume());
723
724 let (_first_broadcast, first) = track(VIDEO);
725 let _first_sub = first.consume();
726 let first = allocator.reserve(&first.demand(), bps(4_000_000));
727 let mut first_share = first.consumer();
728
729 estimate.set(Some(bps(2_000_000))).unwrap();
730 assert_eq!(first_share.changed().await.unwrap(), Some(bps(2_000_000)));
732
733 let (_second_broadcast, second) = track(VIDEO);
736 let _second_sub = second.consume();
737 let second_share = allocator.reserve(&second.demand(), bps(4_000_000));
738
739 assert_eq!(first_share.changed().await.unwrap(), Some(bps(1_000_000)));
740 assert_eq!(second_share.peek(), Some(bps(1_000_000)));
741 }
742
743 #[tokio::test]
746 async fn an_idle_track_claims_nothing() {
747 let estimate = Producer::new();
748 let allocator = Allocator::new(estimate.consume());
749 estimate.set(Some(bps(2_000_000))).unwrap();
750
751 let (_watched_broadcast, watched) = track(VIDEO);
752 let _watched_sub = watched.consume();
753 let watched_share = allocator.reserve(&watched.demand(), bps(4_000_000));
754
755 let (_idle_broadcast, idle) = track(VIDEO);
756 let idle_share = allocator.reserve(&idle.demand(), bps(4_000_000));
757
758 assert_eq!(watched_share.peek(), Some(bps(2_000_000)));
759 assert_eq!(idle_share.peek(), None);
762
763 let _idle_sub = idle.consume();
765 assert_eq!(watched_share.peek(), Some(bps(1_000_000)));
766 assert_eq!(idle_share.peek(), Some(bps(1_000_000)));
767 }
768
769 #[tokio::test]
772 async fn a_share_wakes_when_a_sibling_goes_idle() {
773 let estimate = Producer::new();
774 let allocator = Allocator::new(estimate.consume());
775 estimate.set(Some(bps(2_000_000))).unwrap();
776
777 let (_mine_broadcast, mine) = track(VIDEO);
778 let _mine_sub = mine.consume();
779 let mine_share_reserved = allocator.reserve(&mine.demand(), bps(4_000_000));
780 let mut mine_share = mine_share_reserved.consumer();
781
782 let (_sibling_broadcast, sibling) = track(VIDEO);
783 let sibling_sub = sibling.consume();
784 let _sibling_share = allocator.reserve(&sibling.demand(), bps(4_000_000));
785
786 assert_eq!(mine_share.changed().await.unwrap(), Some(bps(1_000_000)));
787
788 drop(sibling_sub);
790 assert_eq!(mine_share.changed().await.unwrap(), Some(bps(2_000_000)));
791 }
792
793 #[derive(Default)]
797 struct Woken(std::sync::atomic::AtomicBool);
798
799 impl std::task::Wake for Woken {
800 fn wake(self: std::sync::Arc<Self>) {
801 self.wake_by_ref();
802 }
803
804 fn wake_by_ref(self: &std::sync::Arc<Self>) {
805 self.0.store(true, std::sync::atomic::Ordering::SeqCst);
806 }
807 }
808
809 impl Woken {
810 fn new() -> (std::sync::Arc<Self>, kio::Waiter) {
814 let flag = std::sync::Arc::new(Self::default());
815 let waiter = kio::Waiter::new(std::task::Waker::from(flag.clone()));
816 (flag, waiter)
817 }
818
819 fn woken(&self) -> bool {
820 self.0.load(std::sync::atomic::Ordering::SeqCst)
821 }
822 }
823
824 #[tokio::test]
829 async fn an_unchanged_slice_keeps_watching_the_estimate() {
830 let estimate = Producer::new();
831 let allocator = Allocator::new(estimate.consume());
832
833 let (_broadcast, track) = track(VIDEO);
834 let _sub = track.consume();
835 let share_reserved = allocator.reserve(&track.demand(), bps(4_000_000));
836 let mut share = share_reserved.consumer();
837
838 estimate.set(Some(bps(10_000_000))).unwrap();
839 assert_eq!(share.changed().await.unwrap(), Some(bps(4_000_000)));
840
841 let (woken, waiter) = Woken::new();
844 estimate.set(Some(bps(9_000_000))).unwrap();
845 assert!(share.poll_changed(&waiter).is_pending());
846
847 estimate.set(Some(bps(1_000_000))).unwrap();
849 assert!(woken.woken(), "a parked share must be woken by the next estimate");
850 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
851 }
852
853 #[tokio::test]
855 async fn a_parked_share_is_woken_by_sibling_demand() {
856 let estimate = Producer::new();
857 let allocator = Allocator::new(estimate.consume());
858 estimate.set(Some(bps(2_000_000))).unwrap();
859
860 let (_mine_broadcast, mine) = track(VIDEO);
861 let _mine_sub = mine.consume();
862 let share_reserved = allocator.reserve(&mine.demand(), bps(4_000_000));
863 let mut share = share_reserved.consumer();
864
865 let (_sibling_broadcast, sibling) = track(VIDEO);
866 let sibling_sub = sibling.consume();
867 let _sibling_share = allocator.reserve(&sibling.demand(), bps(4_000_000));
868
869 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
870
871 let (woken, waiter) = Woken::new();
872 assert!(share.poll_changed(&waiter).is_pending());
873 drop(sibling_sub);
874 assert!(woken.woken(), "a sibling going idle must wake a parked share");
875 }
876
877 #[tokio::test]
880 async fn a_share_follows_the_estimate_lifecycle() {
881 let estimate = Producer::new();
882 let allocator = Allocator::new(estimate.consume());
883
884 let (_broadcast, track) = track(VIDEO);
885 let _sub = track.consume();
886 let share_reserved = allocator.reserve(&track.demand(), bps(4_000_000));
887 let mut share = share_reserved.consumer();
888
889 estimate.set(Some(bps(2_000_000))).unwrap();
890 assert_eq!(share.changed().await.unwrap(), Some(bps(2_000_000)));
891
892 estimate.set(None).unwrap();
893 assert_eq!(share.changed().await.unwrap(), None);
894
895 estimate.abort(Error::Cancel).unwrap();
896 assert!(share.changed().await.is_err());
897 assert!(share.changed().await.is_err());
898 }
899
900 #[tokio::test]
903 async fn a_closed_track_is_pruned() {
904 let estimate = Producer::new();
905 let allocator = Allocator::new(estimate.consume());
906
907 let (_first_broadcast, first) = track(VIDEO);
910 let _first_share = allocator.reserve(&first.demand(), bps(4_000_000));
911 first.abort(Error::Cancel).unwrap();
912
913 let (_second_broadcast, second) = track(VIDEO);
914 let _second_share = allocator.reserve(&second.demand(), bps(4_000_000));
915
916 assert_eq!(allocator.registry.read().entries.len(), 1);
917 }
918
919 #[tokio::test]
923 async fn dropping_a_reservation_releases_it() {
924 let estimate = Producer::new();
925 let allocator = Allocator::new(estimate.consume());
926 estimate.set(Some(bps(2_000_000))).unwrap();
927
928 let (_first_broadcast, first) = track(VIDEO);
929 let _first_sub = first.consume();
930 let first_reserved = allocator.reserve(&first.demand(), bps(4_000_000));
931
932 let (_second_broadcast, second) = track(VIDEO);
933 let _second_sub = second.consume();
934 let second_reserved = allocator.reserve(&second.demand(), bps(4_000_000));
935 assert_eq!(second_reserved.peek(), Some(bps(1_000_000)));
936
937 let mut orphan = first_reserved.consumer();
940 assert_eq!(orphan.changed().await.unwrap(), Some(bps(1_000_000)));
941
942 drop(first_reserved);
943 assert_eq!(allocator.registry.read().entries.len(), 1);
944 assert_eq!(second_reserved.peek(), Some(bps(2_000_000)));
945
946 assert_eq!(orphan.changed().await.unwrap(), None);
950 assert_eq!(orphan.peek(), None);
951 }
952
953 #[tokio::test]
957 async fn update_changes_the_claim_in_place() {
958 let estimate = Producer::new();
959 let allocator = Allocator::new(estimate.consume());
960 estimate.set(Some(bps(6_000_000))).unwrap();
961
962 let (_small_broadcast, small) = track(VIDEO);
963 let _small_sub = small.consume();
964 let small_reserved = allocator.reserve(&small.demand(), bps(1_000_000));
965
966 let (_large_broadcast, large) = track(VIDEO);
967 let _large_sub = large.consume();
968 let large_reserved = allocator.reserve(&large.demand(), bps(8_000_000));
969
970 assert_eq!(small_reserved.peek(), Some(bps(1_000_000)));
972 assert_eq!(large_reserved.peek(), Some(bps(5_000_000)));
973
974 small_reserved.update(bps(4_000_000));
977 assert_eq!(allocator.registry.read().entries.len(), 2);
978 assert_eq!(small_reserved.peek(), Some(bps(3_000_000)));
979 assert_eq!(large_reserved.peek(), Some(bps(3_000_000)));
980
981 small_reserved.update(bps(1_000_000));
983 assert_eq!(large_reserved.peek(), Some(bps(5_000_000)));
984 }
985
986 #[tokio::test]
989 async fn update_wakes_a_parked_reader() {
990 let estimate = Producer::new();
991 let allocator = Allocator::new(estimate.consume());
992 estimate.set(Some(bps(6_000_000))).unwrap();
993
994 let (_broadcast, producer) = track(VIDEO);
995 let _sub = producer.consume();
996 let reserved = allocator.reserve(&producer.demand(), bps(1_000_000));
997 let mut share = reserved.consumer();
998 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
999
1000 let (woken, waiter) = Woken::new();
1001 assert!(share.poll_changed(&waiter).is_pending());
1002
1003 reserved.update(bps(4_000_000));
1004 assert!(woken.woken(), "raising the ceiling must wake the reader");
1005 assert_eq!(share.changed().await.unwrap(), Some(bps(4_000_000)));
1006 }
1007
1008 #[tokio::test]
1011 async fn a_share_outliving_the_allocator_reports_closed() {
1012 let estimate = Producer::new();
1013 let allocator = Allocator::new(estimate.consume());
1014
1015 let (_broadcast, producer) = track(VIDEO);
1016 let _sub = producer.consume();
1017 let share_reserved = allocator.reserve(&producer.demand(), bps(4_000_000));
1018 let mut share = share_reserved.consumer();
1019
1020 drop(allocator);
1021 drop(estimate);
1022 assert!(share.changed().await.is_err());
1023 }
1024
1025 #[tokio::test]
1030 async fn closed_is_distinct_from_unavailable() {
1031 let producer = Producer::new();
1032 let mut consumer = producer.consume();
1033
1034 producer.set(Some(bps(1_000_000))).unwrap();
1035 assert_eq!(consumer.changed().await.unwrap(), Some(bps(1_000_000)));
1036
1037 producer.set(None).unwrap();
1039 assert_eq!(consumer.changed().await.unwrap(), None);
1040
1041 producer.abort(Error::Cancel).unwrap();
1043 assert!(consumer.changed().await.is_err());
1044 assert!(consumer.changed().await.is_err());
1046 }
1047}