Skip to main content

block_server/
verifier.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 delivery_blob::compression::{ChunkedArchiveError, DataBuffer};
6use fuchsia_sync::Mutex;
7use mapping::{
8    DELIVERY_DATA_SIZE, DeliveryCommand, DeliveryHandler, PENDING_DELIVERY_COMMANDS_CAPACITY,
9    PageRequest, RawDeliveryCommand,
10};
11use std::ops::Range;
12use std::sync::Arc;
13use storage_ptr_slice::MutPtrByteSlice;
14use vmo_fifo::SyncSender;
15
16/// Manages data verification and page supply via the delivery queue VMO.
17pub struct Verifier {
18    sender: Mutex<Option<SyncSender<RawDeliveryCommand>>>,
19}
20
21impl Verifier {
22    pub fn new(delivery_queue: zx::Vmo) -> Self {
23        let sender = if delivery_queue.get_size().unwrap_or(0) > 0 {
24            // Payloads in the delivery queue are transferred to the kernel pager via
25            // `zx_pager_supply_pages`, which requires the source offset to be page-aligned. Set
26            // the alignment of payload offsets to be page aligned.
27            SyncSender::<RawDeliveryCommand>::new(
28                delivery_queue,
29                zx::system_get_page_size() as usize,
30                PENDING_DELIVERY_COMMANDS_CAPACITY,
31            )
32            .ok()
33        } else {
34            None
35        };
36        Self { sender: Mutex::new(sender) }
37    }
38}
39
40impl DeliveryHandler for Verifier {
41    type Request = Buffer;
42
43    /// Returns a [`PageRequest`] implementation for delivering page-in data.
44    fn get_page_request(self: &Arc<Self>, key: u64, original_range: Range<u64>) -> Buffer {
45        Buffer {
46            verifier: Arc::clone(self),
47            key,
48            read_range: original_range,
49            write_cursor: 0,
50            delivered_len: 0,
51            data: Vec::new(),
52        }
53    }
54
55    /// Registers a blob's Merkle tree leaf hashes with the verifier via the delivery queue.
56    fn register_blob(&self, key: u64, merkle_leaves: &[[u8; 32]]) -> Result<(), anyhow::Error> {
57        let mut sender_guard = self.sender.lock();
58        if let Some(sender) = sender_guard.as_mut() {
59            let leaf_bytes: &[u8] = merkle_leaves.as_flattened();
60            let mut payload = sender.reserve_payload(leaf_bytes.len())?;
61            payload.data().copy_from_slice(leaf_bytes);
62            let cmd = DeliveryCommand::RegisterBlob {
63                key,
64                offset: payload.offset(),
65                length: leaf_bytes.len() as u32,
66            };
67            payload.commit(cmd.into())?;
68        }
69        Ok(())
70    }
71}
72
73/// Implementation of [`PageRequest`] for verifier deliveries that forwards
74/// unverified chunks via the delivery queue.
75pub struct Buffer {
76    verifier: Arc<Verifier>,
77    key: u64,
78    read_range: Range<u64>,
79    // Write cursor into `self.data` advanced by calls to [`DataBuffer::commit`].
80    write_cursor: usize,
81    // Bytes delivered to the verifier queue so far.
82    delivered_len: usize,
83    data: Vec<u8>,
84}
85
86impl Buffer {
87    fn deliver_chunk(&mut self, size: usize) -> Result<(), ChunkedArchiveError> {
88        // TODO(https://fxbug.dev/530494057): Optimize buffer management / payload reservations.
89        let chunk_data = &self.data[self.delivered_len..self.delivered_len + size];
90        let offset = self.read_range.start + self.delivered_len as u64;
91
92        let mut sender_guard = self.verifier.sender.lock();
93        if let Some(sender) = sender_guard.as_mut() {
94            let page_size = zx::system_get_page_size() as usize;
95            let aligned_size = size.next_multiple_of(page_size);
96            let mut payload = sender.reserve_payload(aligned_size).map_err(|e| {
97                log::error!(e:?; "Verifier::deliver_chunk: reserve_payload failed");
98                ChunkedArchiveError::IntegrityError
99            })?;
100
101            let payload_data = payload.data();
102            payload_data.subslice_mut(0..size).copy_from_slice(chunk_data);
103            payload_data.subslice_mut(size..aligned_size).fill(0);
104            let cmd = DeliveryCommand::Data {
105                key: self.key,
106                target_offset: offset,
107                length: aligned_size as u32,
108                offset: payload.offset(),
109            };
110            payload.commit(cmd.into()).map_err(|_| ChunkedArchiveError::IntegrityError)?;
111        }
112
113        self.delivered_len += size;
114        Ok(())
115    }
116}
117
118impl DataBuffer for Buffer {
119    fn range(&self) -> Range<u64> {
120        self.read_range.clone()
121    }
122
123    /// Returns a mutable slice into the mapped transfer buffer.
124    ///
125    /// # Panics
126    ///
127    /// Panics if `prepare()` has not been called prior to accessing this method.
128    fn mut_ptr_slice(&mut self) -> MutPtrByteSlice<'_> {
129        assert!(!self.data.is_empty(), "prepare must be called before accessing mut_ptr_slice");
130        MutPtrByteSlice::from(&mut self.data[self.write_cursor..])
131    }
132
133    fn commit(&mut self, size: usize) -> Result<(), ChunkedArchiveError> {
134        self.write_cursor += size;
135
136        // Deliver complete DELIVERY_DATA_SIZE chunks to the verifier queue as data is committed.
137        while self.write_cursor - self.delivered_len >= DELIVERY_DATA_SIZE {
138            self.deliver_chunk(DELIVERY_DATA_SIZE)?;
139        }
140
141        // Deliver any trailing remainder once the entire prepared range has been committed.
142        if self.write_cursor == self.data.len() && self.delivered_len < self.write_cursor {
143            self.deliver_chunk(self.write_cursor - self.delivered_len)?;
144        }
145
146        Ok(())
147    }
148}
149
150impl PageRequest for Buffer {
151    fn prepare(&mut self, read_range: Range<u64>) -> Result<(), ChunkedArchiveError> {
152        assert!(self.data.is_empty(), "prepare must only be called once");
153        self.read_range = read_range;
154        let len = (self.read_range.end - self.read_range.start) as usize;
155        self.data = vec![0u8; len];
156        Ok(())
157    }
158}
159
160#[cfg(test)]
161mod tests {
162    use super::*;
163    use mapping::DELIVERY_DATA_COMMAND;
164
165    #[fuchsia::test]
166    fn test_verifier_buffer_incremental_commit() {
167        let delivery_queue = zx::Vmo::create(65536).unwrap();
168        let mut receiver = vmo_fifo::Receiver::<RawDeliveryCommand>::new(
169            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
170            PENDING_DELIVERY_COMMANDS_CAPACITY,
171        )
172        .unwrap();
173
174        let verifier = Arc::new(Verifier::new(delivery_queue));
175        let key = 42u64;
176
177        let mut request = verifier.get_page_request(key, 0..8192);
178        request.prepare(0..8192).expect("prepare");
179
180        // Fill first page (4096 bytes) with 0xAA
181        let page1_data = vec![0xAAu8; 4096];
182        request.mut_ptr_slice().subslice_mut(0..4096).copy_from_slice(&page1_data);
183        request.commit(4096).expect("commit page 1");
184
185        // Verify mut_ptr_slice advanced to second page
186        assert_eq!(request.mut_ptr_slice().len(), 4096);
187        // Not all data has been committed yet, so no command should be delivered yet.
188        assert!(receiver.is_empty());
189
190        // Fill second page (4096 bytes) with 0xBB
191        let page2_data = vec![0xBBu8; 4096];
192        request.mut_ptr_slice().subslice_mut(0..4096).copy_from_slice(&page2_data);
193        request.commit(4096).expect("commit page 2");
194
195        // Verify remaining space is 0
196        assert_eq!(request.mut_ptr_slice().len(), 0);
197
198        // Check single consolidated message delivered on queue
199        let msg = receiver.peek().expect("msg");
200        assert_eq!(msg.opcode, DELIVERY_DATA_COMMAND);
201        assert_eq!(msg.key, key);
202        assert_eq!(msg.length, 8192);
203        assert_eq!(msg.target_offset, 0);
204        let mut buf = vec![0u8; 8192];
205        msg.payload_slice(msg.offset, msg.length).copy_to_slice(&mut buf);
206        assert_eq!(&buf[..4096], &page1_data[..]);
207        assert_eq!(&buf[4096..], &page2_data[..]);
208        msg.pop().expect("pop msg");
209    }
210
211    #[fuchsia::test]
212    fn test_verifier_buffer_page_aligned_commits() {
213        let delivery_queue = zx::Vmo::create(65536).unwrap();
214        let mut receiver = vmo_fifo::Receiver::<RawDeliveryCommand>::new(
215            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
216            PENDING_DELIVERY_COMMANDS_CAPACITY,
217        )
218        .unwrap();
219
220        let verifier = Arc::new(Verifier::new(delivery_queue));
221        let key = 43u64;
222
223        let mut request = verifier.get_page_request(key, 0..5000);
224        request.prepare(0..5000).expect("prepare");
225
226        let data = vec![0xCCu8; 5000];
227
228        // Commit first 4096 bytes
229        request.mut_ptr_slice().subslice_mut(0..4096).copy_from_slice(&data[..4096]);
230        request.commit(4096).expect("commit 4096 bytes");
231        assert!(receiver.is_empty());
232
233        // Commit remaining 904 bytes (completes the prepared range)
234        request.mut_ptr_slice().subslice_mut(0..904).copy_from_slice(&data[4096..5000]);
235        request.commit(904).expect("commit 904 bytes");
236
237        // Check single message delivered, page-aligned to 8192
238        let msg = receiver.peek().expect("msg");
239        assert_eq!(msg.target_offset, 0);
240        assert_eq!(msg.length, 8192);
241        let mut buf = vec![0u8; 8192];
242        msg.payload_slice(msg.offset, msg.length).copy_to_slice(&mut buf);
243        assert_eq!(&buf[..5000], &data[..]);
244        assert_eq!(&buf[5000..], &[0u8; 3192]);
245        msg.pop().expect("pop msg");
246    }
247
248    #[fuchsia::test]
249    fn test_verifier_buffer_delivery_data_size_chunking() {
250        let total_size = DELIVERY_DATA_SIZE * 2;
251        let delivery_queue = zx::Vmo::create(total_size as u64 + 65536).unwrap();
252        let mut receiver = vmo_fifo::Receiver::<RawDeliveryCommand>::new(
253            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
254            PENDING_DELIVERY_COMMANDS_CAPACITY,
255        )
256        .unwrap();
257
258        let verifier = Arc::new(Verifier::new(delivery_queue));
259        let key = 44u64;
260
261        let mut request = verifier.get_page_request(key, 0..total_size as u64);
262        request.prepare(0..total_size as u64).expect("prepare");
263
264        let chunk_32k = vec![0x55u8; 32 * 1024];
265
266        // Push 3 chunks of 32 KiB (96 KiB) - less than DELIVERY_DATA_SIZE
267        for _ in 0..3 {
268            request.mut_ptr_slice().subslice_mut(0..32 * 1024).copy_from_slice(&chunk_32k);
269            request.commit(32 * 1024).expect("commit");
270        }
271        assert!(receiver.is_empty());
272
273        // Push 4th chunk (now 128 KiB == DELIVERY_DATA_SIZE)
274        request.mut_ptr_slice().subslice_mut(0..32 * 1024).copy_from_slice(&chunk_32k);
275        request.commit(32 * 1024).expect("commit");
276
277        // 1st 128 KiB message delivered
278        let msg1 = receiver.peek().expect("msg1");
279        assert_eq!(msg1.target_offset, 0);
280        assert_eq!(msg1.length, DELIVERY_DATA_SIZE as u32);
281        msg1.pop().expect("pop msg1");
282
283        // Push remaining 4 chunks (another 128 KiB)
284        for _ in 0..4 {
285            request.mut_ptr_slice().subslice_mut(0..32 * 1024).copy_from_slice(&chunk_32k);
286            request.commit(32 * 1024).expect("commit");
287        }
288
289        // 2nd 128 KiB message delivered
290        let msg2 = receiver.peek().expect("msg2");
291        assert_eq!(msg2.target_offset, DELIVERY_DATA_SIZE as u64);
292        assert_eq!(msg2.length, DELIVERY_DATA_SIZE as u32);
293        msg2.pop().expect("pop msg2");
294    }
295
296    #[fuchsia::test]
297    #[should_panic(expected = "prepare must only be called once")]
298    fn test_verifier_buffer_prepare_multiple_times_panics() {
299        let delivery_queue = zx::Vmo::create(4096).unwrap();
300        let verifier = Arc::new(Verifier::new(delivery_queue));
301
302        let key = 46u64;
303        let mut request = verifier.get_page_request(key, 0..8192);
304        request.prepare(0..4096).expect("first prepare");
305        let _ = request.prepare(0..8192);
306    }
307
308    #[fuchsia::test]
309    fn test_verifier_buffer_supplies_pages_via_delivery_receiver() {
310        let pager = Arc::new(zx::Pager::create(zx::PagerOptions::empty()).unwrap());
311        let port = zx::Port::create();
312        let key = 44u64;
313        let paged_vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, key, 8192).unwrap();
314
315        let delivery_queue = zx::Vmo::create(65536).unwrap();
316        let vmo_provider = Arc::new(blob_pager_and_verifier::TestVmoProvider::new(
317            pager.clone(),
318            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
319        ));
320        vmo_provider
321            .register_vmo(key, paged_vmo.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap());
322        let receiver = vmo_fifo::Receiver::<RawDeliveryCommand>::new(
323            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
324            PENDING_DELIVERY_COMMANDS_CAPACITY,
325        )
326        .unwrap();
327        let _processor = blob_pager_and_verifier::DeliveryQueueProcessor::spawn(
328            receiver,
329            vmo_provider,
330            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
331        )
332        .unwrap();
333
334        let verifier = Arc::new(Verifier::new(delivery_queue));
335        let mut request = verifier.get_page_request(key, 0..8192);
336        request.prepare(0..8192).expect("prepare");
337
338        let page1_data = vec![0xAAu8; 4096];
339        request.mut_ptr_slice().subslice_mut(0..4096).copy_from_slice(&page1_data);
340        request.commit(4096).expect("commit page 1");
341
342        let page2_data = vec![0xBBu8; 4096];
343        request.mut_ptr_slice().subslice_mut(0..4096).copy_from_slice(&page2_data);
344        request.commit(4096).expect("commit page 2");
345
346        let mut read_buf = vec![0u8; 8192];
347        paged_vmo.read(&mut read_buf, 0).expect("read paged_vmo");
348        assert_eq!(&read_buf[..4096], &page1_data[..]);
349        assert_eq!(&read_buf[4096..8192], &page2_data[..]);
350    }
351
352    #[fuchsia::test]
353    fn test_verifier_register_blob() {
354        let delivery_queue = zx::Vmo::create(65536).unwrap();
355        let mut receiver = vmo_fifo::Receiver::<RawDeliveryCommand>::new(
356            delivery_queue.duplicate_handle(zx::Rights::SAME_RIGHTS).unwrap(),
357            PENDING_DELIVERY_COMMANDS_CAPACITY,
358        )
359        .unwrap();
360
361        let verifier = Arc::new(Verifier::new(delivery_queue));
362        let key = 42u64;
363        let leaves = [[0xABu8; 32], [0xCDu8; 32]];
364
365        verifier.register_blob(key, &leaves).expect("register_blob failed");
366
367        let msg = receiver.peek().expect("peek msg");
368        assert_eq!(msg.opcode, mapping::DELIVERY_REGISTER_BLOB_COMMAND);
369        assert_eq!(msg.key, key);
370        assert_eq!(msg.length, 64);
371        let mut buf = vec![0u8; 64];
372        msg.payload_slice(msg.offset, msg.length).copy_to_slice(&mut buf);
373        assert_eq!(&buf[..32], &[0xABu8; 32]);
374        assert_eq!(&buf[32..], &[0xCDu8; 32]);
375        msg.pop().expect("pop msg failed");
376    }
377}