use std::fmt;
use std::sync::Arc;
use lgwks_std::hash::{Digest, Hasher};
use crate::effect::RunId;
use crate::rt::clock::Clock;
use crate::rt::sync::CancellationToken;
use super::policy::Policy;
use super::trail::Trail;
use super::{DEFAULT_TRAIL_STEPS, FlowError, MAX_DEPTH, MAX_TENANT_BYTES};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum NameFault {
Empty,
TooLong,
BadChar,
}
pub(crate) fn name_fault(name: &str, max: usize) -> Option<NameFault> {
if name.is_empty() {
return Some(NameFault::Empty);
}
if name.len() > max {
return Some(NameFault::TooLong);
}
let allowed = |byte: u8| byte.is_ascii_alphanumeric() || b"-_.:".contains(&byte);
if !name.bytes().all(allowed) {
return Some(NameFault::BadChar);
}
None
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Tenant(Arc<str>);
impl Tenant {
pub fn new(name: &str) -> Result<Self, FlowError> {
let reason = match name_fault(name, MAX_TENANT_BYTES) {
None => return Ok(Self(Arc::from(name))),
Some(NameFault::Empty) => "the tenant name is empty",
Some(NameFault::TooLong) => "the tenant name is longer than MAX_TENANT_BYTES",
Some(NameFault::BadChar) => {
"the tenant name may hold only ASCII letters, digits, '-', '_', '.' and ':'"
}
};
Err(FlowError::InvalidTenant { reason })
}
pub(crate) fn assumed(name: &'static str) -> Self {
Self(Arc::from(name))
}
#[must_use]
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for Tenant {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct StepKey(Digest);
impl StepKey {
#[must_use]
pub fn as_bytes(&self) -> &[u8; 32] {
self.0.as_bytes()
}
#[must_use]
pub fn to_hex(&self) -> String {
self.0.to_hex()
}
}
pub(crate) fn step_key(tenant: &str, path: &str) -> StepKey {
let mut hasher = Hasher::new();
hasher
.write_framed(tenant.as_bytes())
.write_framed(path.as_bytes());
StepKey(hasher.finalize())
}
impl fmt::Display for StepKey {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
fmt::Display::fmt(&self.0, formatter)
}
}
#[derive(Clone)]
pub struct Scope {
inner: Arc<ScopeInner>,
}
struct ScopeInner {
tenant: Tenant,
path: Arc<str>,
depth: u16,
token: CancellationToken,
policy: Arc<Policy>,
trail: Arc<Trail>,
clock: Clock,
run: Option<RunId>,
}
impl Scope {
#[must_use]
pub fn root(tenant: Tenant) -> Self {
Self::with_token(tenant, CancellationToken::new())
}
#[must_use]
pub fn with_clock(tenant: Tenant, clock: Clock) -> Self {
Self::with_token_and_clock(tenant, CancellationToken::new(), clock)
}
#[must_use]
pub fn with_token(tenant: Tenant, token: CancellationToken) -> Self {
Self::with_token_and_trail(tenant, token, Trail::new(DEFAULT_TRAIL_STEPS))
}
#[must_use]
pub fn with_token_and_clock(tenant: Tenant, token: CancellationToken, clock: Clock) -> Self {
Self::with_token_trail_and_clock(tenant, token, Trail::new(DEFAULT_TRAIL_STEPS), clock)
}
pub(crate) fn with_token_trail_and_clock(
tenant: Tenant,
token: CancellationToken,
trail: Arc<Trail>,
clock: Clock,
) -> Self {
Self::rooted(tenant, token, trail, clock, None)
}
pub(crate) fn rooted(
tenant: Tenant,
token: CancellationToken,
trail: Arc<Trail>,
clock: Clock,
run: Option<RunId>,
) -> Self {
Self {
inner: Arc::new(ScopeInner {
tenant,
path: Arc::from(""),
depth: 0,
token,
policy: Arc::new(Policy::for_this_machine()),
trail,
clock,
run,
}),
}
}
pub(crate) fn with_token_and_trail(
tenant: Tenant,
token: CancellationToken,
trail: Arc<Trail>,
) -> Self {
Self::with_token_trail_and_clock(tenant, token, trail, Clock::wall())
}
pub fn enter(&self, step: &str) -> Result<Self, FlowError> {
self.descend(self.join(step), Stop::Own)
}
pub(super) fn enter_sharing(&self, step: &str) -> Result<Self, FlowError> {
self.descend(self.join(step), Stop::Shared)
}
pub fn item(&self, step: &str, index: usize) -> Result<Self, FlowError> {
self.descend(self.join(&format!("{step}#{index}")), Stop::Own)
}
pub(super) fn numbered(&self, index: usize) -> Result<Self, FlowError> {
self.descend(
Arc::from(format!("{}#{index}", self.inner.path)),
Stop::Shared,
)
}
fn descend(&self, path: Arc<str>, stop: Stop) -> Result<Self, FlowError> {
self.checkpoint()?;
if self.inner.depth >= MAX_DEPTH {
let refusal = Err(FlowError::TooDeep {
at: Arc::clone(&self.inner.path),
limit: MAX_DEPTH,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "descend: returning an error to the caller");
return refusal;
}
self.inner.trail.record(&path);
Ok(Self {
inner: Arc::new(ScopeInner {
tenant: self.inner.tenant.clone(),
path: Arc::clone(&path),
depth: self.inner.depth.saturating_add(1),
token: match stop {
Stop::Own => self.inner.token.child_token(),
Stop::Shared => self.inner.token.clone(),
},
policy: Arc::clone(&self.inner.policy),
trail: Arc::clone(&self.inner.trail),
clock: self.inner.clock.clone(),
run: self.inner.run,
}),
})
}
pub(super) fn policy(&self) -> &Policy {
&self.inner.policy
}
pub(crate) fn shared_path(&self) -> &Arc<str> {
&self.inner.path
}
pub(super) fn join(&self, segment: &str) -> Arc<str> {
if self.inner.path.is_empty() {
Arc::from(segment)
} else {
Arc::from(format!("{}/{segment}", self.inner.path))
}
}
#[must_use]
pub fn tenant(&self) -> &Tenant {
&self.inner.tenant
}
#[must_use]
pub fn run(&self) -> Option<RunId> {
self.inner.run
}
#[must_use]
pub fn path(&self) -> &str {
&self.inner.path
}
#[must_use]
pub fn key(&self) -> StepKey {
step_key(self.inner.tenant.as_str(), &self.inner.path)
}
#[must_use]
pub fn token(&self) -> &CancellationToken {
&self.inner.token
}
#[must_use]
pub fn clock(&self) -> &Clock {
&self.inner.clock
}
pub fn cancel(&self) {
self.inner.token.cancel();
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.inner.token.is_cancelled()
}
pub fn checkpoint(&self) -> Result<(), FlowError> {
if self.is_cancelled() {
let refusal = Err(FlowError::Cancelled {
at: Arc::clone(&self.inner.path),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "checkpoint: returning an error to the caller");
return refusal;
}
Ok(())
}
pub fn require(&self, caps: &[crate::cap::Cap]) -> Result<(), FlowError> {
self.checkpoint()?;
let Some(shortfall) =
super::run_store::authority().and_then(|authority| authority.shortfall(caps))
else {
return Ok(());
};
let Some(deficit) = crate::cap::Deficit::from_shortages(shortfall) else {
return Ok(());
};
Err(FlowError::Blocked {
at: Arc::clone(&self.inner.path),
deficit: Box::new(deficit),
})
}
}
#[derive(Clone, Copy)]
enum Stop {
Own,
Shared,
}
impl fmt::Debug for Scope {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Scope")
.field("tenant", &self.inner.tenant.as_str())
.field("path", &self.path())
.field("cancelled", &self.is_cancelled())
.finish()
}
}