1use 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
16pub 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 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 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 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
73pub struct Buffer {
76 verifier: Arc<Verifier>,
77 key: u64,
78 read_range: Range<u64>,
79 write_cursor: usize,
81 delivered_len: usize,
83 data: Vec<u8>,
84}
85
86impl Buffer {
87 fn deliver_chunk(&mut self, size: usize) -> Result<(), ChunkedArchiveError> {
88 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 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 while self.write_cursor - self.delivered_len >= DELIVERY_DATA_SIZE {
138 self.deliver_chunk(DELIVERY_DATA_SIZE)?;
139 }
140
141 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 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 assert_eq!(request.mut_ptr_slice().len(), 4096);
187 assert!(receiver.is_empty());
189
190 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 assert_eq!(request.mut_ptr_slice().len(), 0);
197
198 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 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 request.mut_ptr_slice().subslice_mut(0..904).copy_from_slice(&data[4096..5000]);
235 request.commit(904).expect("commit 904 bytes");
236
237 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 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 request.mut_ptr_slice().subslice_mut(0..32 * 1024).copy_from_slice(&chunk_32k);
275 request.commit(32 * 1024).expect("commit");
276
277 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 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 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}