1use super::*;
15
16#[cfg(all(test, feature = "owned-frame-transport"))]
17use crate::UOwnedFrame;
18#[cfg(test)]
19use std::io::Read;
20
21#[derive(Clone, Debug, PartialEq)]
23pub struct UVecTxBuffer {
24 metadata: UFrameMetadata,
25 storage: Vec<u8>,
26 payload_offset: usize,
27 payload_len: usize,
28}
29
30#[derive(Clone, Debug)]
32pub struct UVecUninitTxBuffer {
33 metadata: UFrameMetadata,
34 storage: Vec<MaybeUninit<u8>>,
35 payload_offset: usize,
36 payload_len: usize,
37}
38
39#[derive(Clone, Debug, PartialEq)]
41pub struct UVecRxLease {
42 metadata: UFrameMetadata,
43 payload: Option<Vec<u8>>,
44}
45
46impl UVecTxBuffer {
47 #[must_use]
49 pub fn new(metadata: UFrameMetadata, payload_len: usize) -> Self {
50 Self {
51 metadata,
52 storage: vec![0_u8; payload_len],
53 payload_offset: 0,
54 payload_len,
55 }
56 }
57
58 pub fn with_alignment(
64 metadata: UFrameMetadata,
65 payload_len: usize,
66 alignment: usize,
67 ) -> Result<Self, UStatus> {
68 validate_payload_layout(payload_len, alignment)?;
69 if payload_len == 0 {
70 return Ok(Self::new(metadata, payload_len));
71 }
72 let extra = alignment.saturating_sub(1);
73 let storage_len = payload_len.checked_add(extra).ok_or_else(|| {
74 invalid_argument("payload length plus alignment padding overflows usize")
75 })?;
76 let storage = vec![0_u8; storage_len];
77 let payload_offset = aligned_offset(storage.as_ptr() as usize, alignment);
78 Ok(Self {
79 metadata,
80 storage,
81 payload_offset,
82 payload_len,
83 })
84 }
85
86 #[must_use]
88 pub fn into_rx_lease(self) -> UVecRxLease {
89 let payload = if self.metadata.payload_encoding().is_some() {
90 Some(
91 self.storage
92 .get(self.payload_range())
93 .expect("UVecTxBuffer payload range must be in bounds")
94 .to_vec(),
95 )
96 } else {
97 None
98 };
99 UVecRxLease::new_unchecked(self.metadata, payload)
100 }
101
102 fn payload_range(&self) -> std::ops::Range<usize> {
103 let end = self
104 .payload_offset
105 .checked_add(self.payload_len)
106 .expect("UVecTxBuffer payload range overflow");
107 self.payload_offset..end
108 }
109}
110
111impl UTxBuffer for UVecTxBuffer {
112 fn metadata(&self) -> &UFrameMetadata {
113 &self.metadata
114 }
115
116 fn payload(&self) -> &[u8] {
117 self.storage
118 .get(self.payload_range())
119 .expect("UVecTxBuffer payload range must be in bounds")
120 }
121
122 fn payload_mut(&mut self) -> &mut [u8] {
123 let range = self.payload_range();
124 self.storage
125 .get_mut(range)
126 .expect("UVecTxBuffer payload range must be in bounds")
127 }
128}
129
130impl UVecUninitTxBuffer {
131 #[must_use]
133 pub fn new(metadata: UFrameMetadata, payload_len: usize) -> Self {
134 Self {
135 metadata,
136 storage: vec![MaybeUninit::uninit(); payload_len],
137 payload_offset: 0,
138 payload_len,
139 }
140 }
141
142 pub fn with_alignment(
148 metadata: UFrameMetadata,
149 payload_len: usize,
150 alignment: usize,
151 ) -> Result<Self, UStatus> {
152 validate_payload_layout(payload_len, alignment)?;
153 if payload_len == 0 {
154 return Ok(Self::new(metadata, payload_len));
155 }
156 let extra = alignment.saturating_sub(1);
157 let storage_len = payload_len.checked_add(extra).ok_or_else(|| {
158 invalid_argument("payload length plus alignment padding overflows usize")
159 })?;
160 let storage = vec![MaybeUninit::uninit(); storage_len];
161 let payload_offset = aligned_offset(storage.as_ptr() as usize, alignment);
162 Ok(Self {
163 metadata,
164 storage,
165 payload_offset,
166 payload_len,
167 })
168 }
169
170 fn payload_range(&self) -> std::ops::Range<usize> {
171 let end = self
172 .payload_offset
173 .checked_add(self.payload_len)
174 .expect("UVecUninitTxBuffer payload range overflow");
175 self.payload_offset..end
176 }
177}
178
179impl UUninitTxBuffer for UVecUninitTxBuffer {
180 type Initialized = UVecTxBuffer;
181
182 fn metadata(&self) -> &UFrameMetadata {
183 &self.metadata
184 }
185
186 fn payload_len(&self) -> usize {
187 self.payload_len
188 }
189
190 fn payload_uninit_mut(&mut self) -> &mut [MaybeUninit<u8>] {
191 let range = self.payload_range();
192 self.storage
193 .get_mut(range)
194 .expect("UVecUninitTxBuffer payload range must be in bounds")
195 }
196
197 unsafe fn assume_payload_init(self) -> Self::Initialized {
203 let Self {
204 metadata,
205 mut storage,
206 payload_offset,
207 payload_len,
208 } = self;
209 for slot in storage
210 .get_mut(..payload_offset)
211 .expect("UVecUninitTxBuffer prefix range must be in bounds")
212 {
213 slot.write(0);
214 }
215 let payload_end = payload_offset
216 .checked_add(payload_len)
217 .expect("UVecUninitTxBuffer payload range overflow");
218 for slot in storage
219 .get_mut(payload_end..)
220 .expect("UVecUninitTxBuffer suffix range must be in bounds")
221 {
222 slot.write(0);
223 }
224 let len = storage.len();
225 let capacity = storage.capacity();
226 let ptr = storage.as_mut_ptr().cast::<u8>();
227 std::mem::forget(storage);
228 let storage = unsafe { Vec::from_raw_parts(ptr, len, capacity) };
232 UVecTxBuffer {
233 metadata,
234 storage,
235 payload_offset,
236 payload_len,
237 }
238 }
239}
240
241impl UVecRxLease {
242 pub fn new(metadata: UFrameMetadata, payload: Option<Vec<u8>>) -> Result<Self, UStatus> {
248 let frame = Self::new_unchecked(metadata, payload);
249 validate_frame_view_for_transport(&frame)?;
250 Ok(frame)
251 }
252
253 #[must_use]
255 pub fn new_unchecked(metadata: UFrameMetadata, payload: Option<Vec<u8>>) -> Self {
256 Self { metadata, payload }
257 }
258
259 #[must_use]
261 pub fn into_parts(self) -> (UFrameMetadata, Option<Vec<u8>>) {
262 (self.metadata, self.payload)
263 }
264
265 fn payload_bytes(&self) -> &[u8] {
266 self.payload.as_deref().unwrap_or_default()
267 }
268}
269
270impl UFrameView for UVecRxLease {
271 type PayloadReader<'a>
272 = Cursor<&'a [u8]>
273 where
274 Self: 'a;
275 type PayloadSlices<'a>
276 = std::option::IntoIter<&'a [u8]>
277 where
278 Self: 'a;
279
280 fn metadata(&self) -> &UFrameMetadata {
281 &self.metadata
282 }
283
284 fn payload_len(&self) -> usize {
285 self.payload_bytes().len()
286 }
287
288 fn has_payload(&self) -> bool {
289 self.payload.is_some()
290 }
291
292 fn payload_reader(&self) -> Self::PayloadReader<'_> {
293 Cursor::new(self.payload_bytes())
294 }
295
296 fn payload_slices(&self) -> Self::PayloadSlices<'_> {
297 self.payload.as_deref().into_iter()
298 }
299
300 fn try_contiguous_payload(&self) -> Option<&[u8]> {
301 self.payload.as_deref()
302 }
303}
304
305impl UZeroCopyRxLease for UVecRxLease {}
306
307impl ULoanedContiguousZeroCopyRxFrame for UVecRxLease {
308 fn loaned_contiguous_payload(&self) -> Result<LoanedPayload<'_>, UWireError> {
309 let payload = self.payload.as_deref().ok_or(UWireError::MissingPayload)?;
310 Ok(unsafe {
313 LoanedPayload::new_unchecked(payload, PayloadLoanProvenance::OpaqueTransportLoan)
314 })
315 }
316}
317
318#[cfg(any(test, feature = "test-util"))]
319#[derive(Default)]
320struct InMemoryState {
321 sent: Vec<UVecRxLease>,
322 queue: VecDeque<UVecRxLease>,
323 listeners: Vec<Arc<dyn UZeroCopyListener<UVecRxLease>>>,
324}
325
326#[cfg(any(test, feature = "test-util"))]
328#[derive(Clone, Default)]
329pub struct InMemoryZeroCopyTransport {
330 state: Arc<Mutex<InMemoryState>>,
331}
332
333#[cfg(any(test, feature = "test-util"))]
334impl core::fmt::Debug for InMemoryZeroCopyTransport {
335 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
336 f.debug_struct("InMemoryZeroCopyTransport")
337 .finish_non_exhaustive()
338 }
339}
340
341#[cfg(any(test, feature = "test-util"))]
342impl InMemoryZeroCopyTransport {
343 #[must_use]
345 pub fn sent_frames(&self) -> Vec<UVecRxLease> {
346 self.state
347 .lock()
348 .expect("zero-copy state lock poisoned")
349 .sent
350 .clone()
351 }
352
353 pub async fn inject(&self, frame: UVecRxLease) {
355 self.enqueue_and_dispatch(frame).await;
356 }
357
358 async fn enqueue_and_dispatch(&self, frame: UVecRxLease) {
359 let listeners = {
360 let mut state = self.state.lock().expect("zero-copy state lock poisoned");
361 state.queue.push_back(frame.clone());
362 state.listeners.clone()
363 };
364 for listener in listeners {
365 listener.on_receive_zero_copy(frame.clone()).await;
366 }
367 }
368}
369
370#[cfg(any(test, feature = "test-util"))]
371#[async_trait]
372impl UZeroCopyTransportImpl for InMemoryZeroCopyTransport {
373 type Tx = UVecTxBuffer;
374 type Rx = UVecRxLease;
375
376 async fn loan_validated_tx(&self, spec: UTxLoanSpec) -> Result<Self::Tx, UStatus> {
377 UVecTxBuffer::with_alignment(
378 spec.metadata().clone(),
379 spec.payload_len(),
380 spec.payload_alignment(),
381 )
382 }
383
384 async fn send_validated_zero_copy(&self, buffer: Self::Tx) -> Result<(), UStatus> {
385 let frame = buffer.into_rx_lease();
386 {
387 self.state
388 .lock()
389 .expect("zero-copy state lock poisoned")
390 .sent
391 .push(frame.clone());
392 }
393 self.enqueue_and_dispatch(frame).await;
394 Ok(())
395 }
396
397 async fn receive_validated_zero_copy(
398 &self,
399 _source_filter: &UUri,
400 _sink_filter: Option<&UUri>,
401 ) -> Result<Self::Rx, UStatus> {
402 self.state
403 .lock()
404 .expect("zero-copy state lock poisoned")
405 .queue
406 .pop_front()
407 .ok_or_else(|| UStatus::fail_with_code(UCode::NotFound, "no frame available"))
408 }
409
410 async fn register_validated_zero_copy_listener(
411 &self,
412 _source_filter: &UUri,
413 _sink_filter: Option<&UUri>,
414 listener: Arc<dyn UZeroCopyListener<Self::Rx>>,
415 ) -> Result<(), UStatus> {
416 self.state
417 .lock()
418 .expect("zero-copy state lock poisoned")
419 .listeners
420 .push(listener);
421 Ok(())
422 }
423
424 async fn unregister_validated_zero_copy_listener(
425 &self,
426 _source_filter: &UUri,
427 _sink_filter: Option<&UUri>,
428 listener: Arc<dyn UZeroCopyListener<Self::Rx>>,
429 ) -> Result<(), UStatus> {
430 let mut state = self.state.lock().expect("zero-copy state lock poisoned");
431 let Some(index) = state
432 .listeners
433 .iter()
434 .position(|existing| Arc::ptr_eq(existing, &listener))
435 else {
436 return Err(UStatus::fail_with_code(
437 UCode::NotFound,
438 "no such zero-copy listener registered for filters",
439 ));
440 };
441 state.listeners.remove(index);
442 Ok(())
443 }
444}
445
446#[cfg(any(test, feature = "test-util"))]
447#[async_trait]
448impl UZeroCopyUninitTransportImpl for InMemoryZeroCopyTransport {
449 type UninitTx = UVecUninitTxBuffer;
450
451 async fn loan_validated_uninit_tx(&self, spec: UTxLoanSpec) -> Result<Self::UninitTx, UStatus> {
452 UVecUninitTxBuffer::with_alignment(
453 spec.metadata().clone(),
454 spec.payload_len(),
455 spec.payload_alignment(),
456 )
457 }
458}
459#[cfg(test)]
460mod unit_tests {
461 use super::*;
462 use crate::payload::stable::StablePayloadVariant;
463 use crate::test_support::StableTestBytes as StableBytes;
464 use crate::{PayloadEncoding, UMessageBuilder};
465 use std::sync::Mutex as StdMutex;
466
467 #[repr(C)]
468 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
469 struct OtherStableBytes {
470 bytes: [u8; 4],
471 }
472
473 unsafe impl StablePayload for OtherStableBytes {
476 const TYPE_NAME: &'static str = "uprotocol.test.OtherStableBytes";
477 }
478
479 #[repr(C, align(4))]
480 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
481 struct AlignedStableBytes {
482 bytes: [u8; 4],
483 }
484
485 unsafe impl StablePayload for AlignedStableBytes {
488 const TYPE_NAME: &'static str = "uprotocol.test.AlignedStableBytes";
489 }
490
491 fn topic() -> UUri {
492 UUri::try_from_parts("vehicle", 0x4210, 0x01, 0x9000).expect("failed to create topic")
493 }
494
495 fn wildcard_topic_filter() -> UUri {
496 UUri::try_from_parts("vehicle", 0x4210, 0x01, 0xffff).expect("failed to create filter")
497 }
498
499 fn metadata_without_encoding() -> UFrameMetadata {
500 let message = UMessageBuilder::publish(topic()).build().expect("message");
501 crate::frame::metadata::try_project_attributes_to_frame_metadata(message.attributes(), None)
502 .expect("metadata")
503 }
504
505 fn metadata_with_encoding() -> UFrameMetadata {
506 let message = UMessageBuilder::publish(topic()).build().expect("message");
507 crate::frame::metadata::try_project_attributes_to_frame_metadata(
508 message.attributes(),
509 Some(PayloadEncoding::RAW),
510 )
511 .expect("metadata")
512 }
513
514 fn stable_metadata<T: StablePayload>() -> UFrameMetadata {
515 let message = UMessageBuilder::publish(topic()).build().expect("message");
516 crate::frame::metadata::try_project_attributes_to_frame_metadata(
517 message.attributes(),
518 Some(StableContainerPayload::<T>::encoding()),
519 )
520 .expect("metadata")
521 }
522
523 fn stable_encoding_with<T: StablePayload>(
524 type_name: &str,
525 variant: StablePayloadVariant,
526 size: usize,
527 alignment: usize,
528 ) -> PayloadEncoding {
529 PayloadEncoding::custom(
530 StableContainerPayload::<T>::ENCODING_ID,
531 format!(
532 "application/vnd.uprotocol.stable-container;type=\"{type_name}\";variant={variant};size={size};align={alignment}"
533 ),
534 )
535 .expect("stable encoding")
536 }
537
538 #[test]
539 fn tx_loan_no_payload_rejects_encoding() {
540 let error = UTxLoanSpec::no_payload(metadata_with_encoding()).unwrap_err();
541 assert_eq!(error.code(), UCode::InvalidArgument);
542 }
543
544 #[test]
545 fn tx_loan_payload_rejects_missing_encoding() {
546 let error = UTxLoanSpec::payload(metadata_without_encoding(), 8, 1).unwrap_err();
547 assert_eq!(error.code(), UCode::InvalidArgument);
548 }
549
550 #[test]
551 fn present_empty_payload_preserves_encoding() {
552 let spec = UTxLoanSpec::present_empty_payload(metadata_with_encoding()).unwrap();
553 assert!(spec.has_payload());
554 assert_eq!(spec.payload_len(), 0);
555 assert_eq!(spec.payload_alignment(), 1);
556 assert_eq!(
557 spec.metadata().payload_encoding(),
558 metadata_with_encoding().payload_encoding()
559 );
560 }
561
562 #[test]
563 fn payload_alignment_must_be_power_of_two() {
564 let error = UTxPayloadSpec::present(16, 3).unwrap_err();
565 assert_eq!(error.code(), UCode::InvalidArgument);
566 }
567
568 #[test]
569 fn payload_alignment_proof_rejects_zero() {
570 let error = PayloadAlignment::new(0).unwrap_err();
571 assert_eq!(error.code(), UCode::InvalidArgument);
572 }
573
574 #[test]
575 fn payload_alignment_proof_rejects_non_power_of_two() {
576 let error = PayloadAlignment::new(6).unwrap_err();
577 assert_eq!(error.code(), UCode::InvalidArgument);
578 }
579
580 #[test]
581 fn payload_alignment_proof_round_trips_raw_value() {
582 let alignment = PayloadAlignment::new(8).unwrap();
583 let spec = UTxLoanSpec::new(
584 metadata_with_encoding(),
585 UTxPayloadSpec::present_with_alignment(16, alignment),
586 )
587 .unwrap();
588 assert_eq!(alignment.as_usize(), 8);
589 assert_eq!(spec.payload_alignment(), 8);
590 assert_eq!(spec.payload_alignment_proof(), alignment);
591 }
592
593 struct CountingTransport {
594 loan_calls: StdMutex<usize>,
595 send_calls: StdMutex<usize>,
596 receive: StdMutex<Option<UVecRxLease>>,
597 }
598
599 #[async_trait]
600 impl UZeroCopyTransportImpl for CountingTransport {
601 type Tx = UVecTxBuffer;
602 type Rx = UVecRxLease;
603
604 async fn loan_validated_tx(&self, spec: UTxLoanSpec) -> Result<Self::Tx, UStatus> {
605 *self.loan_calls.lock().expect("loan lock") += 1;
606 UVecTxBuffer::with_alignment(
607 spec.metadata().clone(),
608 spec.payload_len(),
609 spec.payload_alignment(),
610 )
611 }
612
613 async fn send_validated_zero_copy(&self, _buffer: Self::Tx) -> Result<(), UStatus> {
614 *self.send_calls.lock().expect("send lock") += 1;
615 Ok(())
616 }
617
618 async fn receive_validated_zero_copy(
619 &self,
620 _source_filter: &UUri,
621 _sink_filter: Option<&UUri>,
622 ) -> Result<Self::Rx, UStatus> {
623 self.receive
624 .lock()
625 .expect("receive lock")
626 .take()
627 .ok_or_else(|| UStatus::fail_with_code(UCode::NotFound, "none"))
628 }
629 }
630
631 #[tokio::test]
632 async fn invalid_spec_rejected_at_the_typestate_transition() {
633 let invalid = UTxLoanSpec::new_unchecked(metadata_with_encoding(), UTxPayloadSpec::Absent);
636 let error = invalid.validate().unwrap_err();
637 assert_eq!(error.code(), UCode::InvalidArgument);
638 }
639
640 #[tokio::test]
641 async fn invalid_tx_buffer_rejected_before_send_implementation() {
642 let transport = CountingTransport {
643 loan_calls: StdMutex::new(0),
644 send_calls: StdMutex::new(0),
645 receive: StdMutex::new(None),
646 };
647 let buffer = UVecTxBuffer::new(metadata_without_encoding(), 4);
648
649 let error = transport.send_zero_copy(buffer).await.unwrap_err();
650
651 assert_eq!(error.code(), UCode::InvalidArgument);
652 assert_eq!(*transport.send_calls.lock().expect("send lock"), 0);
653 }
654
655 #[tokio::test]
656 async fn invalid_receive_lease_rejected_before_return() {
657 let transport = CountingTransport {
658 loan_calls: StdMutex::new(0),
659 send_calls: StdMutex::new(0),
660 receive: StdMutex::new(Some(UVecRxLease::new_unchecked(
661 metadata_with_encoding(),
662 None,
663 ))),
664 };
665
666 let error = transport
667 .receive_zero_copy(&wildcard_topic_filter(), None)
668 .await
669 .unwrap_err();
670
671 assert_eq!(error.code(), UCode::InvalidArgument);
672 }
673
674 struct SegmentedFrame {
675 metadata: UFrameMetadata,
676 first: Vec<u8>,
677 second: Vec<u8>,
678 }
679
680 impl UFrameView for SegmentedFrame {
681 type PayloadReader<'a>
682 = Cursor<Vec<u8>>
683 where
684 Self: 'a;
685 type PayloadSlices<'a>
686 = std::vec::IntoIter<&'a [u8]>
687 where
688 Self: 'a;
689
690 fn metadata(&self) -> &UFrameMetadata {
691 &self.metadata
692 }
693
694 fn payload_len(&self) -> usize {
695 self.first.len() + self.second.len()
696 }
697
698 fn payload_reader(&self) -> Self::PayloadReader<'_> {
699 let mut bytes = Vec::new();
700 bytes.extend_from_slice(&self.first);
701 bytes.extend_from_slice(&self.second);
702 Cursor::new(bytes)
703 }
704
705 fn payload_slices(&self) -> Self::PayloadSlices<'_> {
706 vec![self.first.as_slice(), self.second.as_slice()].into_iter()
707 }
708 }
709
710 #[test]
711 fn frame_view_reader_and_slices_preserve_order() {
712 let frame = SegmentedFrame {
713 metadata: metadata_with_encoding(),
714 first: b"abc".to_vec(),
715 second: b"def".to_vec(),
716 };
717 let mut from_reader = Vec::new();
718 frame
719 .payload_reader()
720 .read_to_end(&mut from_reader)
721 .expect("reader");
722 let from_slices: Vec<u8> = frame.payload_slices().flatten().copied().collect();
723
724 assert_eq!(from_reader, b"abcdef");
725 assert_eq!(from_slices, b"abcdef");
726 validate_frame_view_for_transport(&frame).unwrap();
727 }
728
729 #[test]
730 fn stable_borrow_accepts_loan_backed_contiguous_payload() {
731 let value = StableBytes { bytes: *b"loan" };
732 let frame = UVecRxLease::new(stable_metadata::<StableBytes>(), Some(value.bytes.to_vec()))
733 .expect("stable frame");
734
735 let borrowed = frame.borrow_stable_payload::<StableBytes>().unwrap();
736
737 assert_eq!(borrowed, &value);
738 assert_eq!(
739 frame.payload_loan_provenance().unwrap(),
740 PayloadLoanProvenance::OpaqueTransportLoan
741 );
742 }
743
744 #[test]
745 fn stable_borrow_rejects_wrong_type_metadata() {
746 let value = StableBytes { bytes: *b"loan" };
747 let frame = UVecRxLease::new(
748 stable_metadata::<OtherStableBytes>(),
749 Some(value.bytes.to_vec()),
750 )
751 .expect("stable frame");
752
753 let error = frame.borrow_stable_payload::<StableBytes>().unwrap_err();
754
755 assert!(matches!(error, UWireError::InvalidPayload(_)));
756 }
757
758 #[test]
759 fn stable_borrow_rejects_wrong_size_metadata() {
760 let value = StableBytes { bytes: *b"loan" };
761 let message = UMessageBuilder::publish(topic()).build().expect("message");
762 let metadata = crate::frame::metadata::try_project_attributes_to_frame_metadata(
763 message.attributes(),
764 Some(stable_encoding_with::<StableBytes>(
765 StableBytes::TYPE_NAME,
766 StablePayloadVariant::FixedSize,
767 std::mem::size_of::<StableBytes>() + 1,
768 std::mem::align_of::<StableBytes>(),
769 )),
770 )
771 .expect("metadata");
772 let frame = UVecRxLease::new(metadata, Some(value.bytes.to_vec())).expect("stable frame");
773
774 let error = frame.borrow_stable_payload::<StableBytes>().unwrap_err();
775
776 assert!(matches!(error, UWireError::InvalidPayload(_)));
777 }
778
779 #[test]
780 fn stable_borrow_rejects_insufficient_advertised_alignment() {
781 let message = UMessageBuilder::publish(topic()).build().expect("message");
782 let metadata = crate::frame::metadata::try_project_attributes_to_frame_metadata(
783 message.attributes(),
784 Some(stable_encoding_with::<AlignedStableBytes>(
785 AlignedStableBytes::TYPE_NAME,
786 StablePayloadVariant::FixedSize,
787 std::mem::size_of::<AlignedStableBytes>(),
788 1,
789 )),
790 )
791 .expect("metadata");
792 let frame = UVecRxLease::new(metadata, Some(vec![0_u8; 4])).expect("stable frame");
793
794 let error = frame
795 .borrow_stable_payload::<AlignedStableBytes>()
796 .unwrap_err();
797
798 assert!(matches!(error, UWireError::InvalidPayload(_)));
799 }
800
801 #[test]
802 fn stable_borrow_rejects_payload_length_mismatch() {
803 let frame = UVecRxLease::new(stable_metadata::<StableBytes>(), Some(vec![1, 2, 3]))
804 .expect("stable frame");
805
806 let error = frame.borrow_stable_payload::<StableBytes>().unwrap_err();
807
808 assert!(
809 matches!(error, UWireError::InvalidPayload(message) if message.contains("payload length"))
810 );
811 }
812
813 #[test]
814 fn stable_borrow_rejects_absent_payload() {
815 let frame = UVecRxLease::new_unchecked(stable_metadata::<StableBytes>(), None);
816
817 let error = frame.borrow_stable_payload::<StableBytes>().unwrap_err();
818
819 assert_eq!(error, UWireError::MissingPayload);
820 }
821
822 #[test]
823 fn segmented_frame_is_not_loan_backed_proof() {
824 let frame = SegmentedFrame {
825 metadata: stable_metadata::<StableBytes>(),
826 first: vec![1, 2],
827 second: vec![3, 4],
828 };
829
830 assert_eq!(frame.try_contiguous_payload(), None);
831 validate_frame_view_for_transport(&frame).unwrap();
832 }
833
834 #[tokio::test]
835 async fn in_memory_zero_copy_transport_round_trips_payload() {
836 let transport = InMemoryZeroCopyTransport::default();
837 let spec = UTxLoanSpec::payload(metadata_with_encoding(), 4, 1).unwrap();
838 let mut buffer = transport.loan_tx(spec).await.unwrap();
839 buffer.payload_mut().copy_from_slice(b"test");
840
841 transport.send_zero_copy(buffer).await.unwrap();
842 assert_eq!(transport.sent_frames().len(), 1);
843 let received = transport
844 .receive_zero_copy(&wildcard_topic_filter(), None)
845 .await
846 .unwrap();
847
848 assert_eq!(received.try_contiguous_payload(), Some(b"test".as_slice()));
849
850 transport
851 .inject(UVecRxLease::new(metadata_with_encoding(), Some(b"next".to_vec())).unwrap())
852 .await;
853 let injected = transport
854 .receive_zero_copy(&wildcard_topic_filter(), None)
855 .await
856 .unwrap();
857 assert_eq!(injected.try_contiguous_payload(), Some(b"next".as_slice()));
858 }
859
860 #[tokio::test]
861 async fn stable_initialized_tx_helper_sends_stable_payload() {
862 let transport = InMemoryZeroCopyTransport::default();
863
864 transport
865 .send_loaned_payload_as::<StableContainerPayload<StableBytes>, StableBytes>(
866 stable_metadata::<StableBytes>(),
867 |payload| payload.bytes.copy_from_slice(b"init"),
868 )
869 .await
870 .expect("send initialized stable payload");
871
872 let sent = transport.sent_frames();
873 assert_eq!(sent.len(), 1);
874 let sent = sent.first().expect("one sent frame");
875 assert_eq!(
876 sent.metadata().payload_encoding(),
877 Some(&StableContainerPayload::<StableBytes>::encoding())
878 );
879 assert_eq!(
880 sent.borrow_stable_payload::<StableBytes>().unwrap(),
881 &StableBytes { bytes: *b"init" }
882 );
883 }
884
885 #[cfg(feature = "zero-copy-transport")]
886 #[tokio::test]
887 async fn stable_uninit_tx_helper_sends_byte_backed_payload() {
888 let transport = InMemoryZeroCopyTransport::default();
889
890 transport
891 .send_uninit_loaned_payload_as::<StableContainerPayload<StableBytes>, StableBytes>(
892 stable_metadata::<StableBytes>(),
893 |slot| Ok(slot.write(StableBytes { bytes: *b"noze" })),
894 )
895 .await
896 .expect("send uninit stable payload");
897
898 let sent = transport.sent_frames();
899 assert_eq!(sent.len(), 1);
900 let sent = sent.first().expect("one sent frame");
901 assert_eq!(
902 sent.borrow_stable_payload::<StableBytes>().unwrap(),
903 &StableBytes { bytes: *b"noze" }
904 );
905 }
906
907 #[cfg(feature = "zero-copy-transport")]
908 #[tokio::test]
909 async fn stable_uninit_tx_helper_uses_stable_payload_init_builder() {
910 let transport = InMemoryZeroCopyTransport::default();
911
912 transport
913 .send_uninit_stable_payload_as::<StableBytes>(
914 stable_metadata::<StableBytes>(),
915 |context| context.into_init().bytes_from_array(b"zcpy").finish(),
916 )
917 .await
918 .expect("send stable init payload");
919
920 let sent = transport.sent_frames();
921 assert_eq!(sent.len(), 1);
922 let sent = sent.first().expect("one sent frame");
923 assert_eq!(
924 sent.borrow_stable_payload::<StableBytes>().unwrap(),
925 &StableBytes { bytes: *b"zcpy" }
926 );
927 }
928
929 #[cfg(feature = "perf-diagnostics")]
930 #[tokio::test]
931 async fn phased_stable_uninit_helper_preserves_send_semantics() {
932 let transport = InMemoryZeroCopyTransport::default();
933
934 let phases = transport
935 .send_uninit_stable_payload_as_phased::<StableBytes>(
936 stable_metadata::<StableBytes>(),
937 |context| context.into_init().bytes_from_array(b"perf").finish(),
938 )
939 .await
940 .expect("send phased stable init payload");
941
942 assert!(phases.total >= phases.accounted());
943 assert_eq!(phases.residual, phases.total - phases.accounted());
944 let sent = transport.sent_frames();
945 assert_eq!(sent.len(), 1);
946 assert_eq!(
947 sent.first()
948 .expect("one sent frame")
949 .borrow_stable_payload::<StableBytes>()
950 .unwrap(),
951 &StableBytes { bytes: *b"perf" }
952 );
953 let _clock_overhead = UninitStableSendPhases::calibrate_clock_overhead(32);
954 }
955
956 #[cfg(feature = "owned-frame-transport")]
957 #[test]
958 fn owned_frame_implements_frame_view_when_feature_enabled() {
959 let frame = UOwnedFrame::with_payload(
960 metadata_with_encoding(),
961 bytes::Bytes::from_static(b"owned"),
962 )
963 .expect("owned frame");
964
965 assert_eq!(UFrameView::payload_len(&frame), 5);
966 assert_eq!(frame.try_contiguous_payload(), Some(b"owned".as_slice()));
967 }
968}