Skip to main content

mapping/
pager.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 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
12/// Runs a synchronous event loop that listens on `port` for pager page requests,
13/// resolves the requested blob using `blobs`, allocates a destination buffer using
14/// `buffer_factory`, and invokes `Blob::read_range`.
15///
16/// Notice that buffer management and page supply (such as calling `zx::Pager::supply_pages`
17/// or transmitting over `vmo-fifo`) are entirely delegated to the `DataBuffer` implementation
18/// returned by `buffer_factory`.
19pub 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
48/// A handle to a spawned background thread running `run_pager_loop`.
49/// Dropping this handle sends a shutdown packet to the port and cleanly joins the thread.
50pub struct PagerThread {
51    port: Port,
52    blobs: Arc<Blobs>,
53    thread: Option<std::thread::JoinHandle<()>>,
54}
55
56impl PagerThread {
57    /// Spawns a background thread running `run_pager_loop`.
58    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    /// Returns a reference to the active blob registry.
73    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}