medi_rs/adapters/
lifecycle.rs1use core::future::{Future, poll_fn};
4use core::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
5
6use crate::StartError;
7use core::task::Poll;
8use futures::task::AtomicWaker;
9
10pub struct Lifecycle {
16 started: AtomicBool,
17 running: AtomicUsize,
18 shutdown_requested: AtomicBool,
19 waker: AtomicWaker,
20}
21
22impl Lifecycle {
23 pub const fn new() -> Self {
25 Self {
26 started: AtomicBool::new(false),
27 running: AtomicUsize::new(0),
28 shutdown_requested: AtomicBool::new(false),
29 waker: AtomicWaker::new(),
30 }
31 }
32
33 pub fn start(&self) -> core::result::Result<(), StartError> {
38 self.started
39 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
40 .map(|_| ())
41 .map_err(|_| StartError::AlreadyStarted)
42 }
43
44 pub fn is_started(&self) -> bool {
46 self.started.load(Ordering::Acquire)
47 }
48
49 pub fn request_shutdown(&self) -> bool {
51 self.shutdown_requested
52 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
53 .is_ok()
54 }
55
56 pub fn begin(&self) {
58 self.running.fetch_add(1, Ordering::AcqRel);
59 }
60
61 pub fn finish(&self) {
63 if self.running.fetch_sub(1, Ordering::AcqRel) == 1 {
64 self.waker.wake();
65 }
66 }
67
68 pub fn wait(&self) -> impl Future<Output = ()> + '_ {
70 poll_fn(|cx| {
71 if self.running.load(Ordering::Acquire) == 0 {
72 return Poll::Ready(());
73 }
74 self.waker.register(cx.waker());
75 if self.running.load(Ordering::Acquire) == 0 {
76 Poll::Ready(())
77 } else {
78 Poll::Pending
79 }
80 })
81 }
82}
83
84impl Default for Lifecycle {
85 fn default() -> Self {
86 Self::new()
87 }
88}