Skip to main content

blob_pager_and_verifier/
delivery.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use anyhow::{Error, anyhow, bail};
6use fuchsia_sync::Mutex;
7use mapping::{DELIVERY_VMO_SIZE, DeliveryCommand, RawDeliveryCommand};
8use std::collections::HashMap;
9use std::sync::Arc;
10use storage_ptr_slice::PtrByteSlice;
11use zx;
12
13pub use mapping::DELIVERY_DATA_SIZE;
14
15/// Unverified data pages delivered from the driver to be verified.
16#[derive(Copy, Clone)]
17pub struct UnverifiedPages<'a> {
18    slice: PtrByteSlice<'a>,
19    /// Byte offset within the shared delivery VMO.
20    delivery_offset: u64,
21}
22
23impl<'a> UnverifiedPages<'a> {
24    pub(crate) fn new(slice: PtrByteSlice<'a>, delivery_offset: u64) -> Result<Self, zx::Status> {
25        let page_size = zx::system_get_page_size() as usize;
26        if slice.len() % page_size != 0 || delivery_offset % (page_size as u64) != 0 {
27            return Err(zx::Status::INVALID_ARGS);
28        }
29        Ok(Self { slice, delivery_offset })
30    }
31
32    /// Returns the length of the payload in bytes.
33    pub fn len_in_bytes(&self) -> usize {
34        self.slice.len()
35    }
36
37    /// Returns the length of the payload in pages.
38    pub fn len_in_pages(&self) -> usize {
39        self.slice.len() / (zx::system_get_page_size() as usize)
40    }
41
42    /// Returns the byte offset within the shared delivery VMO.
43    pub fn delivery_offset(&self) -> u64 {
44        self.delivery_offset
45    }
46
47    /// Returns the underlying raw pointer byte slice to the unverified data.
48    pub fn as_ptr_byte_slice(&self) -> PtrByteSlice<'a> {
49        self.slice
50    }
51}
52
53/// Trait providing data delivery and blob metadata registration for incoming delivery commands.
54pub trait DeliveryQueueProvider: Send + Sync + 'static {
55    /// Delivers data directly from the delivery VMO into the target blob.
56    fn deliver_pages(
57        &self,
58        key: u64,
59        target_offset: u64,
60        unverified_pages: UnverifiedPages<'_>,
61    ) -> Result<(), Error>;
62
63    /// Handles a RegisterBlob command using a pointer slice to shared memory.
64    fn register_blob(&self, key: u64, leaf_data: PtrByteSlice<'_>) -> Result<(), Error> {
65        let _ = (key, leaf_data);
66        Ok(())
67    }
68}
69
70/// A simple in-memory `DeliveryQueueProvider` for testing delivery without Merkle verification.
71pub struct TestVmoProvider {
72    pager: Arc<zx::Pager>,
73    delivery_vmo: zx::Vmo,
74    vmos: Mutex<HashMap<u64, zx::Vmo>>,
75}
76
77impl TestVmoProvider {
78    pub fn new(pager: Arc<zx::Pager>, delivery_vmo: zx::Vmo) -> Self {
79        Self { pager, delivery_vmo, vmos: Mutex::new(HashMap::new()) }
80    }
81
82    pub fn delivery_vmo(&self) -> &zx::Vmo {
83        &self.delivery_vmo
84    }
85
86    pub fn register_vmo(&self, key: u64, vmo: zx::Vmo) {
87        self.vmos.lock().insert(key, vmo);
88    }
89
90    pub fn unregister_vmo(&self, key: u64) {
91        self.vmos.lock().remove(&key);
92    }
93}
94
95impl DeliveryQueueProvider for TestVmoProvider {
96    fn deliver_pages(
97        &self,
98        key: u64,
99        target_offset: u64,
100        unverified_pages: UnverifiedPages<'_>,
101    ) -> Result<(), Error> {
102        let vmos = self.vmos.lock();
103        if let Some(target_vmo) = vmos.get(&key) {
104            let length = unverified_pages.len_in_bytes() as u64;
105            if length == 0 {
106                return Ok(());
107            }
108            self.pager.supply_pages(
109                target_vmo,
110                target_offset..target_offset + length,
111                &self.delivery_vmo,
112                unverified_pages.delivery_offset(),
113            )?;
114            Ok(())
115        } else {
116            bail!("Unknown or expired key {key} in deliver_pages");
117        }
118    }
119}
120
121pub struct DeliveryQueueProcessor {
122    delivery_vmo: zx::Vmo,
123    delivery_thread: Option<std::thread::JoinHandle<()>>,
124}
125
126impl DeliveryQueueProcessor {
127    pub fn spawn(
128        mut receiver: vmo_fifo::Receiver<RawDeliveryCommand>,
129        provider: Arc<dyn DeliveryQueueProvider>,
130        delivery_vmo: zx::Vmo,
131    ) -> Result<Self, Error> {
132        let delivery_vmo_clone = delivery_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS)?;
133
134        let thread = std::thread::spawn(move || {
135            while let Ok(msg) = receiver.peek() {
136                let raw_cmd = RawDeliveryCommand {
137                    opcode: msg.opcode,
138                    _padding: msg._padding,
139                    key: msg.key,
140                    target_offset: msg.target_offset,
141                    length: msg.length,
142                    offset: msg.offset,
143                };
144
145                match DeliveryCommand::try_from(raw_cmd) {
146                    Ok(cmd) => {
147                        if let Err(e) = Self::handle_delivery(cmd, &msg, &*provider) {
148                            log::warn!("Failed to handle delivery command: {:?}", e);
149                        }
150                    }
151                    Err(e) => {
152                        log::warn!("Invalid delivery command opcode: {:?}", e);
153                    }
154                }
155                if let Err(e) = msg.pop() {
156                    log::warn!("Failed to pop delivery message: {:?}", e);
157                    break;
158                }
159            }
160        });
161
162        Ok(Self { delivery_vmo: delivery_vmo_clone, delivery_thread: Some(thread) })
163    }
164
165    pub fn delivery_vmo(&self) -> &zx::Vmo {
166        &self.delivery_vmo
167    }
168
169    fn handle_delivery(
170        cmd: DeliveryCommand,
171        raw_msg: &vmo_fifo::Message<'_, RawDeliveryCommand>,
172        provider: &dyn DeliveryQueueProvider,
173    ) -> Result<(), Error> {
174        match cmd {
175            DeliveryCommand::Data { key, target_offset, length, offset } => {
176                if offset.checked_add(length).unwrap_or(u32::MAX) > DELIVERY_VMO_SIZE as u32 {
177                    bail!("Data payload out of bounds");
178                }
179                let vmo_offset = raw_msg.payload_region_offset() as u64 + offset as u64;
180                // The sender is expected to push page-aligned payloads. If the payload is the
181                // final chunk of a blob and the data does not end on a page boundary, the sender
182                // must extend `length` to the next page boundary and pad the trailing bytes with
183                // zeros. This is required because `fuchsia-merkle`'s `verify_aligned` expects
184                // the buffer length to be a multiple of the system page size.
185                let page_size = zx::system_get_page_size() as u32;
186                if offset % page_size != 0 {
187                    bail!("Delivery payload offset must be page aligned: {offset}");
188                }
189                if length % page_size != 0 {
190                    bail!("Delivery payload length must be page aligned: {length}");
191                }
192                if target_offset % (page_size as u64) != 0 {
193                    bail!("Delivery payload target_offset must be page aligned: {target_offset}");
194                }
195                let unverified_pages =
196                    UnverifiedPages::new(raw_msg.payload_slice(offset, length), vmo_offset)
197                        .map_err(|s| anyhow!("Invalid unverified pages: {s}"))?;
198                provider.deliver_pages(key, target_offset, unverified_pages)?;
199            }
200            DeliveryCommand::RegisterBlob { key, offset, length } => {
201                if offset.checked_add(length).unwrap_or(u32::MAX) > DELIVERY_VMO_SIZE as u32 {
202                    bail!("RegisterBlob payload out of bounds");
203                }
204                let leaf_data = raw_msg.payload_slice(offset, length);
205                provider.register_blob(key, leaf_data)?;
206            }
207        }
208        Ok(())
209    }
210}
211
212impl Drop for DeliveryQueueProcessor {
213    fn drop(&mut self) {
214        let _ = self.delivery_vmo.signal(zx::Signals::empty(), vmo_fifo::SIG_SHUTDOWN);
215        if let Some(thread) = self.delivery_thread.take() {
216            let _ = thread.join();
217        }
218    }
219}