Preview of the proposed up-rust native-frame-model branch (up-rust b6b99c6d, up-spec f0e9b17) — not released documentation. branch · write-up

up_rust/zero_copy/
tests.rs

1/********************************************************************************
2 * Copyright (c) 2026 Contributors to the Eclipse Foundation
3 *
4 * See the NOTICE file(s) distributed with this work for additional
5 * information regarding copyright ownership.
6 *
7 * This program and the accompanying materials are made available under the
8 * terms of the Apache License Version 2.0 which is available at
9 * https://www.apache.org/licenses/LICENSE-2.0
10 *
11 * SPDX-License-Identifier: Apache-2.0
12 ********************************************************************************/
13
14use super::*;
15
16#[cfg(all(test, feature = "owned-frame-transport"))]
17use crate::UOwnedFrame;
18#[cfg(test)]
19use std::io::Read;
20
21/// Owned buffer useful for tests, examples, and adapters that emulate a transmit loan.
22#[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/// Owned uninitialized buffer useful for tests and examples.
31#[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/// In-memory receive lease for tests and examples that need receive-lease shape.
40#[derive(Clone, Debug, PartialEq)]
41pub struct UVecRxLease {
42    metadata: UFrameMetadata,
43    payload: Option<Vec<u8>>,
44}
45
46impl UVecTxBuffer {
47    /// Creates an owned transmit buffer with `payload_len` zero-initialized bytes.
48    #[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    /// Creates an owned transmit buffer whose visible payload starts at `alignment`.
59    ///
60    /// # Errors
61    ///
62    /// Returns an error if the requested alignment is invalid or allocation size overflows.
63    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    /// Converts the buffer into an in-memory receive lease, consuming the emulated loan.
87    #[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    /// Creates an owned uninitialized transmit buffer.
132    #[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    /// Creates an owned uninitialized transmit buffer whose visible payload starts at `alignment`.
143    ///
144    /// # Errors
145    ///
146    /// Returns an error if the requested alignment is invalid or allocation size overflows.
147    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    /// # Safety
198    ///
199    /// Caller upholds the documented loan/witness contract for this handle:
200    /// exclusive, layout-valid storage (constructors) or complete
201    /// initialization of every transported byte (witness discharges).
202    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        // SAFETY: `MaybeUninit<u8>` has the same layout as `u8`, prefix/suffix
229        // bytes were initialized above, and the caller guarantees visible
230        // payload bytes are initialized before conversion.
231        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    /// Creates an in-memory receive lease after validation.
243    ///
244    /// # Errors
245    ///
246    /// Returns an error if metadata or payload presence is invalid.
247    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    /// Creates an in-memory receive lease without validation.
254    #[must_use]
255    pub fn new_unchecked(metadata: UFrameMetadata, payload: Option<Vec<u8>>) -> Self {
256        Self { metadata, payload }
257    }
258
259    /// Consumes the lease and returns metadata plus optional payload bytes.
260    #[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        // SAFETY: `UVecRxLease` is the local vector-backed test receive lease
311        // selected in `USR-04B` preflight as the positive fake loan proof.
312        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/// In-memory zero-copy transport for tests and examples.
327#[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    /// Returns frames sent through [`UZeroCopyTransport::send_zero_copy`].
344    #[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    /// Injects a frame into the receive queue and registered zero-copy listeners.
354    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    // SAFETY: upholds the POD layout contract: repr(C), declared padding
474    // only, every byte of a live value initialized.
475    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    // SAFETY: upholds the POD layout contract: repr(C), declared padding
486    // only, every byte of a live value initialized.
487    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        // Passing an unvalidated spec to `loan_tx` is a compile error; the
634        // runtime rejection lives at the typestate transition:
635        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}