1use crate::Blobs;
6use crate::reader::BlockService;
7use delivery_blob::DataBuffer;
8use std::sync::Arc;
9use zx::sys::zx_page_request_command_t::ZX_PAGER_VMO_READ;
10use zx::{Packet, PacketContents, Port, Rights, UserPacket};
11
12pub fn run_pager_loop<B: DataBuffer>(
20 port: &Port,
21 service: Arc<dyn BlockService>,
22 blobs: &Blobs,
23 mut buffer_factory: impl FnMut(u64, usize) -> B + Send + 'static,
24) {
25 loop {
26 match port.wait(zx::MonotonicInstant::INFINITE) {
27 Ok(packet) => match packet.contents() {
28 PacketContents::User(_) => break,
29 PacketContents::Pager(pager_packet) => {
30 let key = packet.key();
31 let command = pager_packet.command();
32 if command == ZX_PAGER_VMO_READ {
33 let offset = pager_packet.range().start;
34 let length = pager_packet.range().end - offset;
35 if let Some(blob) = blobs.get(key) {
36 let dest_buf = buffer_factory(offset, length as usize);
37 blob.read_range(offset..offset + length, service.as_ref(), dest_buf);
38 }
39 }
40 }
41 _ => {}
42 },
43 Err(_) => break,
44 }
45 }
46}
47
48pub struct PagerThread {
51 port: Port,
52 blobs: Arc<Blobs>,
53 thread: Option<std::thread::JoinHandle<()>>,
54}
55
56impl PagerThread {
57 pub fn spawn<B: DataBuffer + 'static>(
59 port: Port,
60 service: Arc<dyn BlockService>,
61 blobs: Arc<Blobs>,
62 buffer_factory: impl FnMut(u64, usize) -> B + Send + 'static,
63 ) -> Self {
64 let thread_port = port.duplicate_handle(Rights::SAME_RIGHTS).expect("duplicate port");
65 let thread_blobs = Arc::clone(&blobs);
66 let thread = std::thread::spawn(move || {
67 run_pager_loop(&thread_port, service, &thread_blobs, buffer_factory);
68 });
69 Self { port, blobs, thread: Some(thread) }
70 }
71
72 pub fn blobs(&self) -> &Arc<Blobs> {
74 &self.blobs
75 }
76}
77
78impl Drop for PagerThread {
79 fn drop(&mut self) {
80 let packet = Packet::from_user_packet(0, 0, UserPacket::from_u8_array([0; 32]));
81 let _ = self.port.queue(&packet);
82 if let Some(thread) = self.thread.take() {
83 let _ = thread.join();
84 }
85 }
86}
87
88#[cfg(test)]
89mod tests {
90 use super::*;
91 use crate::reader::tests::FakeBlockService;
92 use crate::testing::TestVecBuffer;
93
94 #[fuchsia::test]
95 fn test_pager_thread_lifecycle() {
96 let port = Port::create();
97 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
98 let blobs = Arc::new(Blobs::new());
99 let thread =
100 PagerThread::spawn(port, service, blobs, |_offset, _len| TestVecBuffer::new(4096).0);
101 drop(thread);
102 }
103
104 #[fuchsia::test]
105 fn test_pager_packet() {
106 let port = Port::create();
107 let pager = zx::Pager::create(zx::PagerOptions::empty()).expect("create pager");
108 let vmo = pager.create_vmo(zx::VmoOptions::empty(), &port, 1234, 4096).expect("create vmo");
109
110 let vmo_clone = vmo.duplicate_handle(Rights::SAME_RIGHTS).expect("duplicate vmo");
111 let _reader_thread = std::thread::spawn(move || {
112 let mut b = [0u8; 1];
113 let _ = vmo_clone.read(&mut b, 0);
114 });
115
116 let packet = port.wait(zx::MonotonicInstant::INFINITE).expect("wait packet");
117 assert_eq!(packet.key(), 1234);
118 if let PacketContents::Pager(pager_packet) = packet.contents() {
119 assert_eq!(pager_packet.command(), ZX_PAGER_VMO_READ);
120 assert_eq!(pager_packet.range(), 0..4096);
121 } else {
122 panic!("Expected pager packet");
123 }
124 }
125}