1pub 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
14pub(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 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 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 event_sender.signal(&vmo).expect("signal failed");
80 event_receiver.wait(&vmo, SIG_SHUTDOWN).expect("wait failed");
81
82 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 let wait_task = async {
95 event_receiver.wait_async(&vmo, SIG_SHUTDOWN).await.expect("wait failed");
96 };
97
98 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 futures::join!(wait_task, signal_task);
107 }
108}