blob_pager_and_verifier/
delivery.rs1use 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#[derive(Copy, Clone)]
17pub struct UnverifiedPages<'a> {
18 slice: PtrByteSlice<'a>,
19 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 pub fn len_in_bytes(&self) -> usize {
34 self.slice.len()
35 }
36
37 pub fn len_in_pages(&self) -> usize {
39 self.slice.len() / (zx::system_get_page_size() as usize)
40 }
41
42 pub fn delivery_offset(&self) -> u64 {
44 self.delivery_offset
45 }
46
47 pub fn as_ptr_byte_slice(&self) -> PtrByteSlice<'a> {
49 self.slice
50 }
51}
52
53pub trait DeliveryQueueProvider: Send + Sync + 'static {
55 fn deliver_pages(
57 &self,
58 key: u64,
59 target_offset: u64,
60 unverified_pages: UnverifiedPages<'_>,
61 ) -> Result<(), Error>;
62
63 fn register_blob(&self, key: u64, leaf_data: PtrByteSlice<'_>) -> Result<(), Error> {
65 let _ = (key, leaf_data);
66 Ok(())
67 }
68}
69
70pub 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 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}