use std::fmt;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::Duration;
use crate::rt::sync::CancellationToken;
use crate::rt::sync::watch;
use super::error::FlowError;
use super::{Scope, within};
pub const MAX_DEPENDANTS: usize = 65_536;
pub const MAX_NAME_BYTES: usize = 128;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Generation(u64);
impl Generation {
pub const FIRST: Self = Self(0);
#[must_use]
pub fn next() -> Self {
static COUNTER: AtomicU64 = AtomicU64::new(0);
let mut held = COUNTER.load(Ordering::SeqCst);
loop {
let next = held.saturating_add(1);
match COUNTER.compare_exchange_weak(held, next, Ordering::SeqCst, Ordering::SeqCst) {
Ok(_) => return Self(next),
Err(observed) => held = observed,
}
}
}
#[must_use]
pub const fn get(self) -> u64 {
self.0
}
#[must_use]
pub const fn at(number: u64) -> Self {
Self(number)
}
}
impl fmt::Display for Generation {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(formatter, "generation {}", self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ReadinessError {
NotReady {
what: Arc<str>,
generation: Generation,
reason: Arc<str>,
},
StaleGeneration {
what: Arc<str>,
expected: Generation,
got: Generation,
},
UnknownGeneration {
what: Arc<str>,
got: Generation,
},
AlreadyReady {
what: Arc<str>,
generation: Generation,
},
AlreadyFailed {
what: Arc<str>,
generation: Generation,
},
Closed {
what: Arc<str>,
generation: Generation,
},
DependantsFull {
what: Arc<str>,
cap: usize,
},
InvalidName {
reason: &'static str,
},
}
impl ReadinessError {
#[must_use]
pub fn what(&self) -> &str {
match *self {
Self::NotReady { ref what, .. }
| Self::StaleGeneration { ref what, .. }
| Self::UnknownGeneration { ref what, .. }
| Self::AlreadyReady { ref what, .. }
| Self::AlreadyFailed { ref what, .. }
| Self::Closed { ref what, .. }
| Self::DependantsFull { ref what, .. } => what,
Self::InvalidName { .. } => "",
}
}
#[must_use]
pub const fn released(&self) -> bool {
false
}
#[must_use]
pub const fn is_permanent(&self) -> bool {
!matches!(self, Self::NotReady { .. })
}
#[must_use]
pub fn located(self, at: Arc<str>) -> FlowError {
match self {
Self::NotReady {
what,
generation,
reason,
} => FlowError::Transient {
at,
reason: format!("{what} did not become ready ({generation}): {reason}"),
},
Self::StaleGeneration {
what,
expected,
got,
} => FlowError::Failed {
at,
reason: format!("{what} refused a stale signal: expected {expected}, got {got}"),
},
Self::UnknownGeneration { what, got } => FlowError::Failed {
at,
reason: format!("{what} refused a signal from {got}, which it never issued"),
},
Self::AlreadyReady { what, generation } => FlowError::Failed {
at,
reason: format!(
"{what} is already ready ({generation}); a duplicate signal is refused"
),
},
Self::AlreadyFailed { what, generation } => FlowError::Failed {
at,
reason: format!("{what} already failed ({generation}); arm a new readiness"),
},
Self::Closed { what, generation } => FlowError::Failed {
at,
reason: format!("{what} was shut down ({generation}); a signal is refused"),
},
Self::DependantsFull { what, cap } => FlowError::Failed {
at,
reason: format!("{what} is at its declared cap of {cap} dependants"),
},
Self::InvalidName { reason } => FlowError::Failed {
at,
reason: format!("invalid readiness name: {reason}"),
},
}
}
}
impl fmt::Display for ReadinessError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.clone().located(Arc::from("")).to_string())
}
}
impl std::error::Error for ReadinessError {}
impl From<ReadinessError> for FlowError {
fn from(source: ReadinessError) -> Self {
source.located(Arc::from(""))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ReadyOutcome<T> {
Ready {
generation: Generation,
value: Arc<T>,
},
Failed {
generation: Generation,
reason: Arc<str>,
},
Closed {
generation: Generation,
},
}
impl<T> ReadyOutcome<T> {
#[must_use]
pub const fn generation(&self) -> Generation {
match *self {
Self::Ready { generation, .. }
| Self::Failed { generation, .. }
| Self::Closed { generation } => generation,
}
}
#[must_use]
pub const fn is_ready(&self) -> bool {
matches!(*self, Self::Ready { .. })
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Ready<T> {
generation: Generation,
value: Arc<T>,
}
impl<T> Ready<T> {
#[must_use]
pub const fn generation(&self) -> Generation {
self.generation
}
#[must_use]
pub fn awaited(&self) -> &T {
&self.value
}
#[must_use]
pub fn shared(&self) -> Arc<T> {
Arc::clone(&self.value)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct FailAfterReady {
generation: Generation,
reason: Arc<str>,
}
impl FailAfterReady {
#[must_use]
pub const fn generation(&self) -> Generation {
self.generation
}
#[must_use]
pub fn reason(&self) -> &str {
&self.reason
}
}
#[derive(Debug)]
pub struct Readiness<T> {
inner: Arc<State<T>>,
}
#[derive(Debug)]
struct State<T> {
what: Arc<str>,
cap: usize,
status: Mutex<Status>,
admitted: AtomicU64,
released_once: AtomicBool,
dependants: Mutex<Vec<CancellationToken>>,
settled: watch::Sender<Settled<T>>,
}
#[derive(Debug)]
enum Status {
Pending(Generation),
Released(Generation),
Failed {
generation: Generation,
reason: Arc<str>,
},
Closed(Generation),
}
impl Status {
const fn generation(&self) -> Generation {
match *self {
Self::Pending(generation) | Self::Released(generation) | Self::Closed(generation) => {
generation
}
Self::Failed { generation, .. } => generation,
}
}
}
#[derive(Debug)]
enum Settled<T> {
Waiting,
Ready {
generation: Generation,
value: Arc<T>,
},
Failed {
generation: Generation,
reason: Arc<str>,
},
Closed {
generation: Generation,
},
}
impl<T: Clone> Clone for Settled<T> {
fn clone(&self) -> Self {
match *self {
Self::Waiting => Self::Waiting,
Self::Ready {
generation,
ref value,
} => Self::Ready {
generation,
value: Arc::clone(value),
},
Self::Failed {
generation,
ref reason,
} => Self::Failed {
generation,
reason: Arc::clone(reason),
},
Self::Closed { generation } => Self::Closed { generation },
}
}
}
impl<T> Clone for Readiness<T> {
fn clone(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
}
}
}
impl<T: Clone> Readiness<T> {
pub fn new(what: &str) -> Result<Self, ReadinessError> {
Self::named(what, Generation::FIRST, MAX_DEPENDANTS)
}
pub fn named(what: &str, generation: Generation, cap: usize) -> Result<Self, ReadinessError> {
if what.is_empty() {
let refusal = Err(ReadinessError::InvalidName {
reason: "the name is empty",
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "named: returning an error to the caller");
return refusal;
}
if what.len() > MAX_NAME_BYTES {
let refusal = Err(ReadinessError::InvalidName {
reason: "the name is longer than MAX_NAME_BYTES",
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "named: returning an error to the caller");
return refusal;
}
let (settled, _receiver) = watch::channel(Settled::Waiting);
Ok(Self {
inner: Arc::new(State {
what: Arc::from(what),
cap: cap.max(1),
status: Mutex::new(Status::Pending(generation)),
admitted: AtomicU64::new(0),
released_once: AtomicBool::new(false),
dependants: Mutex::new(Vec::new()),
settled,
}),
})
}
#[must_use]
pub fn what(&self) -> &str {
&self.inner.what
}
#[must_use]
pub fn generation(&self) -> Generation {
self.lock().generation()
}
#[must_use]
pub fn is_settled(&self) -> bool {
!matches!(self.settled_value(), Settled::Waiting)
}
#[must_use]
pub fn outcome(&self) -> Option<ReadyOutcome<T>> {
match self.settled_value() {
Settled::Waiting => None,
Settled::Ready { generation, value } => Some(ReadyOutcome::Ready { generation, value }),
Settled::Failed { generation, reason } => {
Some(ReadyOutcome::Failed { generation, reason })
}
Settled::Closed { generation } => Some(ReadyOutcome::Closed { generation }),
}
}
#[must_use]
pub fn admitted(&self) -> u64 {
self.inner.admitted.load(Ordering::SeqCst)
}
#[must_use]
pub fn cap(&self) -> usize {
self.inner.cap
}
#[must_use]
pub fn dependants_live(&self) -> usize {
crate::journal::owner::lock(&self.inner.dependants).len()
}
pub fn ready(&self, generation: Generation, value: T) -> Result<(), ReadinessError> {
let mut status = self.lock();
let current = status.generation();
if current == generation && matches!(&*status, Status::Pending(_)) {
*status = Status::Released(generation);
self.inner.released_once.store(true, Ordering::SeqCst);
drop(status);
self.inner.settled.send_replace(Settled::Ready {
generation,
value: Arc::new(value),
});
return Ok(());
}
match *status {
Status::Released(released) if released == generation => {
Err(ReadinessError::AlreadyReady {
what: Arc::clone(&self.inner.what),
generation: released,
})
}
Status::Pending(_) => Err(self.mismatch(current, generation, false)),
Status::Released(_) => Err(self.mismatch(current, generation, false)),
Status::Closed(_) => Err(self.mismatch(current, generation, true)),
Status::Failed {
generation: failed, ..
} => Err(ReadinessError::AlreadyFailed {
what: Arc::clone(&self.inner.what),
generation: failed,
}),
}
}
pub fn fail(&self, generation: Generation, reason: &str) -> Result<(), ReadinessError> {
let mut status = self.lock();
let current = status.generation();
let closed = matches!(&*status, Status::Closed(_));
let failed = matches!(&*status, Status::Failed { .. });
if current != generation {
let refusal = Err(self.mismatch(current, generation, closed));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "fail: returning an error to the caller");
return refusal;
}
if failed {
let refusal = Err(ReadinessError::AlreadyFailed {
what: Arc::clone(&self.inner.what),
generation: current,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "fail: returning an error to the caller");
return refusal;
}
let reason: Arc<str> = Arc::from(reason);
let released = matches!(*status, Status::Released(_));
*status = Status::Failed {
generation,
reason: Arc::clone(&reason),
};
drop(status);
if released {
self.cancel_dependants();
}
self.inner
.settled
.send_replace(Settled::Failed { generation, reason });
Ok(())
}
pub fn shutdown(&self, generation: Generation) -> Result<(), ReadinessError> {
let mut status = self.lock();
let current = status.generation();
let closed = matches!(&*status, Status::Closed(_));
if current != generation {
let refusal = Err(self.mismatch(current, generation, closed));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "shutdown: returning an error to the caller");
return refusal;
}
if matches!(*status, Status::Failed { .. }) {
let refusal = Err(ReadinessError::AlreadyFailed {
what: Arc::clone(&self.inner.what),
generation: current,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "shutdown: returning an error to the caller");
return refusal;
}
*status = Status::Closed(generation);
drop(status);
self.cancel_dependants();
self.inner
.settled
.send_replace(Settled::Closed { generation });
Ok(())
}
#[must_use]
pub fn stale(&self, generation: Generation) -> bool {
generation < self.lock().generation()
}
#[must_use]
pub fn released_at(&self) -> Option<Generation> {
let status = self.lock();
match *status {
Status::Released(generation) => Some(generation),
Status::Pending(_) | Status::Failed { .. } | Status::Closed(_) => None,
}
}
#[must_use]
pub fn failed_at(&self) -> Option<FailAfterReady> {
let status = self.lock();
match *status {
Status::Failed {
generation,
ref reason,
} if self.released_before_failure() => Some(FailAfterReady {
generation,
reason: Arc::clone(reason),
}),
Status::Pending(_)
| Status::Released(_)
| Status::Closed(_)
| Status::Failed { .. } => None,
}
}
#[must_use]
fn released_before_failure(&self) -> bool {
self.inner.released_once.load(Ordering::SeqCst)
}
pub async fn wait(&self, scope: &Scope, limit: Duration) -> Result<Ready<T>, FlowError> {
self.admit(scope)?;
let settled = self.settled_value();
let finished = match settled {
Settled::Waiting => None,
other => Some(other),
};
if let Some(settled) = finished {
return self.resolve(settled, scope);
}
let limit = limit.max(Duration::from_millis(1));
within(scope, "await-ready", limit, self.subscribe(scope)).await
}
async fn subscribe(&self, scope: &Scope) -> Result<Ready<T>, FlowError> {
let mut settled = self.inner.settled.subscribe();
loop {
let current = settled.borrow_and_update().clone();
match current {
Settled::Waiting => {}
settled_now => return self.resolve(settled_now, scope),
}
let changed = watch::Receiver::changed(&mut settled);
let cancelled = scope.token().cancelled_owned();
lgwks_deps::tokio::select! {
biased;
_ = changed => {}
() = cancelled => {
return Err(FlowError::Cancelled {
at: scope.join("await-ready"),
});
}
}
}
}
fn resolve(&self, settled: Settled<T>, scope: &Scope) -> Result<Ready<T>, FlowError> {
let at = scope.join("await-ready");
match settled {
Settled::Waiting => Err(FlowError::failed(format!(
"{} settled between a dependant's admission and its first read",
self.inner.what
))),
Settled::Ready { generation, value } => Ok(Ready { generation, value }),
Settled::Failed { generation, reason } => Err(ReadinessError::NotReady {
what: Arc::clone(&self.inner.what),
generation,
reason,
}
.located(at)),
Settled::Closed { generation } => Err(FlowError::Cancelled {
at: Arc::from(format!(
"{}: {} was shut down ({generation})",
scope.join("await-ready"),
self.inner.what
)),
}),
}
}
fn admit(&self, scope: &Scope) -> Result<(), FlowError> {
let cap = self.inner.cap;
let charged = self
.inner
.admitted
.try_update(Ordering::SeqCst, Ordering::SeqCst, |count| {
let within_cap = match usize::try_from(count) {
Ok(admitted) => admitted < cap,
Err(_wider_than_this_host_counts) => false,
};
within_cap.then_some(count.saturating_add(1))
});
if charged.is_err() {
let refusal = Err(ReadinessError::DependantsFull {
what: Arc::clone(&self.inner.what),
cap,
}
.located(Arc::from(scope.path())));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "admit: returning an error to the caller");
return refusal;
}
crate::journal::owner::lock(&self.inner.dependants).push(scope.token().clone());
Ok(())
}
fn cancel_dependants(&self) -> usize {
let mut retained = crate::journal::owner::lock(&self.inner.dependants);
let tokens = std::mem::take(&mut *retained);
drop(retained);
let reached = tokens.len();
for token in tokens {
token.cancel();
}
reached
}
fn settled_value(&self) -> Settled<T> {
self.inner.settled.subscribe().borrow_and_update().clone()
}
fn lock(&self) -> MutexGuard<'_, Status> {
crate::journal::owner::lock(&self.inner.status)
}
fn mismatch(&self, current: Generation, got: Generation, closed: bool) -> ReadinessError {
let what = Arc::clone(&self.inner.what);
if closed {
ReadinessError::Closed {
what,
generation: current,
}
} else if got < current {
ReadinessError::StaleGeneration {
what,
expected: current,
got,
}
} else {
ReadinessError::UnknownGeneration { what, got }
}
}
}