#![cfg_attr(
not(test),
expect(
dead_code,
reason = "the handles hold this; tests exercise it meanwhile"
)
)]
use tokio::sync::watch;
use crate::status::PipeStatus;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PeerPath {
Direct,
Relayed,
}
pub(crate) struct Lifecycle {
status: watch::Sender<PipeStatus>,
torn_down: watch::Sender<bool>,
in_flight: watch::Sender<usize>,
}
pub(crate) struct InFlight {
counter: watch::Sender<usize>,
}
impl Drop for InFlight {
fn drop(&mut self) {
self.counter.send_modify(|n| *n = n.saturating_sub(1));
}
}
impl Lifecycle {
pub(crate) fn new() -> Self {
Self {
status: watch::Sender::new(PipeStatus::Idle),
torn_down: watch::Sender::new(false),
in_flight: watch::Sender::new(0),
}
}
pub(crate) fn enter(&self) -> InFlight {
self.in_flight.send_modify(|n| *n += 1);
InFlight {
counter: self.in_flight.clone(),
}
}
pub(crate) fn in_flight(&self) -> usize {
*self.in_flight.borrow()
}
pub(crate) async fn wait_until_drained(&self) {
let mut rx = self.in_flight.subscribe();
if *rx.borrow_and_update() == 0 {
return;
}
let _ = rx.wait_for(|n| *n == 0).await;
}
pub(crate) fn status(&self) -> PipeStatus {
*self.status.borrow()
}
pub(crate) fn set_status(&self, next: PipeStatus) {
self.status.send_if_modified(|current| {
if *current == PipeStatus::Closed || *current == next {
false
} else {
*current = next;
true
}
});
}
pub(crate) async fn changed_since(&self, snapshot: PipeStatus) -> PipeStatus {
let mut rx = self.status.subscribe();
loop {
let current = *rx.borrow_and_update();
if current != snapshot || current == PipeStatus::Closed {
return current;
}
if rx.changed().await.is_err() {
return PipeStatus::Closed;
}
}
}
pub(crate) fn close(&self) {
self.set_status(PipeStatus::Closed);
}
pub(crate) async fn wait_until_closed(&self) {
let mut rx = self.status.subscribe();
if *rx.borrow_and_update() == PipeStatus::Closed {
return;
}
let _ = rx.wait_for(|s| *s == PipeStatus::Closed).await;
}
pub(crate) fn mark_torn_down(&self) {
self.torn_down.send_replace(true);
}
pub(crate) async fn wait_until_torn_down(&self) {
let mut rx = self.torn_down.subscribe();
if *rx.borrow_and_update() {
return;
}
let _ = rx.wait_for(|done| *done).await;
}
}
pub(crate) fn aggregate(peers: &[PeerPath]) -> PipeStatus {
if peers.is_empty() {
PipeStatus::Idle
} else if peers.contains(&PeerPath::Relayed) {
PipeStatus::Relayed
} else {
PipeStatus::Direct
}
}
#[cfg(test)]
#[path = "lifecycle_tests.rs"]
mod lifecycle_tests;