1#![cfg_attr(
15 not(any(
16 feature = "transport-implementer-api",
17 feature = "selected-wire-user-api"
18 )),
19 allow(dead_code)
20)]
21
22#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
52use std::{
53 any::Any,
54 collections::HashMap,
55 sync::{Arc, Mutex},
56};
57use std::{io::Read, marker::PhantomData};
58
59#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
60use async_trait::async_trait;
61#[cfg(feature = "owned-frame-transport")]
62use bytes::Bytes;
63#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
64use tracing::warn;
65
66#[cfg(feature = "zero-copy-transport")]
67use crate::payload::loan::BorrowPayload;
68use crate::payload::{codec::ReadDecodePayload, UWireError};
69use crate::wire::NativePrefixFrameMetadataCodec;
70use crate::wire::{ProtobufWire, StableContainerWireFormat};
71use crate::wire::{UWire, UWireMetadataCodecFor, UWirePayload};
72use crate::{validate_frame_view_for_transport, UFrameMetadata, UFrameView, UStatus};
73#[cfg(feature = "zero-copy-transport")]
74use crate::{
75 LoanedPayload, PayloadAlignment, ULoanedContiguousZeroCopyRxFrame, UTxBuffer, UTxLoanSpec,
76 UUninitTxBuffer, UZeroCopyListener, UZeroCopyRxLease, UZeroCopyTransport,
77 UZeroCopyTransportImpl, UZeroCopyUninitTransportImpl,
78};
79#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
80use crate::{UCode, UUri};
81#[cfg(feature = "owned-frame-transport")]
82use crate::{UOwnedFrame, UOwnedListener, UOwnedTransportImpl};
83
84pub struct UWireTransport<TCore, W, C>
86where
87 W: UWire,
88 C: UWireMetadataCodecFor<W>,
89{
90 core: TCore,
91 wire: W,
92 metadata_codec: C,
93 #[cfg(feature = "zero-copy-transport")]
94 zero_copy_listeners: Mutex<HashMap<WireListenerKey, Arc<dyn Any + Send + Sync>>>,
95 #[cfg(feature = "owned-frame-transport")]
96 owned_listeners: Mutex<HashMap<WireListenerKey, Arc<dyn Any + Send + Sync>>>,
97}
98
99impl<TCore, W, C> core::fmt::Debug for UWireTransport<TCore, W, C>
100where
101 W: UWire,
102 C: UWireMetadataCodecFor<W>,
103{
104 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
105 f.debug_struct("UWireTransport").finish_non_exhaustive()
106 }
107}
108
109impl<TCore, W, C> UWireTransport<TCore, W, C>
110where
111 W: UWire,
112 C: UWireMetadataCodecFor<W>,
113{
114 #[must_use]
116 pub fn new(core: TCore, wire: W, metadata_codec: C) -> Self {
117 Self {
118 core,
119 wire,
120 metadata_codec,
121 #[cfg(feature = "zero-copy-transport")]
122 zero_copy_listeners: Mutex::new(HashMap::new()),
123 #[cfg(feature = "owned-frame-transport")]
124 owned_listeners: Mutex::new(HashMap::new()),
125 }
126 }
127
128 #[must_use]
130 pub fn with_wire_and_metadata_codec(core: TCore, wire: W, metadata_codec: C) -> Self {
131 Self::new(core, wire, metadata_codec)
132 }
133
134 #[must_use]
136 pub fn core(&self) -> &TCore {
137 &self.core
138 }
139
140 #[must_use]
142 pub fn core_mut(&mut self) -> &mut TCore {
143 &mut self.core
144 }
145
146 #[must_use]
148 pub fn metadata_codec(&self) -> &C {
149 &self.metadata_codec
150 }
151
152 #[must_use]
154 pub fn into_parts(self) -> (TCore, W, C) {
155 (self.core, self.wire, self.metadata_codec)
156 }
157}
158
159pub type UNativePrefixWireTransport<TCore, W> =
165 UWireTransport<TCore, W, NativePrefixFrameMetadataCodec>;
166
167pub type ProtobufWireTransport<TCore> = UNativePrefixWireTransport<TCore, ProtobufWire>;
169
170pub type StableContainerWireTransport<TCore> =
172 UNativePrefixWireTransport<TCore, StableContainerWireFormat>;
173
174pub trait UWithNativePrefixWire: Sized {
176 #[must_use]
178 fn into_native_prefix_wire_transport<W>(self, wire: W) -> UNativePrefixWireTransport<Self, W>
179 where
180 W: UWire;
181
182 #[must_use]
184 fn into_protobuf_transport(self) -> ProtobufWireTransport<Self>;
185
186 #[must_use]
188 fn into_stable_container_transport(self) -> StableContainerWireTransport<Self>;
189}
190
191impl<TCore> UWithNativePrefixWire for TCore {
192 fn into_native_prefix_wire_transport<W>(self, wire: W) -> UNativePrefixWireTransport<Self, W>
193 where
194 W: UWire,
195 {
196 UWireTransport::new(self, wire, NativePrefixFrameMetadataCodec)
197 }
198
199 fn into_protobuf_transport(self) -> ProtobufWireTransport<Self> {
200 self.into_native_prefix_wire_transport(ProtobufWire)
201 }
202
203 fn into_stable_container_transport(self) -> StableContainerWireTransport<Self> {
204 self.into_native_prefix_wire_transport(StableContainerWireFormat)
205 }
206}
207
208pub trait UHasWire {
210 type Wire: UWire;
212
213 #[must_use]
215 fn wire(&self) -> &Self::Wire;
216}
217
218impl<TCore, W, C> UHasWire for UWireTransport<TCore, W, C>
219where
220 W: UWire,
221 C: UWireMetadataCodecFor<W>,
222{
223 type Wire = W;
224
225 fn wire(&self) -> &Self::Wire {
226 &self.wire
227 }
228}
229
230#[cfg(feature = "zero-copy-transport")]
231pub trait USelectedWireZeroCopyTransport: UZeroCopyTransport + UHasWire {
233 type MetadataCodec: UWireMetadataCodecFor<Self::Wire>;
235}
236
237#[cfg(feature = "zero-copy-transport")]
238impl<TCore, W, C> USelectedWireZeroCopyTransport for UWireTransport<TCore, W, C>
239where
240 TCore: UZeroCopyTransportCore,
241 W: UWire + Send + Sync + 'static,
242 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
243{
244 type MetadataCodec = C;
245}
246
247#[cfg(feature = "zero-copy-transport")]
248#[derive(Clone, Debug, PartialEq)]
250pub struct PreparedTxLoanSpec {
251 metadata: UFrameMetadata,
252 encoded_metadata: Vec<u8>,
253 payload_len: usize,
254 payload_alignment: PayloadAlignment,
255}
256
257#[cfg(feature = "zero-copy-transport")]
258impl PreparedTxLoanSpec {
259 pub fn from_validated<W, C>(spec: UTxLoanSpec, codec: &C) -> Result<Self, UStatus>
265 where
266 W: UWire,
267 C: UWireMetadataCodecFor<W>,
268 {
269 let encoded_metadata =
270 codec.encode_frame_metadata(W::metadata_context(), spec.metadata())?;
271 Ok(Self {
272 metadata: spec.metadata().clone(),
273 encoded_metadata,
274 payload_len: spec.payload_len(),
275 payload_alignment: spec.payload_alignment_proof(),
276 })
277 }
278
279 pub fn from_encoded_parts(
293 metadata: UFrameMetadata,
294 encoded_metadata: impl Into<Vec<u8>>,
295 payload_len: usize,
296 payload_alignment: usize,
297 ) -> Result<Self, UStatus> {
298 let spec = if metadata.payload_encoding().is_some() {
299 UTxLoanSpec::payload(metadata, payload_len, payload_alignment)?
300 } else {
301 if payload_len != 0 {
302 return Err(UStatus::fail_with_code(
303 UCode::InvalidArgument,
304 "prepared TX spec without payload encoding cannot carry payload bytes",
305 ));
306 }
307 if payload_alignment != 1 {
308 return Err(UStatus::fail_with_code(
309 UCode::InvalidArgument,
310 "prepared TX spec without payload uses alignment 1",
311 ));
312 }
313 UTxLoanSpec::no_payload(metadata)?
314 };
315 Ok(Self {
316 metadata: spec.metadata().clone(),
317 encoded_metadata: encoded_metadata.into(),
318 payload_len: spec.payload_len(),
319 payload_alignment: spec.payload_alignment_proof(),
320 })
321 }
322
323 #[must_use]
325 pub fn metadata(&self) -> &UFrameMetadata {
326 &self.metadata
327 }
328
329 #[must_use]
331 pub fn encoded_metadata(&self) -> &[u8] {
332 &self.encoded_metadata
333 }
334
335 #[must_use]
337 pub fn payload_len(&self) -> usize {
338 self.payload_len
339 }
340
341 #[must_use]
343 pub fn payload_alignment(&self) -> usize {
344 self.payload_alignment.as_usize()
345 }
346
347 #[must_use]
349 pub fn payload_alignment_proof(&self) -> PayloadAlignment {
350 self.payload_alignment
351 }
352
353 #[must_use]
355 pub fn has_payload(&self) -> bool {
356 self.metadata.payload_encoding().is_some()
357 }
358
359 #[must_use]
361 pub fn into_parts(self) -> (UFrameMetadata, Vec<u8>, usize, usize) {
362 (
363 self.metadata,
364 self.encoded_metadata,
365 self.payload_len,
366 self.payload_alignment.as_usize(),
367 )
368 }
369}
370
371pub trait UEncodedRxFrame {
377 type PayloadReader<'a>: Read + 'a
379 where
380 Self: 'a;
381 type PayloadSlices<'a>: Iterator<Item = &'a [u8]> + 'a
383 where
384 Self: 'a;
385
386 fn encoded_metadata(&self) -> &[u8];
388
389 fn payload_len(&self) -> usize;
391
392 fn payload_reader(&self) -> Self::PayloadReader<'_>;
394
395 fn payload_slices(&self) -> Self::PayloadSlices<'_>;
397
398 fn try_contiguous_payload(&self) -> Option<&[u8]> {
400 None
401 }
402}
403
404#[cfg(feature = "zero-copy-transport")]
405pub trait UEncodedLoanedRxFrame: UEncodedRxFrame {
407 fn loaned_contiguous_payload(&self) -> Result<LoanedPayload<'_>, UWireError>;
412}
413
414pub struct UWireRx<Rx, W, C>
416where
417 W: UWire,
418{
419 metadata: UFrameMetadata,
420 raw: Rx,
421 _wire: PhantomData<W>,
422 _metadata_codec: PhantomData<C>,
423}
424
425impl<Rx, W, C> core::fmt::Debug for UWireRx<Rx, W, C>
426where
427 W: UWire,
428{
429 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
430 f.debug_struct("UWireRx").finish_non_exhaustive()
431 }
432}
433
434impl<Rx, W, C> UWireRx<Rx, W, C>
435where
436 Rx: UEncodedRxFrame,
437 W: UWire,
438 C: UWireMetadataCodecFor<W>,
439{
440 pub fn try_from_encoded(raw: Rx, codec: &C) -> Result<Self, UStatus> {
447 let metadata =
448 codec.decode_frame_metadata(W::metadata_context(), raw.encoded_metadata())?;
449 let frame = Self {
450 metadata,
451 raw,
452 _wire: PhantomData,
453 _metadata_codec: PhantomData,
454 };
455 validate_frame_view_for_transport(&frame)?;
456 Ok(frame)
457 }
458
459 #[must_use]
461 pub fn raw(&self) -> &Rx {
462 &self.raw
463 }
464
465 #[must_use]
467 pub fn into_raw(self) -> Rx {
468 self.raw
469 }
470
471 pub fn decode_payload<T>(&self) -> Result<T, UWireError>
483 where
484 W: UWirePayload<T>,
485 <W as UWirePayload<T>>::Codec: ReadDecodePayload<T>,
486 {
487 <<W as UWirePayload<T>>::Codec as crate::payload::codec::PayloadCodec>::verify_encoding(
488 self.metadata.payload_encoding(),
489 )?;
490 if !self.has_payload() {
491 return Err(UWireError::MissingPayload);
492 }
493 <W as UWirePayload<T>>::Codec::decode_payload_from_reader(
494 self.payload_reader(),
495 self.payload_len(),
496 )
497 }
498}
499
500#[cfg(feature = "zero-copy-transport")]
501impl<Rx, W, C> UWireRx<Rx, W, C>
502where
503 Rx: UEncodedLoanedRxFrame,
504 W: UWire,
505 C: UWireMetadataCodecFor<W>,
506{
507 pub fn borrow_payload<T>(&self) -> Result<&T, UWireError>
519 where
520 W: UWirePayload<T>,
521 <W as UWirePayload<T>>::Codec: BorrowPayload<T>,
522 {
523 self.borrow_payload_as::<<W as UWirePayload<T>>::Codec, T>()
524 }
525}
526
527impl<Rx, W, C> UFrameView for UWireRx<Rx, W, C>
528where
529 Rx: UEncodedRxFrame,
530 W: UWire,
531 C: UWireMetadataCodecFor<W>,
532{
533 type PayloadReader<'a>
534 = Rx::PayloadReader<'a>
535 where
536 Self: 'a;
537 type PayloadSlices<'a>
538 = Rx::PayloadSlices<'a>
539 where
540 Self: 'a;
541
542 fn metadata(&self) -> &UFrameMetadata {
543 &self.metadata
544 }
545
546 fn payload_len(&self) -> usize {
547 self.raw.payload_len()
548 }
549
550 fn has_payload(&self) -> bool {
551 self.payload_len() > 0 || self.metadata.payload_encoding().is_some()
552 }
553
554 fn payload_reader(&self) -> Self::PayloadReader<'_> {
555 self.raw.payload_reader()
556 }
557
558 fn payload_slices(&self) -> Self::PayloadSlices<'_> {
559 self.raw.payload_slices()
560 }
561
562 fn try_contiguous_payload(&self) -> Option<&[u8]> {
563 self.raw.try_contiguous_payload()
564 }
565}
566
567#[cfg(feature = "zero-copy-transport")]
568impl<Rx, W, C> UZeroCopyRxLease for UWireRx<Rx, W, C>
569where
570 Rx: UEncodedRxFrame,
571 W: UWire,
572 C: UWireMetadataCodecFor<W>,
573{
574}
575
576#[cfg(feature = "zero-copy-transport")]
577impl<Rx, W, C> ULoanedContiguousZeroCopyRxFrame for UWireRx<Rx, W, C>
578where
579 Rx: UEncodedLoanedRxFrame,
580 W: UWire,
581 C: UWireMetadataCodecFor<W>,
582{
583 fn loaned_contiguous_payload(&self) -> Result<LoanedPayload<'_>, UWireError> {
584 self.raw.loaned_contiguous_payload()
585 }
586}
587
588#[cfg(feature = "zero-copy-transport")]
589#[async_trait]
591pub trait UEncodedZeroCopyListener<Rx>: Send + Sync
592where
593 Rx: UEncodedRxFrame + Send + 'static,
594{
595 async fn on_receive_encoded_zero_copy(&self, frame: Rx);
597}
598
599#[cfg(feature = "zero-copy-transport")]
602#[async_trait]
614pub trait UZeroCopyTransportCore: Send + Sync {
615 type Tx: UTxBuffer + Send;
617
618 type Rx: UEncodedRxFrame + Send + 'static;
620
621 async fn loan_prepared_tx(&self, spec: PreparedTxLoanSpec) -> Result<Self::Tx, UStatus>;
623
624 async fn send_prepared_zero_copy(&self, buffer: Self::Tx) -> Result<(), UStatus>;
626
627 async fn receive_encoded_zero_copy(
629 &self,
630 _source_filter: &UUri,
631 _sink_filter: Option<&UUri>,
632 ) -> Result<Self::Rx, UStatus> {
633 Err(unimplemented())
634 }
635
636 async fn register_encoded_zero_copy_listener(
638 &self,
639 _source_filter: &UUri,
640 _sink_filter: Option<&UUri>,
641 _listener: Arc<dyn UEncodedZeroCopyListener<Self::Rx>>,
642 ) -> Result<(), UStatus> {
643 Err(unimplemented())
644 }
645
646 async fn unregister_encoded_zero_copy_listener(
648 &self,
649 _source_filter: &UUri,
650 _sink_filter: Option<&UUri>,
651 _listener: Arc<dyn UEncodedZeroCopyListener<Self::Rx>>,
652 ) -> Result<(), UStatus> {
653 Err(unimplemented())
654 }
655}
656
657#[cfg(feature = "zero-copy-transport")]
658#[async_trait]
664pub trait UZeroCopyUninitTransportCore: UZeroCopyTransportCore {
665 type UninitTx: UUninitTxBuffer<Initialized = Self::Tx> + Send;
667
668 async fn loan_prepared_uninit_tx(
670 &self,
671 spec: PreparedTxLoanSpec,
672 ) -> Result<Self::UninitTx, UStatus>;
673}
674
675#[cfg(feature = "zero-copy-transport")]
676#[async_trait]
677impl<TCore, W, C> UZeroCopyTransportImpl for UWireTransport<TCore, W, C>
678where
679 TCore: UZeroCopyTransportCore,
680 W: UWire + Send + Sync + 'static,
681 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
682{
683 type Tx = TCore::Tx;
684 type Rx = UWireRx<TCore::Rx, W, C>;
685
686 async fn loan_validated_tx(&self, spec: UTxLoanSpec) -> Result<Self::Tx, UStatus> {
687 self.core
688 .loan_prepared_tx(PreparedTxLoanSpec::from_validated::<W, C>(
689 spec,
690 &self.metadata_codec,
691 )?)
692 .await
693 }
694
695 async fn send_validated_zero_copy(&self, buffer: Self::Tx) -> Result<(), UStatus> {
696 self.core.send_prepared_zero_copy(buffer).await
697 }
698
699 async fn receive_validated_zero_copy(
700 &self,
701 source_filter: &UUri,
702 sink_filter: Option<&UUri>,
703 ) -> Result<Self::Rx, UStatus> {
704 let core_source_filter = selected_wire_core_source_filter_for(source_filter);
705 loop {
706 let frame = self
707 .core
708 .receive_encoded_zero_copy(&core_source_filter, sink_filter)
709 .await?;
710 let frame = UWireRx::try_from_encoded(frame, &self.metadata_codec)?;
711 if wire_frame_matches(&frame, source_filter, sink_filter) {
712 return Ok(frame);
713 }
714 }
715 }
716
717 async fn register_validated_zero_copy_listener(
718 &self,
719 source_filter: &UUri,
720 sink_filter: Option<&UUri>,
721 listener: Arc<dyn UZeroCopyListener<Self::Rx>>,
722 ) -> Result<(), UStatus> {
723 let key = listener_key(
724 source_filter,
725 sink_filter,
726 zero_copy_listener_pointer::<TCore::Rx, W, C>(&listener),
727 );
728 let (listener, inserted) =
729 self.registered_zero_copy_listener(&key, source_filter, sink_filter, listener);
730 let core_source_filter = selected_wire_core_source_filter();
731 let result = self
732 .core
733 .register_encoded_zero_copy_listener(&core_source_filter, sink_filter, listener)
734 .await;
735 if result.is_err() && inserted {
736 self.zero_copy_listeners
737 .lock()
738 .expect("wire zero-copy listener registry lock poisoned")
739 .remove(&key);
740 }
741 result
742 }
743
744 async fn unregister_validated_zero_copy_listener(
745 &self,
746 source_filter: &UUri,
747 sink_filter: Option<&UUri>,
748 listener: Arc<dyn UZeroCopyListener<Self::Rx>>,
749 ) -> Result<(), UStatus> {
750 let key = listener_key(
751 source_filter,
752 sink_filter,
753 zero_copy_listener_pointer::<TCore::Rx, W, C>(&listener),
754 );
755 let listener = self.zero_copy_listener_for_unregister(&key, listener);
756 let core_source_filter = selected_wire_core_source_filter();
757 let result = self
758 .core
759 .unregister_encoded_zero_copy_listener(&core_source_filter, sink_filter, listener)
760 .await;
761 if result.is_ok() {
762 self.zero_copy_listeners
763 .lock()
764 .expect("wire zero-copy listener registry lock poisoned")
765 .remove(&key);
766 }
767 result
768 }
769}
770
771#[cfg(feature = "zero-copy-transport")]
772#[async_trait]
773impl<TCore, W, C> UZeroCopyUninitTransportImpl for UWireTransport<TCore, W, C>
774where
775 TCore: UZeroCopyUninitTransportCore,
776 W: UWire + Send + Sync + 'static,
777 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
778{
779 type UninitTx = TCore::UninitTx;
780
781 async fn loan_validated_uninit_tx(&self, spec: UTxLoanSpec) -> Result<Self::UninitTx, UStatus> {
782 self.core
783 .loan_prepared_uninit_tx(PreparedTxLoanSpec::from_validated::<W, C>(
784 spec,
785 &self.metadata_codec,
786 )?)
787 .await
788 }
789}
790
791#[cfg(feature = "zero-copy-transport")]
792impl<TCore, W, C> UWireTransport<TCore, W, C>
793where
794 TCore: UZeroCopyTransportCore,
795 W: UWire + Send + Sync + 'static,
796 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
797{
798 fn registered_zero_copy_listener(
799 &self,
800 key: &WireListenerKey,
801 source_filter: &UUri,
802 sink_filter: Option<&UUri>,
803 listener: Arc<dyn UZeroCopyListener<UWireRx<TCore::Rx, W, C>>>,
804 ) -> (Arc<dyn UEncodedZeroCopyListener<TCore::Rx>>, bool) {
805 let mut registry = self
806 .zero_copy_listeners
807 .lock()
808 .expect("wire zero-copy listener registry lock poisoned");
809 if let Some(existing) = registry.get(key) {
810 if let Ok(existing) = existing
811 .clone()
812 .downcast::<WireZeroCopyListener<TCore::Rx, W, C>>()
813 {
814 return (existing, false);
815 }
816 }
817
818 let wrapped = Arc::new(WireZeroCopyListener::<TCore::Rx, W, C> {
819 source_filter: source_filter.clone(),
820 sink_filter: sink_filter.cloned(),
821 listener,
822 metadata_codec: self.metadata_codec.clone(),
823 _wire: PhantomData,
824 });
825 registry.insert(key.clone(), wrapped.clone());
826 (wrapped, true)
827 }
828
829 fn zero_copy_listener_for_unregister(
830 &self,
831 key: &WireListenerKey,
832 fallback: Arc<dyn UZeroCopyListener<UWireRx<TCore::Rx, W, C>>>,
833 ) -> Arc<dyn UEncodedZeroCopyListener<TCore::Rx>> {
834 self.zero_copy_listeners
835 .lock()
836 .expect("wire zero-copy listener registry lock poisoned")
837 .get(key)
838 .and_then(|listener| {
839 listener
840 .clone()
841 .downcast::<WireZeroCopyListener<TCore::Rx, W, C>>()
842 .ok()
843 })
844 .unwrap_or_else(|| {
845 Arc::new(WireZeroCopyListener::<TCore::Rx, W, C> {
846 source_filter: key.source_filter.clone(),
847 sink_filter: key.sink_filter.clone(),
848 listener: fallback,
849 metadata_codec: self.metadata_codec.clone(),
850 _wire: PhantomData,
851 })
852 })
853 }
854}
855
856#[cfg(feature = "zero-copy-transport")]
857struct WireZeroCopyListener<Rx, W, C>
858where
859 Rx: UEncodedRxFrame + Send + 'static,
860 W: UWire,
861 C: UWireMetadataCodecFor<W>,
862{
863 source_filter: UUri,
864 sink_filter: Option<UUri>,
865 listener: Arc<dyn UZeroCopyListener<UWireRx<Rx, W, C>>>,
866 metadata_codec: C,
867 _wire: PhantomData<W>,
868}
869
870#[cfg(feature = "zero-copy-transport")]
871#[async_trait]
872impl<Rx, W, C> UEncodedZeroCopyListener<Rx> for WireZeroCopyListener<Rx, W, C>
873where
874 Rx: UEncodedRxFrame + Send + 'static,
875 W: UWire + Send + Sync + 'static,
876 C: UWireMetadataCodecFor<W> + Send + Sync + 'static,
877{
878 async fn on_receive_encoded_zero_copy(&self, frame: Rx) {
879 match UWireRx::<Rx, W, C>::try_from_encoded(frame, &self.metadata_codec) {
880 Ok(frame)
881 if wire_frame_matches(&frame, &self.source_filter, self.sink_filter.as_ref()) =>
882 {
883 self.listener.on_receive_zero_copy(frame).await;
884 }
885 Ok(_) => {}
886 Err(error) => warn!(%error, "dropping invalid selected-wire zero-copy frame"),
887 }
888 }
889}
890
891#[cfg(feature = "zero-copy-transport")]
892fn wire_frame_matches<Rx, W, C>(
893 frame: &UWireRx<Rx, W, C>,
894 source_filter: &UUri,
895 sink_filter: Option<&UUri>,
896) -> bool
897where
898 Rx: UEncodedRxFrame,
899 W: UWire,
900 C: UWireMetadataCodecFor<W>,
901{
902 source_filter.matches(frame.metadata().source())
903 && sink_filter.is_none_or(|filter| {
904 frame
905 .metadata()
906 .sink()
907 .is_some_and(|sink| filter.matches(sink))
908 })
909}
910
911#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
912fn selected_wire_core_source_filter() -> UUri {
913 UUri::try_from_parts("*", u32::MAX, u8::MAX, u16::MAX)
914 .expect("valid selected-wire core wildcard source filter")
915}
916
917#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
918fn selected_wire_core_source_filter_for(source_filter: &UUri) -> UUri {
919 if source_filter.verify_no_wildcards().is_ok() {
920 source_filter.clone()
921 } else {
922 selected_wire_core_source_filter()
923 }
924}
925
926#[cfg(feature = "owned-frame-transport")]
927#[derive(Clone, Debug, PartialEq)]
929pub struct PreparedOwnedFrame {
930 metadata: UFrameMetadata,
931 encoded_metadata: Vec<u8>,
932 payload: Option<Bytes>,
933}
934
935#[cfg(feature = "owned-frame-transport")]
936impl PreparedOwnedFrame {
937 pub fn from_validated<W, C>(frame: UOwnedFrame, codec: &C) -> Result<Self, UStatus>
943 where
944 W: UWire,
945 C: UWireMetadataCodecFor<W>,
946 {
947 let (metadata, payload) = frame.into_parts();
948 let encoded_metadata = codec.encode_frame_metadata(W::metadata_context(), &metadata)?;
949 Ok(Self {
950 metadata,
951 encoded_metadata,
952 payload,
953 })
954 }
955
956 #[must_use]
958 pub fn metadata(&self) -> &UFrameMetadata {
959 &self.metadata
960 }
961
962 #[must_use]
964 pub fn encoded_metadata(&self) -> &[u8] {
965 &self.encoded_metadata
966 }
967
968 #[must_use]
970 pub fn payload(&self) -> Option<&Bytes> {
971 self.payload.as_ref()
972 }
973
974 #[must_use]
976 pub fn into_parts(self) -> (UFrameMetadata, Vec<u8>, Option<Bytes>) {
977 (self.metadata, self.encoded_metadata, self.payload)
978 }
979}
980
981#[cfg(feature = "owned-frame-transport")]
982#[derive(Clone, Debug, PartialEq)]
984pub struct EncodedOwnedFrame {
985 encoded_metadata: Vec<u8>,
986 payload: Option<Bytes>,
987}
988
989#[cfg(feature = "owned-frame-transport")]
990impl EncodedOwnedFrame {
991 #[must_use]
993 pub fn new(encoded_metadata: impl Into<Vec<u8>>, payload: Option<Bytes>) -> Self {
994 Self {
995 encoded_metadata: encoded_metadata.into(),
996 payload,
997 }
998 }
999
1000 #[must_use]
1002 pub fn encoded_metadata(&self) -> &[u8] {
1003 &self.encoded_metadata
1004 }
1005
1006 #[must_use]
1008 pub fn payload(&self) -> Option<&Bytes> {
1009 self.payload.as_ref()
1010 }
1011
1012 pub fn decode<W, C>(self, codec: &C) -> Result<UOwnedFrame, UStatus>
1018 where
1019 W: UWire,
1020 C: UWireMetadataCodecFor<W>,
1021 {
1022 let metadata =
1023 codec.decode_frame_metadata(W::metadata_context(), &self.encoded_metadata)?;
1024 UOwnedFrame::new(metadata, self.payload).map_err(invalid_metadata)
1025 }
1026
1027 #[must_use]
1029 pub fn into_parts(self) -> (Vec<u8>, Option<Bytes>) {
1030 (self.encoded_metadata, self.payload)
1031 }
1032}
1033
1034#[cfg(feature = "owned-frame-transport")]
1035#[async_trait]
1037pub trait UEncodedOwnedListener: Send + Sync {
1038 async fn on_receive_encoded_owned(&self, frame: EncodedOwnedFrame);
1040}
1041
1042#[cfg(feature = "owned-frame-transport")]
1045#[async_trait]
1053pub trait UOwnedTransportCore: Send + Sync {
1054 async fn send_prepared_owned(&self, frame: PreparedOwnedFrame) -> Result<(), UStatus>;
1056
1057 async fn receive_encoded_owned(
1059 &self,
1060 _source_filter: &UUri,
1061 _sink_filter: Option<&UUri>,
1062 ) -> Result<EncodedOwnedFrame, UStatus> {
1063 Err(unimplemented())
1064 }
1065
1066 async fn register_encoded_owned_listener(
1068 &self,
1069 _source_filter: &UUri,
1070 _sink_filter: Option<&UUri>,
1071 _listener: Arc<dyn UEncodedOwnedListener>,
1072 ) -> Result<(), UStatus> {
1073 Err(unimplemented())
1074 }
1075
1076 async fn unregister_encoded_owned_listener(
1078 &self,
1079 _source_filter: &UUri,
1080 _sink_filter: Option<&UUri>,
1081 _listener: Arc<dyn UEncodedOwnedListener>,
1082 ) -> Result<(), UStatus> {
1083 Err(unimplemented())
1084 }
1085}
1086
1087#[cfg(feature = "owned-frame-transport")]
1088#[async_trait]
1089impl<TCore, W, C> UOwnedTransportImpl for UWireTransport<TCore, W, C>
1090where
1091 TCore: UOwnedTransportCore,
1092 W: UWire + Send + Sync + 'static,
1093 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
1094{
1095 async fn send_validated_owned(&self, frame: UOwnedFrame) -> Result<(), UStatus> {
1096 self.core
1097 .send_prepared_owned(PreparedOwnedFrame::from_validated::<W, C>(
1098 frame,
1099 &self.metadata_codec,
1100 )?)
1101 .await
1102 }
1103
1104 async fn receive_validated_owned(
1105 &self,
1106 source_filter: &UUri,
1107 sink_filter: Option<&UUri>,
1108 ) -> Result<UOwnedFrame, UStatus> {
1109 let core_source_filter = selected_wire_core_source_filter_for(source_filter);
1110 loop {
1111 let frame = self
1112 .core
1113 .receive_encoded_owned(&core_source_filter, sink_filter)
1114 .await?
1115 .decode::<W, C>(&self.metadata_codec)?;
1116 if owned_frame_matches(&frame, source_filter, sink_filter) {
1117 return Ok(frame);
1118 }
1119 }
1120 }
1121
1122 async fn register_validated_owned_listener(
1123 &self,
1124 source_filter: &UUri,
1125 sink_filter: Option<&UUri>,
1126 listener: Arc<dyn UOwnedListener>,
1127 ) -> Result<(), UStatus> {
1128 let key = listener_key(
1129 source_filter,
1130 sink_filter,
1131 owned_listener_pointer(&listener),
1132 );
1133 let (listener, inserted) =
1134 self.registered_owned_listener(&key, source_filter, sink_filter, listener);
1135 let core_source_filter = selected_wire_core_source_filter();
1136 let result = self
1137 .core
1138 .register_encoded_owned_listener(&core_source_filter, sink_filter, listener)
1139 .await;
1140 if result.is_err() && inserted {
1141 self.owned_listeners
1142 .lock()
1143 .expect("wire owned listener registry lock poisoned")
1144 .remove(&key);
1145 }
1146 result
1147 }
1148
1149 async fn unregister_validated_owned_listener(
1150 &self,
1151 source_filter: &UUri,
1152 sink_filter: Option<&UUri>,
1153 listener: Arc<dyn UOwnedListener>,
1154 ) -> Result<(), UStatus> {
1155 let key = listener_key(
1156 source_filter,
1157 sink_filter,
1158 owned_listener_pointer(&listener),
1159 );
1160 let listener = self.owned_listener_for_unregister(&key, listener);
1161 let core_source_filter = selected_wire_core_source_filter();
1162 let result = self
1163 .core
1164 .unregister_encoded_owned_listener(&core_source_filter, sink_filter, listener)
1165 .await;
1166 if result.is_ok() {
1167 self.owned_listeners
1168 .lock()
1169 .expect("wire owned listener registry lock poisoned")
1170 .remove(&key);
1171 }
1172 result
1173 }
1174}
1175
1176#[cfg(feature = "owned-frame-transport")]
1177impl<TCore, W, C> UWireTransport<TCore, W, C>
1178where
1179 TCore: UOwnedTransportCore,
1180 W: UWire + Send + Sync + 'static,
1181 C: UWireMetadataCodecFor<W> + Clone + Send + Sync + 'static,
1182{
1183 fn registered_owned_listener(
1184 &self,
1185 key: &WireListenerKey,
1186 source_filter: &UUri,
1187 sink_filter: Option<&UUri>,
1188 listener: Arc<dyn UOwnedListener>,
1189 ) -> (Arc<dyn UEncodedOwnedListener>, bool) {
1190 let mut registry = self
1191 .owned_listeners
1192 .lock()
1193 .expect("wire owned listener registry lock poisoned");
1194 if let Some(existing) = registry.get(key) {
1195 if let Ok(existing) = existing.clone().downcast::<WireOwnedListener<W, C>>() {
1196 return (existing, false);
1197 }
1198 }
1199
1200 let wrapped = Arc::new(WireOwnedListener::<W, C> {
1201 source_filter: source_filter.clone(),
1202 sink_filter: sink_filter.cloned(),
1203 listener,
1204 metadata_codec: self.metadata_codec.clone(),
1205 _wire: PhantomData,
1206 });
1207 registry.insert(key.clone(), wrapped.clone());
1208 (wrapped, true)
1209 }
1210
1211 fn owned_listener_for_unregister(
1212 &self,
1213 key: &WireListenerKey,
1214 fallback: Arc<dyn UOwnedListener>,
1215 ) -> Arc<dyn UEncodedOwnedListener> {
1216 self.owned_listeners
1217 .lock()
1218 .expect("wire owned listener registry lock poisoned")
1219 .get(key)
1220 .and_then(|listener| listener.clone().downcast::<WireOwnedListener<W, C>>().ok())
1221 .unwrap_or_else(|| {
1222 Arc::new(WireOwnedListener::<W, C> {
1223 source_filter: key.source_filter.clone(),
1224 sink_filter: key.sink_filter.clone(),
1225 listener: fallback,
1226 metadata_codec: self.metadata_codec.clone(),
1227 _wire: PhantomData,
1228 })
1229 })
1230 }
1231}
1232
1233#[cfg(feature = "owned-frame-transport")]
1234struct WireOwnedListener<W, C>
1235where
1236 W: UWire,
1237 C: UWireMetadataCodecFor<W>,
1238{
1239 source_filter: UUri,
1240 sink_filter: Option<UUri>,
1241 listener: Arc<dyn UOwnedListener>,
1242 metadata_codec: C,
1243 _wire: PhantomData<W>,
1244}
1245
1246#[cfg(feature = "owned-frame-transport")]
1247#[async_trait]
1248impl<W, C> UEncodedOwnedListener for WireOwnedListener<W, C>
1249where
1250 W: UWire + Send + Sync + 'static,
1251 C: UWireMetadataCodecFor<W> + Send + Sync + 'static,
1252{
1253 async fn on_receive_encoded_owned(&self, frame: EncodedOwnedFrame) {
1254 match frame.decode::<W, C>(&self.metadata_codec) {
1255 Ok(frame)
1256 if owned_frame_matches(&frame, &self.source_filter, self.sink_filter.as_ref()) =>
1257 {
1258 self.listener.on_receive_owned(frame).await;
1259 }
1260 Ok(_) => {}
1261 Err(error) => warn!(%error, "dropping invalid selected-wire owned frame"),
1262 }
1263 }
1264}
1265
1266#[cfg(feature = "owned-frame-transport")]
1267fn owned_frame_matches(
1268 frame: &UOwnedFrame,
1269 source_filter: &UUri,
1270 sink_filter: Option<&UUri>,
1271) -> bool {
1272 source_filter.matches(frame.metadata().source())
1273 && sink_filter.is_none_or(|filter| {
1274 frame
1275 .metadata()
1276 .sink()
1277 .is_some_and(|sink| filter.matches(sink))
1278 })
1279}
1280
1281#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
1282#[derive(Clone, Debug, Eq, Hash, PartialEq)]
1283struct WireListenerKey {
1284 source_filter: UUri,
1285 sink_filter: Option<UUri>,
1286 listener: usize,
1287}
1288
1289#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
1290fn listener_key(
1291 source_filter: &UUri,
1292 sink_filter: Option<&UUri>,
1293 listener: usize,
1294) -> WireListenerKey {
1295 WireListenerKey {
1296 source_filter: source_filter.clone(),
1297 sink_filter: sink_filter.cloned(),
1298 listener,
1299 }
1300}
1301
1302#[cfg(feature = "zero-copy-transport")]
1303fn zero_copy_listener_pointer<Rx, W, C>(
1304 listener: &Arc<dyn UZeroCopyListener<UWireRx<Rx, W, C>>>,
1305) -> usize
1306where
1307 Rx: UEncodedRxFrame,
1308 W: UWire,
1309 C: UWireMetadataCodecFor<W>,
1310{
1311 let ptr = Arc::as_ptr(listener);
1312 let thin_ptr = ptr as *const ();
1313 thin_ptr as usize
1314}
1315
1316#[cfg(feature = "owned-frame-transport")]
1317fn owned_listener_pointer(listener: &Arc<dyn UOwnedListener>) -> usize {
1318 let ptr = Arc::as_ptr(listener);
1319 let thin_ptr = ptr as *const ();
1320 thin_ptr as usize
1321}
1322
1323#[cfg(any(feature = "zero-copy-transport", feature = "owned-frame-transport"))]
1324fn unimplemented() -> UStatus {
1325 UStatus::fail_with_code(UCode::Unimplemented, "not implemented")
1326}
1327
1328#[cfg(feature = "owned-frame-transport")]
1329fn invalid_metadata(error: crate::UFrameMetadataError) -> UStatus {
1330 UStatus::fail_with_code(UCode::InvalidArgument, error.to_string())
1331}
1332
1333#[cfg(test)]
1334mod tests {
1335 use std::{collections::VecDeque, io::Cursor, sync::Arc, sync::Mutex as StdMutex};
1336
1337 use protobuf::well_known_types::wrappers::StringValue;
1338
1339 use super::*;
1340 use crate::{
1341 test_support::StableTestBytes as WireStableBytes, EncodePayload, PayloadEncoding,
1342 PayloadLoanProvenance, ProtobufPayload, ProtobufWire, StableContainerPayload,
1343 StableContainerWireFormat, StablePayload, UMessageBuilder, UProtocolNativeWire,
1344 UTxPayloadSpec, UVecRxLease, UVecTxBuffer, UVecUninitTxBuffer, UWireMetadataCodec,
1345 UZeroCopyTransportExt,
1346 };
1347
1348 #[derive(Clone)]
1349 struct RawRx {
1350 encoded_metadata: Vec<u8>,
1351 payload: Vec<u8>,
1352 }
1353
1354 impl UEncodedRxFrame for RawRx {
1355 type PayloadReader<'a>
1356 = Cursor<&'a [u8]>
1357 where
1358 Self: 'a;
1359 type PayloadSlices<'a>
1360 = std::iter::Once<&'a [u8]>
1361 where
1362 Self: 'a;
1363
1364 fn encoded_metadata(&self) -> &[u8] {
1365 &self.encoded_metadata
1366 }
1367
1368 fn payload_len(&self) -> usize {
1369 self.payload.len()
1370 }
1371
1372 fn payload_reader(&self) -> Self::PayloadReader<'_> {
1373 Cursor::new(&self.payload)
1374 }
1375
1376 fn payload_slices(&self) -> Self::PayloadSlices<'_> {
1377 std::iter::once(self.payload.as_slice())
1378 }
1379
1380 fn try_contiguous_payload(&self) -> Option<&[u8]> {
1381 Some(&self.payload)
1382 }
1383 }
1384
1385 impl UEncodedLoanedRxFrame for RawRx {
1386 fn loaned_contiguous_payload(&self) -> Result<LoanedPayload<'_>, UWireError> {
1387 Ok(unsafe {
1390 LoanedPayload::new_unchecked(
1391 self.payload.as_slice(),
1392 PayloadLoanProvenance::OpaqueTransportLoan,
1393 )
1394 })
1395 }
1396 }
1397
1398 #[derive(Default)]
1399 struct RecordingCore {
1400 prepared: StdMutex<Vec<PreparedTxLoanSpec>>,
1401 sent: StdMutex<Vec<UVecTxBuffer>>,
1402 received: StdMutex<VecDeque<RawRx>>,
1403 receive_filters: StdMutex<Vec<(UUri, Option<UUri>)>>,
1404 }
1405
1406 #[async_trait]
1407 impl UZeroCopyTransportCore for RecordingCore {
1408 type Tx = UVecTxBuffer;
1409 type Rx = RawRx;
1410
1411 async fn loan_prepared_tx(&self, spec: PreparedTxLoanSpec) -> Result<Self::Tx, UStatus> {
1412 self.prepared.lock().unwrap().push(spec.clone());
1413 UVecTxBuffer::with_alignment(
1414 spec.metadata().clone(),
1415 spec.payload_len(),
1416 spec.payload_alignment(),
1417 )
1418 }
1419
1420 async fn send_prepared_zero_copy(&self, buffer: Self::Tx) -> Result<(), UStatus> {
1421 self.sent.lock().unwrap().push(buffer);
1422 Ok(())
1423 }
1424
1425 async fn receive_encoded_zero_copy(
1426 &self,
1427 source_filter: &UUri,
1428 sink_filter: Option<&UUri>,
1429 ) -> Result<Self::Rx, UStatus> {
1430 self.receive_filters
1431 .lock()
1432 .unwrap()
1433 .push((source_filter.clone(), sink_filter.cloned()));
1434 self.received.lock().unwrap().pop_front().ok_or_else(|| {
1435 UStatus::fail_with_code(UCode::NotFound, "no test encoded frame available")
1436 })
1437 }
1438 }
1439
1440 #[cfg(feature = "owned-frame-transport")]
1441 #[async_trait]
1442 impl UOwnedTransportCore for RecordingCore {
1443 async fn send_prepared_owned(&self, _frame: PreparedOwnedFrame) -> Result<(), UStatus> {
1444 Ok(())
1445 }
1446
1447 async fn receive_encoded_owned(
1448 &self,
1449 source_filter: &UUri,
1450 sink_filter: Option<&UUri>,
1451 ) -> Result<EncodedOwnedFrame, UStatus> {
1452 self.receive_filters
1453 .lock()
1454 .unwrap()
1455 .push((source_filter.clone(), sink_filter.cloned()));
1456 self.received
1457 .lock()
1458 .unwrap()
1459 .pop_front()
1460 .map(|raw| EncodedOwnedFrame::new(raw.encoded_metadata, Some(raw.payload.into())))
1461 .ok_or_else(|| {
1462 UStatus::fail_with_code(UCode::NotFound, "no test owned frame available")
1463 })
1464 }
1465 }
1466
1467 #[derive(Default)]
1468 struct RecordingZeroCopyListener {
1469 frames: StdMutex<Vec<UWireRx<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>>>,
1470 }
1471
1472 #[async_trait]
1473 impl UZeroCopyListener<UWireRx<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>>
1474 for RecordingZeroCopyListener
1475 {
1476 async fn on_receive_zero_copy(
1477 &self,
1478 frame: UWireRx<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>,
1479 ) {
1480 self.frames.lock().unwrap().push(frame);
1481 }
1482 }
1483
1484 struct CompileCore;
1485
1486 #[async_trait]
1487 impl UZeroCopyTransportCore for CompileCore {
1488 type Tx = UVecTxBuffer;
1489 type Rx = RawRx;
1490
1491 async fn loan_prepared_tx(&self, spec: PreparedTxLoanSpec) -> Result<Self::Tx, UStatus> {
1492 UVecTxBuffer::with_alignment(
1493 spec.metadata().clone(),
1494 spec.payload_len(),
1495 spec.payload_alignment(),
1496 )
1497 }
1498
1499 async fn send_prepared_zero_copy(&self, _buffer: Self::Tx) -> Result<(), UStatus> {
1500 Ok(())
1501 }
1502 }
1503
1504 #[async_trait]
1505 impl UZeroCopyUninitTransportCore for CompileCore {
1506 type UninitTx = UVecUninitTxBuffer;
1507
1508 async fn loan_prepared_uninit_tx(
1509 &self,
1510 spec: PreparedTxLoanSpec,
1511 ) -> Result<Self::UninitTx, UStatus> {
1512 UVecUninitTxBuffer::with_alignment(
1513 spec.metadata().clone(),
1514 spec.payload_len(),
1515 spec.payload_alignment(),
1516 )
1517 }
1518 }
1519
1520 #[cfg(feature = "owned-frame-transport")]
1521 #[async_trait]
1522 impl UOwnedTransportCore for CompileCore {
1523 async fn send_prepared_owned(&self, _frame: PreparedOwnedFrame) -> Result<(), UStatus> {
1524 Ok(())
1525 }
1526 }
1527
1528 fn metadata_with_payload() -> UFrameMetadata {
1529 metadata_with_payload_encoding(PayloadEncoding::RAW)
1530 }
1531
1532 fn metadata_with_payload_encoding(payload_encoding: PayloadEncoding) -> UFrameMetadata {
1533 metadata_with_topic_and_payload_encoding(0x9000, payload_encoding)
1534 }
1535
1536 fn metadata_with_topic_and_payload_encoding(
1537 resource_id: u16,
1538 payload_encoding: PayloadEncoding,
1539 ) -> UFrameMetadata {
1540 let topic = UUri::try_from_parts("vehicle", 0x4210, 0x01, resource_id).expect("topic URI");
1541 let message = UMessageBuilder::publish(topic).build().expect("message");
1542 crate::frame::metadata::try_project_attributes_to_frame_metadata(
1543 message.attributes(),
1544 Some(payload_encoding),
1545 )
1546 .expect("metadata")
1547 }
1548
1549 fn stable_metadata<T: StablePayload>() -> UFrameMetadata {
1550 let topic = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).expect("topic URI");
1551 let message = UMessageBuilder::publish(topic).build().expect("message");
1552 crate::frame::metadata::try_project_attributes_to_frame_metadata(
1553 message.attributes(),
1554 Some(StableContainerPayload::<T>::encoding()),
1555 )
1556 .expect("metadata")
1557 }
1558
1559 fn raw_frame_for_topic(resource_id: u16, payload: &[u8]) -> RawRx {
1560 let metadata = metadata_with_topic_and_payload_encoding(resource_id, PayloadEncoding::RAW);
1561 RawRx {
1562 encoded_metadata: encode_metadata::<UProtocolNativeWire>(&metadata),
1563 payload: payload.to_vec(),
1564 }
1565 }
1566
1567 fn encode_metadata<W>(metadata: &UFrameMetadata) -> Vec<u8>
1568 where
1569 W: UWire,
1570 {
1571 NativePrefixFrameMetadataCodec
1572 .encode_frame_metadata(W::metadata_context(), metadata)
1573 .unwrap()
1574 }
1575
1576 fn decode_metadata<W>(encoded: &[u8]) -> UFrameMetadata
1577 where
1578 W: UWire,
1579 {
1580 NativePrefixFrameMetadataCodec
1581 .decode_frame_metadata(W::metadata_context(), encoded)
1582 .unwrap()
1583 }
1584
1585 #[test]
1586 fn wire_rx_decodes_metadata_and_delegates_payload() {
1587 let metadata = metadata_with_payload();
1588 let encoded_metadata = encode_metadata::<UProtocolNativeWire>(&metadata);
1589 let raw = RawRx {
1590 encoded_metadata,
1591 payload: b"abc".to_vec(),
1592 };
1593
1594 let rx = UWireRx::<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>::try_from_encoded(
1595 raw,
1596 &NativePrefixFrameMetadataCodec,
1597 )
1598 .unwrap();
1599
1600 assert_eq!(rx.metadata(), &metadata);
1601 assert_eq!(rx.payload_len(), 3);
1602 assert_eq!(rx.try_contiguous_payload(), Some(&b"abc"[..]));
1603 }
1604
1605 #[test]
1606 fn wire_rx_decodes_payload_with_selected_wire() {
1607 let value = StringValue {
1608 value: "selected-wire".to_string(),
1609 special_fields: Default::default(),
1610 };
1611 let payload = ProtobufWire::encode_payload_owned(&value).unwrap();
1612 let metadata = metadata_with_payload_encoding(ProtobufPayload::encoding());
1613 let encoded_metadata = encode_metadata::<ProtobufWire>(&metadata);
1614 let raw = RawRx {
1615 encoded_metadata,
1616 payload: payload.to_vec(),
1617 };
1618
1619 let rx = UWireRx::<RawRx, ProtobufWire, NativePrefixFrameMetadataCodec>::try_from_encoded(
1620 raw,
1621 &NativePrefixFrameMetadataCodec,
1622 )
1623 .unwrap();
1624 let decoded: StringValue = rx.decode_payload().unwrap();
1625
1626 assert_eq!(decoded.value, "selected-wire");
1627 }
1628
1629 #[test]
1630 fn stable_borrow_accepts_wire_rx_with_loaned_raw_frame() {
1631 fn assert_loaned_rx<T: ULoanedContiguousZeroCopyRxFrame>() {}
1632 assert_loaned_rx::<UWireRx<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>>();
1633
1634 let value = WireStableBytes { bytes: *b"wire" };
1635 let metadata = stable_metadata::<WireStableBytes>();
1636 let encoded_metadata = encode_metadata::<UProtocolNativeWire>(&metadata);
1637 let raw = RawRx {
1638 encoded_metadata,
1639 payload: value.bytes.to_vec(),
1640 };
1641
1642 let rx = UWireRx::<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec>::try_from_encoded(
1643 raw,
1644 &NativePrefixFrameMetadataCodec,
1645 )
1646 .unwrap();
1647 let borrowed = rx.borrow_stable_payload::<WireStableBytes>().unwrap();
1648
1649 assert_eq!(borrowed, &value);
1650 assert_eq!(
1651 rx.payload_loan_provenance().unwrap(),
1652 PayloadLoanProvenance::OpaqueTransportLoan
1653 );
1654 }
1655
1656 #[tokio::test]
1657 async fn stable_initialized_tx_helper_sends_through_selected_wire_transport() {
1658 let transport =
1659 RecordingCore::default().into_native_prefix_wire_transport(StableContainerWireFormat);
1660
1661 transport
1662 .send_loaned_payload::<WireStableBytes>(
1663 stable_metadata::<WireStableBytes>(),
1664 |payload| payload.bytes.copy_from_slice(b"wire"),
1665 )
1666 .await
1667 .expect("send initialized stable payload through selected wire");
1668
1669 let prepared = transport.core().prepared.lock().unwrap();
1670 assert_eq!(prepared.len(), 1);
1671 let prepared_frame = prepared.first().expect("one prepared frame");
1672 let decoded =
1673 decode_metadata::<StableContainerWireFormat>(prepared_frame.encoded_metadata());
1674 assert_eq!(
1675 decoded.payload_encoding(),
1676 Some(&StableContainerPayload::<WireStableBytes>::encoding())
1677 );
1678 drop(prepared);
1679
1680 let sent = transport.core().sent.lock().unwrap();
1681 assert_eq!(sent.len(), 1);
1682 let sent_frame = sent.first().expect("one sent frame");
1683 let frame = UVecRxLease::new(
1684 sent_frame.metadata().clone(),
1685 Some(sent_frame.payload().to_vec()),
1686 )
1687 .expect("sent stable frame");
1688 assert_eq!(
1689 frame.borrow_stable_payload::<WireStableBytes>().unwrap(),
1690 &WireStableBytes { bytes: *b"wire" }
1691 );
1692 }
1693
1694 #[tokio::test]
1695 async fn receive_filters_after_selected_wire_decode() {
1696 let core = RecordingCore::default();
1697 core.received.lock().unwrap().extend([
1698 raw_frame_for_topic(0x9001, b"drop"),
1699 raw_frame_for_topic(0x9000, b"keep"),
1700 ]);
1701 let transport = core.into_native_prefix_wire_transport(UProtocolNativeWire);
1702 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).unwrap();
1703
1704 let frame = transport
1705 .receive_validated_zero_copy(&source_filter, None)
1706 .await
1707 .expect("matching decoded frame");
1708
1709 assert_eq!(frame.try_contiguous_payload(), Some(&b"keep"[..]));
1710 assert!(transport.core().received.lock().unwrap().is_empty());
1711 let filters = transport.core().receive_filters.lock().unwrap();
1712 assert_eq!(filters.len(), 2);
1713 let filter = filters.first().expect("first receive filter");
1714 assert_eq!(filter.0, source_filter);
1715 assert_eq!(filter.1, None);
1716 }
1717
1718 #[tokio::test]
1719 async fn receive_rejects_invalid_selected_wire_frame_before_later_frames() {
1720 let core = RecordingCore::default();
1721 core.received.lock().unwrap().extend([
1722 RawRx {
1723 encoded_metadata: b"invalid selected-wire metadata".to_vec(),
1724 payload: b"drop".to_vec(),
1725 },
1726 raw_frame_for_topic(0x9000, b"keep"),
1727 ]);
1728 let transport = core.into_native_prefix_wire_transport(UProtocolNativeWire);
1729 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).unwrap();
1730
1731 let status = match transport
1732 .receive_validated_zero_copy(&source_filter, None)
1733 .await
1734 {
1735 Ok(_) => panic!("invalid selected-wire metadata must be rejected"),
1736 Err(status) => status,
1737 };
1738 assert_eq!(status.code(), UCode::InvalidArgument);
1739 assert!(status
1740 .message()
1741 .is_some_and(|message| message.contains("wrong native-prefix metadata magic")));
1742 assert_eq!(transport.core().received.lock().unwrap().len(), 1);
1743
1744 let frame = transport
1745 .receive_validated_zero_copy(&source_filter, None)
1746 .await
1747 .expect("later matching decoded frame");
1748
1749 assert_eq!(frame.try_contiguous_payload(), Some(&b"keep"[..]));
1750 assert!(transport.core().received.lock().unwrap().is_empty());
1751 let filters = transport.core().receive_filters.lock().unwrap();
1752 assert_eq!(filters.len(), 2);
1753 let filter = filters.first().expect("first receive filter");
1754 assert_eq!(filter.0, source_filter);
1755 assert_eq!(filter.1, None);
1756 }
1757
1758 #[tokio::test]
1759 async fn receive_uses_wildcard_core_filter_for_wildcard_source_filter() {
1760 let core = RecordingCore::default();
1761 core.received
1762 .lock()
1763 .unwrap()
1764 .push_back(raw_frame_for_topic(0x9000, b"keep"));
1765 let transport = core.into_native_prefix_wire_transport(UProtocolNativeWire);
1766 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, u16::MAX).unwrap();
1767
1768 let frame = transport
1769 .receive_validated_zero_copy(&source_filter, None)
1770 .await
1771 .expect("matching decoded frame");
1772
1773 assert_eq!(frame.try_contiguous_payload(), Some(&b"keep"[..]));
1774 let filters = transport.core().receive_filters.lock().unwrap();
1775 assert_eq!(filters.len(), 1);
1776 let filter = filters.first().expect("one receive filter");
1777 assert_eq!(filter.0, selected_wire_core_source_filter());
1778 assert_eq!(filter.1, None);
1779 }
1780
1781 #[cfg(feature = "owned-frame-transport")]
1782 #[tokio::test]
1783 async fn owned_receive_uses_exact_core_filter_for_exact_source_filter() {
1784 let core = RecordingCore::default();
1785 core.received
1786 .lock()
1787 .unwrap()
1788 .push_back(raw_frame_for_topic(0x9000, b"owned"));
1789 let transport = core.into_native_prefix_wire_transport(UProtocolNativeWire);
1790 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).unwrap();
1791
1792 let frame = transport
1793 .receive_validated_owned(&source_filter, None)
1794 .await
1795 .expect("matching owned frame");
1796
1797 assert_eq!(frame.payload_bytes(), b"owned");
1798 let filters = transport.core().receive_filters.lock().unwrap();
1799 assert_eq!(filters.len(), 1);
1800 let filter = filters.first().expect("one receive filter");
1801 assert_eq!(filter.0, source_filter);
1802 assert_eq!(filter.1, None);
1803 }
1804
1805 #[cfg(feature = "owned-frame-transport")]
1806 #[tokio::test]
1807 async fn owned_receive_uses_wildcard_core_filter_for_wildcard_source_filter() {
1808 let core = RecordingCore::default();
1809 core.received
1810 .lock()
1811 .unwrap()
1812 .push_back(raw_frame_for_topic(0x9000, b"owned"));
1813 let transport = core.into_native_prefix_wire_transport(UProtocolNativeWire);
1814 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, u16::MAX).unwrap();
1815
1816 let frame = transport
1817 .receive_validated_owned(&source_filter, None)
1818 .await
1819 .expect("matching owned frame");
1820
1821 assert_eq!(frame.payload_bytes(), b"owned");
1822 let filters = transport.core().receive_filters.lock().unwrap();
1823 assert_eq!(filters.len(), 1);
1824 let filter = filters.first().expect("one receive filter");
1825 assert_eq!(filter.0, selected_wire_core_source_filter());
1826 assert_eq!(filter.1, None);
1827 }
1828
1829 #[test]
1830 fn selected_wire_core_source_filter_uses_exact_source_when_safe() {
1831 let exact = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).unwrap();
1832 let wildcard = UUri::try_from_parts("vehicle", 0x4210, 0x01, u16::MAX).unwrap();
1833
1834 assert_eq!(selected_wire_core_source_filter_for(&exact), exact);
1835 assert_eq!(
1836 selected_wire_core_source_filter_for(&wildcard),
1837 selected_wire_core_source_filter()
1838 );
1839 }
1840
1841 #[tokio::test]
1842 async fn listener_filters_after_selected_wire_decode() {
1843 let source_filter = UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).unwrap();
1844 let listener = Arc::new(RecordingZeroCopyListener::default());
1845 let wire_listener =
1846 WireZeroCopyListener::<RawRx, UProtocolNativeWire, NativePrefixFrameMetadataCodec> {
1847 source_filter,
1848 sink_filter: None,
1849 listener: listener.clone(),
1850 metadata_codec: NativePrefixFrameMetadataCodec,
1851 _wire: PhantomData,
1852 };
1853
1854 wire_listener
1855 .on_receive_encoded_zero_copy(raw_frame_for_topic(0x9001, b"drop"))
1856 .await;
1857 wire_listener
1858 .on_receive_encoded_zero_copy(raw_frame_for_topic(0x9000, b"keep"))
1859 .await;
1860
1861 let frames = listener.frames.lock().unwrap();
1862 assert_eq!(frames.len(), 1);
1863 let frame = frames.first().expect("one received frame");
1864 assert_eq!(frame.try_contiguous_payload(), Some(&b"keep"[..]));
1865 }
1866
1867 #[test]
1868 fn wire_transport_fits_zero_copy_blanket_boundaries() {
1869 fn assert_zero_copy_impl<T: UZeroCopyTransportImpl>() {}
1870 fn assert_uninit_impl<T: UZeroCopyUninitTransportImpl>() {}
1871
1872 type Transport =
1873 UWireTransport<CompileCore, UProtocolNativeWire, NativePrefixFrameMetadataCodec>;
1874 assert_zero_copy_impl::<Transport>();
1875 assert_uninit_impl::<Transport>();
1876 }
1877
1878 #[test]
1879 fn prepared_tx_spec_carries_metadata_bytes_and_layout() {
1880 let metadata = metadata_with_payload();
1881 let spec = crate::UTxLoanSpec::new(
1882 metadata.clone(),
1883 UTxPayloadSpec::Present {
1884 len: 4,
1885 alignment: crate::PayloadAlignment::new(2).unwrap(),
1886 },
1887 )
1888 .unwrap();
1889
1890 let prepared = PreparedTxLoanSpec::from_validated::<
1891 UProtocolNativeWire,
1892 NativePrefixFrameMetadataCodec,
1893 >(spec, &NativePrefixFrameMetadataCodec)
1894 .unwrap();
1895
1896 assert_eq!(prepared.metadata(), &metadata);
1897 assert_eq!(prepared.payload_len(), 4);
1898 assert_eq!(prepared.payload_alignment(), 2);
1899 assert_eq!(
1900 prepared.payload_alignment_proof(),
1901 crate::PayloadAlignment::new(2).unwrap()
1902 );
1903 assert!(!prepared.encoded_metadata().is_empty());
1904 }
1905
1906 #[cfg(feature = "owned-frame-transport")]
1907 #[test]
1908 fn wire_transport_fits_owned_blanket_boundary() {
1909 fn assert_owned_impl<T: UOwnedTransportImpl>() {}
1910
1911 type Transport =
1912 UWireTransport<CompileCore, UProtocolNativeWire, NativePrefixFrameMetadataCodec>;
1913 assert_owned_impl::<Transport>();
1914 }
1915
1916 #[test]
1917 fn with_wire_constructs_adapter() {
1918 let transport = CompileCore.into_native_prefix_wire_transport(UProtocolNativeWire);
1919 let _: &UProtocolNativeWire = transport.wire();
1920 let _: &CompileCore = transport.core();
1921 }
1922
1923 #[test]
1924 fn existing_public_rx_lease_still_compiles_independently() {
1925 fn assert_public_rx<T: UZeroCopyRxLease>() {}
1926
1927 assert_public_rx::<UVecRxLease>();
1928 }
1929}