Skip to main content

vmo_fifo/
signal.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
5// Use Zircon user signals to coordinate queue states between producer and consumer.
6// We use a pair of alternating signals for each event type (data available vs. space available)
7// to eliminate the extra syscall required to clear signal bits after waking up.
8pub const SIG_DATA_AVAILABLE_0: zx::Signals = zx::Signals::USER_0;
9pub const SIG_DATA_AVAILABLE_1: zx::Signals = zx::Signals::USER_1;
10pub const SIG_SPACE_AVAILABLE_0: zx::Signals = zx::Signals::USER_2;
11pub const SIG_SPACE_AVAILABLE_1: zx::Signals = zx::Signals::USER_3;
12pub const SIG_SHUTDOWN: zx::Signals = zx::Signals::USER_4;
13
14// A synchronization helper that coordinates event waiting and signaling on a VMO using a pair of
15// Zircon user signals and toggling between them.
16pub(crate) struct EventSignal {
17    set_next: zx::Signals,
18    clear_next: zx::Signals,
19}
20
21impl EventSignal {
22    pub(crate) const fn new(sig_0: zx::Signals, sig_1: zx::Signals) -> Self {
23        Self { set_next: sig_0, clear_next: sig_1 }
24    }
25
26    // Blocks until the currently expected signal bit (or `shutdown_sig`) is asserted.
27    pub(crate) fn wait(
28        &mut self,
29        vmo: &zx::Vmo,
30        shutdown_sig: zx::Signals,
31    ) -> Result<(), zx::Status> {
32        let mask = self.set_next | shutdown_sig;
33
34        let observed = vmo.wait_one(mask, zx::MonotonicInstant::INFINITE).to_result()?;
35        if observed.contains(shutdown_sig) {
36            return Err(zx::Status::CANCELED);
37        }
38
39        std::mem::swap(&mut self.set_next, &mut self.clear_next);
40        Ok(())
41    }
42
43    pub(crate) async fn wait_async(
44        &mut self,
45        vmo: &zx::Vmo,
46        shutdown_sig: zx::Signals,
47    ) -> Result<(), zx::Status> {
48        let mask = self.set_next | shutdown_sig;
49
50        let observed = fuchsia_async::OnSignals::new(vmo, mask).await?;
51        if observed.contains(shutdown_sig) {
52            return Err(zx::Status::CANCELED);
53        }
54
55        std::mem::swap(&mut self.set_next, &mut self.clear_next);
56        Ok(())
57    }
58
59    pub(crate) fn signal(&mut self, vmo: &zx::Vmo) -> Result<(), zx::Status> {
60        // Asserts the active signal bit and clears the alternate bit.
61        vmo.signal(self.clear_next, self.set_next)?;
62        std::mem::swap(&mut self.set_next, &mut self.clear_next);
63        Ok(())
64    }
65}
66
67#[cfg(test)]
68mod tests {
69    use super::*;
70    use fuchsia_async::{DurationExt, Timer};
71
72    #[fuchsia::test]
73    fn test_event_signal() {
74        let vmo = zx::Vmo::create(4096).expect("VMO creation failed");
75        let mut event_sender = EventSignal::new(SIG_DATA_AVAILABLE_0, SIG_DATA_AVAILABLE_1);
76        let mut event_receiver = EventSignal::new(SIG_DATA_AVAILABLE_0, SIG_DATA_AVAILABLE_1);
77
78        // First cycle
79        event_sender.signal(&vmo).expect("signal failed");
80        event_receiver.wait(&vmo, SIG_SHUTDOWN).expect("wait failed");
81
82        // Second cycle
83        event_sender.signal(&vmo).expect("signal failed");
84        event_receiver.wait(&vmo, SIG_SHUTDOWN).expect("wait failed");
85    }
86
87    #[fuchsia::test]
88    async fn test_event_signal_async_concurrency() {
89        let vmo = zx::Vmo::create(4096).expect("VMO creation failed");
90        let mut event_sender = EventSignal::new(SIG_DATA_AVAILABLE_0, SIG_DATA_AVAILABLE_1);
91        let mut event_receiver = EventSignal::new(SIG_DATA_AVAILABLE_0, SIG_DATA_AVAILABLE_1);
92
93        // This task will wait for the signal before it is sent.
94        let wait_task = async {
95            event_receiver.wait_async(&vmo, SIG_SHUTDOWN).await.expect("wait failed");
96        };
97
98        // This task will delay briefly, then send the signal.
99        let signal_task = async {
100            Timer::new(zx::MonotonicDuration::from_millis(1).after_now()).await;
101            event_sender.signal(&vmo).expect("signal failed");
102        };
103
104        // Run the tasks concurrently. If the wait_task blocks the executor thread synchronously,
105        // this join will deadlock.
106        futures::join!(wait_task, signal_task);
107    }
108}