1use alloc::collections::VecDeque;
2use alloc::sync::Arc;
3use core::future::Future;
4use core::pin::Pin;
5use core::task::{Context, Poll, Waker};
6
7use spin::Mutex;
8
9use crate::element::{OutputSink, PushOutcome, QosMessage, Reconfigure};
10use crate::error::G2gError;
11use crate::frame::PipelinePacket;
12use crate::link::LinkPolicy;
13use crate::runtime::instrument::{EdgeCounters, Probe};
14
15pub fn bounded<T>(capacity: usize) -> (Sender<T>, Receiver<T>) {
16 assert!(capacity > 0, "channel capacity must be > 0");
17 let inner = Arc::new(Mutex::new(Inner {
18 queue: VecDeque::with_capacity(capacity),
19 capacity,
20 send_waker: None,
21 recv_waker: None,
22 senders: 1,
23 receivers: 1,
24 }));
25 (
26 Sender {
27 inner: inner.clone(),
28 },
29 Receiver { inner },
30 )
31}
32
33#[derive(Debug)]
34struct Inner<T> {
35 queue: VecDeque<T>,
36 capacity: usize,
37 send_waker: Option<Waker>,
38 recv_waker: Option<Waker>,
39 senders: usize,
40 receivers: usize,
41}
42
43#[derive(Debug)]
44pub struct Sender<T> {
45 inner: Arc<Mutex<Inner<T>>>,
46}
47
48#[derive(Debug)]
49pub struct Receiver<T> {
50 inner: Arc<Mutex<Inner<T>>>,
51}
52
53impl<T> Clone for Sender<T> {
54 fn clone(&self) -> Self {
55 self.inner.lock().senders += 1;
56 Self {
57 inner: self.inner.clone(),
58 }
59 }
60}
61
62impl<T> Drop for Sender<T> {
63 fn drop(&mut self) {
64 let mut g = self.inner.lock();
65 g.senders -= 1;
66 if g.senders == 0 {
67 if let Some(w) = g.recv_waker.take() {
68 w.wake();
69 }
70 }
71 }
72}
73
74impl<T> Drop for Receiver<T> {
75 fn drop(&mut self) {
76 let mut g = self.inner.lock();
77 g.receivers -= 1;
78 if g.receivers == 0 {
79 if let Some(w) = g.send_waker.take() {
80 w.wake();
81 }
82 }
83 }
84}
85
86#[derive(Debug, Clone, Copy, PartialEq, Eq)]
87pub enum SendError {
88 Closed,
90 Full,
92}
93
94impl<T> Sender<T> {
95 pub fn try_send(&self, value: T) -> Result<(), (T, SendError)> {
98 let mut g = self.inner.lock();
99 if g.receivers == 0 {
100 return Err((value, SendError::Closed));
101 }
102 if g.queue.len() >= g.capacity {
103 return Err((value, SendError::Full));
104 }
105 g.queue.push_back(value);
106 if let Some(w) = g.recv_waker.take() {
107 w.wake();
108 }
109 Ok(())
110 }
111
112 pub fn send(&self, value: T) -> SendFuture<'_, T> {
113 SendFuture {
114 sender: self,
115 value: Some(value),
116 }
117 }
118
119 pub fn poll_send(
123 &self,
124 cx: &mut Context<'_>,
125 value: &mut Option<T>,
126 ) -> Poll<Result<(), SendError>> {
127 let mut g = self.inner.lock();
128 if g.receivers == 0 {
129 return Poll::Ready(Err(SendError::Closed));
130 }
131 if g.queue.len() < g.capacity {
132 let v = value.take().expect("poll_send called without a value");
133 g.queue.push_back(v);
134 if let Some(w) = g.recv_waker.take() {
135 w.wake();
136 }
137 return Poll::Ready(Ok(()));
138 }
139 g.send_waker = Some(cx.waker().clone());
140 Poll::Pending
141 }
142
143 pub(crate) fn evict_front_matching(&self, pred: impl Fn(&T) -> bool) -> Option<T> {
149 let mut g = self.inner.lock();
150 let idx = g.queue.iter().position(pred)?;
151 g.queue.remove(idx)
152 }
153}
154
155#[allow(missing_debug_implementations)]
156pub struct SendFuture<'a, T> {
157 sender: &'a Sender<T>,
158 value: Option<T>,
159}
160
161impl<'a, T: Unpin> Future for SendFuture<'a, T> {
162 type Output = Result<(), SendError>;
163
164 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
165 let this = self.get_mut();
166 this.sender.poll_send(cx, &mut this.value)
167 }
168}
169
170impl<T> Receiver<T> {
171 pub fn recv(&self) -> RecvFuture<'_, T> {
172 RecvFuture { receiver: self }
173 }
174
175 pub fn fill_percent(&self) -> u8 {
178 let g = self.inner.lock();
179 ((g.queue.len() * 100) / g.capacity) as u8
180 }
181
182 pub fn try_recv(&self) -> Option<T> {
185 let mut g = self.inner.lock();
186 let v = g.queue.pop_front();
187 if v.is_some() {
188 if let Some(w) = g.send_waker.take() {
189 w.wake();
190 }
191 }
192 v
193 }
194}
195
196#[allow(missing_debug_implementations)]
197pub struct RecvFuture<'a, T> {
198 receiver: &'a Receiver<T>,
199}
200
201impl<'a, T> Future for RecvFuture<'a, T> {
202 type Output = Option<T>;
203
204 fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
205 let this = self.get_mut();
206 let mut g = this.receiver.inner.lock();
207 if let Some(v) = g.queue.pop_front() {
208 if let Some(w) = g.send_waker.take() {
209 w.wake();
210 }
211 return Poll::Ready(Some(v));
212 }
213 if g.senders == 0 {
214 return Poll::Ready(None);
215 }
216 g.recv_waker = Some(cx.waker().clone());
217 Poll::Pending
218 }
219}
220
221#[derive(Debug, Clone, Default)]
225pub struct ReconfigureSlot {
226 inner: Arc<Mutex<Option<Reconfigure>>>,
227}
228
229impl ReconfigureSlot {
230 pub fn store(&self, value: Reconfigure) {
231 *self.inner.lock() = Some(value);
232 }
233
234 pub fn take(&self) -> Option<Reconfigure> {
235 self.inner.lock().take()
236 }
237}
238
239#[derive(Debug, Clone, Default)]
244pub struct QosSlot {
245 inner: Arc<Mutex<Option<QosMessage>>>,
246}
247
248impl QosSlot {
249 pub fn store(&self, value: QosMessage) {
250 *self.inner.lock() = Some(value);
251 }
252
253 pub fn take(&self) -> Option<QosMessage> {
254 self.inner.lock().take()
255 }
256}
257
258#[derive(Debug, Clone, Default)]
263pub struct BitrateSlot {
264 inner: Arc<Mutex<Option<u32>>>,
265}
266
267impl BitrateSlot {
268 pub fn store(&self, value: u32) {
269 *self.inner.lock() = Some(value);
270 }
271
272 pub fn take(&self) -> Option<u32> {
273 self.inner.lock().take()
274 }
275}
276
277#[derive(Debug, Clone)]
283pub struct LinkSender {
284 pub(crate) data: Sender<PipelinePacket>,
285 pub(crate) reconfigure: ReconfigureSlot,
286 pub(crate) qos: QosSlot,
289 pub(crate) bitrate: BitrateSlot,
292 pub(crate) policy: LinkPolicy,
295 pub(crate) dropped: Option<Arc<Mutex<u64>>>,
299 pub(crate) probe: ProbeSlot,
304 pub(crate) transit: Option<TransitRing>,
308 pub(crate) counters: Option<Arc<EdgeCounters>>,
312}
313
314pub(crate) type TransitRing = Arc<Mutex<VecDeque<u64>>>;
319
320#[inline]
322fn stamp_now_ns() -> u64 {
323 #[cfg(feature = "std")]
324 {
325 crate::metrics::monotonic_ns()
326 }
327 #[cfg(not(feature = "std"))]
328 {
329 0
330 }
331}
332
333impl LinkSender {
334 #[cfg(feature = "std")]
338 pub(crate) fn set_policy(&mut self, policy: LinkPolicy) {
339 self.policy = policy;
340 }
341
342 #[cfg(feature = "std")]
344 pub(crate) fn set_drop_counter(&mut self, counter: Arc<Mutex<u64>>) {
345 self.dropped = Some(counter);
346 }
347
348 #[cfg(feature = "std")]
351 pub(crate) fn set_counters(&mut self, counters: Arc<EdgeCounters>) {
352 self.counters = Some(counters);
353 }
354
355 fn record_drop(&self) {
357 if let Some(c) = &self.dropped {
358 *c.lock() += 1;
359 }
360 if let Some(c) = &self.counters {
361 c.record_drop();
362 }
363 }
364
365 fn record_sent(&self, bytes: u64, blocked_since: Option<u64>) {
369 if let Some(c) = &self.counters {
370 let blocked = blocked_since.map_or(0, |t0| stamp_now_ns().saturating_sub(t0));
371 c.record_packet(bytes, blocked);
372 }
373 }
374}
375
376pub(crate) fn packet_bytes(packet: &PipelinePacket) -> u64 {
380 match packet {
381 PipelinePacket::DataFrame(f) => match &f.domain {
382 crate::memory::MemoryDomain::System(s) => s.as_slice().len() as u64,
383 #[cfg(feature = "alloc")]
384 crate::memory::MemoryDomain::SystemView(v) => v.backing().len() as u64,
385 #[cfg(feature = "alloc")]
386 _ => 0,
387 },
388 _ => 0,
389 }
390}
391
392#[derive(Debug)]
397pub struct LinkReceiver {
398 pub(crate) data: Receiver<PipelinePacket>,
399 pub(crate) reconfigure: ReconfigureSlot,
400 pub(crate) qos: QosSlot,
401 pub(crate) bitrate: BitrateSlot,
402 pub(crate) transit: Option<TransitRing>,
405}
406
407impl LinkReceiver {
408 pub fn recv(&self) -> RecvFuture<'_, PipelinePacket> {
409 self.data.recv()
410 }
411
412 pub fn try_recv(&self) -> Option<PipelinePacket> {
414 self.data.try_recv()
415 }
416
417 pub fn fill_percent(&self) -> u8 {
419 self.data.fill_percent()
420 }
421
422 pub fn pop_transit_ns(&self) -> Option<u64> {
429 let ring = self.transit.as_ref()?;
430 let sent = ring.lock().pop_front()?;
431 #[cfg(feature = "std")]
432 {
433 Some(crate::metrics::monotonic_ns().saturating_sub(sent))
434 }
435 #[cfg(not(feature = "std"))]
436 {
437 let _ = sent;
438 Some(0)
439 }
440 }
441
442 pub fn request_reconfigure(&self, r: Reconfigure) {
446 self.reconfigure.store(r);
447 }
448
449 pub fn request_qos(&self, q: QosMessage) {
453 self.qos.store(q);
454 }
455
456 pub fn request_bitrate(&self, bps: u32) {
460 self.bitrate.store(bps);
461 }
462
463 pub(crate) fn qos_slot(&self) -> QosSlot {
468 self.qos.clone()
469 }
470
471 pub(crate) fn reconfigure_slot(&self) -> ReconfigureSlot {
474 self.reconfigure.clone()
475 }
476
477 pub(crate) fn bitrate_slot(&self) -> BitrateSlot {
479 self.bitrate.clone()
480 }
481}
482
483pub fn link(capacity: usize) -> (LinkSender, LinkReceiver) {
486 build_link(capacity, None)
487}
488
489#[cfg(feature = "std")]
494pub(crate) fn link_with_transit(capacity: usize) -> (LinkSender, LinkReceiver) {
495 build_link(capacity, Some(Arc::new(Mutex::new(VecDeque::new()))))
496}
497
498fn build_link(capacity: usize, transit: Option<TransitRing>) -> (LinkSender, LinkReceiver) {
499 let (data_tx, data_rx) = bounded::<PipelinePacket>(capacity);
500 let slot = ReconfigureSlot::default();
501 let qos = QosSlot::default();
502 let bitrate = BitrateSlot::default();
503 (
504 LinkSender {
505 data: data_tx,
506 reconfigure: slot.clone(),
507 qos: qos.clone(),
508 bitrate: bitrate.clone(),
509 policy: LinkPolicy::Block,
510 dropped: None,
511 probe: ProbeSlot::default(),
512 transit: transit.clone(),
513 counters: None,
514 },
515 LinkReceiver {
516 data: data_rx,
517 reconfigure: slot,
518 qos,
519 bitrate,
520 transit,
521 },
522 )
523}
524
525#[derive(Debug, Clone, Copy, PartialEq, Eq)]
527pub enum ProbeAction {
528 Pass,
530 Drop,
532}
533
534pub trait LinkInterceptor {
538 fn on_packet(&self, packet: &PipelinePacket) -> ProbeAction;
539}
540
541#[derive(Clone, Default)]
546pub struct ProbeSlot {
547 inner: Arc<Mutex<Option<Arc<dyn LinkInterceptor + Send + Sync>>>>,
548}
549
550impl core::fmt::Debug for ProbeSlot {
551 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
552 f.debug_struct("ProbeSlot").finish_non_exhaustive()
553 }
554}
555
556impl ProbeSlot {
557 pub fn install(&self, probe: Arc<dyn LinkInterceptor + Send + Sync>) {
559 *self.inner.lock() = Some(probe);
560 }
561
562 pub fn remove(&self) {
564 *self.inner.lock() = None;
565 }
566
567 pub(crate) fn action(&self, packet: &PipelinePacket) -> ProbeAction {
571 match self.inner.lock().as_ref() {
572 Some(probe) => probe.on_packet(packet),
573 None => ProbeAction::Pass,
574 }
575 }
576}
577
578#[derive(Debug)]
599pub struct SenderSink {
600 link: LinkSender,
601 probe: ProbeSlot,
602 upstream_qos: Option<QosSlot>,
609 upstream_reconfigure: Option<ReconfigureSlot>,
612 reconfigure_answered: ReconfigureAnswered,
618 upstream_bitrate: Option<BitrateSlot>,
620 #[cfg(feature = "metadata")]
628 meta_stash: Option<crate::meta::FrameMetaSet>,
629 eos_forwarded: bool,
634 push_wait_probe: Probe,
640 push_phase: PushPhase,
643}
644
645pub(crate) fn advertise_orientation(in_rx: &LinkReceiver, absorbs: bool) {
653 if absorbs {
654 in_rx.request_reconfigure(crate::element::Reconfigure::AbsorbOrientation);
655 }
656}
657
658#[derive(Debug, Clone, Copy, PartialEq, Eq)]
666pub(crate) struct ReconfigureAnswered {
667 pub keyframe: bool,
669 pub orientation: bool,
671}
672
673impl Default for ReconfigureAnswered {
674 fn default() -> Self {
675 ReconfigureAnswered {
679 keyframe: true,
680 orientation: false,
681 }
682 }
683}
684
685#[derive(Debug, Clone, Copy)]
687enum PushPhase {
688 Idle,
690 Sending {
695 bytes: u64,
696 blocked_since: Option<u64>,
697 stamped: bool,
698 },
699}
700
701impl SenderSink {
702 pub fn new(link: LinkSender) -> Self {
703 let probe = link.probe.clone();
707 Self {
708 link,
709 probe,
710 upstream_qos: None,
711 upstream_reconfigure: None,
712 reconfigure_answered: ReconfigureAnswered::default(),
713 upstream_bitrate: None,
714 #[cfg(feature = "metadata")]
715 meta_stash: None,
716 eos_forwarded: false,
717 push_wait_probe: None,
718 push_phase: PushPhase::Idle,
719 }
720 }
721
722 pub(crate) fn set_push_wait_probe(&mut self, probe: Probe) {
727 self.push_wait_probe = probe;
728 }
729
730 fn record_push_wait(&self, since: Option<u64>) {
734 if let (Some(probe), Some(t0)) = (&self.push_wait_probe, since) {
735 probe.add_push_wait(stamp_now_ns().saturating_sub(t0));
736 }
737 }
738
739 fn wants_blocked_stamp(&self) -> bool {
742 self.link.counters.is_some() || self.push_wait_probe.is_some()
743 }
744
745 pub(crate) fn eos_forwarded(&self) -> bool {
747 self.eos_forwarded
748 }
749
750 #[cfg(all(feature = "metadata", feature = "std"))]
756 pub(crate) fn set_meta_stash(&mut self, meta: Option<crate::meta::FrameMetaSet>) {
757 self.meta_stash = meta;
758 }
759
760 pub fn probe(&self) -> ProbeSlot {
763 self.probe.clone()
764 }
765
766 pub(crate) fn relay_qos_to(&mut self, upstream: QosSlot) {
770 self.upstream_qos = Some(upstream);
771 }
772
773 pub(crate) fn relay_reconfigure_to(
778 &mut self,
779 upstream: ReconfigureSlot,
780 answered: ReconfigureAnswered,
781 ) {
782 self.upstream_reconfigure = Some(upstream);
783 self.reconfigure_answered = answered;
784 }
785
786 pub(crate) fn relay_bitrate_to(&mut self, upstream: BitrateSlot) {
788 self.upstream_bitrate = Some(upstream);
789 }
790
791 fn take_reconfigure_or_relay(&self, pre_send: bool) -> Option<crate::element::Reconfigure> {
803 use crate::element::Reconfigure;
804 let r = self.link.reconfigure.take()?;
805 let answered = match &r {
806 Reconfigure::ForceKeyframe => self.reconfigure_answered.keyframe,
807 Reconfigure::AbsorbOrientation => self.reconfigure_answered.orientation,
808 Reconfigure::Propose(_) | Reconfigure::Renegotiate => true,
809 };
810 if answered {
811 if !pre_send && matches!(r, Reconfigure::AbsorbOrientation) {
816 self.link.reconfigure.store(r);
817 return None;
818 }
819 return Some(r);
820 }
821 if let Some(upstream) = &self.upstream_reconfigure {
822 upstream.store(r);
823 }
824 None
825 }
826
827 fn post_send_outcome(&self) -> PushOutcome {
828 if let Some(r) = self.take_reconfigure_or_relay(false) {
829 return PushOutcome::Reconfigure(r);
830 }
831 if let Some(q) = self.link.qos.take() {
832 match &self.upstream_qos {
833 Some(upstream) => upstream.store(q),
834 None => return PushOutcome::Qos(q),
835 }
836 }
837 if let Some(bps) = self.link.bitrate.take() {
838 match &self.upstream_bitrate {
841 Some(upstream) => upstream.store(bps),
842 None => return PushOutcome::Bitrate(bps),
843 }
844 }
845 PushOutcome::Accepted
846 }
847}
848
849impl SenderSink {
850 fn poll_blocking_send(
855 &mut self,
856 cx: &mut core::task::Context<'_>,
857 packet: &mut Option<PipelinePacket>,
858 bytes: u64,
859 blocked_since: Option<u64>,
860 stamped: bool,
861 ) -> Poll<Result<PushOutcome, G2gError>> {
862 match self.link.data.poll_send(cx, packet) {
863 Poll::Pending => Poll::Pending,
864 Poll::Ready(Ok(())) => {
868 self.push_phase = PushPhase::Idle;
869 self.link.record_sent(bytes, blocked_since);
870 self.record_push_wait(blocked_since);
871 Poll::Ready(Ok(self.post_send_outcome()))
872 }
873 Poll::Ready(Err(SendError::Closed)) => {
874 self.push_phase = PushPhase::Idle;
875 if stamped {
876 if let Some(ring) = &self.link.transit {
877 ring.lock().pop_back();
878 }
879 }
880 packet.take();
883 Poll::Ready(Err(G2gError::Shutdown))
884 }
885 Poll::Ready(Err(SendError::Full)) => unreachable!("poll_send never returns Full"),
886 }
887 }
888}
889
890impl OutputSink for SenderSink {
891 fn begin_push(&mut self) {
892 self.push_phase = PushPhase::Idle;
895 }
896
897 fn poll_push(
898 &mut self,
899 cx: &mut core::task::Context<'_>,
900 packet_slot: &mut Option<PipelinePacket>,
901 ) -> Poll<Result<PushOutcome, G2gError>> {
902 if let PushPhase::Sending {
903 bytes,
904 blocked_since,
905 stamped,
906 } = self.push_phase
907 {
908 return self.poll_blocking_send(cx, packet_slot, bytes, blocked_since, stamped);
909 }
910 let packet = packet_slot
911 .as_mut()
912 .expect("poll_push called without a packet");
913 #[cfg(feature = "metadata")]
918 if let (Some(stash), PipelinePacket::DataFrame(frame)) = (&self.meta_stash, &mut *packet) {
919 if frame.meta.is_empty() {
920 frame.meta = stash.clone();
921 }
922 }
923 if self.probe.action(packet) == ProbeAction::Drop {
925 packet_slot.take();
926 return Poll::Ready(Ok(PushOutcome::Accepted));
927 }
928 if !matches!(packet, PipelinePacket::Eos) {
937 if let Some(r) = self.take_reconfigure_or_relay(true) {
938 packet_slot.take();
939 return Poll::Ready(Ok(PushOutcome::Reconfigure(r)));
940 }
941 }
942 if matches!(packet, PipelinePacket::Eos) {
945 self.eos_forwarded = true;
946 }
947 if let (PipelinePacket::CapsChanged(caps), Some(c)) = (&*packet, &self.link.counters) {
950 c.record_caps(caps);
951 }
952 #[cfg(feature = "std")]
955 if let (Some(probe), PipelinePacket::DataFrame(frame)) = (&self.push_wait_probe, &*packet) {
956 if frame.timing.arrival_ns != 0 {
957 probe.record_age_at_emit(stamp_now_ns().saturating_sub(frame.timing.arrival_ns));
958 }
959 }
960 let is_data = matches!(packet, PipelinePacket::DataFrame(_));
965 let bytes = packet_bytes(packet);
967 if is_data && self.link.policy != LinkPolicy::Block {
968 let taken = packet_slot.take().expect("packet checked above");
969 match self.link.policy {
970 LinkPolicy::DropNewest => match self.link.data.try_send(taken) {
971 Ok(()) => self.link.record_sent(bytes, None),
972 Err((_dropped, SendError::Full)) => self.link.record_drop(),
974 Err((_v, SendError::Closed)) => return Poll::Ready(Err(G2gError::Shutdown)),
975 },
976 LinkPolicy::DropOldest => match self.link.data.try_send(taken) {
977 Ok(()) => self.link.record_sent(bytes, None),
978 Err((returned, SendError::Full)) => {
979 if self
983 .link
984 .data
985 .evict_front_matching(|p| matches!(p, PipelinePacket::DataFrame(_)))
986 .is_some()
987 {
988 self.link.record_drop();
989 match self.link.data.try_send(returned) {
990 Ok(()) => self.link.record_sent(bytes, None),
991 Err((_v, SendError::Closed)) => {
992 return Poll::Ready(Err(G2gError::Shutdown))
993 }
994 Err((_v, SendError::Full)) => {
995 unreachable!("a slot was just freed by eviction")
996 }
997 }
998 } else {
999 *packet_slot = Some(returned);
1000 let blocked_since = self.wants_blocked_stamp().then(stamp_now_ns);
1001 self.push_phase = PushPhase::Sending {
1002 bytes,
1003 blocked_since,
1004 stamped: false,
1005 };
1006 return self.poll_blocking_send(
1007 cx,
1008 packet_slot,
1009 bytes,
1010 blocked_since,
1011 false,
1012 );
1013 }
1014 }
1015 Err((_v, SendError::Closed)) => return Poll::Ready(Err(G2gError::Shutdown)),
1016 },
1017 LinkPolicy::Block => unreachable!("guarded by policy != Block"),
1018 }
1019 return Poll::Ready(Ok(self.post_send_outcome()));
1020 }
1021 let stamped = is_data && self.link.transit.is_some();
1025 if stamped {
1026 if let Some(ring) = &self.link.transit {
1027 ring.lock().push_back(stamp_now_ns());
1028 }
1029 }
1030 let blocked_since = self.wants_blocked_stamp().then(stamp_now_ns);
1034 self.push_phase = PushPhase::Sending {
1035 bytes,
1036 blocked_since,
1037 stamped,
1038 };
1039 self.poll_blocking_send(cx, packet_slot, bytes, blocked_since, stamped)
1040 }
1041}
1042
1043#[cfg(test)]
1044mod link_tests {
1045 use super::*;
1046 use crate::caps::{Caps, Dim, Rate, VideoCodec};
1047 use crate::element::OutputSinkExt;
1048 use crate::frame::{Frame, FrameTiming};
1049 use crate::memory::{MemoryDomain, SystemSlice};
1050 use alloc::boxed::Box;
1051 use alloc::vec::Vec;
1052 use core::pin::Pin;
1053 use core::task::{Context, Poll, RawWaker, RawWakerVTable, Waker};
1054
1055 static NOOP_VTABLE: RawWakerVTable = RawWakerVTable::new(
1059 |_| RawWaker::new(core::ptr::null(), &NOOP_VTABLE),
1060 |_| {},
1061 |_| {},
1062 |_| {},
1063 );
1064
1065 fn noop_waker() -> Waker {
1066 unsafe { Waker::from_raw(RawWaker::new(core::ptr::null(), &NOOP_VTABLE)) }
1069 }
1070
1071 fn run_to_ready<F: core::future::Future>(mut fut: F) -> F::Output {
1072 let waker = noop_waker();
1073 let mut cx = Context::from_waker(&waker);
1074 let mut pinned = unsafe { Pin::new_unchecked(&mut fut) };
1077 match pinned.as_mut().poll(&mut cx) {
1078 Poll::Ready(v) => v,
1079 Poll::Pending => panic!("link_tests::run_to_ready saw Pending"),
1080 }
1081 }
1082
1083 fn dummy_frame() -> PipelinePacket {
1084 PipelinePacket::DataFrame(Frame {
1085 domain: MemoryDomain::System(SystemSlice::from_boxed(Box::new([0u8; 4]))),
1086 timing: FrameTiming::default(),
1087 sequence: 0,
1088 meta: Default::default(),
1089 })
1090 }
1091
1092 fn proposed_caps() -> Caps {
1093 Caps::CompressedVideo {
1094 codec: VideoCodec::H264,
1095 width: Dim::Fixed(1280),
1096 height: Dim::Fixed(720),
1097 framerate: Rate::Any,
1098 }
1099 }
1100
1101 #[test]
1102 fn push_returns_accepted_when_no_reconfigure_pending() {
1103 let (tx, _rx) = link(2);
1104 let mut sink = SenderSink::new(tx);
1105 let outcome = run_to_ready(sink.push(dummy_frame())).expect("send ok");
1106 assert_eq!(outcome, PushOutcome::Accepted);
1107 }
1108
1109 #[test]
1110 fn request_reconfigure_surfaces_on_next_push() {
1111 let (tx, rx) = link(2);
1112 let mut sink = SenderSink::new(tx);
1113
1114 rx.request_reconfigure(Reconfigure::Propose(proposed_caps()));
1116
1117 let outcome = run_to_ready(sink.push(dummy_frame())).expect("push ok");
1122 match outcome {
1123 PushOutcome::Reconfigure(Reconfigure::Propose(c)) => {
1124 assert_eq!(c, proposed_caps());
1125 }
1126 other => panic!("expected Reconfigure::Propose, got {other:?}"),
1127 }
1128
1129 assert!(
1131 rx.try_recv().is_none(),
1132 "packet must not enqueue when reconfigure pending"
1133 );
1134 }
1135
1136 #[test]
1137 fn second_push_returns_accepted_after_reconfigure_drained() {
1138 let (tx, rx) = link(2);
1139 let mut sink = SenderSink::new(tx);
1140
1141 rx.request_reconfigure(Reconfigure::Renegotiate);
1142 let first = run_to_ready(sink.push(dummy_frame())).unwrap();
1143 assert!(matches!(first, PushOutcome::Reconfigure(_)));
1144
1145 let second = run_to_ready(sink.push(dummy_frame())).unwrap();
1146 assert_eq!(second, PushOutcome::Accepted);
1147 }
1148
1149 #[test]
1150 fn request_qos_surfaces_after_the_packet_is_sent() {
1151 let (tx, rx) = link(2);
1152 let mut sink = SenderSink::new(tx);
1153
1154 rx.request_qos(QosMessage {
1157 jitter_ns: 5_000_000,
1158 running_time_ns: 100,
1159 });
1160 let outcome = run_to_ready(sink.push(dummy_frame())).expect("push ok");
1161 match outcome {
1162 PushOutcome::Qos(q) => {
1163 assert_eq!(q.jitter_ns, 5_000_000);
1164 assert_eq!(q.running_time_ns, 100);
1165 }
1166 other => panic!("expected Qos, got {other:?}"),
1167 }
1168 assert!(
1170 rx.try_recv().is_some(),
1171 "QoS is advisory; the frame still flowed"
1172 );
1173 }
1174
1175 #[test]
1176 fn reconfigure_takes_priority_over_qos() {
1177 let (tx, rx) = link(2);
1178 let mut sink = SenderSink::new(tx);
1179
1180 rx.request_qos(QosMessage {
1182 jitter_ns: 1_000,
1183 running_time_ns: 0,
1184 });
1185 rx.request_reconfigure(Reconfigure::Renegotiate);
1186 let first = run_to_ready(sink.push(dummy_frame())).unwrap();
1187 assert!(
1188 matches!(first, PushOutcome::Reconfigure(_)),
1189 "reconfigure first"
1190 );
1191
1192 let second = run_to_ready(sink.push(dummy_frame())).unwrap();
1193 assert!(
1194 matches!(second, PushOutcome::Qos(_)),
1195 "QoS surfaces once reconfigure drained"
1196 );
1197 }
1198
1199 #[test]
1200 fn try_recv_returns_value_then_none() {
1201 let (tx, rx) = bounded::<u32>(2);
1202 assert_eq!(rx.try_recv(), None, "empty queue");
1203 tx.try_send(7).unwrap();
1204 assert_eq!(rx.try_recv(), Some(7));
1205 assert_eq!(rx.try_recv(), None, "drained");
1206 }
1207
1208 #[test]
1209 fn try_recv_drains_then_none_after_senders_drop() {
1210 let (tx, rx) = bounded::<u32>(2);
1211 tx.try_send(1).unwrap();
1212 drop(tx);
1213 assert_eq!(rx.try_recv(), Some(1), "remaining value still drains");
1214 assert_eq!(rx.try_recv(), None, "empty and closed");
1215 }
1216
1217 #[test]
1221 fn relay_is_decided_per_variant() {
1222 let (up_tx, up_rx) = link(2);
1223 let (down_tx, down_rx) = link(2);
1224 drop(up_tx);
1226 let mut adapter = SenderSink::new(down_tx);
1227 adapter.relay_reconfigure_to(
1228 up_rx.reconfigure_slot(),
1229 ReconfigureAnswered {
1230 keyframe: true,
1231 orientation: false,
1232 },
1233 );
1234
1235 down_rx.request_reconfigure(Reconfigure::ForceKeyframe);
1236 let outcome = run_to_ready(adapter.push(dummy_frame())).expect("push ok");
1237 assert!(
1238 matches!(
1239 outcome,
1240 PushOutcome::Reconfigure(Reconfigure::ForceKeyframe)
1241 ),
1242 "an answered variant surfaces to the producer, got {outcome:?}"
1243 );
1244 assert!(
1245 up_rx.reconfigure.take().is_none(),
1246 "an answered variant must not also travel upstream"
1247 );
1248
1249 down_rx.request_reconfigure(Reconfigure::AbsorbOrientation);
1250 let outcome = run_to_ready(adapter.push(dummy_frame())).expect("push ok");
1251 assert_eq!(
1252 outcome,
1253 PushOutcome::Accepted,
1254 "an unanswered variant is relayed, not surfaced"
1255 );
1256 assert!(
1257 matches!(
1258 up_rx.reconfigure.take(),
1259 Some(Reconfigure::AbsorbOrientation)
1260 ),
1261 "the advertisement must reach the upstream link"
1262 );
1263 assert!(
1264 down_rx.try_recv().is_some(),
1265 "a relayed variant does not hold the packet back"
1266 );
1267 }
1268
1269 #[test]
1272 fn an_answered_orientation_surfaces_while_a_keyframe_relays() {
1273 let (up_tx, up_rx) = link(2);
1274 let (down_tx, down_rx) = link(2);
1275 drop(up_tx);
1276 let mut adapter = SenderSink::new(down_tx);
1277 adapter.relay_reconfigure_to(
1278 up_rx.reconfigure_slot(),
1279 ReconfigureAnswered {
1280 keyframe: false,
1281 orientation: true,
1282 },
1283 );
1284
1285 down_rx.request_reconfigure(Reconfigure::AbsorbOrientation);
1286 let outcome = run_to_ready(adapter.push(dummy_frame())).expect("push ok");
1287 assert!(
1288 matches!(
1289 outcome,
1290 PushOutcome::Reconfigure(Reconfigure::AbsorbOrientation)
1291 ),
1292 "the flip has to see the advertisement, got {outcome:?}"
1293 );
1294 assert!(
1295 down_rx.try_recv().is_none(),
1296 "the pre-send check holds the packet back for the producer to resend"
1297 );
1298
1299 down_rx.request_reconfigure(Reconfigure::ForceKeyframe);
1300 let outcome = run_to_ready(adapter.push(dummy_frame())).expect("push ok");
1301 assert_eq!(outcome, PushOutcome::Accepted);
1302 assert!(matches!(
1303 up_rx.reconfigure.take(),
1304 Some(Reconfigure::ForceKeyframe)
1305 ));
1306 }
1307
1308 #[test]
1312 fn an_eos_crosses_even_with_a_reconfigure_pending() {
1313 let (tx, rx) = link(2);
1314 let mut sink = SenderSink::new(tx);
1315 rx.request_reconfigure(Reconfigure::AbsorbOrientation);
1316 let outcome = run_to_ready(sink.push(PipelinePacket::Eos)).expect("push ok");
1317 assert_eq!(outcome, PushOutcome::Accepted);
1318 assert!(
1319 matches!(rx.try_recv(), Some(PipelinePacket::Eos)),
1320 "the end of stream must still reach the consumer"
1321 );
1322 }
1323
1324 #[test]
1329 fn an_unanswered_variant_without_a_relay_target_is_dropped() {
1330 let (tx, rx) = link(2);
1331 let mut adapter = SenderSink::new(tx);
1332 adapter.reconfigure_answered = ReconfigureAnswered {
1333 keyframe: true,
1334 orientation: false,
1335 };
1336
1337 rx.request_reconfigure(Reconfigure::AbsorbOrientation);
1338 let outcome = run_to_ready(adapter.push(dummy_frame())).expect("push ok");
1339 assert_eq!(outcome, PushOutcome::Accepted);
1340 assert!(rx.try_recv().is_some(), "the frame still crossed");
1341 }
1342
1343 #[test]
1344 fn latest_reconfigure_overwrites_older_pending() {
1345 let (tx, rx) = link(2);
1346 let mut sink = SenderSink::new(tx);
1347
1348 rx.request_reconfigure(Reconfigure::Renegotiate);
1350 rx.request_reconfigure(Reconfigure::Propose(proposed_caps()));
1351
1352 let outcome = run_to_ready(sink.push(dummy_frame())).unwrap();
1353 match outcome {
1354 PushOutcome::Reconfigure(Reconfigure::Propose(c)) => {
1355 assert_eq!(c, proposed_caps(), "newest proposal must win");
1356 }
1357 other => panic!("expected newest Propose, got {other:?}"),
1358 }
1359 }
1360
1361 fn frame_seq(seq: u64) -> PipelinePacket {
1362 PipelinePacket::DataFrame(Frame {
1363 domain: MemoryDomain::System(SystemSlice::from_boxed(Box::new([0u8; 4]))),
1364 timing: FrameTiming::default(),
1365 sequence: seq,
1366 meta: Default::default(),
1367 })
1368 }
1369
1370 struct DropOdd;
1372 impl LinkInterceptor for DropOdd {
1373 fn on_packet(&self, packet: &PipelinePacket) -> ProbeAction {
1374 match packet {
1375 PipelinePacket::DataFrame(f) if f.sequence % 2 == 1 => ProbeAction::Drop,
1376 _ => ProbeAction::Pass,
1377 }
1378 }
1379 }
1380
1381 #[test]
1382 fn installed_probe_drops_selected_packets() {
1383 let (tx, rx) = link(8);
1384 let mut sink = SenderSink::new(tx);
1385 sink.probe().install(Arc::new(DropOdd));
1386
1387 for seq in 0..4 {
1388 run_to_ready(sink.push(frame_seq(seq))).unwrap();
1389 }
1390
1391 let mut got = Vec::new();
1392 while let Some(PipelinePacket::DataFrame(f)) = rx.try_recv() {
1393 got.push(f.sequence);
1394 }
1395 assert_eq!(got, [0, 2], "odd-sequence frames dropped by the probe");
1396 }
1397
1398 #[test]
1399 fn removed_probe_lets_packets_pass_again() {
1400 let (tx, rx) = link(8);
1401 let mut sink = SenderSink::new(tx);
1402 let probe = sink.probe();
1403
1404 probe.install(Arc::new(DropOdd));
1405 run_to_ready(sink.push(frame_seq(1))).unwrap(); probe.remove();
1407 run_to_ready(sink.push(frame_seq(3))).unwrap(); let mut got = Vec::new();
1410 while let Some(PipelinePacket::DataFrame(f)) = rx.try_recv() {
1411 got.push(f.sequence);
1412 }
1413 assert_eq!(got, [3], "after remove(), the odd frame passes");
1414 }
1415
1416 #[cfg(feature = "std")]
1417 fn drained_sequences(rx: &LinkReceiver) -> Vec<u64> {
1418 let mut got = Vec::new();
1419 while let Some(PipelinePacket::DataFrame(f)) = rx.try_recv() {
1420 got.push(f.sequence);
1421 }
1422 got
1423 }
1424
1425 #[cfg(feature = "std")]
1427 #[test]
1428 fn drop_newest_discards_incoming_when_full() {
1429 let (mut tx, rx) = link(2);
1430 tx.set_policy(LinkPolicy::DropNewest);
1431 let counter = Arc::new(Mutex::new(0u64));
1432 tx.set_drop_counter(counter.clone());
1433 let mut sink = SenderSink::new(tx);
1434
1435 for seq in 0..2 {
1438 assert_eq!(
1439 run_to_ready(sink.push(frame_seq(seq))).unwrap(),
1440 PushOutcome::Accepted
1441 );
1442 }
1443 assert_eq!(
1444 run_to_ready(sink.push(frame_seq(2))).unwrap(),
1445 PushOutcome::Accepted
1446 );
1447
1448 assert_eq!(
1449 drained_sequences(&rx),
1450 [0, 1],
1451 "drop-newest keeps the oldest"
1452 );
1453 assert_eq!(*counter.lock(), 1);
1454 }
1455
1456 #[cfg(feature = "std")]
1457 #[test]
1458 fn drop_oldest_evicts_front_when_full() {
1459 let (mut tx, rx) = link(2);
1460 tx.set_policy(LinkPolicy::DropOldest);
1461 let counter = Arc::new(Mutex::new(0u64));
1462 tx.set_drop_counter(counter.clone());
1463 let mut sink = SenderSink::new(tx);
1464
1465 for seq in 0..2 {
1466 run_to_ready(sink.push(frame_seq(seq))).unwrap();
1467 }
1468 assert_eq!(
1470 run_to_ready(sink.push(frame_seq(2))).unwrap(),
1471 PushOutcome::Accepted
1472 );
1473
1474 assert_eq!(
1475 drained_sequences(&rx),
1476 [1, 2],
1477 "drop-oldest keeps the newest"
1478 );
1479 assert_eq!(*counter.lock(), 1);
1480 }
1481
1482 #[test]
1483 fn fill_percent_tracks_link_occupancy() {
1484 let (tx, rx) = link(4);
1485 assert_eq!(rx.fill_percent(), 0, "empty link reads 0%");
1486 let mut sink = SenderSink::new(tx);
1487 run_to_ready(sink.push(frame_seq(0))).unwrap();
1488 run_to_ready(sink.push(frame_seq(1))).unwrap();
1489 assert_eq!(rx.fill_percent(), 50, "2 of 4 slots = 50%");
1490 run_to_ready(sink.push(frame_seq(2))).unwrap();
1491 run_to_ready(sink.push(frame_seq(3))).unwrap();
1492 assert_eq!(rx.fill_percent(), 100, "full link reads 100%");
1493 rx.try_recv();
1494 assert_eq!(rx.fill_percent(), 75, "after one drain, 3 of 4 = 75%");
1495 }
1496
1497 #[cfg(feature = "std")]
1498 #[test]
1499 fn leaky_links_never_drop_control_packets() {
1500 let (mut tx, rx) = link(1);
1503 tx.set_policy(LinkPolicy::DropNewest);
1504 let mut sink = SenderSink::new(tx);
1505 run_to_ready(sink.push(frame_seq(0))).unwrap();
1506
1507 let waker = noop_waker();
1508 let mut cx = Context::from_waker(&waker);
1509 let mut fut = core::pin::pin!(sink.push(PipelinePacket::CapsChanged(proposed_caps())));
1510 assert!(
1511 matches!(fut.as_mut().poll(&mut cx), Poll::Pending),
1512 "a control packet blocks on a full leaky link, never dropped"
1513 );
1514
1515 assert_eq!(drained_sequences(&rx), [0]);
1517 }
1518}