use std::fmt;
use std::future::Future;
use std::io;
use std::num::NonZeroUsize;
use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, MutexGuard, OnceLock};
use std::time::{Duration, Instant};
use lgwks_std::hash::{Digest, Hasher};
use crate::cap::{Cap, Deficit, Demand, Shortage, uncovered};
use crate::effect::{InputIdentity, RunId};
use crate::gate::GrantSet;
use crate::journal::frame::SaturatingFrom;
use crate::journal::owner::lock;
use crate::rt::clock::Clock;
use crate::rt::runtime::{Handle, Runtime};
use crate::rt::sync::{CancellationToken, OwnedSemaphorePermit, Semaphore};
use crate::rt::task_local;
use crate::script::run_store::{
Authority as RecordsAuthority, AuthorityCheck, Records, StoredValue,
};
use crate::script::scope::step_key as reserved_step_key;
use crate::script::trail::Trail;
use crate::script::{
Appended, DEFAULT_TRAIL_STEPS, Durable, FlowError, MAX_IN_FLIGHT, Scope, Tenant, within,
};
use self::name::TaskName;
use self::request::{
BINDING_STEP, TERMINAL_STEP, TerminalRecord, derive_run as derive_request_run, digest_of_record,
};
mod definition;
mod ledger;
mod name;
pub mod repair;
mod request;
mod store;
pub use definition::{DefinitionIdentity, Drift, UNVERSIONED_CODEC};
pub use ledger::{Control, LeaseRefusal, RunLedger};
pub use repair::{MAX_TICKET_NEEDS, RepairError, RepairTicket};
pub use request::{
InFlight, InputDigest, MAX_REQUEST_KEY_BYTES, RequestConflict, RequestError, RequestKey,
Submission,
};
pub use store::{
MAX_RECORD_BYTES, MAX_RECORDS_PER_RUN, MAX_STORE_BYTES, RunStore, StoreError, StoreLimitKind,
};
pub const MAX_ADMITTED_TASKS: usize = MAX_IN_FLIGHT;
pub const MAX_TASK_DEADLINE: Duration = Duration::from_secs(24 * 60 * 60);
pub const MAX_PROGRESS_STEPS: usize = MAX_IN_FLIGHT;
const BODY_STEP: &str = "body";
static HOST_IDENTITIES: AtomicU64 = AtomicU64::new(1);
task_local! {
static HELD_PERMITS: Vec<u64>;
}
type HostIdentity = u64;
fn own_fault() -> &'static FlowError {
OWN_FAULT
.get_or_init(|| FlowError::failed("a report with neither an output nor a located failure"))
}
static OWN_FAULT: OnceLock<FlowError> = OnceLock::new();
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum Disposition {
Succeeded,
Failed,
Cancelled,
DeadlineExceeded,
Refused,
Blocked,
}
impl Disposition {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Succeeded => "Succeeded",
Self::Failed => "Failed",
Self::Cancelled => "Cancelled",
Self::DeadlineExceeded => "DeadlineExceeded",
Self::Refused => "Refused",
Self::Blocked => "Blocked",
}
}
#[must_use]
pub const fn is_success(self) -> bool {
matches!(self, Self::Succeeded)
}
#[must_use]
pub const fn was_admitted(self) -> bool {
!matches!(self, Self::Refused)
}
}
fn terminal_for<O: Durable>(report: &Report<O>) -> Result<Option<TerminalRecord>, RequestError> {
Ok(match report.disposition() {
Disposition::Succeeded => {
let output = report
.output()
.ok_or(RequestError::Record(FlowError::failed(
"a successful run reported no output to record",
)))?;
Some(TerminalRecord::succeeded(output)?)
}
Disposition::Failed | Disposition::DeadlineExceeded => Some(TerminalRecord::stopped(
report.disposition(),
report.error().map_or("", FlowError::at),
report.error().map_or_else(String::new, ToString::to_string),
)),
Disposition::Cancelled | Disposition::Refused | Disposition::Blocked => None,
})
}
impl fmt::Display for Disposition {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.label())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum EffectKnowledge {
None,
StepRecords {
records: usize,
},
}
pub struct Task<F> {
name: TaskName,
body: F,
needs: Vec<Cap>,
}
impl<F> Task<F> {
#[must_use]
pub fn name(&self) -> &str {
self.name.as_str()
}
#[must_use]
pub fn needs(&self) -> &[Cap] {
&self.needs
}
#[must_use]
pub fn requiring(mut self, caps: &[Cap]) -> Self {
self.needs = caps.to_vec();
self
}
pub(crate) fn name_owned(&self) -> TaskName {
self.name.clone()
}
}
impl<F: fmt::Debug> fmt::Debug for Task<F> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Task")
.field("name", &self.name.as_str())
.field("needs", &self.needs)
.field("body", &self.body)
.finish()
}
}
pub fn task<F>(name: &str, body: F) -> Result<Task<F>, FlowError> {
Ok(Task {
name: TaskName::new(name)?,
body,
needs: Vec::new(),
})
}
#[derive(Debug)]
pub struct Report<O> {
disposition: Disposition,
output: Option<O>,
error: Option<FlowError>,
elapsed: Duration,
steps: Vec<Arc<str>>,
dropped_steps: usize,
progress_capacity: usize,
effects: EffectKnowledge,
tenant: Tenant,
task: TaskName,
run: Option<RunId>,
ticket: Option<Ticket>,
needs: Option<Deficit>,
repair: Option<RepairTicket>,
}
impl<O> Report<O> {
#[must_use]
pub const fn disposition(&self) -> Disposition {
self.disposition
}
#[must_use]
pub fn output(&self) -> Option<&O> {
self.output.as_ref()
}
#[must_use]
pub const fn error(&self) -> Option<&FlowError> {
self.error.as_ref()
}
pub fn result(&self) -> Result<&O, &FlowError> {
match (self.output.as_ref(), self.error.as_ref()) {
(Some(output), _) if self.disposition.is_success() => Ok(output),
(_, Some(error)) => Err(error),
_ => Err(own_fault()),
}
}
pub fn into_result(self) -> Result<O, FlowError> {
let (output, error) = (self.output, self.error);
match (output, error) {
(Some(output), _) if self.disposition.is_success() => Ok(output),
(_, Some(error)) => Err(error),
_ => Err(FlowError::failed(
"a report with neither an output nor a located failure",
)),
}
}
#[must_use]
pub fn into_output(self) -> Option<O> {
self.output
}
#[must_use]
pub const fn elapsed(&self) -> Duration {
self.elapsed
}
#[must_use]
pub fn steps(&self) -> &[Arc<str>] {
&self.steps
}
#[must_use]
pub const fn dropped_steps(&self) -> usize {
self.dropped_steps
}
#[must_use]
pub const fn progress_capacity(&self) -> usize {
self.progress_capacity
}
#[must_use]
pub const fn effects(&self) -> EffectKnowledge {
self.effects
}
#[must_use]
pub fn tenant(&self) -> &Tenant {
&self.tenant
}
#[must_use]
pub fn task_name(&self) -> &str {
self.task.as_str()
}
#[must_use]
pub fn ended_at(&self) -> Option<&str> {
self.error
.as_ref()
.map(FlowError::at)
.filter(|path| !path.is_empty())
}
#[must_use]
pub const fn run_id(&self) -> Option<RunId> {
self.run
}
#[must_use]
pub const fn ticket(&self) -> Option<&Ticket> {
self.ticket.as_ref()
}
#[must_use]
pub const fn needs(&self) -> Option<&Deficit> {
self.needs.as_ref()
}
#[must_use]
pub const fn repair(&self) -> Option<&RepairTicket> {
self.repair.as_ref()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct Ticket {
run: RunId,
tenant: Tenant,
task: TaskName,
disposition: Disposition,
at: Option<Arc<str>>,
}
impl Ticket {
#[must_use]
pub const fn run(&self) -> RunId {
self.run
}
#[must_use]
pub fn tenant(&self) -> &Tenant {
&self.tenant
}
#[must_use]
pub fn task(&self) -> &str {
self.task.as_str()
}
#[must_use]
pub const fn disposition(&self) -> Disposition {
self.disposition
}
#[must_use]
pub fn stopped_at(&self) -> Option<&str> {
self.at.as_deref()
}
}
impl fmt::Display for Ticket {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
formatter,
"{}/{}: {}",
self.tenant.as_str(),
self.run.id().to_hex(),
self.disposition.label()
)?;
if let Some(at) = self.at.as_deref() {
write!(formatter, " at {at}")?;
}
Ok(())
}
}
fn recorded_receipt(
records: &Records,
tenant: &str,
run: RunId,
binding_key: crate::script::StepKey,
) -> Result<StoredValue, RequestError> {
match records.lookup(tenant, run, binding_key)? {
Some(held) => Ok(held),
None => {
let refusal = Err(RequestError::IdCollision);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "recorded_receipt: the claim reported a record the lookup cannot find");
refusal
}
}
}
struct Existing<'a> {
store: &'a RunStore,
records: &'a Records,
tenant: &'a str,
run: RunId,
key: &'a RequestKey,
held: &'a StoredValue,
digest: InputDigest,
}
#[derive(Clone)]
pub struct Host {
inner: Arc<Installation>,
}
struct Installation {
identity: HostIdentity,
tenant: Tenant,
token: CancellationToken,
limits: HostLimits,
admission: Arc<Semaphore>,
in_flight: LiveRuns,
peak_in_flight: HighWater,
admitted: AtomicU64,
refused: AtomicU64,
clock: Clock,
store: Option<store::RunStore>,
codec: String,
ledger: Option<RunLedger>,
grants: GrantSet,
repair_attempts: u64,
repair_spend: u64,
reactor: Mutex<Option<Runtime>>,
}
#[derive(Debug, Clone)]
enum Plan {
Fresh,
Declared(DefinitionIdentity),
Resume(RunId),
ResumeUnder(RunId, DefinitionIdentity),
}
impl Plan {
fn declared(&self) -> Option<&DefinitionIdentity> {
match self {
&Self::Declared(ref identity) | &Self::ResumeUnder(_, ref identity) => Some(identity),
&Self::Fresh | &Self::Resume(_) => None,
}
}
}
impl Plan {
const fn run(&self) -> Option<RunId> {
match *self {
Self::Fresh | Self::Declared(_) => None,
Self::Resume(run) | Self::ResumeUnder(run, _) => Some(run),
}
}
}
#[cfg(test)]
mod tests {
use super::Plan;
use crate::effect::RunId;
use crate::task::{Host, name::TaskName};
type TestResult = Result<(), String>;
fn run_id(nibble: char) -> Result<RunId, String> {
let literal = format!("{nibble}{}", "0".repeat(31));
RunId::from_hex(&literal).map_err(|error| error.to_string())
}
#[test]
fn only_the_resume_arms_adopt_a_run() -> TestResult {
let host = host()?;
let task = TaskName::new("probe").map_err(|error| error.to_string())?;
let run = run_id('1')?;
let identity = host.definition_for(task.as_str(), None);
assert!(Plan::Fresh.run().is_none(), "a fresh plan adopts nothing");
assert!(
Plan::Declared(identity.clone()).run().is_none(),
"a declared fresh plan adopts nothing either"
);
assert_eq!(
Plan::Resume(run).run(),
Some(run),
"a resume adopts its run"
);
assert_eq!(
Plan::ResumeUnder(run, identity).run(),
Some(run),
"a resume under a declared identity adopts the same run"
);
Ok(())
}
#[test]
fn a_derived_identity_is_stable_per_run_and_distinct_across_runs() -> TestResult {
let host = host()?;
let first = run_id('1')?;
let second = run_id('2')?;
assert_eq!(
host.definition_for("probe", Some(&first)).digest(),
host.definition_for("probe", Some(&first)).digest(),
"one run's derived identity must not move between two readings"
);
assert_ne!(
host.definition_for("probe", Some(&first)).digest(),
host.definition_for("probe", Some(&second)).digest(),
"two runs must not share a derived identity"
);
assert_ne!(
host.definition_for("probe", None).digest(),
host.definition_for("probe", Some(&first)).digest(),
"a run that has not adopted an id is not the run that has"
);
Ok(())
}
fn host() -> Result<Host, String> {
Host::builder("acme")
.map_err(|error| error.to_string())?
.build()
.map_err(|error| error.to_string())
}
}
impl Host {
pub fn builder(tenant: &str) -> Result<HostBuilder, FlowError> {
Ok(HostBuilder {
tenant: Tenant::new(tenant)?,
max_concurrent: default_max_concurrent(),
deadline: default_deadline(),
progress: DEFAULT_TRAIL_STEPS,
clock: Clock::wall(),
store: None,
codec: UNVERSIONED_CODEC.to_owned(),
ledger: None,
grants: GrantSet::empty(),
repair_attempts: default_repair_attempts(),
repair_spend: default_repair_spend(),
})
}
#[must_use]
pub fn tenant(&self) -> &Tenant {
&self.inner.tenant
}
#[must_use]
pub fn limits(&self) -> &HostLimits {
&self.inner.limits
}
#[must_use]
pub fn clock(&self) -> &Clock {
&self.inner.clock
}
#[must_use]
pub fn admission(&self) -> Admission<'_> {
Admission { inner: &self.inner }
}
#[must_use]
pub fn token(&self) -> &CancellationToken {
&self.inner.token
}
#[must_use]
pub fn run_store(&self) -> Option<&store::RunStore> {
self.inner.store.as_ref()
}
fn identity_for(&self, run: RunId, declared: &DefinitionIdentity) -> DefinitionIdentity {
match self
.inner
.store
.as_ref()
.and_then(|store| store.definition_of(run))
{
Some(recorded) => recorded,
None => declared.clone(),
}
}
fn definition_for(&self, task_name: &str, run: Option<&RunId>) -> DefinitionIdentity {
let mut hasher = Hasher::new();
hasher.write_framed(b"lgwks.bot.definition.host-default");
hasher.write_framed(self.inner.tenant.as_str().as_bytes());
hasher.write_framed(task_name.as_bytes());
match run {
Some(run) => hasher.write_framed(run.id().to_hex().as_bytes()),
None => hasher.write_framed(b"fresh"),
};
DefinitionIdentity::new(task_name, 0, hasher.finalize(), 1).with_codec(&self.inner.codec)
}
#[must_use]
pub fn definition(
&self,
task_name: &str,
definition_revision: u64,
input_digest: Option<Digest>,
steps: usize,
) -> DefinitionIdentity {
let input = match input_digest {
Some(declared) => declared,
None => {
let mut hasher = Hasher::new();
hasher.write_framed(b"lgwks.bot.input.undeclared");
hasher.write_framed(self.inner.codec.as_bytes());
hasher.finalize()
}
};
DefinitionIdentity::new(task_name, definition_revision, input, steps)
.with_codec(&self.inner.codec)
}
#[must_use]
pub fn run_ledger(&self) -> Option<&RunLedger> {
self.inner.ledger.as_ref()
}
#[must_use]
pub fn grants(&self) -> &GrantSet {
&self.inner.grants
}
#[must_use]
pub fn repair_bounds(&self) -> (u64, u64) {
(self.inner.repair_attempts, self.inner.repair_spend)
}
pub fn cancel(&self) {
self.inner.token.cancel();
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.inner.token.is_cancelled()
}
pub async fn run<I, O, F, Fut>(&self, task: &Task<F>, input: I) -> Report<O>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
self.run_plan(task, input, Plan::Fresh).await
}
async fn run_plan<I, O, F, Fut>(&self, task: &Task<F>, input: I, plan: Plan) -> Report<O>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
self.execute(task, input, plan, None, &GrantSet::empty(), 1)
.await
}
pub async fn run_under<I, O, F, Fut>(
&self,
definition: &DefinitionIdentity,
task: &Task<F>,
input: I,
) -> Report<O>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
self.run_plan(task, input, Plan::Declared(definition.clone()))
.await
}
pub async fn resume<I, O, F, Fut>(&self, run: RunId, task: &Task<F>, input: I) -> Report<O>
where
O: Durable,
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
let started = Instant::now();
let report = self.run_plan(task, input, Plan::Resume(run)).await;
self.settle_request(run, task.name_owned(), started, report)
.await
}
pub async fn resume_under<I, O, F, Fut>(
&self,
run: RunId,
definition: &DefinitionIdentity,
task: &Task<F>,
input: I,
) -> Report<O>
where
O: Durable,
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
let started = Instant::now();
let report = self
.run_plan(task, input, Plan::ResumeUnder(run, definition.clone()))
.await;
self.settle_request(run, task.name_owned(), started, report)
.await
}
pub async fn resume_ticket<I, O, F, Fut>(
&self,
ticket: &Ticket,
task: &Task<F>,
input: I,
) -> Result<Report<O>, FlowError>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
if ticket.tenant() != &self.inner.tenant {
let refusal = Err(FlowError::Failed {
at: Arc::from(ticket.task()),
reason: format!(
"the ticket names tenant {:?}, not {:?}; refusing to resume another \
tenant's run",
ticket.tenant().as_str(),
self.inner.tenant.as_str()
),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "resume_ticket: returning an error to the caller");
return refusal;
}
Ok(self.run_plan(task, input, Plan::Resume(ticket.run())).await)
}
pub async fn submit<I, O, F, Fut>(
&self,
key: &RequestKey,
task: &Task<F>,
input: I,
) -> Result<Submission<O>, RequestError>
where
I: InputIdentity,
O: Durable,
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
let Some(store) = self.inner.store.as_ref() else {
let refusal = Err(RequestError::NoStore);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "submit: returning an error to the caller");
return refusal;
};
let tenant = self.inner.tenant.as_str();
let run = derive_request_run(tenant, key)?;
let records = self.records(Some(run)).ok_or(RequestError::NoStore)?;
let digest = InputDigest::of(&input);
let binding_key = reserved_step_key(tenant, BINDING_STEP);
let definition = self.definition_for(task.name(), Some(&run));
if let Some(held) = records.lookup(tenant, run, binding_key)? {
let existing = Existing {
store,
records: &records,
tenant,
run,
key,
held: &held,
digest,
};
return self.observe_existing(existing, task);
}
let claimed = records
.claim(
tenant,
run,
binding_key,
BINDING_STEP,
&definition,
digest.as_bytes().to_vec(),
)
.await?;
match claimed {
Appended::Recorded => {}
Appended::AlreadyRecorded => {
let held = recorded_receipt(&records, tenant, run, binding_key)?;
let existing = Existing {
store,
records: &records,
tenant,
run,
key,
held: &held,
digest,
};
return self.observe_existing(existing, task);
}
Appended::Conflicting => {
let held = recorded_receipt(&records, tenant, run, binding_key)?;
let existing = digest_of_record(held.bytes())?;
let refusal = Err(RequestError::Conflict(RequestConflict::new(
key.clone(),
existing,
digest,
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "submit: returning an error to the caller");
return refusal;
}
}
let report = self
.execute(task, input, Plan::Resume(run), None, &GrantSet::empty(), 1)
.await;
if let Some(terminal) = terminal_for(&report)? {
records
.record(
tenant,
run,
reserved_step_key(tenant, TERMINAL_STEP),
TERMINAL_STEP,
&definition,
terminal.encode()?,
)
.await?;
}
Ok(Submission::Executed(report))
}
fn observe_existing<O, F>(
&self,
existing: Existing<'_>,
task: &Task<F>,
) -> Result<Submission<O>, RequestError>
where
O: Durable,
{
let Existing {
store,
records,
tenant,
run,
key,
held,
digest,
} = existing;
let recorded = digest_of_record(held.bytes())?;
if recorded != digest {
let refusal = Err(RequestError::Conflict(RequestConflict::new(
key.clone(),
recorded,
digest,
)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "observe_existing: returning an error to the caller");
return refusal;
}
let terminal_key = reserved_step_key(tenant, TERMINAL_STEP);
if let Some(head) = records.lookup(tenant, run, terminal_key)? {
let (disposition, output, error) =
TerminalRecord::decode(head.bytes())?.parts::<O>()?;
let report = self.report(Terminal {
started: Instant::now(),
task: task.name_owned(),
disposition,
output,
error,
trail: TrailSnapshot::Empty,
run: Some(run),
needs: None,
repair: None,
});
return Ok(Submission::Reattached(report));
}
Ok(Submission::InFlight(InFlight::new(
run,
task.name().to_owned(),
store.record_count(run),
)))
}
async fn settle_request<O>(
&self,
run: RunId,
task: TaskName,
started: Instant,
report: Report<O>,
) -> Report<O>
where
O: Durable,
{
if matches!(
report.disposition(),
Disposition::Refused | Disposition::Cancelled
) {
return report;
}
let Some(records) = self.records(Some(run)) else {
return report;
};
let tenant = self.inner.tenant.as_str();
match self.unsettled(&records, tenant, run).await {
Ok(false) => report,
Err(error) => self.report(Self::unrecorded(task, started, run, error)),
Ok(true) => {
self.settle(records, tenant, run, task, started, report)
.await
}
}
}
async fn unsettled(
&self,
records: &Records,
tenant: &str,
run: RunId,
) -> Result<bool, FlowError> {
if records
.lookup(tenant, run, reserved_step_key(tenant, BINDING_STEP))?
.is_none()
{
return Ok(false);
}
Ok(records
.lookup(tenant, run, reserved_step_key(tenant, TERMINAL_STEP))?
.is_none())
}
async fn settle<O>(
&self,
records: Records,
tenant: &str,
run: RunId,
task: TaskName,
started: Instant,
report: Report<O>,
) -> Report<O>
where
O: Durable,
{
let Some(terminal) = terminal_for(&report).ok().flatten() else {
return report;
};
let definition = self.definition_for(task.as_str(), Some(&run));
let bytes = match terminal.encode() {
Ok(bytes) => bytes,
Err(error) => return self.report(Self::unrecorded(task, started, run, error)),
};
let stored = records
.record(
tenant,
run,
reserved_step_key(tenant, TERMINAL_STEP),
TERMINAL_STEP,
&definition,
bytes,
)
.await;
match stored {
Ok(()) => report,
Err(error) => self.report(Self::unrecorded(task, started, run, error)),
}
}
fn unrecorded<O>(
task: TaskName,
started: Instant,
run: RunId,
cause: FlowError,
) -> Terminal<O> {
let at = Arc::from(task.as_str());
let error = FlowError::Failed {
at,
reason: format!("the request's outcome was not recorded: {cause}"),
};
Terminal {
started,
task,
disposition: disposition_of(&error),
output: None,
error: Some(error),
trail: TrailSnapshot::Empty,
run: Some(run),
needs: None,
repair: None,
}
}
pub async fn repair<I, O, F, Fut>(
&self,
ticket: &RepairTicket,
grant: &GrantSet,
task: &Task<F>,
input: I,
spend: u64,
) -> Result<Report<O>, RepairError>
where
O: Durable,
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
let started = Instant::now();
let ledger = self.inner.ledger.as_ref().ok_or(RepairError::NoLedger)?;
if ticket.tenant() != self.inner.tenant.as_str() {
let refusal = Err(RepairError::ForeignTenant {
ticket: ticket.tenant().to_owned(),
host: self.inner.tenant.as_str().to_owned(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "repair: returning an error to the caller");
return refusal;
}
ticket.check_grant(grant)?;
let control = ledger
.control(ticket.run())
.ok_or_else(|| RepairError::UnknownRun {
run: ticket.run().id().to_hex(),
})?;
if control.epoch() != ticket.epoch() {
let refusal = Err(RepairError::StaleEpoch {
current: control.epoch(),
offered: ticket.epoch(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "repair: returning an error to the caller");
return refusal;
}
let delta = ticket.delta();
let report = self
.execute(
task,
input,
Plan::Resume(ticket.run()),
Some(ticket),
&delta,
spend,
)
.await;
Ok(self
.settle_request(ticket.run(), task.name_owned(), started, report)
.await)
}
async fn execute<I, O, F, Fut>(
&self,
task: &Task<F>,
input: I,
plan: Plan,
repair: Option<&RepairTicket>,
grant_delta: &GrantSet,
spend: u64,
) -> Report<O>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
let started = Instant::now();
let task_name = task.name_owned();
let resume = plan.run();
let declared = plan.declared().cloned();
let definition = match declared.clone() {
Some(identity) => identity,
None => self.definition_for(task_name.as_str(), resume.as_ref()),
};
if let Some(run) = resume
&& let Some(store) = self.inner.store.as_ref()
&& let Some(owner) = store.tenant_of(run)
&& owner != self.inner.tenant.as_str()
{
let error = StoreError::ForeignTenant {
owner,
asked: self.inner.tenant.as_str().to_owned(),
};
return self.refuse(started, task_name, error.to_string());
}
let recorded_definition = match resume {
Some(adopted) => self.identity_for(adopted, &definition),
None => definition.clone(),
};
if let (Some(run), Some(store), Some(declared)) =
(resume, self.inner.store.as_ref(), declared.as_ref())
&& let Err(error) = store.check_definition(run, declared)
{
let located = match error {
StoreError::Incompatible { ref drift, .. } => {
FlowError::incompatible(task_name.as_str(), drift.clone())
}
other => FlowError::Store {
at: Arc::from(task_name.as_str()),
source: Box::new(other),
},
};
self.inner.refused.fetch_add(1, Ordering::Relaxed);
return self.report(Terminal {
started,
task: task_name,
disposition: Disposition::Refused,
output: None,
error: Some(located),
trail: TrailSnapshot::Empty,
run: Some(run),
needs: None,
repair: None,
});
}
if let Some(run) = resume
&& self.inner.store.is_none()
{
return self.refuse(
started,
task_name,
format!(
"this host has no run store, so run {} has nothing to resume from",
run.id().to_hex()
),
);
}
let run = self.adopt_run(resume);
if run.is_none() && self.inner.store.is_some() {
return self.refuse(
started,
task_name,
String::from(
"this host has a run store but could not mint a run identity \
(the `ephemeral` feature is off or its entropy source failed); \
no step ran unrecorded",
),
);
}
let covers = |cap: &Cap| self.inner.grants.grants(cap) || grant_delta.grants(cap);
let shortfall = Deficit::from_shortages(uncovered(
task.needs(),
covers,
Some(&Demand::new(task_name.as_str())),
));
if let Some(deficit) = shortfall {
return self.blocked(started, task_name, run, deficit);
}
let permit = match self.admit(started, &task_name).await {
Ok(permit) => permit,
Err(failure) => {
let disposition = failure.disposition();
if disposition == Disposition::Refused {
self.inner.refused.fetch_add(1, Ordering::Relaxed);
}
return self.report(Terminal {
started,
task: task_name,
disposition,
output: None,
error: Some(failure.into_error()),
trail: TrailSnapshot::Empty,
run,
needs: None,
repair: None,
});
}
};
let _admission = self.charge(permit);
if let Some(run) = run
&& let Some(ledger) = self.inner.ledger.as_ref()
&& let Err(refusal) = ledger
.charge(
self.inner.tenant.as_str(),
run,
repair.map(RepairTicket::stamp).as_ref(),
spend,
self.inner.repair_attempts,
self.inner.repair_spend,
)
.await
{
return self.report(Terminal {
started,
task: task_name,
disposition: Disposition::Refused,
output: None,
error: Some(FlowError::failed(format_args!(
"the run's root budget refused this attempt: {refusal}"
))),
trail: TrailSnapshot::Empty,
run: Some(run),
needs: None,
repair: None,
});
}
let trail = Trail::new(self.inner.limits.progress_capacity());
let root = Scope::rooted(
self.inner.tenant.clone(),
self.inner.token.child_token(),
Arc::clone(&trail),
self.inner.clock.clone(),
run,
);
let scope = match root.enter(task_name.as_str()) {
Ok(scope) => scope,
Err(error) => {
return self.report(Terminal {
started,
task: task_name,
disposition: disposition_of(&error),
output: None,
error: Some(error),
trail: TrailSnapshot::Empty,
run,
needs: None,
repair: None,
});
}
};
let held = match HELD_PERMITS.try_with(Clone::clone) {
Ok(held) => held,
Err(not_entered) => {
lgwks_std::trace::debug!(
?not_entered,
"submit: no enclosing scope holds a permit set"
);
Vec::new()
}
};
let charged = if held.contains(&self.inner.identity) {
held
} else {
let mut charged = held;
charged.push(self.inner.identity);
charged
};
let body = (task.body)(scope.clone(), input);
let deadline = self.inner.limits.default_deadline();
let outcome = Box::pin(crate::script::run_store::with_authority(
Some(self.authority(grant_delta)),
crate::script::run_store::within(
self.records(run),
&recorded_definition,
HELD_PERMITS.scope(charged, within(&scope, BODY_STEP, deadline, body)),
),
))
.await;
let snapshot = trail.snapshot();
let (disposition, output, error, needs) = match outcome {
Ok(value) => (Disposition::Succeeded, Some(value), None, None),
Err(error) => {
let located = error.located(&scope);
let deficit = located.deficit().cloned();
match deficit {
Some(deficit) => (Disposition::Blocked, None, Some(located), Some(deficit)),
None => (disposition_of(&located), None, Some(located), None),
}
}
};
let repair = match (
needs.as_ref(),
run,
self.inner.ledger.as_ref(),
repair.is_some(),
) {
(Some(deficit), Some(run), Some(ledger), false) => {
let epoch = ledger.control(run).map_or(0, |control| control.epoch());
Some(RepairTicket::of_deficit(
run,
self.inner.tenant.as_str(),
epoch,
deficit,
))
}
_ => None,
};
self.report(Terminal {
started,
task: task_name,
disposition,
output,
error,
trail: TrailSnapshot::Taken(snapshot),
run,
needs,
repair,
})
}
fn not_admitted<O>(&self, started: Instant, task: TaskName, declined: Declined) -> Report<O> {
let Declined {
disposition,
reason,
run,
needs,
repair,
} = declined;
let at = Arc::from(task.as_str());
self.report(Terminal {
started,
task,
disposition,
output: None,
error: Some(FlowError::Failed { at, reason }),
trail: TrailSnapshot::Empty,
run,
needs,
repair,
})
}
fn refuse<O>(&self, started: Instant, task: TaskName, reason: String) -> Report<O> {
self.inner.refused.fetch_add(1, Ordering::Relaxed);
self.not_admitted(started, task, Declined::plain(Disposition::Refused, reason))
}
fn blocked<O>(
&self,
started: Instant,
task: TaskName,
run: Option<RunId>,
deficit: Deficit,
) -> Report<O> {
let repair = match (run, self.inner.ledger.as_ref()) {
(Some(run), Some(_ledger)) => {
let epoch = self
.inner
.ledger
.as_ref()
.and_then(|ledger| ledger.control(run))
.map_or(0, |control| control.epoch());
Some(RepairTicket::of_deficit(
run,
self.inner.tenant.as_str(),
epoch,
&deficit,
))
}
_ => None,
};
self.not_admitted(
started,
task,
Declined {
disposition: Disposition::Blocked,
reason: String::from("this host's authority does not cover what this task needs"),
run,
needs: Some(deficit),
repair,
},
)
}
fn adopt_run(&self, resume: Option<RunId>) -> Option<RunId> {
match (self.inner.store.is_some(), resume) {
(true, Some(run)) => Some(run),
(true, None) => self.mint_run(),
(false, _) => None,
}
}
#[cfg(feature = "ephemeral")]
fn mint_run(&self) -> Option<RunId> {
RunId::mint().ok()
}
#[cfg(not(feature = "ephemeral"))]
fn mint_run(&self) -> Option<RunId> {
None
}
fn authority(&self, grant_delta: &GrantSet) -> RecordsAuthority {
RecordsAuthority(Arc::new(RunAuthority {
base: self.inner.grants.clone(),
delta: grant_delta.clone(),
}))
}
fn records(&self, run: Option<RunId>) -> Option<Records> {
let store = self.inner.store.as_ref()?;
run?;
Some(Records(Arc::new(store.clone())))
}
pub fn block_on<I, O, F, Fut>(&self, task: &Task<F>, input: I) -> Result<Report<O>, HostError>
where
F: Fn(Scope, I) -> Fut,
Fut: Future<Output = Result<O, FlowError>>,
{
if lgwks_deps::tokio::runtime::Handle::try_current().is_ok() {
let refusal = Err(HostError::InsideRuntime);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "block_on: returning an error to the caller");
return refusal;
}
let handle = self.reactor()?;
Ok(handle.block_on(self.run(task, input)))
}
fn reactor(&self) -> Result<Handle, HostError> {
let mut slot = self.reactor_slot();
if let Some(runtime) = slot.as_ref() {
return Ok(runtime.handle());
}
let runtime = Runtime::new().map_err(|cause| HostError::Runtime { cause })?;
let handle = runtime.handle();
*slot = Some(runtime);
Ok(handle)
}
fn reactor_slot(&self) -> MutexGuard<'_, Option<Runtime>> {
lock(&self.inner.reactor)
}
async fn admit(&self, started: Instant, task: &TaskName) -> Result<Permit, AdmissionFailure> {
let nested = matches!(
HELD_PERMITS.try_with(|held| held.contains(&self.inner.identity)),
Ok(true)
);
if nested {
return Ok(Permit::Charged);
}
if self.inner.token.is_cancelled() {
let refusal = Err(AdmissionFailure::Refused(Refused {
at: Arc::from(task.as_str()),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "admit: returning an error to the caller");
return refusal;
}
let deadline = self.inner.limits.default_deadline();
let remaining = deadline.saturating_sub(started.elapsed());
let waiting = self
.inner
.token
.run_until_cancelled(Semaphore::acquire_owned(Arc::clone(&self.inner.admission)));
match crate::rt::time::timeout(remaining, waiting).await {
Ok(Some(Ok(permit))) if !self.inner.token.is_cancelled() => Ok(Permit::Held(permit)),
Ok(Some(Ok(_after_the_stop))) => Err(AdmissionFailure::Refused(Refused {
at: Arc::from(task.as_str()),
})),
Ok(Some(Err(_closed))) => Err(AdmissionFailure::Failed(FlowError::failed(
"the host's admission budget is closed",
))),
Ok(None) => Err(AdmissionFailure::Refused(Refused {
at: Arc::from(task.as_str()),
})),
Err(_elapsed) => Err(AdmissionFailure::DeadlineExceeded(FlowError::TimedOut {
at: Arc::from(BODY_STEP),
after: deadline,
})),
}
}
fn report<O>(&self, terminal: Terminal<O>) -> Report<O> {
let Terminal {
started,
task,
disposition,
output,
error,
trail,
run,
needs,
repair,
} = terminal;
let (steps, dropped_steps) = match trail {
TrailSnapshot::Taken((paths, dropped)) => (paths, dropped),
TrailSnapshot::Empty => (Vec::new(), 0),
};
let effects = match (run, self.inner.store.as_ref()) {
(Some(run), Some(store)) => EffectKnowledge::StepRecords {
records: store.record_count(run),
},
_ => EffectKnowledge::None,
};
let ticket = run.and_then(|run| {
if disposition.is_success() {
return None;
}
Some(Ticket {
run,
tenant: self.inner.tenant.clone(),
task: task.clone(),
disposition,
at: error
.as_ref()
.map(FlowError::at)
.filter(|path| !path.is_empty())
.map(Arc::from),
})
});
Report {
disposition,
output,
error,
elapsed: started.elapsed(),
steps,
dropped_steps,
progress_capacity: self.inner.limits.progress_capacity(),
effects,
tenant: self.inner.tenant.clone(),
task,
run,
ticket,
needs,
repair,
}
}
fn charge(&self, permit: Permit) -> AdmissionCharge<'_> {
let held = matches!(permit, Permit::Held(_));
if held {
let live = self.inner.in_flight.raise();
self.inner.peak_in_flight.raise_to(live);
}
self.inner.admitted.fetch_add(1, Ordering::Relaxed);
AdmissionCharge {
installation: &self.inner,
permit,
held,
}
}
}
impl Default for LiveRuns {
fn default() -> Self {
Self(AtomicUsize::new(0))
}
}
struct HighWater(AtomicUsize);
impl HighWater {
fn raise_to(&self, live: usize) {
let _previous = self
.0
.try_update(Ordering::Relaxed, Ordering::Relaxed, |peak| {
(live > peak).then_some(live)
});
}
fn get(&self) -> usize {
self.0.load(Ordering::Relaxed)
}
}
struct AdmissionCharge<'installation> {
installation: &'installation Installation,
permit: Permit,
held: bool,
}
impl Drop for AdmissionCharge<'_> {
fn drop(&mut self) {
if let Permit::Held(permit) = std::mem::replace(&mut self.permit, Permit::Charged) {
drop(permit);
}
if self.held {
self.installation.in_flight.release_one();
}
}
}
impl Default for HighWater {
fn default() -> Self {
Self(AtomicUsize::new(0))
}
}
struct LiveRuns(AtomicUsize);
impl LiveRuns {
fn raise(&self) -> usize {
self.0.fetch_add(1, Ordering::Relaxed).saturating_add(1)
}
fn release_one(&self) {
let _previous = self
.0
.try_update(Ordering::Relaxed, Ordering::Relaxed, |live| {
live.checked_sub(1)
});
}
fn get(&self) -> usize {
self.0.load(Ordering::Relaxed)
}
}
#[derive(Debug)]
struct Refused {
at: Arc<str>,
}
#[derive(Debug)]
enum AdmissionFailure {
Refused(Refused),
Failed(FlowError),
DeadlineExceeded(FlowError),
}
impl AdmissionFailure {
const fn disposition(&self) -> Disposition {
match *self {
Self::Refused(_) => Disposition::Refused,
Self::DeadlineExceeded(_) => Disposition::DeadlineExceeded,
Self::Failed(_) => Disposition::Failed,
}
}
fn into_error(self) -> FlowError {
match self {
Self::Refused(Refused { at }) => FlowError::Failed {
at,
reason: String::from("the host refused to admit this run"),
},
Self::Failed(error) | Self::DeadlineExceeded(error) => error,
}
}
}
fn disposition_of(error: &FlowError) -> Disposition {
match *error {
FlowError::Cancelled { .. } => Disposition::Cancelled,
FlowError::TimedOut { .. } => Disposition::DeadlineExceeded,
_ => Disposition::Failed,
}
}
struct Terminal<O> {
started: Instant,
task: TaskName,
disposition: Disposition,
output: Option<O>,
error: Option<FlowError>,
trail: TrailSnapshot,
run: Option<RunId>,
needs: Option<Deficit>,
repair: Option<RepairTicket>,
}
struct Declined {
disposition: Disposition,
reason: String,
run: Option<RunId>,
needs: Option<Deficit>,
repair: Option<RepairTicket>,
}
impl Declined {
fn plain(disposition: Disposition, reason: String) -> Self {
Self {
disposition,
reason,
run: None,
needs: None,
repair: None,
}
}
}
enum TrailSnapshot {
Taken((Vec<Arc<str>>, usize)),
Empty,
}
struct RunAuthority {
base: GrantSet,
delta: GrantSet,
}
impl AuthorityCheck for RunAuthority {
fn uncovered(&self, required: &[Cap]) -> Vec<Shortage> {
uncovered(
required,
|cap| self.base.grants(cap) || self.delta.grants(cap),
None,
)
}
}
enum Permit {
Held(OwnedSemaphorePermit),
Charged,
}
#[must_use = "a HostBuilder builds nothing until `build` is called"]
#[derive(Debug)]
pub struct HostBuilder {
tenant: Tenant,
max_concurrent: NonZeroUsize,
deadline: Duration,
progress: usize,
clock: Clock,
grants: GrantSet,
store: Option<store::RunStore>,
codec: String,
ledger: Option<RunLedger>,
repair_attempts: u64,
repair_spend: u64,
}
impl HostBuilder {
pub fn max_concurrent_tasks(mut self, tasks: NonZeroUsize) -> Self {
self.max_concurrent = tasks;
self
}
pub fn default_deadline(mut self, deadline: Duration) -> Self {
self.deadline = deadline;
self
}
pub fn clock(mut self, clock: Clock) -> Self {
self.clock = clock;
self
}
pub fn progress_capacity(mut self, steps: usize) -> Self {
self.progress = steps;
self
}
pub fn run_store(mut self, dir: impl AsRef<std::path::Path>) -> Result<Self, StoreError> {
self.store = Some(store::RunStore::open_in(
dir.as_ref(),
self.tenant.as_str(),
)?);
Ok(self)
}
pub fn store(mut self, store: store::RunStore) -> Self {
self.store = Some(store);
self
}
#[must_use = "a builder that declares nothing is the host the caller already had"]
pub fn durable_codec(mut self, codec: &str) -> Self {
codec.clone_into(&mut self.codec);
self
}
pub fn grants(mut self, grants: GrantSet) -> Self {
self.grants = grants;
self
}
pub fn repair_ledger(mut self, dir: impl AsRef<std::path::Path>) -> Result<Self, StoreError> {
self.ledger = Some(RunLedger::open_in(dir.as_ref(), self.tenant.as_str())?);
Ok(self)
}
pub fn ledger(mut self, ledger: RunLedger) -> Self {
self.ledger = Some(ledger);
self
}
pub fn repair_bounds(mut self, attempts: u64, spend: u64) -> Self {
self.repair_attempts = attempts;
self.repair_spend = spend;
self
}
pub fn build(self) -> Result<Host, HostError> {
let tasks = u64::saturating_from(self.max_concurrent.get());
let tasks_ceiling = u64::saturating_from(MAX_ADMITTED_TASKS);
if tasks > tasks_ceiling {
let refusal = Err(HostError::Bound {
what: "the admission ceiling",
value: tasks,
max: tasks_ceiling,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "build: returning an error to the caller");
return refusal;
}
if self.deadline.is_zero() {
let refusal = Err(HostError::Bound {
what: "the default deadline",
value: 0,
max: MAX_TASK_DEADLINE.as_secs(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "build: returning an error to the caller");
return refusal;
}
if self.deadline > MAX_TASK_DEADLINE {
let refusal = Err(HostError::Bound {
what: "the default deadline",
value: self.deadline.as_secs(),
max: MAX_TASK_DEADLINE.as_secs(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "build: returning an error to the caller");
return refusal;
}
let steps = u64::saturating_from(self.progress);
let steps_ceiling = u64::saturating_from(MAX_PROGRESS_STEPS);
if steps > steps_ceiling {
let refusal = Err(HostError::Bound {
what: "the progress capacity",
value: steps,
max: steps_ceiling,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "build: returning an error to the caller");
return refusal;
}
Ok(Host {
inner: Arc::new(Installation {
identity: HOST_IDENTITIES.fetch_add(1, Ordering::Relaxed),
tenant: self.tenant,
token: CancellationToken::new(),
limits: HostLimits {
max_concurrent: self.max_concurrent,
deadline: self.deadline,
progress: self.progress,
},
admission: Arc::new(Semaphore::new(self.max_concurrent.get())),
in_flight: LiveRuns::default(),
peak_in_flight: HighWater::default(),
admitted: AtomicU64::new(0),
refused: AtomicU64::new(0),
clock: self.clock,
store: self.store,
codec: self.codec,
ledger: self.ledger,
grants: self.grants,
repair_attempts: self.repair_attempts,
repair_spend: self.repair_spend,
reactor: Mutex::new(None),
}),
})
}
}
fn default_max_concurrent() -> NonZeroUsize {
const PER_CORE: NonZeroUsize = NonZeroUsize::MIN.saturating_add(63);
const CEILING: NonZeroUsize =
NonZeroUsize::MIN.saturating_add(MAX_ADMITTED_TASKS.saturating_sub(1));
let cores = match std::thread::available_parallelism() {
Ok(cores) => cores,
Err(uncounted) => {
lgwks_std::trace::debug!(?uncounted, "default_max_concurrent: sizing for one core");
NonZeroUsize::MIN
}
};
cores.saturating_mul(PER_CORE).min(CEILING)
}
fn default_deadline() -> Duration {
Duration::from_secs(30)
}
fn default_repair_attempts() -> u64 {
8
}
fn default_repair_spend() -> u64 {
default_repair_attempts().saturating_mul(8)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HostLimits {
max_concurrent: NonZeroUsize,
deadline: Duration,
progress: usize,
}
impl HostLimits {
#[must_use]
pub const fn max_concurrent_tasks(&self) -> usize {
self.max_concurrent.get()
}
#[must_use]
pub const fn default_deadline(&self) -> Duration {
self.deadline
}
#[must_use]
pub const fn progress_capacity(&self) -> usize {
self.progress
}
}
#[derive(Clone, Copy)]
pub struct Admission<'installation> {
inner: &'installation Installation,
}
impl fmt::Debug for Admission<'_> {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_struct("Admission")
.field("available_permits", &self.available_permits())
.field("in_flight", &self.in_flight())
.field("peak_in_flight", &self.peak_in_flight())
.field("admitted", &self.admitted())
.field("refused", &self.refused())
.finish()
}
}
impl Admission<'_> {
#[must_use]
pub fn available_permits(&self) -> usize {
self.inner.admission.available_permits()
}
#[must_use]
pub fn in_flight(&self) -> usize {
self.inner.in_flight.get()
}
#[must_use]
pub fn peak_in_flight(&self) -> usize {
self.inner.peak_in_flight.get()
}
#[must_use]
pub fn admitted(&self) -> u64 {
self.inner.admitted.load(Ordering::Relaxed)
}
#[must_use]
pub fn refused(&self) -> u64 {
self.inner.refused.load(Ordering::Relaxed)
}
}
#[derive(Debug)]
#[non_exhaustive]
pub enum HostError {
Bound {
what: &'static str,
value: u64,
max: u64,
},
InsideRuntime,
Runtime {
cause: io::Error,
},
}
impl fmt::Display for HostError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Bound { what, value, max } => {
write!(formatter, "{what}: {value} is outside 1..={max}")
}
Self::InsideRuntime => formatter.write_str(
"the synchronous entry cannot park a thread an async runtime is driving; \
await `Host::run` on that runtime, or call it outside it",
),
Self::Runtime { ref cause } => {
write!(formatter, "the host could not build a runtime: {cause}")
}
}
}
}
impl std::error::Error for HostError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Bound { .. } | Self::InsideRuntime => None,
Self::Runtime { ref cause } => Some(cause),
}
}
}