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]
673 fn one_tier_still_serves_audio_before_video() {
674 let flat = [want(0, 0, 128_000), want(1, 0, 4_000_000)];
675 let tiered = [want(0, AUDIO, 128_000), want(1, VIDEO, 4_000_000)];
676
677 for wants in [flat, tiered] {
678 assert_eq!(allocate(bps(2_000_000), &wants, 0), Some(bps(128_000)));
679 assert_eq!(allocate(bps(2_000_000), &wants, 1), Some(bps(1_872_000)));
680 }
681 }
682
683 #[test]
684 fn an_even_tier_splits_evenly() {
685 let wants = [want(0, VIDEO, 4_000_000), want(1, VIDEO, 4_000_000)];
686
687 assert_eq!(allocate(bps(6_000_000), &wants, 0), Some(bps(3_000_000)));
688 assert_eq!(allocate(bps(6_000_000), &wants, 1), Some(bps(3_000_000)));
689 }
690
691 #[test]
694 fn a_small_share_frees_what_it_does_not_want() {
695 let wants = [want(0, VIDEO, 1_000_000), want(1, VIDEO, 8_000_000)];
696
697 assert_eq!(allocate(bps(6_000_000), &wants, 0), Some(bps(1_000_000)));
698 assert_eq!(allocate(bps(6_000_000), &wants, 1), Some(bps(5_000_000)));
699 }
700
701 #[test]
705 fn surplus_is_left_unclaimed() {
706 assert_eq!(
707 allocate(bps(10_000_000), &[want(0, VIDEO, 4_000_000)], 0),
708 Some(bps(4_000_000))
709 );
710 }
711
712 #[test]
713 fn an_unregistered_share_has_no_grant() {
714 assert_eq!(allocate(bps(1_000_000), &[], 0), None);
715 assert_eq!(allocate(bps(1_000_000), &[want(0, VIDEO, 1_000)], 7), None);
716 }
717
718 #[tokio::test]
721 async fn concurrent_tracks_split_the_estimate() {
722 let estimate = Producer::new();
723 let allocator = Allocator::new(estimate.consume());
724
725 let (_first_broadcast, first) = track(VIDEO);
726 let _first_sub = first.consume();
727 let first = allocator.reserve(&first.demand(), bps(4_000_000));
728 let mut first_share = first.consumer();
729
730 estimate.set(Some(bps(2_000_000))).unwrap();
731 assert_eq!(first_share.changed().await.unwrap(), Some(bps(2_000_000)));
733
734 let (_second_broadcast, second) = track(VIDEO);
737 let _second_sub = second.consume();
738 let second_share = allocator.reserve(&second.demand(), bps(4_000_000));
739
740 assert_eq!(first_share.changed().await.unwrap(), Some(bps(1_000_000)));
741 assert_eq!(second_share.peek(), Some(bps(1_000_000)));
742 }
743
744 #[tokio::test]
747 async fn an_idle_track_claims_nothing() {
748 let estimate = Producer::new();
749 let allocator = Allocator::new(estimate.consume());
750 estimate.set(Some(bps(2_000_000))).unwrap();
751
752 let (_watched_broadcast, watched) = track(VIDEO);
753 let _watched_sub = watched.consume();
754 let watched_share = allocator.reserve(&watched.demand(), bps(4_000_000));
755
756 let (_idle_broadcast, idle) = track(VIDEO);
757 let idle_share = allocator.reserve(&idle.demand(), bps(4_000_000));
758
759 assert_eq!(watched_share.peek(), Some(bps(2_000_000)));
760 assert_eq!(idle_share.peek(), None);
763
764 let _idle_sub = idle.consume();
766 assert_eq!(watched_share.peek(), Some(bps(1_000_000)));
767 assert_eq!(idle_share.peek(), Some(bps(1_000_000)));
768 }
769
770 #[tokio::test]
773 async fn a_share_wakes_when_a_sibling_goes_idle() {
774 let estimate = Producer::new();
775 let allocator = Allocator::new(estimate.consume());
776 estimate.set(Some(bps(2_000_000))).unwrap();
777
778 let (_mine_broadcast, mine) = track(VIDEO);
779 let _mine_sub = mine.consume();
780 let mine_share_reserved = allocator.reserve(&mine.demand(), bps(4_000_000));
781 let mut mine_share = mine_share_reserved.consumer();
782
783 let (_sibling_broadcast, sibling) = track(VIDEO);
784 let sibling_sub = sibling.consume();
785 let _sibling_share = allocator.reserve(&sibling.demand(), bps(4_000_000));
786
787 assert_eq!(mine_share.changed().await.unwrap(), Some(bps(1_000_000)));
788
789 drop(sibling_sub);
791 assert_eq!(mine_share.changed().await.unwrap(), Some(bps(2_000_000)));
792 }
793
794 #[derive(Default)]
798 struct Woken(std::sync::atomic::AtomicBool);
799
800 impl std::task::Wake for Woken {
801 fn wake(self: std::sync::Arc<Self>) {
802 self.wake_by_ref();
803 }
804
805 fn wake_by_ref(self: &std::sync::Arc<Self>) {
806 self.0.store(true, std::sync::atomic::Ordering::SeqCst);
807 }
808 }
809
810 impl Woken {
811 fn new() -> (std::sync::Arc<Self>, kio::Waiter) {
815 let flag = std::sync::Arc::new(Self::default());
816 let waiter = kio::Waiter::new(std::task::Waker::from(flag.clone()));
817 (flag, waiter)
818 }
819
820 fn woken(&self) -> bool {
821 self.0.load(std::sync::atomic::Ordering::SeqCst)
822 }
823 }
824
825 #[tokio::test]
830 async fn an_unchanged_slice_keeps_watching_the_estimate() {
831 let estimate = Producer::new();
832 let allocator = Allocator::new(estimate.consume());
833
834 let (_broadcast, track) = track(VIDEO);
835 let _sub = track.consume();
836 let share_reserved = allocator.reserve(&track.demand(), bps(4_000_000));
837 let mut share = share_reserved.consumer();
838
839 estimate.set(Some(bps(10_000_000))).unwrap();
840 assert_eq!(share.changed().await.unwrap(), Some(bps(4_000_000)));
841
842 let (woken, waiter) = Woken::new();
845 estimate.set(Some(bps(9_000_000))).unwrap();
846 assert!(share.poll_changed(&waiter).is_pending());
847
848 estimate.set(Some(bps(1_000_000))).unwrap();
850 assert!(woken.woken(), "a parked share must be woken by the next estimate");
851 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
852 }
853
854 #[tokio::test]
856 async fn a_parked_share_is_woken_by_sibling_demand() {
857 let estimate = Producer::new();
858 let allocator = Allocator::new(estimate.consume());
859 estimate.set(Some(bps(2_000_000))).unwrap();
860
861 let (_mine_broadcast, mine) = track(VIDEO);
862 let _mine_sub = mine.consume();
863 let share_reserved = allocator.reserve(&mine.demand(), bps(4_000_000));
864 let mut share = share_reserved.consumer();
865
866 let (_sibling_broadcast, sibling) = track(VIDEO);
867 let sibling_sub = sibling.consume();
868 let _sibling_share = allocator.reserve(&sibling.demand(), bps(4_000_000));
869
870 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
871
872 let (woken, waiter) = Woken::new();
873 assert!(share.poll_changed(&waiter).is_pending());
874 drop(sibling_sub);
875 assert!(woken.woken(), "a sibling going idle must wake a parked share");
876 }
877
878 #[tokio::test]
881 async fn a_share_follows_the_estimate_lifecycle() {
882 let estimate = Producer::new();
883 let allocator = Allocator::new(estimate.consume());
884
885 let (_broadcast, track) = track(VIDEO);
886 let _sub = track.consume();
887 let share_reserved = allocator.reserve(&track.demand(), bps(4_000_000));
888 let mut share = share_reserved.consumer();
889
890 estimate.set(Some(bps(2_000_000))).unwrap();
891 assert_eq!(share.changed().await.unwrap(), Some(bps(2_000_000)));
892
893 estimate.set(None).unwrap();
894 assert_eq!(share.changed().await.unwrap(), None);
895
896 estimate.abort(Error::Cancel).unwrap();
897 assert!(share.changed().await.is_err());
898 assert!(share.changed().await.is_err());
899 }
900
901 #[tokio::test]
904 async fn a_closed_track_is_pruned() {
905 let estimate = Producer::new();
906 let allocator = Allocator::new(estimate.consume());
907
908 let (_first_broadcast, first) = track(VIDEO);
911 let _first_share = allocator.reserve(&first.demand(), bps(4_000_000));
912 first.abort(Error::Cancel).unwrap();
913
914 let (_second_broadcast, second) = track(VIDEO);
915 let _second_share = allocator.reserve(&second.demand(), bps(4_000_000));
916
917 assert_eq!(allocator.registry.read().entries.len(), 1);
918 }
919
920 #[tokio::test]
924 async fn dropping_a_reservation_releases_it() {
925 let estimate = Producer::new();
926 let allocator = Allocator::new(estimate.consume());
927 estimate.set(Some(bps(2_000_000))).unwrap();
928
929 let (_first_broadcast, first) = track(VIDEO);
930 let _first_sub = first.consume();
931 let first_reserved = allocator.reserve(&first.demand(), bps(4_000_000));
932
933 let (_second_broadcast, second) = track(VIDEO);
934 let _second_sub = second.consume();
935 let second_reserved = allocator.reserve(&second.demand(), bps(4_000_000));
936 assert_eq!(second_reserved.peek(), Some(bps(1_000_000)));
937
938 let mut orphan = first_reserved.consumer();
941 assert_eq!(orphan.changed().await.unwrap(), Some(bps(1_000_000)));
942
943 drop(first_reserved);
944 assert_eq!(allocator.registry.read().entries.len(), 1);
945 assert_eq!(second_reserved.peek(), Some(bps(2_000_000)));
946
947 assert_eq!(orphan.changed().await.unwrap(), None);
951 assert_eq!(orphan.peek(), None);
952 }
953
954 #[tokio::test]
958 async fn update_changes_the_claim_in_place() {
959 let estimate = Producer::new();
960 let allocator = Allocator::new(estimate.consume());
961 estimate.set(Some(bps(6_000_000))).unwrap();
962
963 let (_small_broadcast, small) = track(VIDEO);
964 let _small_sub = small.consume();
965 let small_reserved = allocator.reserve(&small.demand(), bps(1_000_000));
966
967 let (_large_broadcast, large) = track(VIDEO);
968 let _large_sub = large.consume();
969 let large_reserved = allocator.reserve(&large.demand(), bps(8_000_000));
970
971 assert_eq!(small_reserved.peek(), Some(bps(1_000_000)));
973 assert_eq!(large_reserved.peek(), Some(bps(5_000_000)));
974
975 small_reserved.update(bps(4_000_000));
978 assert_eq!(allocator.registry.read().entries.len(), 2);
979 assert_eq!(small_reserved.peek(), Some(bps(3_000_000)));
980 assert_eq!(large_reserved.peek(), Some(bps(3_000_000)));
981
982 small_reserved.update(bps(1_000_000));
984 assert_eq!(large_reserved.peek(), Some(bps(5_000_000)));
985 }
986
987 #[tokio::test]
990 async fn update_wakes_a_parked_reader() {
991 let estimate = Producer::new();
992 let allocator = Allocator::new(estimate.consume());
993 estimate.set(Some(bps(6_000_000))).unwrap();
994
995 let (_broadcast, producer) = track(VIDEO);
996 let _sub = producer.consume();
997 let reserved = allocator.reserve(&producer.demand(), bps(1_000_000));
998 let mut share = reserved.consumer();
999 assert_eq!(share.changed().await.unwrap(), Some(bps(1_000_000)));
1000
1001 let (woken, waiter) = Woken::new();
1002 assert!(share.poll_changed(&waiter).is_pending());
1003
1004 reserved.update(bps(4_000_000));
1005 assert!(woken.woken(), "raising the ceiling must wake the reader");
1006 assert_eq!(share.changed().await.unwrap(), Some(bps(4_000_000)));
1007 }
1008
1009 #[tokio::test]
1012 async fn a_share_outliving_the_allocator_reports_closed() {
1013 let estimate = Producer::new();
1014 let allocator = Allocator::new(estimate.consume());
1015
1016 let (_broadcast, producer) = track(VIDEO);
1017 let _sub = producer.consume();
1018 let share_reserved = allocator.reserve(&producer.demand(), bps(4_000_000));
1019 let mut share = share_reserved.consumer();
1020
1021 drop(allocator);
1022 drop(estimate);
1023 assert!(share.changed().await.is_err());
1024 }
1025
1026 #[tokio::test]
1031 async fn closed_is_distinct_from_unavailable() {
1032 let producer = Producer::new();
1033 let mut consumer = producer.consume();
1034
1035 producer.set(Some(bps(1_000_000))).unwrap();
1036 assert_eq!(consumer.changed().await.unwrap(), Some(bps(1_000_000)));
1037
1038 producer.set(None).unwrap();
1040 assert_eq!(consumer.changed().await.unwrap(), None);
1041
1042 producer.abort(Error::Cancel).unwrap();
1044 assert!(consumer.changed().await.is_err());
1045 assert!(consumer.changed().await.is_err());
1047 }
1048}