use std::any::Any;
use std::collections::{HashMap, HashSet};
use std::fmt;
use std::future::Future;
use std::num::{NonZeroU32, NonZeroU128};
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::task::{Context, Poll, Waker};
use std::thread;
use std::time::Duration;
use crate::effect::InputIdentity;
use lgwks_deps::bevy_ecs::{
self,
prelude::{Changed, Component, Entity, Resource, World},
schedule::{
IntoScheduleConfigs, LogLevel, Schedule, ScheduleBuildSettings, SingleThreadedExecutor,
},
};
use lgwks_std::hash::{Digest, Hasher};
use crate::journal::frame::SaturatingFrom;
use crate::journal::owner::{lock, wait_timeout};
use super::broker::{Authority, Broker, DispatchError, prepare_dispatch};
use super::cap::{Deficit, Demand, Shortage};
#[cfg(feature = "ephemeral")]
use super::effect::MintError;
use super::effect::{
ActionDigest, ActionId, AttemptId, EffectIdentity, EffectKey, EnvironmentEpoch, FlowRevision,
Id128,
};
use super::error::{BotError, DispatchCertainty, Escaped, RetryClass};
use super::gate::GrantSet;
#[cfg(feature = "ephemeral")]
use super::journal::MemoryJournal;
use super::journal::{
AttemptStatus, DurabilityPromise, DurableAck, EffectEvent, EffectJournal, EventKind,
JournalError, JournalPosition, recover,
};
use super::registry::{DomainRegistry, Source};
use super::spec::{
ActionSpec, Admission, BotSpec, ChainEntry, ChainSpec, Erased, Need, NeedSet, ObserveAny,
Witness, typed_entry,
};
use super::verb::EffectLifetime;
use super::verb::RefreshReason;
use super::verb::{Evaluate, Execute, Observe};
use crate::clock::Clock;
#[cfg(feature = "profile")]
mod profile;
#[cfg(feature = "profile")]
use profile::Charge;
#[cfg(feature = "profile")]
pub use profile::{TickProfile, TickStage};
const ACTION_ID_DOMAIN: &[u8] = b"lgwks.bot.action-id.v2";
const ACTION_DIGEST_DOMAIN: &[u8] = b"lgwks.bot.action-digest.v2";
const OUTPUT_IDENTITY_DOMAIN: &[u8] = b"lgwks.bot.input-identity.v2";
fn changed_outcome(key: EffectKey) -> JournalError {
JournalError::OutOfOrder {
key: Box::new(key),
expected: Some(EventKind::Verified),
attempted: EventKind::OutcomeObserved,
}
}
fn hash_parts(parts: &[&[u8]]) -> Digest {
let mut hasher = Hasher::new();
for part in parts {
hasher.write_framed(part);
}
hasher.finalize()
}
fn portable_index(index: usize) -> u64 {
u64::saturating_from(index)
}
fn derive_input_stamp(stamp: u64) -> [u8; 16] {
let digest = hash_parts(&[b"lgwks.bot.input-stamp.v1", &stamp.to_le_bytes()]);
let mut wide = [0_u8; 16];
wide.copy_from_slice(&digest.as_bytes()[..16]);
wide
}
fn id_from_digest(digest: &Digest) -> Id128 {
let bytes = digest.as_bytes();
let mut wide = [0_u8; 16];
for (slot, byte) in wide.iter_mut().zip(bytes.iter()) {
*slot = *byte;
}
match NonZeroU128::new(u128::from_be_bytes(wide)) {
Some(identity) => Id128::from_nonzero(identity),
None => Id128::from_nonzero(NonZeroU128::MIN),
}
}
fn derive_action_id(bot: &str, chain: usize, entry: usize, domain: &str) -> ActionId {
ActionId::new(id_from_digest(&hash_parts(&[
ACTION_ID_DOMAIN,
bot.as_bytes(),
&portable_index(chain).to_le_bytes(),
&portable_index(entry).to_le_bytes(),
domain.as_bytes(),
])))
}
fn derive_action_digest(
flow: FlowRevision,
chain: usize,
entry: usize,
input: &[u8; 16],
) -> ActionDigest {
ActionDigest::new(hash_parts(&[
ACTION_DIGEST_DOMAIN,
flow.digest().as_bytes(),
&portable_index(chain).to_le_bytes(),
&portable_index(entry).to_le_bytes(),
input,
]))
}
pub struct EffectScope {
identity: EffectIdentity,
broker: Broker,
journal: Box<dyn EffectJournal>,
}
#[cfg(feature = "ephemeral")]
#[derive(Debug)]
#[non_exhaustive]
pub enum EphemeralError {
Mint(MintError),
Broker(super::broker::BrokerError),
}
#[cfg(feature = "ephemeral")]
impl fmt::Display for EphemeralError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::Mint(ref cause) => write!(f, "could not mint an ephemeral identity: {cause}"),
Self::Broker(ref cause) => {
write!(f, "could not register the ephemeral environment: {cause}")
}
}
}
}
#[cfg(feature = "ephemeral")]
impl std::error::Error for EphemeralError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
match *self {
Self::Mint(ref cause) => Some(cause),
Self::Broker(ref cause) => Some(cause),
}
}
}
#[cfg(feature = "ephemeral")]
impl From<MintError> for EphemeralError {
fn from(cause: MintError) -> Self {
Self::Mint(cause)
}
}
impl fmt::Debug for EffectScope {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("EffectScope")
.field("identity", &self.identity)
.field("broker", &self.broker)
.field("durability", &self.journal.durability())
.field("tail", &self.journal.tail())
.finish()
}
}
impl EffectScope {
#[must_use]
pub fn new(identity: EffectIdentity, broker: Broker, journal: Box<dyn EffectJournal>) -> Self {
Self {
identity,
broker,
journal,
}
}
#[cfg(feature = "ephemeral")]
pub fn ephemeral() -> Result<Self, EphemeralError> {
let identity = EffectIdentity::ephemeral()?;
let mut broker = Broker::new();
broker
.register(identity.environment())
.map_err(EphemeralError::Broker)?;
Ok(Self::new(identity, broker, Box::new(MemoryJournal::new())))
}
#[must_use]
pub const fn identity(&self) -> EffectIdentity {
self.identity
}
#[must_use]
pub const fn broker(&self) -> &Broker {
&self.broker
}
#[must_use]
pub fn journal(&self) -> &dyn EffectJournal {
&*self.journal
}
#[must_use]
pub fn journal_mut(&mut self) -> &mut dyn EffectJournal {
&mut *self.journal
}
#[must_use]
pub fn into_journal(self) -> Box<dyn EffectJournal> {
self.journal
}
}
struct Effects {
scope: EffectScope,
tail: JournalPosition,
unsettled: Vec<EffectKey>,
recording: Vec<RecordedOutcome>,
requirements: Vec<(EffectKey, DurabilityPromise)>,
applied_latest: HashMap<ActionId, EffectKey>,
applied_events: HashSet<(ActionId, ActionDigest)>,
attempted: Vec<(ActionId, AttemptId)>,
}
struct RecordedOutcome {
key: EffectKey,
evidence: EffectEvidence,
cause: JournalError,
}
impl Effects {
fn new(scope: EffectScope, tail: JournalPosition) -> Self {
Self {
scope,
tail,
unsettled: Vec::new(),
recording: Vec::new(),
requirements: Vec::new(),
applied_latest: HashMap::new(),
applied_events: HashSet::new(),
attempted: Vec::new(),
}
}
fn note_applied(&mut self, key: EffectKey) {
self.applied_events.insert((key.action(), key.digest()));
self.applied_latest.insert(key.action(), key);
}
fn recording_for(&self, action: ActionId) -> Option<&RecordedOutcome> {
self.recording
.iter()
.find(|entry| entry.key.action() == action)
}
const fn identity(&self) -> EffectIdentity {
self.scope.identity()
}
const fn broker(&self) -> &Broker {
self.scope.broker()
}
fn into_journal(self) -> Box<dyn EffectJournal> {
self.scope.into_journal()
}
async fn prepare(
&mut self,
key: EffectKey,
lifetime: EffectLifetime,
) -> Result<Authority, DispatchError> {
let required = match lifetime {
EffectLifetime::Local => DurabilityPromise::Ephemeral,
EffectLifetime::External => {
self.scope
.journal()
.admit_external_handoff()
.map_err(DispatchError::Journal)?;
self.scope
.journal()
.reserve_handoff_capacity(3)
.map_err(DispatchError::Journal)?;
DurabilityPromise::ProcessCrash
}
};
let intent = EffectEvent::IntentAdmitted { key };
let intent_ack = self.append_async(&intent).await?;
self.note_journal_attempt(key.action(), key.attempt());
if !intent_ack.promise().meets(required) {
let refusal = Err(DispatchError::Journal(JournalError::PromiseUnmet {
required,
offered: intent_ack.promise(),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "prepare: returning an error to the caller");
return refusal;
}
self.accept_position(&intent, intent_ack.position())?;
let expected_tail = self.tail;
let scope = &mut self.scope;
let (authority, ack) =
prepare_dispatch(&scope.broker, &mut *scope.journal, expected_tail, key)
.await?
.into_parts();
if !ack.promise().meets(required) {
let refusal = Err(DispatchError::Journal(JournalError::PromiseUnmet {
required,
offered: ack.promise(),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "prepare: returning an error to the caller");
return refusal;
}
self.accept_position(&EffectEvent::DispatchPrepared { key }, ack.position())?;
self.requirements.push((key, required));
Ok(authority)
}
fn epoch(&self) -> Option<EnvironmentEpoch> {
self.scope.broker().epoch(self.identity().environment())
}
fn key(
&self,
action: ActionId,
chain: usize,
entry: usize,
input: [u8; 16],
attempt: AttemptId,
) -> Option<EffectKey> {
Some(self.identity().key(
action,
attempt,
derive_action_digest(self.identity().flow(), chain, entry, &input),
self.epoch()?,
))
}
fn continue_if_due(&mut self) -> Result<(), JournalError> {
if !self.scope.journal().continuation_watermark()?.is_due() {
return Ok(());
}
let successor = match self.scope.journal_mut().continue_as_new()? {
Some(successor) => successor,
None => return Ok(()),
};
let adopted = successor.tail();
self.scope.journal = successor;
self.tail = adopted;
Ok(())
}
fn append(&mut self, event: &EffectEvent) -> Result<DurableAck, JournalError> {
self.continue_if_due()?;
self.scope
.journal_mut()
.compare_and_append(self.tail, event)
}
async fn append_async(&mut self, event: &EffectEvent) -> Result<DurableAck, JournalError> {
self.continue_if_due()?;
self.scope
.journal_mut()
.compare_and_append_async(self.tail, event)
.await
}
fn accept_position(
&mut self,
event: &EffectEvent,
position: JournalPosition,
) -> Result<(), JournalError> {
self.event_position(event, position)?;
let actual = self.scope.journal().tail();
if actual != position {
let refusal = Err(JournalError::TailMismatch {
expected: position,
actual,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "accept_position: returning an error to the caller");
return refusal;
}
self.tail = position;
Ok(())
}
fn blocks(&self, action: ActionId) -> bool {
self.unsettled.iter().any(|key| key.action() == action)
|| self.recording_for(action).is_some()
}
fn unsettled_for(&self, action: ActionId) -> Option<EffectKey> {
self.unsettled
.iter()
.find(|key| key.action() == action)
.copied()
}
fn applied_in(
&self,
action: ActionId,
chain: usize,
entry: usize,
input: [u8; 16],
event: bool,
) -> bool {
let digest = derive_action_digest(self.identity().flow(), chain, entry, &input);
if event {
return self.applied_events.contains(&(action, digest));
}
self.applied_latest
.get(&action)
.is_some_and(|key| key.digest() == digest)
}
fn recorded_attempt(&self, action: ActionId) -> Option<AttemptId> {
self.attempted
.iter()
.find(|entry| entry.0 == action)
.map(|entry| entry.1)
}
fn note_journal_attempt(&mut self, action: ActionId, attempt: AttemptId) {
match self.attempted.iter_mut().find(|entry| entry.0 == action) {
Some(slot) => {
if attempt.get() > slot.1.get() {
slot.1 = attempt;
}
}
None => self.attempted.push((action, attempt)),
}
}
fn settle_recovered(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
) -> Result<(), JournalError> {
self.ensure_outcome(key, evidence)?;
self.unsettled.retain(|held| *held != key);
Ok(())
}
fn settle_recording(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
) -> Result<(), JournalError> {
self.confirm_recorded_outcome(key, evidence)?;
self.recording.retain(|entry| entry.key != key);
self.fold_outcome(key, evidence);
Ok(())
}
fn outcome_position(
&self,
key: EffectKey,
evidence: EffectEvidence,
) -> Result<JournalPosition, JournalError> {
let event = EffectEvent::OutcomeObserved { key, evidence };
let Some((position, recorded)) = self.scope.journal().outcome_at(key)? else {
let refusal = Err(JournalError::OutOfOrder {
key: Box::new(key),
expected: Some(EventKind::OutcomeObserved),
attempted: EventKind::OutcomeObserved,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "outcome_position: returning an error to the caller");
return refusal;
};
if recorded != evidence {
let refusal = Err(changed_outcome(key));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "outcome_position: the journal holds a different outcome");
return refusal;
}
self.event_position(&event, position)?;
Ok(position)
}
fn event_position(
&self,
expected: &EffectEvent,
position: JournalPosition,
) -> Result<(), JournalError> {
let Some(entry) = self.scope.journal().committed_entry(position)? else {
let refusal = Err(JournalError::OutOfOrder {
key: Box::new(expected.key()),
expected: Some(expected.kind()),
attempted: expected.kind(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "event_position: returning an error to the caller");
return refusal;
};
if entry.event() != expected {
let refusal = Err(JournalError::EntryMismatch {
position,
expected: Box::new(*expected),
actual: Box::new(*entry.event()),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "event_position: returning an error to the caller");
return refusal;
}
Ok(())
}
fn outcome_for(&self, key: &EffectKey) -> Result<Option<EffectEvidence>, JournalError> {
Ok(self
.scope
.journal()
.outcome_at(*key)?
.map(|(_, evidence)| evidence))
}
fn confirm_outcome(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
required: DurabilityPromise,
acknowledgment: DurableAck,
) -> Result<(), JournalError> {
let position = self.outcome_position(key, evidence)?;
if acknowledgment.position() != position {
let refusal = Err(JournalError::ReceiptMismatch {
expected: position,
actual: acknowledgment.position(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "confirm_outcome: returning an error to the caller");
return refusal;
}
if acknowledgment.promise().meets(required) {
return Ok(());
}
self.obtain_receipt(key, evidence, position, required)
}
fn obtain_receipt(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
position: JournalPosition,
required: DurabilityPromise,
) -> Result<(), JournalError> {
let receipt = self
.scope
.journal_mut()
.confirm_outcome(key, evidence, position, required)?;
if receipt.position() != position {
let refusal = Err(JournalError::ReceiptMismatch {
expected: position,
actual: receipt.position(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "obtain_receipt: the receipt names another position");
return refusal;
}
if receipt.promise().meets(required) {
Ok(())
} else {
Err(JournalError::PromiseUnmet {
required,
offered: receipt.promise(),
})
}
}
fn fold_outcome(&mut self, key: EffectKey, evidence: EffectEvidence) {
self.requirements.retain(|entry| entry.0 != key);
self.note_journal_attempt(key.action(), key.attempt());
if evidence == EffectEvidence::Applied {
self.note_applied(key);
}
}
fn confirm_recorded_outcome(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
) -> Result<(), JournalError> {
let required = self.scope.journal().durability();
if required == DurabilityPromise::Ephemeral {
return Ok(());
}
let position = self.outcome_position(key, evidence)?;
self.obtain_receipt(key, evidence, position, required)
}
fn settle_committed_outcome(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
required: DurabilityPromise,
) -> Result<(), JournalError> {
let offered = self.scope.journal().durability();
if !offered.meets(required) {
let refusal = Err(JournalError::ReceiptUnavailable { required });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "settle_committed_outcome: returning an error to the caller");
return refusal;
}
if offered != DurabilityPromise::Ephemeral {
self.confirm_recorded_outcome(key, evidence)?;
}
self.accept_position(
&EffectEvent::OutcomeObserved { key, evidence },
self.outcome_position(key, evidence)?,
)?;
self.fold_outcome(key, evidence);
Ok(())
}
fn ensure_outcome(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
) -> Result<(), JournalError> {
let required = match self.requirements.iter().find(|entry| entry.0 == key) {
Some(&(_, declared)) => declared,
None => self.scope.journal().durability(),
};
if let Some(recorded) = self.outcome_for(&key)? {
if recorded != evidence {
let refusal = Err(changed_outcome(key));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "ensure_outcome: the journal holds a different outcome");
return refusal;
}
return self.settle_committed_outcome(key, evidence, required);
}
match self.append(&EffectEvent::OutcomeObserved { key, evidence }) {
Ok(acknowledgment) => self.record_new_outcome(key, evidence, required, acknowledgment),
Err(JournalError::OutOfOrder { expected, .. })
if expected == Some(EventKind::Verified) || expected.is_none() =>
{
match self.outcome_for(&key)? {
Some(recorded) if recorded == evidence => {
self.settle_committed_outcome(key, evidence, required)
}
_ => Err(JournalError::OutOfOrder {
key: Box::new(key),
expected,
attempted: EventKind::OutcomeObserved,
}),
}
}
Err(cause @ JournalError::OutcomeUnknown { .. }) => {
match self.outcome_for(&key)? {
Some(recorded) if recorded == evidence => {
self.settle_committed_outcome(key, evidence, required)
}
Some(_) => Err(JournalError::OutOfOrder {
key: Box::new(key),
expected: Some(EventKind::Verified),
attempted: EventKind::OutcomeObserved,
}),
None => Err(cause),
}
}
Err(cause) => {
Err(cause)
}
}
}
}
impl Effects {
fn record_new_outcome(
&mut self,
key: EffectKey,
evidence: EffectEvidence,
required: DurabilityPromise,
acknowledgment: DurableAck,
) -> Result<(), JournalError> {
self.confirm_outcome(key, evidence, required, acknowledgment)?;
let position = self.outcome_position(key, evidence)?;
self.accept_position(&EffectEvent::OutcomeObserved { key, evidence }, position)?;
self.fold_outcome(key, evidence);
Ok(())
}
}
impl fmt::Debug for Effects {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Effects")
.field("identity", &self.identity())
.field("journal", &self.scope.journal().durability())
.field("unsettled", &self.unsettled.len())
.field("applied", &self.applied_events.len())
.field("attempted", &self.attempted.len())
.finish()
}
}
#[derive(Component, Debug, Clone)]
pub(crate) struct SourceId {
pub(crate) chain: usize,
domain: String,
}
impl SourceId {
pub(crate) fn domain(&self) -> &str {
&self.domain
}
}
#[derive(Component, Debug, Clone, Copy, PartialEq, Eq, Default)]
pub(crate) struct Revision(u64);
#[derive(Resource, Debug)]
struct Grants(GrantSet);
#[derive(Resource, Debug, Default)]
struct Fired(usize);
#[derive(Resource, Debug, Default)]
struct TickError(Option<BotError>);
#[derive(Resource, Debug, Default)]
struct Order(Vec<Entity>);
#[derive(Resource, Debug, Default)]
struct PollScratch {
declared: Vec<Option<RefreshReason>>,
stalled: Vec<usize>,
slots: Vec<usize>,
}
#[derive(Debug, Default)]
struct ChangeTicks {
seen: Vec<Option<u64>>,
polled: Vec<Option<u64>>,
}
impl ChangeTicks {
fn quiet(&self, chain: usize, reported: Option<u64>) -> bool {
reported.is_some() && self.seen.get(chain).copied().flatten() == reported
}
fn commit(&mut self, chain: usize, polled: Option<u64>) {
if let Some(slot) = self.seen.get_mut(chain) {
*slot = polled;
}
}
}
#[derive(Default)]
struct AdmitScratch {
moving: Vec<bool>,
kinds: Vec<AdmitKind>,
inputs: Vec<Option<AdmittedInput>>,
}
impl AdmitScratch {
fn take(world: &mut World, count: usize) -> Self {
let mut scratch = std::mem::take(&mut *world.non_send_mut::<AdmitScratch>());
scratch.moving.clear();
scratch.moving.resize(count, false);
scratch.kinds.clear();
scratch.kinds.resize(count, AdmitKind::Idle);
scratch.inputs.clear();
scratch.inputs.resize(count, None);
scratch
}
}
#[derive(Resource, Debug)]
struct Policy(RetryPolicy);
#[derive(Debug, Default)]
struct Invalidated {
reasons: Vec<Option<RefreshReason>>,
}
#[derive(Debug, Clone, Copy, Resource)]
struct PollBudget(Duration);
impl Invalidated {
fn mark(&mut self, chain: usize, reason: RefreshReason) {
if let Some(slot) = self.reasons.get_mut(chain)
&& slot.is_none()
{
*slot = Some(reason);
}
}
fn is_invalid(&self, chain: usize) -> bool {
self.reasons
.get(chain)
.is_some_and(|reason| reason.is_some())
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ForcedRefresh {
chain: usize,
domain: String,
reason: RefreshReason,
}
impl ForcedRefresh {
#[must_use]
pub const fn chain(&self) -> usize {
self.chain
}
#[must_use]
pub fn domain(&self) -> &str {
&self.domain
}
#[must_use]
pub const fn reason(&self) -> RefreshReason {
self.reason
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct StalledSource {
chain: usize,
domain: String,
deadline: Duration,
}
impl StalledSource {
#[must_use]
pub const fn chain(&self) -> usize {
self.chain
}
#[must_use]
pub fn domain(&self) -> &str {
&self.domain
}
#[must_use]
pub const fn deadline(&self) -> Duration {
self.deadline
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Resource)]
pub struct TickReport {
fired: usize,
forced: Vec<ForcedRefresh>,
superseded: Vec<SupersededObservation>,
stalled: Vec<StalledSource>,
watchdogs: u32,
}
impl TickReport {
#[must_use]
pub const fn fired(&self) -> usize {
self.fired
}
#[must_use]
pub fn forced(&self) -> &[ForcedRefresh] {
&self.forced
}
#[must_use]
pub fn superseded(&self) -> &[SupersededObservation] {
&self.superseded
}
#[must_use]
pub fn forced_any(&self) -> bool {
!self.forced.is_empty()
}
#[must_use]
pub fn superseded_any(&self) -> bool {
!self.superseded.is_empty()
}
#[must_use]
pub fn stalled(&self) -> &[StalledSource] {
&self.stalled
}
#[must_use]
pub fn stalled_any(&self) -> bool {
!self.stalled.is_empty()
}
#[must_use]
pub const fn watchdogs(&self) -> u32 {
self.watchdogs
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct SupersededObservation {
chain: usize,
revision: u64,
}
impl SupersededObservation {
#[must_use]
pub const fn chain(&self) -> usize {
self.chain
}
#[must_use]
pub const fn revision(&self) -> u64 {
self.revision
}
}
#[derive(Debug, Default)]
struct Committed {
slots: Vec<SlotAdmission>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
enum SlotState {
#[default]
Nothing,
Admitted,
Unacted,
}
#[derive(Debug, Clone, Copy, Default)]
struct SlotAdmission {
state: SlotState,
revision: u64,
}
impl Committed {
fn commit(&mut self, chain: usize, revision: u64) {
if let Some(slot) = self.slots.get_mut(chain) {
slot.state = SlotState::Unacted;
slot.revision = revision;
}
}
fn admit(&mut self, chain: usize) {
if let Some(slot) = self.slots.get_mut(chain) {
slot.state = SlotState::Admitted;
}
}
fn slot(&self, chain: usize) -> Option<SlotAdmission> {
self.slots.get(chain).copied()
}
}
pub(crate) struct EcsChain {
source: Box<dyn ObserveAny>,
same: fn(&dyn Any, &dyn Any) -> bool,
identify: fn(&dyn Any) -> AdmittedInput,
entries: Vec<ChainEntry>,
witness: Witness,
}
impl EcsChain {
fn from_registry(source: Source, entries: Vec<ChainEntry>) -> Self {
let (inner, same, identify, witness) = source.into_parts();
Self {
source: inner,
same,
identify,
witness,
entries,
}
}
}
#[derive(Default)]
struct Chains(Vec<EcsChain>);
#[derive(Default)]
struct Observed(Vec<Option<Erased>>);
#[derive(Default)]
struct Polled(Vec<Result<Option<Erased>, BotError>>);
#[derive(Default)]
struct Moved(Vec<bool>);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Decision {
Attempt,
Skip,
}
struct Step {
chain: usize,
entry: usize,
decision: Decision,
}
#[derive(Default)]
struct Plan {
steps: Vec<Step>,
failures: Vec<Option<BotError>>,
}
const MAX_IN_FLIGHT_POLLS: usize = 32;
pub const DEFAULT_POLL_DEADLINE: Duration = Duration::from_secs(30);
pub const MAX_POLL_DEADLINE: Duration = Duration::from_secs(600);
#[derive(Debug)]
struct PollWatchdog {
state: Arc<Mutex<WatchdogState>>,
settled: Arc<Condvar>,
reaper: Mutex<Option<thread::JoinHandle<()>>>,
}
impl PollWatchdog {
fn new(deadline: Duration, members: usize) -> Self {
Self {
state: Arc::new(Mutex::new(WatchdogState {
members,
resolved: 0,
expired: false,
released: false,
deadline,
pending: vec![false; members],
wakers: (0..members).map(|_| None).collect(),
})),
settled: Arc::new(Condvar::new()),
reaper: Mutex::new(None),
}
}
fn poll_finished(&self) {
let mut state = lock(&self.state);
state.resolved = state.resolved.saturating_add(1);
if state.resolved >= state.members {
state.released = true;
self.settled.notify_all();
}
}
fn poll_pending(&self, slot: usize) {
let mut state = lock(&self.state);
if let Some(pending) = state.pending.get_mut(slot) {
*pending = true;
}
}
fn expired(&self) -> bool {
lock(&self.state).expired
}
fn start(&self, clock: &Clock, armed: &AtomicBool) -> bool {
let mut owned = lock(&self.reaper);
if owned.is_some() {
return true;
}
{
let state = lock(&self.state);
if state.expired || state.released || !state.pending.iter().any(|held| *held) {
return true;
}
}
let shared = WatchdogShared {
deadline: lock(&self.state).deadline,
state: Arc::clone(&self.state),
settled: Arc::clone(&self.settled),
};
let built = thread::Builder::new()
.name("lgwks-poll-deadline".into())
.spawn({
let clock = clock.clone();
move || reaper(&clock, &shared)
});
match built {
Ok(handle) => {
*owned = Some(handle);
armed.store(true, Ordering::Relaxed);
true
}
Err(_) => false,
}
}
fn join(&self) {
let handle = lock(&self.reaper).take();
if let Some(handle) = handle {
let _joined = handle.join();
}
}
}
#[derive(Debug, Default)]
struct WatchdogState {
expired: bool,
released: bool,
resolved: usize,
members: usize,
pending: Vec<bool>,
deadline: Duration,
wakers: Vec<Option<Waker>>,
}
struct WatchdogShared {
deadline: Duration,
state: Arc<Mutex<WatchdogState>>,
settled: Arc<Condvar>,
}
fn reaper(clock: &Clock, watchdog: &WatchdogShared) {
let fired = clock.wall_watchdog();
let mut state = lock(&watchdog.state);
loop {
if state.released {
return;
}
let remaining = watchdog.deadline.saturating_sub(fired.elapsed());
if remaining.is_zero() {
state.expired = true;
for slot in &mut state.wakers {
if let Some(waker) = slot.take() {
waker.wake();
}
}
return;
}
let (woken, _) = wait_timeout(&watchdog.settled, state, remaining);
state = woken;
if state.expired || state.released {
return;
}
}
}
struct WavePoll<'wave> {
chain: usize,
watchdog: &'wave PollWatchdog,
slot: usize,
installed: bool,
poll: Pin<Box<dyn Future<Output = Result<Option<Erased>, BotError>> + 'wave>>,
}
impl WavePoll<'_> {
fn stalled(&self, deadline: Duration) -> Result<Option<Erased>, BotError> {
Err(BotError::PollStalled {
chain: self.chain,
deadline,
})
}
fn turn(
&mut self,
cx: &mut Context<'_>,
clock: &Clock,
armed: &AtomicBool,
deadline: Duration,
) -> Poll<Result<Option<Erased>, BotError>> {
{
let mut state = lock(&self.watchdog.state);
if state.expired {
return Poll::Ready(self.stalled(deadline));
}
let stale = if self.installed {
state
.wakers
.get(self.slot)
.and_then(|slot| slot.as_ref())
.is_none_or(|installed| !installed.will_wake(cx.waker()))
} else {
true
};
if let (true, Some(slot)) = (stale, state.wakers.get_mut(self.slot)) {
*slot = Some(cx.waker().clone());
self.installed = true;
}
}
match self.poll.as_mut().poll(cx) {
Poll::Ready(outcome) => Poll::Ready(outcome),
Poll::Pending => {
self.watchdog.poll_pending(self.slot);
if !self.watchdog.start(clock, armed) {
return Poll::Ready(self.stalled(deadline));
}
if self.watchdog.expired() {
Poll::Ready(self.stalled(deadline))
} else {
Poll::Pending
}
}
}
}
}
impl Drop for WavePoll<'_> {
fn drop(&mut self) {
self.watchdog.poll_finished();
}
}
async fn bounded_wave(
deadline: Duration,
watchdog: &PollWatchdog,
polls: impl IntoIterator<Item = WavePoll<'_>>,
) -> (Vec<Result<Option<Erased>, BotError>>, bool) {
let clock = Clock::wall();
let armed = AtomicBool::new(false);
let joined = lgwks_std::task::join_all_boxed(polls.into_iter().map(|mut poll| {
let clock = clock.clone();
let armed = &armed;
Box::pin(std::future::poll_fn(move |cx: &mut Context<'_>| {
poll.turn(cx, &clock, armed, deadline)
}))
}));
let results = joined.await;
watchdog.join();
(results, armed.load(Ordering::Relaxed))
}
#[derive(Clone, Copy)]
pub(crate) struct AdmittedInput {
identity: [u8; 16],
event: bool,
}
pub(crate) fn identify_output<S>(value: &dyn Any) -> AdmittedInput
where
S: Observe,
S::Output: InputIdentity + 'static,
{
let mut hasher = Hasher::new();
hasher
.write_framed(OUTPUT_IDENTITY_DOMAIN)
.write_framed(<S::Output as InputIdentity>::SCHEMA_ID);
let event = match value.downcast_ref::<S::Output>() {
Some(typed) => {
typed.write_identity(&mut hasher);
typed.names_an_event()
}
None => {
hasher.write_framed(b"downcast-miss");
false
}
};
let digest = hasher.finalize();
let mut wide = [0_u8; 16];
wide.copy_from_slice(&digest.as_bytes()[..16]);
AdmittedInput {
identity: wide,
event,
}
}
pub(crate) fn same_output<S>(left: &dyn Any, right: &dyn Any) -> bool
where
S: Observe + 'static,
S::Output: PartialEq + InputIdentity + 'static,
{
match (
left.downcast_ref::<S::Output>(),
right.downcast_ref::<S::Output>(),
) {
(Some(left), Some(right)) => left == right,
_ => false,
}
}
fn parked(world: &World) -> bool {
world.get_resource::<TickError>().is_some_and(|error| {
error
.0
.as_ref()
.is_some_and(|error| !matches!(error, BotError::PollStalled { .. }))
})
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub struct WorkId {
chain: usize,
entry: usize,
}
impl WorkId {
#[must_use]
pub const fn chain(self) -> usize {
self.chain
}
#[must_use]
pub const fn entry(self) -> usize {
self.entry
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum AbandonReason {
AttemptsExhausted {
attempts: u32,
},
Terminal,
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum TransitionHold {
NotStarted,
Unrecorded,
UnsettledByRecovery {
action: ActionId,
},
Failed {
attempts: u32,
budget: u32,
cause: String,
},
OutcomeUnknown {
attempts: u32,
cause: String,
},
RecordingFailed {
evidence: EffectEvidence,
cause: String,
},
Abandoned {
reason: AbandonReason,
cause: String,
},
}
impl fmt::Display for TransitionHold {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match *self {
Self::NotStarted => f.write_str("not attempted: an entry ahead of it is unresolved"),
Self::Unrecorded => {
f.write_str("an attempt began and no outcome was recorded: the effect may be live")
}
Self::UnsettledByRecovery { action } => write!(
f,
"may have happened: an earlier attempt on action {action} was dispatched \
and nothing recorded its outcome, so nothing is sent again until it is \
settled"
),
Self::Failed {
attempts,
budget,
ref cause,
} => write!(
f,
"did not happen: attempt {attempts} of {budget} failed with {}",
Escaped(cause)
),
Self::OutcomeUnknown {
attempts,
ref cause,
} => write!(
f,
"may have happened: attempt {attempts} was indeterminate ({})",
Escaped(cause)
),
Self::RecordingFailed {
evidence,
ref cause,
} => write!(
f,
"the effect is {evidence} and only its record is missing: {}. \
The journal append is retried; the action is not",
Escaped(cause)
),
Self::Abandoned { reason, ref cause } => match reason {
AbandonReason::AttemptsExhausted { attempts } => write!(
f,
"given up on after {attempts} attempts: {}",
Escaped(cause)
),
AbandonReason::Terminal => {
write!(f, "given up on, no retry can fix: {}", Escaped(cause))
}
},
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PendingWork {
id: WorkId,
key: Option<Box<EffectKey>>,
hold: TransitionHold,
}
impl PendingWork {
#[must_use]
pub const fn id(&self) -> WorkId {
self.id
}
#[must_use]
pub fn key(&self) -> Option<EffectKey> {
self.key.as_deref().copied()
}
#[must_use]
pub const fn hold(&self) -> &TransitionHold {
&self.hold
}
}
impl fmt::Display for PendingWork {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "chain {} entry {} (", self.id.chain(), self.id.entry())?;
match self.key.as_deref() {
Some(key) => write!(
f,
"attempt {}, binding {}",
key.attempt().get(),
key.digest()
)?,
None => f.write_str("no attempt yet")?,
}
write!(f, "): {}", self.hold)
}
}
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Hash,
lgwks_std::wire::Archive,
lgwks_std::wire::Serialize,
lgwks_std::wire::Deserialize,
)]
#[rkyv(attr(non_exhaustive), crate = lgwks_std::wire::rkyv, compare(PartialEq), derive(Debug))]
#[non_exhaustive]
pub enum EffectEvidence {
Applied,
NotApplied,
}
impl fmt::Display for EffectEvidence {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match *self {
Self::Applied => "applied",
Self::NotApplied => "not applied",
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct RetryPolicy {
max_attempts: NonZeroU32,
}
impl RetryPolicy {
pub const ONE_ATTEMPT: Self = Self {
max_attempts: NonZeroU32::MIN,
};
pub const DEFAULT: Self = Self {
max_attempts: match NonZeroU32::new(3) {
Some(three) => three,
None => NonZeroU32::MIN,
},
};
#[must_use]
pub const fn new(max_attempts: NonZeroU32) -> Self {
Self { max_attempts }
}
#[must_use]
pub const fn max_attempts(self) -> u32 {
self.max_attempts.get()
}
}
#[derive(Debug, Clone)]
enum EntryState {
NotStarted,
Unrecorded,
Succeeded,
Skipped,
DefinitelyFailed {
attempts: u32,
cause: String,
},
OutcomeUnknown {
attempts: u32,
cause: String,
},
RecordingFailed {
key: EffectKey,
evidence: EffectEvidence,
cause: String,
attempts: u32,
},
Abandoned {
reason: AbandonReason,
cause: String,
},
}
impl EntryState {
fn is_open(&self) -> bool {
matches!(
self,
Self::NotStarted
| Self::Unrecorded
| Self::DefinitelyFailed { .. }
| Self::OutcomeUnknown { .. }
| Self::RecordingFailed { .. }
)
}
fn is_abandoned(&self) -> bool {
matches!(self, Self::Abandoned { .. })
}
fn hold(&self, budget: u32) -> Option<TransitionHold> {
match *self {
Self::Succeeded | Self::Skipped => None,
Self::NotStarted => Some(TransitionHold::NotStarted),
Self::Unrecorded => Some(TransitionHold::Unrecorded),
Self::DefinitelyFailed {
attempts,
ref cause,
} => Some(TransitionHold::Failed {
attempts,
budget,
cause: cause.clone(),
}),
Self::OutcomeUnknown {
attempts,
ref cause,
} => Some(TransitionHold::OutcomeUnknown {
attempts,
cause: cause.clone(),
}),
Self::RecordingFailed {
ref evidence,
ref cause,
..
} => Some(TransitionHold::RecordingFailed {
evidence: *evidence,
cause: cause.clone(),
}),
Self::Abandoned { reason, ref cause } => Some(TransitionHold::Abandoned {
reason,
cause: cause.clone(),
}),
}
}
}
#[derive(Debug, Clone, Copy, Default)]
struct AttemptRecord {
begun: Option<AttemptId>,
settled: Option<SettledAttempt>,
}
#[derive(Debug, Clone, Copy)]
struct SettledAttempt {
attempt: AttemptId,
evidence: EffectEvidence,
}
struct Transition {
revision: u64,
input: [u8; 16],
event: bool,
entries: Vec<EntryState>,
attempts: Vec<AttemptRecord>,
value: Option<Erased>,
}
impl fmt::Debug for Transition {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("Transition")
.field("revision", &self.revision)
.field("entries", &self.entries)
.field("attempts", &self.attempts)
.field("bound", &self.value.is_some())
.finish()
}
}
impl Transition {
fn attempt_of(&self, entry: usize) -> Option<AttemptId> {
self.attempts.get(entry).and_then(|record| record.begun)
}
fn opened(input: AdmittedInput, revision: u64, entries: usize, value: Option<Erased>) -> Self {
Self {
revision,
input: input.identity,
event: input.event,
entries: vec![EntryState::NotStarted; entries],
attempts: vec![AttemptRecord::default(); entries],
value,
}
}
fn resumed(
input: AdmittedInput,
revision: u64,
previous: &Self,
value: Option<Erased>,
) -> Self {
Self {
revision,
input: input.identity,
event: input.event,
entries: previous
.entries
.iter()
.map(|state| match *state {
EntryState::Abandoned { reason, ref cause } => EntryState::Abandoned {
reason,
cause: cause.clone(),
},
_ => EntryState::NotStarted,
})
.collect(),
attempts: vec![AttemptRecord::default(); previous.entries.len()],
value,
}
}
fn has_open(&self) -> bool {
self.entries.iter().any(EntryState::is_open)
}
fn is_retained(&self) -> bool {
self.entries
.iter()
.any(|state| state.is_open() || state.is_abandoned())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Settled {
Decided,
Duplicate,
Superseded {
id: WorkId,
current: ActionDigest,
},
StaleAttempt {
id: WorkId,
reported: AttemptId,
outstanding: AttemptId,
},
Contradicted {
id: WorkId,
settled: EffectEvidence,
},
NoSuchWork {
id: Option<WorkId>,
},
}
#[derive(Debug)]
struct Ledger {
transitions: Vec<Option<Transition>>,
next_input: u64,
actions: Vec<Vec<ActionId>>,
effects: Effects,
}
impl Ledger {
fn new(actions: Vec<Vec<ActionId>>, effects: Effects, next_input: u64) -> Self {
Self {
transitions: (0..actions.len()).map(|_| None).collect(),
actions,
effects,
next_input,
}
}
fn mint_input(&mut self) -> [u8; 16] {
let stamp = self.next_input;
self.next_input = self.next_input.saturating_add(1);
derive_input_stamp(stamp)
}
fn action_of(&self, id: WorkId) -> Option<ActionId> {
self.actions.get(id.chain())?.get(id.entry()).copied()
}
fn into_journal(self) -> Box<dyn EffectJournal> {
self.effects.into_journal()
}
fn generation(&self, chain: usize) -> Option<([u8; 16], bool)> {
self.transitions
.get(chain)
.and_then(Option::as_ref)
.map(|transition| (transition.input, transition.event))
}
fn locate(&self, action: ActionId) -> Option<WorkId> {
for (chain, entries) in self.actions.iter().enumerate() {
if let Some(entry) = entries.iter().position(|declared| *declared == action) {
return Some(WorkId { chain, entry });
}
}
None
}
fn take(&mut self, chain: usize) -> Option<Transition> {
self.transitions.get_mut(chain).and_then(Option::take)
}
fn put(&mut self, chain: usize, transition: Option<Transition>) {
if let Some(slot) = self.transitions.get_mut(chain) {
*slot = transition;
}
}
fn entry_state(&self, id: WorkId) -> Option<&EntryState> {
self.transitions
.get(id.chain())
.and_then(Option::as_ref)
.and_then(|transition| transition.entries.get(id.entry()))
}
fn entry_mut(&mut self, id: WorkId) -> Option<&mut EntryState> {
self.transitions
.get_mut(id.chain())
.and_then(Option::as_mut)
.and_then(|transition| transition.entries.get_mut(id.entry()))
}
async fn begin_attempt(
&mut self,
id: WorkId,
lifetime: EffectLifetime,
) -> Result<(EffectKey, Authority), BotError> {
let no_such_work = || BotError::NoSuchWork { work: id };
let action = self.action_of(id).ok_or_else(no_such_work)?;
if self.effects.blocks(action) {
let refusal = Err(BotError::EffectUnsettled { action });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "begin_attempt: returning an error to the caller");
return refusal;
}
let transition = self
.transitions
.get_mut(id.chain())
.and_then(Option::as_mut)
.ok_or_else(no_such_work)?;
let state = transition
.entries
.get_mut(id.entry())
.ok_or_else(no_such_work)?;
if !matches!(
*state,
EntryState::NotStarted | EntryState::DefinitelyFailed { .. }
) {
let refusal = Err(no_such_work());
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "begin_attempt: returning an error to the caller");
return refusal;
}
let record = transition
.attempts
.get_mut(id.entry())
.ok_or_else(no_such_work)?;
let attempt = match record.begun {
Some(previous) => previous
.checked_next()
.ok_or(BotError::EffectUnsettled { action })?,
None => match self.effects.recorded_attempt(action) {
Some(recorded) => recorded
.checked_next()
.ok_or(BotError::EffectUnsettled { action })?,
None => AttemptId::FIRST,
},
};
let input = transition.input;
let key = self
.effects
.key(action, id.chain(), id.entry(), input, attempt)
.ok_or(BotError::EffectUnsettled { action })?;
let authority = self
.effects
.prepare(key, lifetime)
.await
.map_err(|cause| BotError::EffectRefused { cause })?;
let transition = self
.transitions
.get_mut(id.chain())
.and_then(Option::as_mut)
.ok_or_else(no_such_work)?;
let record = transition
.attempts
.get_mut(id.entry())
.ok_or_else(no_such_work)?;
record.begun = Some(attempt);
if let Some(state) = transition.entries.get_mut(id.entry()) {
*state = EntryState::Unrecorded;
}
Ok((key, authority))
}
fn skip(&mut self, id: WorkId) -> bool {
let Some(state) = self.entry_mut(id) else {
return false;
};
if !matches!(
*state,
EntryState::NotStarted | EntryState::DefinitelyFailed { .. }
) {
return false;
}
*state = EntryState::Skipped;
true
}
fn succeed(&mut self, id: WorkId) -> bool {
let Some(state) = self.entry_mut(id) else {
return false;
};
if !matches!(*state, EntryState::Unrecorded) {
return false;
}
*state = EntryState::Succeeded;
true
}
fn recording_failed(
&self,
id: WorkId,
) -> Option<(crate::effect::EffectKey, EffectEvidence, u32)> {
match *self.entry_state(id)? {
EntryState::RecordingFailed {
key,
evidence,
attempts,
..
} => Some((key, evidence, attempts)),
_ => None,
}
}
fn hold_recording(&mut self, id: WorkId, evidence: EffectEvidence, error: &BotError) -> bool {
let Some(state) = self.entry_mut(id) else {
return false;
};
let EntryState::RecordingFailed {
ref mut cause,
evidence: ref mut held,
..
} = *state
else {
return false;
};
*cause = error.to_string();
*held = evidence;
true
}
fn finish_recording(
&mut self,
id: WorkId,
evidence: EffectEvidence,
attempts: u32,
budget: u32,
) -> bool {
let Some(state) = self.entry_mut(id) else {
return false;
};
if !matches!(*state, EntryState::RecordingFailed { .. }) {
return false;
}
*state = match evidence {
EffectEvidence::Applied => EntryState::Succeeded,
EffectEvidence::NotApplied => {
if attempts < budget {
EntryState::DefinitelyFailed {
attempts,
cause: "effect not applied; its record landed on retry".into(),
}
} else {
EntryState::Abandoned {
reason: AbandonReason::AttemptsExhausted { attempts },
cause: "effect not applied; its record landed on retry".into(),
}
}
}
};
true
}
fn fail(&mut self, id: WorkId, error: &BotError, attempt: AttemptId, budget: u32) -> bool {
let attempts = u32::saturating_from(attempt.get());
let Some(state) = self.entry_mut(id) else {
return false;
};
if !matches!(*state, EntryState::Unrecorded) {
return false;
}
*state = failure_state(error, attempts, budget);
true
}
fn classify_settlement(
&self,
key: &EffectKey,
evidence: EffectEvidence,
) -> Result<Settled, JournalError> {
let Some(id) = self.locate(key.action()) else {
return Ok(Settled::NoSuchWork { id: None });
};
let identity = self.effects.identity();
if key.run() != identity.run()
|| key.flow() != identity.flow()
|| key.environment() != identity.environment()
{
return Ok(Settled::NoSuchWork { id: None });
}
let Some(transition) = self.transitions.get(id.chain()).and_then(Option::as_ref) else {
return self.durable_or_missing(key, id, evidence);
};
let live = derive_action_digest(identity.flow(), id.chain(), id.entry(), &transition.input);
if key.digest() != live {
return Ok(Settled::Superseded { id, current: live });
}
let Some(record) = transition.attempts.get(id.entry()).copied() else {
return self.durable_or_missing(key, id, evidence);
};
if let Some(recorded) = self.effects.outcome_for(key)?
&& recorded != evidence
{
return Ok(Settled::Contradicted {
id,
settled: recorded,
});
}
let settleable = matches!(
transition.entries.get(id.entry()),
Some(
EntryState::Unrecorded
| EntryState::OutcomeUnknown { .. }
| EntryState::Abandoned { .. }
)
);
if !settleable {
return match record.settled {
Some(previous)
if previous.attempt == key.attempt() && previous.evidence == evidence =>
{
Ok(Settled::Duplicate)
}
Some(previous) if previous.attempt == key.attempt() => Ok(Settled::Contradicted {
id,
settled: previous.evidence,
}),
Some(_) | None => self.durable_or_missing(key, id, evidence),
};
}
let Some(outstanding) = record.begun else {
return Ok(Settled::NoSuchWork { id: Some(id) });
};
if key.attempt() != outstanding {
return Ok(Settled::StaleAttempt {
id,
reported: key.attempt(),
outstanding,
});
}
Ok(Settled::Decided)
}
fn durable_or_missing(
&self,
key: &EffectKey,
id: WorkId,
evidence: EffectEvidence,
) -> Result<Settled, JournalError> {
Ok(match self.effects.outcome_for(key)? {
Some(recorded) if recorded == evidence => Settled::Duplicate,
Some(recorded) => Settled::Contradicted {
id,
settled: recorded,
},
None => Settled::NoSuchWork { id: Some(id) },
})
}
fn apply_settlement(&mut self, key: &EffectKey, evidence: EffectEvidence) {
let Some(id) = self.locate(key.action()) else {
return;
};
let Some(transition) = self
.transitions
.get_mut(id.chain())
.and_then(Option::as_mut)
else {
return;
};
let Some(record) = transition.attempts.get(id.entry()).copied() else {
return;
};
let Some(outstanding) = record.begun else {
return;
};
let next = match evidence {
EffectEvidence::Applied => EntryState::Succeeded,
EffectEvidence::NotApplied => EntryState::NotStarted,
};
if let Some(state) = transition.entries.get_mut(id.entry()) {
*state = next;
}
if let Some(slot) = transition.attempts.get_mut(id.entry()) {
slot.settled = Some(SettledAttempt {
attempt: outstanding,
evidence,
});
}
}
fn bound(&self, chain: usize) -> Option<&Erased> {
self.transitions
.get(chain)
.and_then(Option::as_ref)
.and_then(|transition| transition.value.as_ref())
}
fn key_of(&self, chain: usize, entry: usize, transition: &Transition) -> Option<EffectKey> {
let id = WorkId { chain, entry };
self.effects.key(
self.action_of(id)?,
chain,
entry,
transition.input,
transition.attempt_of(entry)?,
)
}
fn unresolved(&self, budget: u32) -> Vec<PendingWork> {
let mut unresolved = Vec::new();
for chain in 0..self.actions.len().max(self.transitions.len()) {
let declared = self.actions.get(chain).map_or(0, Vec::len);
let held = self
.transitions
.get(chain)
.and_then(Option::as_ref)
.map_or(0, |transition| transition.entries.len());
for entry in 0..declared.max(held) {
let id = WorkId { chain, entry };
if let Some(key) = self
.action_of(id)
.and_then(|action| self.effects.unsettled_for(action))
{
let action = key.action();
unresolved.push(PendingWork {
id,
key: Some(Box::new(key)),
hold: TransitionHold::UnsettledByRecovery { action },
});
continue;
}
if let Some(recording) = self
.action_of(id)
.and_then(|action| self.effects.recording_for(action))
{
unresolved.push(PendingWork {
id,
key: Some(Box::new(recording.key)),
hold: TransitionHold::RecordingFailed {
evidence: recording.evidence,
cause: recording.cause.to_string(),
},
});
continue;
}
let Some(transition) = self.transitions.get(chain).and_then(Option::as_ref) else {
continue;
};
let Some(state) = transition.entries.get(entry) else {
continue;
};
if !(state.is_open() || state.is_abandoned()) {
continue;
}
let Some(hold) = state.hold(budget) else {
continue;
};
unresolved.push(PendingWork {
id,
key: self.key_of(chain, entry, transition).map(Box::new),
hold,
});
}
}
unresolved
}
fn first_unresolved(&self, budget: u32) -> Option<PendingWork> {
self.unresolved(budget).into_iter().next()
}
fn pending(&self, budget: u32) -> Vec<PendingWork> {
self.unresolved(budget)
}
fn unresolved_count(&self, budget: u32) -> usize {
self.unresolved(budget).len()
}
}
fn revision_of(world: &World, chain: usize) -> u64 {
world
.resource::<Order>()
.0
.get(chain)
.and_then(|entity| world.get::<Revision>(*entity))
.map_or(0, |revision| revision.0)
}
fn take_observed(world: &mut World, chain: usize) -> Option<Erased> {
let taken = world
.non_send_mut::<Observed>()
.0
.get_mut(chain)
.and_then(Option::take);
if taken.is_some() {
world.non_send_mut::<Committed>().admit(chain);
}
taken
}
fn admits(world: &World, chain: usize, bound: Option<&Erased>) -> bool {
let seen = world.non_send::<Observed>();
let Some(next) = seen.0.get(chain).and_then(|slot| slot.as_ref()) else {
return false;
};
let Some(bound) = bound else {
return true;
};
let chains = world.non_send::<Chains>();
chains
.0
.get(chain)
.is_some_and(|chain| !(chain.same)(bound.as_any(), next.as_any()))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AdmitKind {
Keep,
Resume,
Open,
Idle,
}
#[derive(Clone, Copy)]
enum Admit {
Keep,
Bind(AdmittedInput),
}
impl Admit {
fn bind(kind: AdmitKind, input: Option<AdmittedInput>) -> Self {
match (kind, input) {
(AdmitKind::Resume | AdmitKind::Open, Some(input)) => Self::Bind(input),
_ => Self::Keep,
}
}
}
fn plan_admissions(world: &World, moving: &[bool], kinds: &mut [AdmitKind]) {
for (chain, kind) in kinds.iter_mut().enumerate() {
let ledger = world.non_send::<Ledger>();
let held = ledger.transitions.get(chain).and_then(Option::as_ref);
*kind = match held {
Some(transition) if transition.has_open() => AdmitKind::Keep,
Some(transition) => {
if admits(world, chain, transition.value.as_ref()) {
AdmitKind::Resume
} else {
AdmitKind::Keep
}
}
None if moving.get(chain).copied().is_some_and(|held| held)
&& world
.non_send::<Observed>()
.0
.get(chain)
.is_some_and(|slot| slot.is_some()) =>
{
AdmitKind::Open
}
None => AdmitKind::Idle,
};
}
}
fn plan_inputs(world: &mut World, kinds: &[AdmitKind], inputs: &mut [Option<AdmittedInput>]) {
let binds = |kind: &AdmitKind| matches!(kind, AdmitKind::Resume | AdmitKind::Open);
{
let chains = world.non_send::<Chains>();
let observed = world.non_send::<Observed>();
for (chain, kind) in kinds.iter().enumerate() {
if !binds(kind) {
continue;
}
let (Some(held), Some(value)) = (
chains.0.get(chain),
observed.0.get(chain).and_then(|slot| slot.as_ref()),
) else {
continue;
};
if let Some(slot) = inputs.get_mut(chain) {
*slot = Some((held.identify)(value.as_any()));
}
}
}
let mut ledger = world.non_send_mut::<Ledger>();
for (chain, kind) in kinds.iter().enumerate() {
if !binds(kind) {
continue;
}
if inputs.get(chain).is_some_and(Option::is_none) {
let minted = AdmittedInput {
identity: ledger.mint_input(),
event: false,
};
if let Some(slot) = inputs.get_mut(chain) {
*slot = Some(minted);
}
}
}
}
fn resume(
world: &mut World,
chain: usize,
admit: Admit,
held: Option<Transition>,
) -> Option<Transition> {
match (admit, held) {
(Admit::Keep, Some(transition)) => Some(transition),
(Admit::Keep, None) => None,
(Admit::Bind(input), Some(transition)) => {
let value = take_observed(world, chain);
Some(Transition::resumed(
input,
revision_of(world, chain),
&transition,
value,
))
}
(Admit::Bind(input), None) => {
let entries = world
.non_send::<Chains>()
.0
.get(chain)
.map_or(0, |chain| chain.entries.len());
let value = take_observed(world, chain);
Some(Transition::opened(
input,
revision_of(world, chain),
entries,
value,
))
}
}
}
fn failure_state(error: &BotError, attempts: u32, budget: u32) -> EntryState {
if let BotError::EffectUnrecorded {
ref key,
ref evidence,
ref cause,
} = *error
{
return EntryState::RecordingFailed {
key: **key,
evidence: *evidence,
cause: cause.to_string(),
attempts,
};
}
let cause = error.to_string();
match error.retry_class() {
RetryClass::Safe if attempts < budget => EntryState::DefinitelyFailed { attempts, cause },
RetryClass::Safe => EntryState::Abandoned {
reason: AbandonReason::AttemptsExhausted { attempts },
cause,
},
RetryClass::RequiresEvidence => EntryState::OutcomeUnknown { attempts, cause },
RetryClass::Never => EntryState::Abandoned {
reason: AbandonReason::Terminal,
cause,
},
}
}
fn observe_fold(world: &mut World) {
if parked(world) {
return;
}
let mut polled = std::mem::take(&mut world.non_send_mut::<Polled>().0);
let mut first_error = None;
let mut first_stall = None;
for (index, result) in polled.iter_mut().enumerate() {
let stalled = matches!(result, Err(BotError::PollStalled { .. }));
if stalled {
first_stall.get_or_insert(index);
*result = Ok(None);
} else if result.is_err() && first_error.is_none() {
first_error = std::mem::replace(result, Ok(None)).err();
}
}
match (first_stall, first_error) {
(Some(chain), _) => {
world.resource_mut::<TickError>().0 = Some(BotError::PollStalled {
chain,
deadline: world.resource::<PollBudget>().0,
});
}
(None, Some(error)) => {
world.resource_mut::<TickError>().0 = Some(error);
put_polled(world, polled);
return;
}
(None, None) => {}
}
{
let chains = world.non_send::<Chains>();
let mismatch = polled.iter().enumerate().find_map(|(index, result)| {
let next = result.as_ref().ok()?.as_ref()?;
let chain = chains.0.get(index)?;
(!chain.witness.agrees_with(next.witness)).then_some(BotError::TypeMismatch {
site: "observe_fold rendezvous",
chain: Some(index),
expected: chain.witness.name(),
observed: next.witness.name(),
})
});
if let Some(error) = mismatch {
world.resource_mut::<TickError>().0 = Some(error);
put_polled(world, polled);
return;
}
}
let mut changed = std::mem::take(&mut world.non_send_mut::<Moved>().0);
changed.clear();
changed.resize(polled.len(), false);
let mut superseded: Vec<SupersededObservation> = Vec::new();
let revisions = std::mem::take(&mut world.non_send_mut::<ChangeTicks>().polled);
{
let mut observed = world.non_send_mut::<Observed>();
for (index, result) in polled.iter_mut().enumerate() {
let moved = result
.as_mut()
.ok()
.and_then(Option::take)
.is_some_and(|value| {
if let Some(slot) = observed.0.get_mut(index) {
*slot = Some(value);
return true;
}
false
});
if let Some(flag) = changed.get_mut(index) {
*flag = moved;
}
}
}
for index in 0..changed.len() {
if changed.get(index).copied().is_some_and(|held| held) {
let reported = revisions.get(index).copied().flatten();
world.non_send_mut::<ChangeTicks>().commit(index, reported);
}
}
world.non_send_mut::<ChangeTicks>().polled = revisions;
put_polled(world, polled);
for (index, moved) in changed.iter().copied().enumerate() {
if !moved {
continue;
}
let Some(entity) = world.resource::<Order>().0.get(index).copied() else {
continue;
};
if let Some(mut revision) = world.get_mut::<Revision>(entity) {
revision.0 = revision.0.wrapping_add(1);
}
}
for index in changed
.iter()
.copied()
.enumerate()
.filter_map(|(index, moved)| moved.then_some(index))
{
if let Some(previous) = world.non_send::<Committed>().slot(index)
&& previous.state == SlotState::Unacted
{
superseded.push(SupersededObservation {
chain: index,
revision: previous.revision,
});
}
let revision = revision_of(world, index);
world.non_send_mut::<Committed>().commit(index, revision);
if let Some(slot) = world.non_send_mut::<Invalidated>().reasons.get_mut(index) {
*slot = None;
}
}
world.resource_mut::<TickReport>().superseded = superseded;
world.non_send_mut::<Moved>().0 = changed;
}
fn put_polled(world: &mut World, polled: Vec<Result<Option<Erased>, BotError>>) {
world.non_send_mut::<Polled>().0 = polled;
}
fn fire_plan(world: &mut World) {
let (mut steps, mut failures) = take_plan_buffers(world);
steps.clear();
failures.clear();
if parked(world) {
put_plan_buffers(world, steps, failures);
return;
}
let count = world.non_send::<Chains>().0.len();
let mut scratch = AdmitScratch::take(world, count);
let AdmitScratch {
ref mut moving,
ref mut kinds,
ref mut inputs,
} = scratch;
#[cfg(feature = "profile")]
let compare_charge = Charge::new(TickStage::Compare);
{
let mut query = world.query_filtered::<&SourceId, Changed<Revision>>();
for id in query.iter(world) {
if let Some(flag) = moving.get_mut(id.chain) {
*flag = true;
}
}
}
plan_admissions(world, moving, kinds);
#[cfg(feature = "profile")]
drop(compare_charge);
#[cfg(feature = "profile")]
let fingerprint_charge = Charge::new(TickStage::Fingerprint);
plan_inputs(world, kinds, inputs);
#[cfg(feature = "profile")]
drop(fingerprint_charge);
failures.resize_with(count, || None);
#[cfg(feature = "profile")]
let decide_charge = Charge::new(TickStage::Decide);
for index in 0..count {
let held = world.non_send_mut::<Ledger>().take(index);
let admit = Admit::bind(
match kinds.get(index).copied() {
Some(kind) => kind,
None => AdmitKind::Idle,
},
inputs.get(index).copied().flatten(),
);
let Some(mut transition) = resume(world, index, admit, held) else {
continue;
};
{
let chains = world.non_send::<Chains>();
let ledger = world.non_send::<Ledger>();
if let (Some(chain), Some(value), Some(failure)) = (
chains.0.get(index),
transition.value.as_ref(),
failures.get_mut(index),
) {
plan_chain(
chain,
&transition,
value,
index,
ledger,
&mut steps,
failure,
);
}
}
let retained = transition.is_retained();
let returned = if retained {
None
} else {
std::mem::take(&mut transition.value)
};
world
.non_send_mut::<Ledger>()
.put(index, if retained { Some(transition) } else { None });
if let Some(returned) = returned {
let mut observed = world.non_send_mut::<Observed>();
if let Some(slot) = observed.0.get_mut(index).filter(|slot| slot.is_none()) {
*slot = Some(returned);
}
}
}
#[cfg(feature = "profile")]
drop(decide_charge);
*world.non_send_mut::<AdmitScratch>() = scratch;
put_plan_buffers(world, steps, failures);
}
fn take_plan_buffers(world: &mut World) -> (Vec<Step>, Vec<Option<BotError>>) {
let mut plan = world.non_send_mut::<Plan>();
(
std::mem::take(&mut plan.steps),
std::mem::take(&mut plan.failures),
)
}
fn put_plan_buffers(world: &mut World, steps: Vec<Step>, failures: Vec<Option<BotError>>) {
let mut plan = world.non_send_mut::<Plan>();
plan.steps = steps;
plan.failures = failures;
}
fn plan_chain(
chain: &EcsChain,
transition: &Transition,
value: &Erased,
index: usize,
ledger: &Ledger,
steps: &mut Vec<Step>,
failure: &mut Option<BotError>,
) {
for (entry_index, entry) in chain.entries.iter().enumerate() {
let Some(state) = transition.entries.get(entry_index) else {
continue;
};
if ledger
.action_of(WorkId {
chain: index,
entry: entry_index,
})
.is_some_and(|action| ledger.effects.blocks(action))
{
break;
}
match *state {
EntryState::Succeeded | EntryState::Skipped => continue,
EntryState::Abandoned { .. } => break,
EntryState::Unrecorded | EntryState::OutcomeUnknown { .. } => break,
EntryState::RecordingFailed { .. } => {}
EntryState::NotStarted | EntryState::DefinitelyFailed { .. } => {}
}
match entry.condition.check_any(value) {
Ok(true) => {}
Ok(false) => {
steps.push(Step {
chain: index,
entry: entry_index,
decision: Decision::Skip,
});
continue;
}
Err(error) => {
*failure = Some(error);
break;
}
}
steps.push(Step {
chain: index,
entry: entry_index,
decision: Decision::Attempt,
});
}
}
fn schedule() -> Schedule {
let mut schedule = Schedule::default();
schedule.set_executor(SingleThreadedExecutor::default());
schedule.set_build_settings(ScheduleBuildSettings {
ambiguity_detection: LogLevel::Error,
..Default::default()
});
schedule.add_systems((observe_fold, fire_plan).chain());
schedule
}
fn validate(schedule: &mut Schedule, world: &mut World) -> Result<(), BotError> {
schedule
.initialize(world)
.map(|_| ())
.map_err(|error| BotError::DomainError {
domain: "ecs::schedule".into(),
certainty: DispatchCertainty::Refused,
cause: error.to_string(schedule.graph(), world),
})
}
pub struct EcsBot {
name: String,
world: World,
schedule: Schedule,
poll_deadline: Duration,
}
impl EcsBot {
#[must_use]
pub fn builder(name: impl Into<String>) -> EcsBuilder {
EcsBuilder {
name: name.into(),
chains: Vec::new(),
policy: RetryPolicy::DEFAULT,
effects: None,
poll_deadline: DEFAULT_POLL_DEADLINE,
}
}
pub fn from_spec(
spec: &BotSpec,
registry: &DomainRegistry,
grants: &GrantSet,
effects: EffectScope,
) -> Result<Self, Admission> {
registry.validate().map_err(Admission::Refused)?;
if spec.version != BotSpec::CURRENT_VERSION {
let refusal = Err(Admission::Refused(BotError::UnsupportedSpecVersion {
found: spec.version,
supported: BotSpec::CURRENT_VERSION,
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "from_spec: returning an error to the caller");
return refusal;
}
if spec.chains.is_empty() {
let refusal = Err(Admission::Refused(BotError::IncompleteSpec {
field: "chains",
cause: String::from("the document declares no observation chains at all"),
}));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "from_spec: returning an error to the caller");
return refusal;
}
let mut needs: Vec<Need> = Vec::new();
let mut chains: Vec<EcsChain> = Vec::with_capacity(spec.chains.len());
for (chain_index, chain) in spec.chains.iter().enumerate() {
if let Some(built) = materialize_chain(chain_index, chain, registry, grants, &mut needs)
{
chains.push(built);
}
}
if !needs.is_empty() {
let refusal = Err(Admission::Needs(NeedSet::new(needs)));
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "from_spec: returning an error to the caller");
return refusal;
}
let mut builder = Self::builder(spec.name.clone());
for chain in chains {
builder = builder.chain(chain);
}
builder
.with_effects(effects)
.build(grants)
.map_err(Admission::Refused)
}
#[must_use]
pub fn into_journal(mut self) -> Option<Box<dyn EffectJournal>> {
self.world
.remove_non_send::<Ledger>()
.map(Ledger::into_journal)
}
#[must_use]
pub fn name(&self) -> &str {
&self.name
}
#[must_use]
pub fn fired(&self) -> usize {
self.world.resource::<Fired>().0
}
#[must_use]
pub fn world(&self) -> &World {
&self.world
}
#[must_use]
pub fn revisions(&self) -> Vec<u64> {
let order = self.world.resource::<Order>().0.clone();
order
.iter()
.filter_map(|entity| self.world.get::<Revision>(*entity))
.map(|revision| revision.0)
.collect()
}
#[must_use]
pub fn source_domains(&self) -> Vec<String> {
let order = self.world.resource::<Order>().0.clone();
order
.iter()
.filter_map(|entity| self.world.get::<SourceId>(*entity))
.map(|id| id.domain().to_owned())
.collect()
}
pub async fn tick_async(&mut self) -> Result<usize, BotError> {
self.begin_tick();
let mut polled = std::mem::take(&mut self.world.non_send_mut::<Polled>().0);
polled.clear();
let mut revisions = std::mem::take(&mut self.world.non_send_mut::<ChangeTicks>().polled);
let mut scratch = {
let mut guard = self.world.non_send_mut::<PollScratch>();
std::mem::take(&mut *guard)
};
#[cfg(feature = "profile")]
let poll_charge = Charge::new(TickStage::Poll);
let watchdogs = self
.poll_sources(&mut polled, &mut revisions, &mut scratch)
.await;
#[cfg(feature = "profile")]
drop(poll_charge);
self.world.non_send_mut::<Polled>().0 = polled;
self.world.non_send_mut::<ChangeTicks>().polled = revisions;
self.record_invalidations(&scratch.declared);
self.publish_refreshes();
self.publish_stalls(&scratch.stalled, watchdogs);
*self.world.non_send_mut::<PollScratch>() = scratch;
#[cfg(feature = "profile")]
let schedule_charge = Charge::new(TickStage::Schedule);
self.schedule.run(&mut self.world);
#[cfg(feature = "profile")]
drop(schedule_charge);
let plan = {
let mut guard = self.world.non_send_mut::<Plan>();
std::mem::take(&mut *guard)
};
let Plan {
mut steps,
mut failures,
} = plan;
let budget = self.world.resource::<Policy>().0.max_attempts();
#[cfg(feature = "profile")]
let act_charge = Charge::new(TickStage::Act);
let (fired, failure) = self.run_steps(&steps, &mut failures, budget).await;
#[cfg(feature = "profile")]
drop(act_charge);
steps.clear();
failures.clear();
put_plan_buffers(&mut self.world, steps, failures);
self.world.resource_mut::<Fired>().0 = fired;
self.world.resource_mut::<TickReport>().fired = fired;
if let Some(error) = failure {
self.world.resource_mut::<TickError>().0 = Some(error);
}
if let Some(error) = self.world.resource_mut::<TickError>().0.take() {
let refusal = Err(error);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "tick_async: returning an error to the caller");
return refusal;
}
let ledger = self.world.non_send::<Ledger>();
match ledger.first_unresolved(budget) {
Some(work) => Err(BotError::PendingTransition {
work,
outstanding: ledger.unresolved_count(budget),
}),
None => Ok(self.world.resource::<Fired>().0),
}
}
#[cfg(feature = "profile")]
pub async fn tick_profiled(&mut self) -> (Result<usize, BotError>, TickProfile) {
profile::begin();
let outcome = self.tick_async().await;
(outcome, profile::end())
}
pub fn tick(&mut self) -> Result<usize, BotError> {
#[cfg(feature = "rt")]
if lgwks_deps::tokio::runtime::Handle::try_current().is_ok() {
let refusal = Err(BotError::TickInsideRuntime);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "tick: returning an error to the caller");
return refusal;
}
lgwks_std::task::block_on(self.tick_async())
}
fn begin_tick(&mut self) {
self.world.resource_mut::<Fired>().0 = 0;
self.world.resource_mut::<TickError>().0 = None;
*self.world.resource_mut::<TickReport>() = TickReport::default();
}
#[must_use]
pub fn tick_report(&self) -> TickReport {
self.world.resource::<TickReport>().clone()
}
async fn poll_sources(
&self,
polled: &mut Vec<Result<Option<Erased>, BotError>>,
revisions: &mut Vec<Option<u64>>,
scratch: &mut PollScratch,
) -> u32 {
let count = self.world.non_send::<Chains>().0.len();
polled.clear();
polled.resize_with(count, || Ok(None));
revisions.clear();
revisions.resize(count, None);
self.declared_refresh(&mut scratch.declared, count);
scratch.stalled.clear();
let PollScratch {
ref mut declared,
ref mut stalled,
ref mut slots,
} = *scratch;
let mut watchdogs = 0_u32;
let chains = self.world.non_send::<Chains>();
for wave_index in 0..count.div_ceil(MAX_IN_FLIGHT_POLLS) {
let base = wave_index.saturating_mul(MAX_IN_FLIGHT_POLLS);
let width = count.saturating_sub(base).min(MAX_IN_FLIGHT_POLLS);
slots.clear();
{
let ticks = self.world.non_send::<ChangeTicks>();
for slot in 0..width {
let index = base.saturating_add(slot);
let Some(chain) = chains.0.get(index) else {
continue;
};
let unsound = declared
.get(index)
.copied()
.flatten()
.is_some_and(|reason| reason.invalidates_baseline());
let reported = chain.source.revision();
if let Some(slot) = revisions.get_mut(index) {
*slot = reported;
}
if !(ticks.quiet(index, reported) && !unsound) {
slots.push(slot);
}
}
}
if slots.is_empty() {
continue;
}
let wave_watchdog = PollWatchdog::new(self.poll_deadline, slots.len());
let (results, armed) = {
let mut polls: Vec<WavePoll<'_>> = Vec::new();
for &slot in slots.iter() {
{
let index = base.saturating_add(slot);
let chain = &chains.0[index];
let seen = self.world.non_send::<Observed>();
let ledger = self.world.non_send::<Ledger>();
let grants = self.world.resource::<Grants>();
let unsound = declared
.get(index)
.copied()
.flatten()
.is_some_and(|reason| reason.invalidates_baseline());
let baseline = if unsound {
None
} else {
seen.0
.get(index)
.and_then(|slot| slot.as_ref())
.or_else(|| ledger.bound(index))
};
polls.push(WavePoll {
chain: index,
slot,
installed: false,
poll: Box::pin(chain.source.poll_any(&grants.0, baseline)),
watchdog: &wave_watchdog,
});
}
}
bounded_wave(self.poll_deadline, &wave_watchdog, polls.drain(..)).await
};
watchdogs = watchdogs.saturating_add(u32::from(armed));
for (&slot, result) in slots.iter().zip(results) {
let index = base.saturating_add(slot);
let Some(slot) = polled.get_mut(index) else {
continue;
};
if matches!(result, Err(BotError::PollStalled { .. })) {
stalled.push(index);
}
*slot = result;
}
}
for (index, reason) in declared.iter_mut().enumerate() {
let Some(chain) = chains.0.get(index) else {
continue;
};
match (
chain.source.cache_state(),
polled.get(index).map(Result::is_ok),
) {
(Some(declared_now), Some(true)) => *reason = Some(declared_now),
(_, Some(false)) if reason.is_none() => {
*reason = Some(RefreshReason::Disconnected);
}
_ => {}
}
}
stalled.sort_unstable();
watchdogs
}
fn declared_refresh(&self, declared: &mut Vec<Option<RefreshReason>>, count: usize) {
let chains = self.world.non_send::<Chains>();
let invalidated = self.world.non_send::<Invalidated>();
declared.clear();
declared.resize(count, None);
for (index, chain) in chains.0.iter().enumerate() {
if invalidated.is_invalid(index) {
continue;
}
if let Some(reason) = chain.source.cache_state() {
declared[index] = Some(reason);
}
}
}
fn record_invalidations(&mut self, declared: &[Option<RefreshReason>]) {
let mut invalidated = self.world.non_send_mut::<Invalidated>();
for (chain, reason) in declared.iter().enumerate() {
if let Some(reason) = *reason {
invalidated.mark(chain, reason);
}
}
}
fn publish_refreshes(&mut self) {
let forced: Vec<ForcedRefresh> = {
let invalidated = self.world.non_send::<Invalidated>();
let order = self.world.resource::<Order>();
invalidated
.reasons
.iter()
.enumerate()
.filter_map(|(chain, reason)| {
let reason = (*reason)?;
let domain = order
.0
.get(chain)
.and_then(|entity| self.world.get::<SourceId>(*entity))
.map_or_else(String::new, |id| id.domain().to_owned());
Some(ForcedRefresh {
chain,
domain,
reason,
})
})
.collect()
};
self.world.resource_mut::<TickReport>().forced = forced;
}
fn publish_stalls(&mut self, stalled: &[usize], watchdogs: u32) {
let rows: Vec<StalledSource> = {
let order = self.world.resource::<Order>();
stalled
.iter()
.map(|chain| {
let domain = order
.0
.get(*chain)
.and_then(|entity| self.world.get::<SourceId>(*entity))
.map_or_else(String::new, |id| id.domain().to_owned());
StalledSource {
chain: *chain,
domain,
deadline: self.poll_deadline,
}
})
.collect()
};
let mut report = self.world.resource_mut::<TickReport>();
report.stalled = rows;
report.watchdogs = watchdogs;
}
async fn run_steps(
&mut self,
steps: &[Step],
failures: &mut [Option<BotError>],
budget: u32,
) -> (usize, Option<BotError>) {
let mut fired: usize = 0;
let mut failure: Option<BotError> = None;
let mut start = 0;
for chain in 0..failures.len() {
let end = steps[start..]
.iter()
.position(|step| step.chain != chain)
.map_or(steps.len(), |offset| start.saturating_add(offset));
let run = match steps.get(start) {
Some(step) if step.chain == chain => &steps[start..end],
_ => &[],
};
let condition_failure = failures.get_mut(chain).and_then(Option::take);
let (chain_fired, chain_failure) = self.run_chain(run, condition_failure, budget).await;
fired = fired.saturating_add(chain_fired);
if failure.is_none() {
failure = chain_failure;
}
start = end;
}
(fired, failure)
}
async fn run_chain(
&mut self,
steps: &[Step],
condition_failure: Option<BotError>,
budget: u32,
) -> (usize, Option<BotError>) {
let mut fired: usize = 0;
let mut failure: Option<BotError> = None;
for step in steps {
let work = WorkId {
chain: step.chain,
entry: step.entry,
};
if matches!(step.decision, Decision::Skip) {
let blocked = {
let ledger = self.world.non_send::<Ledger>();
ledger
.action_of(work)
.is_some_and(|action| ledger.effects.blocks(action))
};
if blocked {
break;
}
self.world.non_send_mut::<Ledger>().skip(work);
continue;
}
if let Some((key, evidence, attempts)) =
self.world.non_send::<Ledger>().recording_failed(work)
{
let result = self
.world
.non_send_mut::<Ledger>()
.effects
.ensure_outcome(key, evidence);
match result {
Ok(()) => {
if self
.world
.non_send_mut::<Ledger>()
.finish_recording(work, evidence, attempts, budget)
&& evidence == EffectEvidence::Applied
{
fired = fired.saturating_add(1);
}
continue;
}
Err(cause) => {
let error = BotError::EffectUnrecorded {
key: Box::new(key),
evidence,
cause: Box::new(DispatchError::Journal(cause)),
};
self.world
.non_send_mut::<Ledger>()
.hold_recording(work, evidence, &error);
failure = Some(error);
break;
}
}
}
if let Some(action) = self.world.non_send::<Ledger>().action_of(work) {
let ledger = self.world.non_send::<Ledger>();
if ledger.generation(step.chain).is_some_and(|(input, event)| {
ledger
.effects
.applied_in(action, step.chain, step.entry, input, event)
}) {
self.world.non_send_mut::<Ledger>().skip(work);
continue;
}
if ledger.effects.blocks(action) {
break;
}
}
let lifetime = self
.world
.non_send::<Chains>()
.0
.get(step.chain)
.and_then(|chain| chain.entries.get(step.entry))
.map_or(EffectLifetime::External, |entry| {
entry.action.effect_lifetime()
});
let (key, authority) = match self
.world
.non_send_mut::<Ledger>()
.begin_attempt(work, lifetime)
.await
{
Ok(pair) => pair,
Err(error) => {
failure = Some(error);
break;
}
};
if let Err(cause) = self
.world
.non_send::<Ledger>()
.effects
.broker()
.revalidate(&authority)
{
let error = BotError::EffectRefused {
cause: DispatchError::Broker(cause),
};
self.world
.non_send_mut::<Ledger>()
.fail(work, &error, key.attempt(), budget);
failure = Some(error);
break;
}
let outcome = {
let chains = self.world.non_send::<Chains>();
let grants = self.world.resource::<Grants>();
let bound = self.world.non_send::<Ledger>().bound(step.chain);
match (chains.0.get(step.chain), bound) {
(Some(chain), Some(value)) => match chain.entries.get(step.entry) {
Some(entry) => Some(entry.action.run_any(&grants.0, value).await),
None => None,
},
_ => None,
}
};
let Some(outcome) = outcome else {
break;
};
let observed = match outcome {
Ok(_) => Some(EffectEvidence::Applied),
Err(ref error) => match error.dispatch_certainty() {
DispatchCertainty::Refused | DispatchCertainty::NotDelivered => {
Some(EffectEvidence::NotApplied)
}
DispatchCertainty::Unsettled => None,
DispatchCertainty::Occurred => None,
},
};
if let Some(evidence) = observed
&& let Err(cause) = self
.world
.non_send_mut::<Ledger>()
.effects
.ensure_outcome(key, evidence)
{
let error = BotError::EffectUnrecorded {
key: Box::new(key),
evidence,
cause: Box::new(DispatchError::Journal(cause)),
};
self.world
.non_send_mut::<Ledger>()
.fail(work, &error, key.attempt(), budget);
failure = Some(error);
break;
}
match outcome {
Ok(_) => {
if self.world.non_send_mut::<Ledger>().succeed(work) {
fired = fired.saturating_add(1);
}
}
Err(error) => {
self.world
.non_send_mut::<Ledger>()
.fail(work, &error, key.attempt(), budget);
failure = Some(error);
break;
}
}
}
(fired, failure.or(condition_failure))
}
#[must_use]
pub fn pending(&self) -> Vec<PendingWork> {
let budget = self.world.resource::<Policy>().0.max_attempts();
self.world.non_send::<Ledger>().pending(budget)
}
pub fn resolve_effect(
&mut self,
key: &EffectKey,
evidence: EffectEvidence,
) -> Result<(), BotError> {
let recovered = self
.world
.non_send::<Ledger>()
.effects
.unsettled_for(key.action())
.is_some_and(|held| held == *key);
if recovered {
return self
.world
.non_send_mut::<Ledger>()
.effects
.settle_recovered(*key, evidence)
.map_err(|cause| BotError::EffectRefused {
cause: DispatchError::Journal(cause),
});
}
let recording = self
.world
.non_send::<Ledger>()
.effects
.recording_for(key.action())
.is_some_and(|held| held.key == *key && held.evidence == evidence);
if recording {
return self
.world
.non_send_mut::<Ledger>()
.effects
.settle_recording(*key, evidence)
.map_err(|cause| BotError::EffectUnrecorded {
key: Box::new(*key),
evidence,
cause: Box::new(DispatchError::Journal(cause)),
});
}
match self
.world
.non_send::<Ledger>()
.classify_settlement(key, evidence)
{
Err(cause) => Err(BotError::EffectRefused {
cause: DispatchError::Journal(cause),
}),
Ok(verdict @ (Settled::Decided | Settled::Duplicate)) => {
self.world
.non_send_mut::<Ledger>()
.effects
.ensure_outcome(*key, evidence)
.map_err(|cause| BotError::EffectUnrecorded {
key: Box::new(*key),
evidence,
cause: Box::new(DispatchError::Journal(cause)),
})?;
if verdict == Settled::Decided {
self.world
.non_send_mut::<Ledger>()
.apply_settlement(key, evidence);
}
Ok(())
}
Ok(Settled::Superseded { id, current }) => Err(BotError::EvidenceSuperseded {
work: id,
named: key.digest(),
current,
}),
Ok(Settled::StaleAttempt {
id,
reported,
outstanding,
}) => Err(BotError::EvidenceStaleAttempt {
work: id,
reported,
outstanding,
}),
Ok(Settled::Contradicted { id, settled }) => Err(BotError::EvidenceContradicted {
work: id,
settled,
submitted: evidence,
}),
Ok(Settled::NoSuchWork { id: Some(id) }) => Err(BotError::NoSuchWork { work: id }),
Ok(Settled::NoSuchWork { id: None }) => Err(BotError::ActionNotDeclared {
action: key.action(),
}),
}
}
}
fn materialize_chain(
chain_index: usize,
chain: &ChainSpec,
registry: &DomainRegistry,
grants: &GrantSet,
needs: &mut Vec<Need>,
) -> Option<EcsChain> {
let before = needs.len();
let Some(source_ctor) = registry.source(&chain.source) else {
needs.push(Need::UnknownSource {
chain: chain_index,
domain: chain.source.clone(),
});
return None;
};
let source = match source_ctor(&chain.target) {
Ok(source) => source,
Err(cause) => {
needs.push(Need::SourceTargetRejected {
chain: chain_index,
domain: chain.source.clone(),
cause: cause.to_string(),
});
return None;
}
};
for shortage in grants.uncovered(source.required_caps(), &Demand::new(chain.source.clone())) {
needs.push(Need::MissingCapability {
chain: chain_index,
action: None,
domain: chain.source.clone(),
capability: shortage.required().clone(),
});
}
let mut entries: Vec<ChainEntry> = Vec::with_capacity(chain.on.len());
for (action_index, entry) in chain.on.iter().enumerate() {
let condition_id: &String = &entry.0;
let action_spec: &ActionSpec = &entry.1;
let action = match registry.action(&action_spec.domain) {
Some(ctor) => match ctor(&action_spec.target) {
Ok(action) => Some(action),
Err(cause) => {
needs.push(Need::ActionTargetRejected {
chain: chain_index,
action: action_index,
domain: action_spec.domain.clone(),
cause: cause.to_string(),
});
None
}
},
None => {
needs.push(Need::UnknownAction {
chain: chain_index,
action: action_index,
domain: action_spec.domain.clone(),
});
None
}
};
if let Some(ref action) = action {
for shortage in grants.uncovered(
action.required_caps(),
&Demand::new(action_spec.domain.clone()),
) {
needs.push(Need::MissingCapability {
chain: chain_index,
action: Some(action_index),
domain: action_spec.domain.clone(),
capability: shortage.required().clone(),
});
}
}
let condition = match source.condition(condition_id) {
Ok(condition) => Some(condition),
Err(_) => {
needs.push(Need::UnknownCondition {
chain: chain_index,
action: action_index,
condition: condition_id.clone(),
});
None
}
};
if let (Some(action), Some(condition)) = (action, condition) {
entries.push(ChainEntry::erased(
condition.into_evaluate_any(),
action.into_execute_any(),
));
}
}
if needs.len() == before {
Some(EcsChain::from_registry(source, entries))
} else {
None
}
}
pub struct EcsBuilder {
name: String,
chains: Vec<EcsChain>,
policy: RetryPolicy,
effects: Option<EffectScope>,
poll_deadline: Duration,
}
impl EcsBuilder {
#[must_use]
pub fn with_poll_deadline(mut self, deadline: Duration) -> Self {
self.poll_deadline = deadline;
self
}
#[must_use]
pub fn with_retry_policy(mut self, policy: RetryPolicy) -> Self {
self.policy = policy;
self
}
#[must_use]
pub fn with_effects(mut self, effects: EffectScope) -> Self {
self.effects = Some(effects);
self
}
#[must_use]
pub(crate) fn chain(mut self, chain: EcsChain) -> Self {
self.chains.push(chain);
self
}
#[must_use]
pub fn observe<S>(self, source: S) -> EcsObserveBuilder<S>
where
S: Observe,
{
EcsObserveBuilder {
name: self.name,
prior: self.chains,
source,
entries: Vec::new(),
policy: self.policy,
effects: self.effects,
poll_deadline: self.poll_deadline,
}
}
pub fn build(self, grants: &GrantSet) -> Result<EcsBot, BotError> {
EcsBot::assemble(
self.name,
self.chains,
grants,
self.policy,
self.effects,
self.poll_deadline,
)
}
}
impl<S: Observe> EcsObserveBuilder<S> {
#[must_use]
pub fn with_effects(mut self, effects: EffectScope) -> Self {
self.effects = Some(effects);
self
}
#[must_use]
pub fn on<C, A>(mut self, condition: C, action: A) -> Self
where
S::Output: 'static,
C: Evaluate<S::Output> + 'static,
A: Execute<Input = S::Output> + 'static,
A::Output: 'static,
{
self.entries
.push(typed_entry::<C, A, S::Output>(condition, action));
self
}
#[must_use]
pub fn observe<U>(self, source: U) -> EcsObserveBuilder<U>
where
S: 'static,
S::Output: PartialEq + InputIdentity + 'static,
U: Observe,
{
let Self {
name,
mut prior,
source: previous,
entries,
policy,
effects,
poll_deadline,
} = self;
prior.push(EcsChain {
source: Box::new(previous),
same: same_output::<S>,
identify: identify_output::<S>,
witness: Witness::of::<S::Output>(),
entries,
});
EcsObserveBuilder {
name,
prior,
source,
entries: Vec::new(),
policy,
effects,
poll_deadline,
}
}
pub fn build(self, grants: &GrantSet) -> Result<EcsBot, BotError>
where
S: 'static,
S::Output: PartialEq + InputIdentity + 'static,
{
let Self {
name,
mut prior,
source,
entries,
policy,
effects,
poll_deadline,
} = self;
prior.push(EcsChain {
source: Box::new(source),
same: same_output::<S>,
identify: identify_output::<S>,
witness: Witness::of::<S::Output>(),
entries,
});
EcsBot::assemble(name, prior, grants, policy, effects, poll_deadline)
}
}
pub struct EcsObserveBuilder<S> {
name: String,
effects: Option<EffectScope>,
prior: Vec<EcsChain>,
source: S,
entries: Vec<ChainEntry>,
policy: RetryPolicy,
poll_deadline: Duration,
}
impl EcsBot {
fn assemble(
name: String,
chains: Vec<EcsChain>,
grants: &GrantSet,
policy: RetryPolicy,
effects: Option<EffectScope>,
poll_deadline: Duration,
) -> Result<Self, BotError> {
if name.is_empty() {
let refusal = Err(BotError::IncompleteSpec {
field: "name",
cause: String::from("the bot was given an empty name"),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
if poll_deadline.is_zero() {
let refusal = Err(BotError::PollDeadlineUnbounded {
deadline: poll_deadline,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
if poll_deadline > MAX_POLL_DEADLINE {
let refusal = Err(BotError::PollDeadlineExceeded {
deadline: poll_deadline,
ceiling: MAX_POLL_DEADLINE,
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
let Some(effects) = effects else {
let refusal = Err(BotError::IncompleteSpec {
field: "effects",
cause: String::from(
"no effect scope was supplied, so a dispatch could not be recorded",
),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
};
let committed = effects
.journal()
.committed()
.map_err(|cause| BotError::EffectRefused {
cause: DispatchError::Journal(cause),
})?;
let tail = effects.journal().tail();
let recovered_events = match u64::try_from(committed.len()) {
Ok(count) => count,
Err(_too_many_events) => {
let refusal = Err(BotError::EffectRefused {
cause: DispatchError::Journal(JournalError::Exhausted),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
};
if recovered_events != tail.sequence() {
let refusal = Err(BotError::EffectRefused {
cause: DispatchError::Journal(JournalError::SnapshotStale {
recovered_events,
committed_events: tail.sequence(),
}),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
let recovered = recover(committed.iter());
let mut shortages: Vec<Shortage> = Vec::new();
for chain in &chains {
shortages.extend(grants.uncovered(
chain.source.required_caps(),
&Demand::new(chain.source.domain_id()),
));
for entry in &chain.entries {
shortages.extend(grants.uncovered(
entry.action.required_caps(),
&Demand::new(entry.action.domain_id()),
));
}
}
if let Some(deficit) = Deficit::from_shortages(shortages) {
let refusal = Err(BotError::CapabilityDenied { deficit });
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
let mut world = World::new();
world.insert_resource(Grants(grants.clone()));
world.insert_resource(Fired::default());
world.insert_resource(TickError::default());
world.insert_resource(Policy(policy));
world.insert_resource(PollBudget(poll_deadline));
let mut order = Vec::with_capacity(chains.len());
for (index, chain) in chains.iter().enumerate() {
order.push(
world
.spawn((
SourceId {
chain: index,
domain: chain.source.domain_id().to_owned(),
},
Revision::default(),
))
.id(),
);
}
world.insert_resource(Order(order));
let count = chains.len();
let actions: Vec<Vec<ActionId>> = chains
.iter()
.enumerate()
.map(|(chain, held)| {
held.entries
.iter()
.enumerate()
.map(|(entry, declared)| {
derive_action_id(&name, chain, entry, declared.action.domain_id())
})
.collect()
})
.collect();
let next_input = tail.sequence().saturating_add(1);
let mut ledger = Ledger::new(actions, Effects::new(effects, tail), next_input);
for attempt in recovered.attempts() {
let key = attempt.key();
let Some(id) = ledger.locate(key.action()) else {
let refusal = Err(BotError::ActionNotDeclared {
action: key.action(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
};
let identity = ledger.effects.identity();
if key.run() != identity.run()
|| key.flow() != identity.flow()
|| key.environment() != identity.environment()
{
let refusal = Err(BotError::ActionNotDeclared {
action: key.action(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "assemble: returning an error to the caller");
return refusal;
}
let _ = id;
ledger
.effects
.note_journal_attempt(key.action(), key.attempt());
match attempt.status() {
AttemptStatus::OutcomeUnknown => ledger.effects.unsettled.push(key),
AttemptStatus::Applied
| AttemptStatus::Verified
| AttemptStatus::VerificationFailed => {
if ledger.effects.scope.journal().durability() == DurabilityPromise::Ephemeral {
ledger.effects.note_applied(key);
} else if let Err(cause) = ledger
.effects
.confirm_recorded_outcome(key, EffectEvidence::Applied)
{
ledger.effects.recording.push(RecordedOutcome {
key,
evidence: EffectEvidence::Applied,
cause,
});
} else {
ledger.effects.note_applied(key);
}
}
AttemptStatus::Prepared => {}
AttemptStatus::NotApplied => {
if ledger.effects.scope.journal().durability() != DurabilityPromise::Ephemeral
&& let Err(cause) = ledger
.effects
.confirm_recorded_outcome(key, EffectEvidence::NotApplied)
{
ledger.effects.recording.push(RecordedOutcome {
key,
evidence: EffectEvidence::NotApplied,
cause,
});
}
}
}
}
world.insert_non_send(Chains(chains));
world.insert_non_send(Observed((0..count).map(|_| None).collect()));
world.insert_non_send(ledger);
world.insert_non_send(Polled::default());
world.insert_non_send(ChangeTicks {
seen: vec![None; count],
polled: vec![None; count],
});
world.insert_non_send(AdmitScratch::default());
world.insert_non_send(Moved::default());
world.insert_non_send(Plan::default());
world.insert_non_send(PollScratch::default());
world.insert_non_send(Invalidated {
reasons: vec![None; count],
});
world.insert_non_send(Committed {
slots: vec![SlotAdmission::default(); count],
});
world.insert_resource(TickReport::default());
let mut schedule = schedule();
validate(&mut schedule, &mut world)?;
Ok(Self {
name,
world,
schedule,
poll_deadline,
})
}
}
impl core::fmt::Debug for EcsBot {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("EcsBot")
.field("name", &self.name)
.finish_non_exhaustive()
}
}
impl core::fmt::Debug for EcsBuilder {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("EcsBuilder")
.field("name", &self.name)
.field("chains", &self.chains.len())
.field("policy", &self.policy)
.finish()
}
}
impl<S> core::fmt::Debug for EcsObserveBuilder<S> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("EcsObserveBuilder")
.field("name", &self.name)
.field("prior", &self.prior.len())
.field("entries", &self.entries.len())
.field("policy", &self.policy)
.finish_non_exhaustive()
}
}
#[cfg(test)]
pub(super) mod tests {
use std::cell::{Cell, RefCell};
use std::future::Future;
use std::marker::PhantomData;
use std::rc::Rc;
use std::task::{Context, Poll, Waker};
use super::*;
use crate::EventId;
use crate::cap::{Auth, Cap, Demand};
use crate::effect::{EnvironmentId, RunId};
use crate::journal::{DurabilityPromise, DurableAck, JournalPosition, MemoryJournal};
use std::io;
type TestResult = Result<(), Box<dyn std::error::Error>>;
fn domain_refusal<T>(domain: &str, cause: &str) -> Result<T, BotError> {
let refusal = Err(BotError::DomainError {
domain: domain.into(),
certainty: DispatchCertainty::NotDelivered,
cause: cause.into(),
});
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "a test source refused its poll");
refusal
}
struct Script {
values: Rc<RefCell<Vec<u16>>>,
cursor: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Script {
fn new(values: Vec<u16>) -> Self {
Self {
values: Rc::new(RefCell::new(values)),
cursor: Rc::new(Cell::new(0)),
caps: vec![Cap::net()],
}
}
}
impl Observe for Script {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
let index = self.cursor.get();
self.cursor.set(index.saturating_add(1));
match self.values.borrow().get(index).copied() {
Some(value) => Ok(value),
None => domain_refusal("test::script", "polled past the end of its script"),
}
}
fn domain_id(&self) -> &str {
"test::script"
}
}
struct Holds {
value: u16,
caps: Vec<Cap>,
}
impl Holds {
fn new(value: u16) -> Self {
Self {
value,
caps: vec![Cap::net()],
}
}
}
impl Observe for Holds {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
Ok(self.value)
}
fn domain_id(&self) -> &str {
"test::holds"
}
}
struct Counted {
value: Rc<Cell<u16>>,
polls: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Counted {
fn new(value: u16, polls: Rc<Cell<usize>>) -> Self {
Self {
value: Rc::new(Cell::new(value)),
polls,
caps: vec![Cap::net()],
}
}
}
impl Observe for Counted {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
self.polls.set(self.polls.get().saturating_add(1));
Ok(self.value.get())
}
fn fingerprint(&self) -> Option<u128> {
Some(u128::from(self.value.get()))
}
fn domain_id(&self) -> &str {
"test::counted"
}
}
struct Rev {
value: Rc<Cell<u64>>,
rev: Rc<Cell<u64>>,
polls: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Observe for Rev {
type Output = u64;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u64, BotError> {
call.0.check(&self.caps)?;
self.polls.set(self.polls.get().saturating_add(1));
Ok(self.value.get())
}
fn revision(&self) -> Option<u64> {
Some(self.rev.get())
}
fn domain_id(&self) -> &str {
"test::rev"
}
}
struct CountU64(Rc<Cell<usize>>);
impl Execute for CountU64 {
type Input = u64;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u64)) -> Result<(), BotError> {
call.0.check(&[])?;
self.0.set(self.0.get().saturating_add(1));
Ok(())
}
fn domain_id(&self) -> &'static str {
"test::count_u64"
}
}
#[test]
fn a_change_tick_source_fires_once_per_movement_and_is_not_polled_otherwise() -> TestResult {
let value = Rc::new(Cell::new(0u64));
let rev = Rc::new(Cell::new(0u64));
let polls = Rc::new(Cell::new(0usize));
let counter = Rc::new(Cell::new(0usize));
let mut bot = EcsBot::builder("revs")
.observe(Rev {
value: Rc::clone(&value),
rev: Rc::clone(&rev),
polls: Rc::clone(&polls),
caps: Vec::new(),
})
.on(
|value: &u64| value.is_multiple_of(2),
CountU64(Rc::clone(&counter)),
)
.with_effects(test_effects()?)
.build(&GrantSet::empty())?;
let mut fired = Vec::new();
for step in 0..6u64 {
let held = step
.checked_div(2)
.ok_or("halving a u64 by two cannot fail")?;
value.set(held);
rev.set(held);
fired.push(bot.tick()?);
}
assert_eq!(
fired,
vec![1, 0, 0, 0, 1, 0],
"one fire per even movement and nothing while the revision holds: {fired:?}"
);
assert_eq!(
polls.get(),
3,
"a skipped chain is not polled at all, so three movements are three polls"
);
assert_eq!(
bot.pending(),
Vec::new(),
"a value that moves must be acted on, and a chain with nothing outstanding \
must not be reported"
);
Ok(())
}
struct Exhausting {
remaining: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Observe for Exhausting {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
let left = self.remaining.get();
self.remaining.set(left.saturating_sub(1));
if left == 0 {
return domain_refusal("test::exhausting", "script exhausted");
}
Ok(200)
}
fn domain_id(&self) -> &str {
"test::exhausting"
}
}
struct Switched {
value: Rc<Cell<u16>>,
fail: Rc<Cell<bool>>,
polls: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Switched {
fn new(value: u16, fail: bool, polls: Rc<Cell<usize>>) -> Self {
Self {
value: Rc::new(Cell::new(value)),
fail: Rc::new(Cell::new(fail)),
polls,
caps: vec![Cap::net()],
}
}
}
impl Observe for Switched {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
self.polls.set(self.polls.get().saturating_add(1));
if self.fail.get() {
return domain_refusal("test::switched", "switch is off");
}
Ok(self.value.get())
}
fn fingerprint(&self) -> Option<u128> {
if self.fail.get() {
None
} else {
Some(u128::from(self.value.get()))
}
}
fn domain_id(&self) -> &str {
"test::switched"
}
}
struct Yielding {
value: u16,
yielded: Rc<Cell<bool>>,
polls: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl Yielding {
fn new(value: u16, yielded: Rc<Cell<bool>>, polls: Rc<Cell<usize>>) -> Self {
Self {
value,
yielded,
polls,
caps: vec![Cap::net()],
}
}
}
impl Observe for Yielding {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
self.polls.set(self.polls.get().saturating_add(1));
let flag = Rc::clone(&self.yielded);
std::future::poll_fn(move |cx| {
if flag.get() {
Poll::Ready(())
} else {
flag.set(true);
cx.waker().wake_by_ref();
Poll::Pending
}
})
.await;
Ok(self.value)
}
fn fingerprint(&self) -> Option<u128> {
Some(u128::from(self.value))
}
fn domain_id(&self) -> &str {
"test::yielding"
}
}
struct AbaSource {
value: Rc<Cell<u16>>,
phase: Rc<Cell<u8>>,
polls: Rc<Cell<usize>>,
caps: Vec<Cap>,
}
impl AbaSource {
fn new(value: Rc<Cell<u16>>, phase: Rc<Cell<u8>>, polls: Rc<Cell<usize>>) -> Self {
Self {
value,
phase,
polls,
caps: vec![Cap::net()],
}
}
}
impl Observe for AbaSource {
type Output = u16;
fn required_caps(&self) -> &[Cap] {
&self.caps
}
async fn poll(&self, call: (Auth, ())) -> Result<u16, BotError> {
call.0.check(&self.caps)?;
self.polls.set(self.polls.get().saturating_add(1));
let value = Rc::clone(&self.value);
let phase = Rc::clone(&self.phase);
let captured = Rc::new(Cell::new(0));
let captured = std::future::poll_fn(move |cx| match phase.get() {
0 => {
phase.set(1);
cx.waker().wake_by_ref();
Poll::Pending
}
1 => {
captured.set(value.get());
phase.set(2);
cx.waker().wake_by_ref();
Poll::Pending
}
_ => {
phase.set(0);
Poll::Ready(captured.get())
}
})
.await;
Ok(captured)
}
fn fingerprint(&self) -> Option<u128> {
Some(u128::from(self.value.get()))
}
fn domain_id(&self) -> &str {
"test::aba"
}
}
struct Count(Rc<Cell<usize>>);
impl Execute for Count {
type Input = u16;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u16)) -> Result<(), BotError> {
call.0.check(&[])?;
self.0.set(self.0.get().saturating_add(1));
Ok(())
}
fn domain_id(&self) -> &str {
"test::count"
}
}
struct Record(Rc<RefCell<Vec<u16>>>);
impl Execute for Record {
type Input = u16;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u16)) -> Result<(), BotError> {
call.0.check(&[])?;
self.0.borrow_mut().push(*call.1);
Ok(())
}
fn domain_id(&self) -> &str {
"test::record"
}
}
struct Flaky {
remaining: Rc<Cell<usize>>,
log: Rc<RefCell<Vec<u16>>>,
kind: FlakyKind,
}
#[derive(Clone, Copy)]
enum FlakyKind {
Refused,
Indeterminate,
}
impl Execute for Flaky {
type Input = u16;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u16)) -> Result<(), BotError> {
call.0.check(&[])?;
let left = self.remaining.get();
if left == 0 {
self.log.borrow_mut().push(*call.1);
return Ok(());
}
self.remaining.set(left.saturating_sub(1));
Err(match self.kind {
FlakyKind::Refused => BotError::DomainError {
domain: "test::flaky".into(),
certainty: DispatchCertainty::NotDelivered,
cause: "refused, and the effect did not happen".into(),
},
FlakyKind::Indeterminate => BotError::EffectIndeterminate {
domain: "test::flaky".into(),
cause: "the acknowledgment never arrived".into(),
},
})
}
fn domain_id(&self) -> &str {
"test::flaky"
}
}
fn identity(work: &PendingWork) -> (usize, usize) {
(work.id().chain(), work.id().entry())
}
fn identities(bot: &EcsBot) -> Vec<(usize, usize)> {
bot.pending().iter().map(identity).collect()
}
fn first_hold(bot: &EcsBot) -> Option<TransitionHold> {
bot.pending().first().map(|work| work.hold().clone())
}
struct Refuses {
attempts: Rc<Cell<usize>>,
domain: &'static str,
certainty: DispatchCertainty,
cause: &'static str,
}
impl Refuses {
fn undelivered(attempts: Rc<Cell<usize>>) -> Self {
Self {
attempts,
domain: "test::refuses",
certainty: DispatchCertainty::NotDelivered,
cause: "refused",
}
}
fn permanently(attempts: Rc<Cell<usize>>) -> Self {
Self {
attempts,
domain: "test::permanent",
certainty: DispatchCertainty::Refused,
cause: "no binding is installed for this domain",
}
}
}
impl Execute for Refuses {
type Input = u16;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u16)) -> Result<(), BotError> {
call.0.check(&[])?;
self.attempts.set(self.attempts.get().saturating_add(1));
Err(BotError::DomainError {
domain: self.domain.into(),
certainty: self.certainty,
cause: self.cause.into(),
})
}
fn domain_id(&self) -> &str {
self.domain
}
}
struct RefusesToEvaluate;
impl Evaluate<u16> for RefusesToEvaluate {
fn check(&self, observed: &u16) -> Result<bool, BotError> {
Err(BotError::EvaluateError {
cause: format!("this condition cannot decide about {observed}"),
})
}
fn condition_id(&self) -> &str {
"test::refuses"
}
}
struct NeedsCaps {
caps: Vec<Cap>,
domain: &'static str,
}
impl Execute for NeedsCaps {
type Input = u16;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&self.caps
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u16)) -> Result<(), BotError> {
call.0.check(&self.caps)?;
Ok(())
}
fn domain_id(&self) -> &str {
self.domain
}
}
fn net_grants() -> GrantSet {
GrantSet::empty().grant(Cap::net())
}
fn held_key(work: &PendingWork) -> Result<EffectKey, Box<dyn std::error::Error>> {
work.key().ok_or_else(|| {
format!(
"chain {} entry {} has no attempt",
work.id().chain(),
work.id().entry()
)
.into()
})
}
const TEST_RUN: &str = "0102030405060708090a0b0c0d0e0f10";
const TEST_ENV: &str = "2122232425262728292a2b2c2d2e2f30";
const TEST_FLOW: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
pub(crate) fn test_effects() -> Result<EffectScope, Box<dyn std::error::Error>> {
test_effects_with(Box::new(MemoryJournal::new()))
}
struct FaultJournal {
inner: Rc<RefCell<MemoryJournal>>,
durable: bool,
refuse_intent: Rc<Cell<bool>>,
refuse_prepare: Rc<Cell<bool>>,
refuse_outcome: Rc<Cell<bool>>,
lose_reply_once: Rc<Cell<bool>>,
unknown_once: Rc<Cell<bool>>,
outcome_appends: Rc<Cell<usize>>,
}
impl FaultJournal {
fn new() -> Self {
Self::graded(false)
}
fn durable() -> Self {
Self::graded(true)
}
fn graded(durable: bool) -> Self {
Self {
inner: Rc::new(RefCell::new(MemoryJournal::new())),
durable,
refuse_intent: Rc::new(Cell::new(false)),
refuse_prepare: Rc::new(Cell::new(false)),
refuse_outcome: Rc::new(Cell::new(false)),
lose_reply_once: Rc::new(Cell::new(false)),
unknown_once: Rc::new(Cell::new(false)),
outcome_appends: Rc::new(Cell::new(0)),
}
}
fn refuse_outcome(&self) -> Rc<Cell<bool>> {
Rc::clone(&self.refuse_outcome)
}
fn refuse_intent(&self) -> Rc<Cell<bool>> {
Rc::clone(&self.refuse_intent)
}
fn refuse_prepare(&self) -> Rc<Cell<bool>> {
Rc::clone(&self.refuse_prepare)
}
fn outcome_appends(&self) -> Rc<Cell<usize>> {
Rc::clone(&self.outcome_appends)
}
fn lose_reply_once(&self) -> Rc<Cell<bool>> {
Rc::clone(&self.lose_reply_once)
}
fn unknown_once(&self) -> Rc<Cell<bool>> {
Rc::clone(&self.unknown_once)
}
fn store(&self) -> Rc<RefCell<MemoryJournal>> {
Rc::clone(&self.inner)
}
}
fn injected<T>(error: JournalError) -> Result<T, JournalError> {
let refusal = Err(error);
lgwks_std::trace::debug!(error = ?refusal.as_ref().err(), "a test journal injected a fault");
refusal
}
impl EffectJournal for FaultJournal {
fn durability(&self) -> DurabilityPromise {
if self.durable {
DurabilityPromise::ProcessCrash
} else {
self.inner.borrow().durability()
}
}
fn tail(&self) -> JournalPosition {
self.inner.borrow().tail()
}
fn committed(&self) -> Result<Vec<EffectEvent>, JournalError> {
EffectJournal::committed(&*self.inner.borrow())
}
fn committed_entries(&self) -> Result<Vec<crate::journal::JournalEntry>, JournalError> {
Ok(self.inner.borrow().committed().to_vec())
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError> {
let is_outcome = matches!(event, EffectEvent::OutcomeObserved { .. });
let refuse = match *event {
EffectEvent::IntentAdmitted { .. } => self.refuse_intent.get(),
EffectEvent::DispatchPrepared { .. } => self.refuse_prepare.get(),
EffectEvent::OutcomeObserved { .. } => {
self.outcome_appends
.set(self.outcome_appends.get().saturating_add(1));
self.refuse_outcome.get()
}
EffectEvent::Verified { .. } => false,
};
if refuse {
return injected(JournalError::Storage(io::Error::other(
"injected append refusal",
)));
}
if is_outcome && self.unknown_once.take() {
return injected(JournalError::OutcomeUnknown {
cause: io::Error::other("injected unknown outcome, nothing written"),
});
}
let ack = self
.inner
.borrow_mut()
.compare_and_append(expected_tail, event)?;
if is_outcome && self.lose_reply_once.take() {
return injected(JournalError::OutcomeUnknown {
cause: io::Error::other("injected post-commit lost reply"),
});
}
Ok(ack)
}
fn confirm_outcome(
&mut self,
key: crate::effect::EffectKey,
evidence: EffectEvidence,
position: JournalPosition,
required: DurabilityPromise,
) -> Result<DurableAck, JournalError> {
let held = self.durable
&& EffectJournal::committed(&*self.inner.borrow())?.iter().any(
|event| matches!(*event, EffectEvent::OutcomeObserved { key: held, evidence: held_evidence }
if held == key && held_evidence == evidence),
);
if !held {
return injected(JournalError::ReceiptUnavailable { required });
}
Ok(DurableAck::new(position, self.durability()))
}
}
fn refused(
ticked: Result<usize, BotError>,
reported: impl FnOnce(usize) -> String,
) -> Result<BotError, Box<dyn std::error::Error>> {
match ticked {
Ok(fired) => Err(reported(fired).into()),
Err(error) => Ok(error),
}
}
fn indeterminate_once_bot(
name: &str,
log: &Rc<RefCell<Vec<u16>>>,
) -> Result<EcsBot, Box<dyn std::error::Error>> {
Ok(EcsBot::builder(name)
.observe(Script::new(vec![200, 200, 200]))
.on(
|value: &u16| *value >= 200,
Flaky {
remaining: Rc::new(Cell::new(1)),
log: Rc::clone(log),
kind: FlakyKind::Indeterminate,
},
)
.with_effects(test_effects()?)
.build(&net_grants())?)
}
fn counting_bot(
name: &str,
runs: &Rc<Cell<usize>>,
effects: EffectScope,
) -> Result<EcsBot, Box<dyn std::error::Error>> {
Ok(EcsBot::builder(name)
.observe(Holds::new(200))
.on(|value: &u16| *value >= 200, Count(Rc::clone(runs)))
.with_effects(effects)
.build(&net_grants())?)
}
#[test]
fn hash_parts_frames_distinct_variable_length_splits() -> TestResult {
let left = hash_parts(&[b"ab", b"c"]);
let right = hash_parts(&[b"a", b"bc"]);
assert_ne!(left, right, "distinct field splits must not collide");
assert_eq!(left, hash_parts(&[b"ab", b"c"]));
Ok(())
}
#[test]
fn an_ambiguous_commit_then_error_is_idempotent() -> TestResult {
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::new();
let lose_reply = journal.lose_reply_once();
let store = journal.store();
let mut bot = counting_bot(
"ambiguous-commit",
&runs,
test_effects_with(Box::new(journal))?,
)?;
lose_reply.set(true);
assert_eq!(bot.tick()?, 1);
assert_eq!(runs.get(), 1);
let committed = EffectJournal::committed(&*store.borrow())?;
assert!(matches!(
committed.last(),
Some(EffectEvent::OutcomeObserved {
evidence: EffectEvidence::Applied,
..
})
));
let key = committed[0].key();
bot.resolve_effect(&key, EffectEvidence::Applied)?;
assert_eq!(bot.tick()?, 0);
assert_eq!(runs.get(), 1, "retry must not dispatch the action again");
assert!(bot.pending().is_empty());
Ok(())
}
struct SwappedEntryJournal {
inner: MemoryJournal,
swap: EffectEvent,
}
impl EffectJournal for SwappedEntryJournal {
fn durability(&self) -> DurabilityPromise {
self.inner.durability()
}
fn tail(&self) -> JournalPosition {
self.inner.tail()
}
fn committed(&self) -> Result<Vec<EffectEvent>, JournalError> {
EffectJournal::committed(&self.inner)
}
fn committed_entries(&self) -> Result<Vec<crate::journal::JournalEntry>, JournalError> {
Ok(self.inner.committed().to_vec())
}
fn committed_entry(
&self,
position: JournalPosition,
) -> Result<Option<crate::journal::JournalEntry>, JournalError> {
Ok(Some(crate::journal::JournalEntry::new(position, self.swap)))
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError> {
self.inner.compare_and_append(expected_tail, event)
}
}
#[test]
fn a_position_holding_a_different_event_is_a_mismatch_not_an_acknowledgement() -> TestResult {
let key = crate::effect::EffectIdentity::new(
RunId::from_hex(TEST_RUN)?,
EnvironmentId::from_hex(TEST_ENV)?,
crate::effect::FlowRevision::from_tagged("blake3_256", TEST_FLOW)?,
)
.key(
crate::effect::ActionId::from_hex("1112131415161718191a1b1c1d1e1f20")?,
crate::effect::AttemptId::from_decimal("1")?,
crate::effect::ActionDigest::from_tagged(
"blake3_256",
"f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef",
)?,
crate::effect::EnvironmentEpoch::from_decimal("1")?,
);
let admitted = EffectEvent::IntentAdmitted { key };
let mut backing = MemoryJournal::new();
let ack = backing.compare_and_append(backing.tail(), &admitted)?;
let swapped = EffectEvent::DispatchPrepared { key };
let scope = test_effects_with(Box::new(SwappedEntryJournal {
inner: backing,
swap: swapped,
}))?;
let mut effects = Effects::new(scope, ack.position());
match effects.accept_position(&admitted, ack.position()) {
Err(JournalError::EntryMismatch {
position,
expected,
actual,
}) => {
assert_eq!(position, ack.position());
assert_eq!(*expected, admitted);
assert_eq!(*actual, swapped);
}
other => {
return Err(format!(
"a swapped positioned entry must be EntryMismatch, got {other:?}"
)
.into());
}
}
Ok(())
}
#[test]
fn a_durable_retry_after_a_lost_reply_settles_instead_of_stalling() -> TestResult {
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::durable();
let ambiguous_once = journal.lose_reply_once();
let store = journal.store();
let mut bot = counting_bot("lost-reply", &runs, test_effects_with(Box::new(journal))?)?;
ambiguous_once.set(true);
assert_eq!(bot.tick()?, 1);
assert_eq!(runs.get(), 1);
let committed = EffectJournal::committed(&*store.borrow())?;
assert!(
matches!(
committed.last(),
Some(EffectEvent::OutcomeObserved {
evidence: EffectEvidence::Applied,
..
})
),
"the store must hold the outcome whose reply was lost"
);
assert!(bot.pending().is_empty());
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::durable();
let unknown_once = journal.unknown_once();
let store = journal.store();
let mut bot = counting_bot(
"unknown-refused",
&runs,
test_effects_with(Box::new(journal))?,
)?;
unknown_once.set(true);
match bot.tick() {
Err(error @ BotError::EffectUnrecorded { .. }) => {
assert_eq!(
error.dispatch_certainty(),
DispatchCertainty::Occurred,
"the effect happened; only its record's reply was lost: {error:?}"
);
}
other => return Err(format!("expected a lost reply, got {other:?}").into()),
}
assert_eq!(runs.get(), 1);
let committed = EffectJournal::committed(&*store.borrow())?;
assert!(
!committed
.iter()
.any(|event| matches!(event, EffectEvent::OutcomeObserved { .. })),
"an unknown outcome without a commit must leave the ladder free"
);
assert_eq!(
bot.tick()?,
1,
"the record lands on the recovery tick and counts as fired"
);
assert_eq!(runs.get(), 1, "a lost reply must never re-run the action");
assert!(bot.pending().is_empty());
Ok(())
}
fn test_effects_with(
journal: Box<dyn EffectJournal>,
) -> Result<EffectScope, Box<dyn std::error::Error>> {
let environment = EnvironmentId::from_hex(TEST_ENV)?;
let mut broker = Broker::new();
broker.register(environment)?;
Ok(EffectScope::new(
EffectIdentity::new(
RunId::from_hex(TEST_RUN)?,
environment,
FlowRevision::from_tagged("blake3_256", TEST_FLOW)?,
),
broker,
journal,
))
}
#[test]
fn a_post_applied_append_failure_is_a_recording_failure_not_a_refusal() -> TestResult {
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::new();
let refuse_outcome = journal.refuse_outcome();
let outcome_appends = journal.outcome_appends();
let mut bot = counting_bot(
"unrecorded-applied",
&runs,
test_effects_with(Box::new(journal))?,
)?;
refuse_outcome.set(true);
let error = refused(bot.tick(), |fired| {
format!("a refused outcome append was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectUnrecorded { .. }),
"expected EffectUnrecorded, got {error:?}"
);
assert_eq!(
error.dispatch_certainty(),
DispatchCertainty::Occurred,
"the effect happened: the certainty must not be Refused or NotDelivered"
);
assert_eq!(
error.retry_class(),
RetryClass::Never,
"the action is never re-entered for a recording failure"
);
assert_eq!(runs.get(), 1, "the action ran exactly once");
assert_eq!(
outcome_appends.get(),
1,
"the outcome append was attempted once and refused"
);
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::RecordingFailed { .. })
),
"the entry is held for an append-only retry: {:?}",
first_hold(&bot)
);
refuse_outcome.set(false);
assert_eq!(
bot.tick()?,
1,
"the append lands and the recovered effect is counted as fired"
);
assert_eq!(
runs.get(),
1,
"a retry of recording must not re-enter the action"
);
assert_eq!(
outcome_appends.get(),
2,
"the outcome append was retried exactly once"
);
assert!(
bot.pending().is_empty(),
"the entry is resolved once its record lands: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn a_post_failure_append_failure_keeps_the_not_applied_fact() -> TestResult {
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::new();
let refuse_outcome = journal.refuse_outcome();
let mut bot = EcsBot::builder("unrecorded-not-applied")
.observe(Holds::new(200))
.on(
|value: &u16| *value >= 200,
Refuses::undelivered(Rc::clone(&runs)),
)
.with_effects(test_effects_with(Box::new(journal))?)
.build(&net_grants())?;
refuse_outcome.set(true);
let error = refused(bot.tick(), |fired| {
format!("a refused outcome append was reported as {fired} fired")
})?;
assert!(
matches!(
error,
BotError::EffectUnrecorded {
evidence: EffectEvidence::NotApplied,
..
}
),
"expected EffectUnrecorded carrying NotApplied, got {error:?}"
);
assert_eq!(
error.dispatch_certainty(),
DispatchCertainty::NotDelivered,
"the effect definitely did not happen"
);
assert_eq!(
error.retry_class(),
RetryClass::Never,
"a recording failure never re-enters the action"
);
assert_eq!(runs.get(), 1, "the action ran exactly once");
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::RecordingFailed { .. })
),
"the entry is held for an append-only retry: {:?}",
first_hold(&bot)
);
refuse_outcome.set(false);
let error = refused(bot.tick(), |fired| {
format!("an open entry was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::PendingTransition { .. }),
"expected the still-open entry, got {error:?}"
);
assert_eq!(
runs.get(),
1,
"the append retry must not re-enter the action"
);
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::Failed { attempts: 1, .. })
),
"the entry is eligible again under the budget: {:?}",
first_hold(&bot)
);
Ok(())
}
#[test]
fn an_append_failure_before_the_action_is_still_a_pre_dispatch_refusal() -> TestResult {
let runs = Rc::new(Cell::new(0));
let journal = FaultJournal::new();
let refuse_intent = journal.refuse_intent();
let mut bot = counting_bot(
"refuse-intent",
&runs,
test_effects_with(Box::new(journal))?,
)?;
refuse_intent.set(true);
let error = refused(bot.tick(), |fired| {
format!("a refused intent append was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectRefused { .. }),
"expected EffectRefused before intent, got {error:?}"
);
assert_eq!(error.dispatch_certainty(), DispatchCertainty::Refused);
assert_eq!(runs.get(), 0, "the action never ran");
let journal = FaultJournal::new();
let refuse_prepare = journal.refuse_prepare();
let mut bot = counting_bot(
"refuse-prepare",
&runs,
test_effects_with(Box::new(journal))?,
)?;
refuse_prepare.set(true);
let error = refused(bot.tick(), |fired| {
format!("a refused prepare append was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectRefused { .. }),
"expected EffectRefused before preparation, got {error:?}"
);
assert_eq!(error.dispatch_certainty(), DispatchCertainty::Refused);
assert_eq!(runs.get(), 0, "the action never ran");
Ok(())
}
const SCRIPT: [u16; 5] = [200, 200, 503, 503, 200];
#[test]
fn a_condition_fires_only_on_the_tick_its_source_moves() -> TestResult {
let counter = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("scripted")
.observe(Script::new(SCRIPT.to_vec()))
.on(|value: &u16| *value >= 500, Count(Rc::clone(&counter)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let mut fired = Vec::new();
let mut revisions = Vec::new();
for _ in 0..5 {
fired.push(bot.tick()?);
revisions.push(
bot.revisions()
.first()
.copied()
.ok_or("a bot with one chain reports one revision")?,
);
}
assert_eq!(
revisions,
vec![1, 1, 2, 2, 3],
"Revision must track value movement, not system execution"
);
assert_eq!(fired, vec![0, 0, 1, 0, 0]);
assert_eq!(counter.get(), 1);
Ok(())
}
#[test]
fn a_settled_chain_does_not_fire_again_on_every_later_tick() -> TestResult {
let counter = Rc::new(Cell::new(0));
let mut bot = counting_bot("held", &counter, test_effects()?)?;
assert_eq!(bot.tick()?, 1, "the movement opens the chain");
for tick in 0_u32..20 {
assert_eq!(
bot.tick()?,
0,
"tick {}: the source held still, so the chain had nothing to run",
tick.saturating_add(2)
);
}
assert_eq!(
counter.get(),
1,
"the effect is owed once per movement, not once per tick"
);
assert_eq!(
bot.revisions().first().copied(),
Some(1),
"Revision counts movements, so a settled chain's is still 1"
);
Ok(())
}
#[test]
fn a_source_with_a_detached_digest_is_polled_while_it_holds_still() -> TestResult {
let polls = Rc::new(Cell::new(0));
let counter = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("lazy")
.observe(Counted::new(200, Rc::clone(&polls)))
.on(|value: &u16| *value >= 200, Count(Rc::clone(&counter)))
.with_effects(test_effects()?)
.build(&net_grants())?;
assert_eq!(bot.tick()?, 1, "the first tick polls");
assert_eq!(polls.get(), 1, "and polls exactly once");
for tick in 0_u32..20 {
assert_eq!(
bot.tick()?,
0,
"tick {}: nothing moved",
tick.saturating_add(2)
);
}
assert_eq!(
polls.get(),
21,
"a detached digest cannot suppress a later async poll"
);
assert_eq!(counter.get(), 1, "and must not fire the effect again");
Ok(())
}
#[test]
fn a_detached_digest_cannot_suppress_the_final_stable_aba_state() -> TestResult {
let value = Rc::new(Cell::new(0));
let phase = Rc::new(Cell::new(0));
let polls = Rc::new(Cell::new(0));
let seen = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("aba")
.observe(AbaSource::new(
Rc::clone(&value),
Rc::clone(&phase),
Rc::clone(&polls),
))
.on(|_: &u16| true, Record(Rc::clone(&seen)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let mut tick = Box::pin(bot.tick_async());
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
if !matches!(tick.as_mut().poll(&mut cx), Poll::Pending) {
return Err("the first poll must pause before sampling B".into());
}
value.set(1);
if !matches!(tick.as_mut().poll(&mut cx), Poll::Pending) {
return Err("the source must pause after sampling B".into());
}
value.set(0);
match tick.as_mut().poll(&mut cx) {
Poll::Ready(Ok(fired)) => assert_eq!(fired, 1, "the sampled B action fires"),
Poll::Ready(Err(error)) => return Err(format!("the B tick failed: {error}").into()),
Poll::Pending => return Err("the third poll must admit the sampled B value".into()),
}
drop(tick);
assert_eq!(*seen.borrow(), vec![1], "the first tick delivers sampled B");
assert_eq!(bot.tick()?, 1, "the next tick must observe stable A");
assert_eq!(
*seen.borrow(),
vec![1, 0],
"stable A is delivered instead of being suppressed by B's stale digest"
);
assert_eq!(polls.get(), 2, "both B and the final stable A were polled");
Ok(())
}
#[test]
fn a_source_value_that_moves_is_polled_again_and_fires() -> TestResult {
let polls = Rc::new(Cell::new(0));
let counter = Rc::new(Cell::new(0));
let source = Counted::new(200, Rc::clone(&polls));
let value = Rc::clone(&source.value);
let mut bot = EcsBot::builder("lazy-moving")
.observe(source)
.on(|value: &u16| *value >= 200, Count(Rc::clone(&counter)))
.with_effects(test_effects()?)
.build(&net_grants())?;
assert_eq!(bot.tick()?, 1, "the first value fires");
assert_eq!(bot.tick()?, 0, "and holds");
assert_eq!(
polls.get(),
2,
"the second tick polls and compares its value"
);
value.set(503);
assert_eq!(bot.tick()?, 1, "the movement is seen, not skipped");
assert_eq!(
polls.get(),
3,
"the tick paid for a value to see the movement"
);
assert_eq!(bot.tick()?, 0, "and settles again");
assert_eq!(polls.get(), 4, "the stable value is polled and compared");
assert_eq!(counter.get(), 2, "two movements, two effects");
Ok(())
}
#[test]
fn a_source_without_a_digest_is_polled_every_tick() -> TestResult {
let counter = Rc::new(Cell::new(0));
let mut bot = counting_bot("unguarded", &counter, test_effects()?)?;
assert_eq!(bot.tick()?, 1);
for _ in 0..5 {
assert_eq!(bot.tick()?, 0);
}
assert_eq!(
counter.get(),
1,
"the value decides the chain, not the digest"
);
Ok(())
}
#[test]
fn building_without_the_grant_is_refused() -> TestResult {
let counter = Rc::new(Cell::new(0));
match EcsBot::builder("ungranted")
.observe(Script::new(SCRIPT.to_vec()))
.on(|value: &u16| *value >= 500, Count(counter))
.with_effects(test_effects()?)
.build(&GrantSet::empty())
{
Ok(_) => Err("a source requiring `bot.net` built without it".into()),
Err(error) => {
assert!(
matches!(error, BotError::CapabilityDenied { .. }),
"expected CapabilityDenied, got {error:?}"
);
Ok(())
}
}
}
#[test]
fn admission_names_every_unmet_requirement_and_the_domain_that_declared_it() -> TestResult {
match EcsBot::builder("short")
.observe(Holds::new(1))
.on(
|value: &u16| *value > 0,
NeedsCaps {
caps: vec![Cap::fs()],
domain: "test::writes",
},
)
.with_effects(test_effects()?)
.build(&GrantSet::empty())
{
Ok(_) => Err("a bot requiring two ungranted capabilities built".into()),
Err(BotError::CapabilityDenied { deficit }) => {
let named: Vec<(&str, Option<&str>)> = deficit
.shortages()
.map(|shortage| {
(
shortage.required().as_str(),
shortage.demand().map(Demand::domain),
)
})
.collect();
assert_eq!(
named,
vec![
(Cap::NET, Some("test::holds")),
(Cap::FS, Some("test::writes")),
],
"one refusal must name both requirements and both domains: {deficit}"
);
assert!(
deficit
.to_grant_set()
.admit(&[Cap::net(), Cap::fs()])
.is_ok(),
"the deficit must derive the whole repair: {deficit}"
);
Ok(())
}
Err(other) => Err(format!("expected a capability denial, got {other:?}").into()),
}
}
#[test]
fn a_poll_error_is_returned_and_commits_nothing() -> TestResult {
let counter = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("exhausting")
.observe(Exhausting {
remaining: Rc::new(Cell::new(1)),
caps: vec![Cap::net()],
})
.on(|value: &u16| *value >= 500, Count(Rc::clone(&counter)))
.with_effects(test_effects()?)
.build(&net_grants())?;
assert_eq!(
bot.tick()?,
0,
"the first poll answers 200 and fires nothing"
);
let after_success = bot.revisions();
let error = refused(bot.tick(), |_| {
String::from("a failing poll was reported as a successful tick")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the observer's own error, got {error:?}"
);
assert_eq!(
bot.revisions(),
after_success,
"a failed tick must not commit a new Revision"
);
assert_eq!(
counter.get(),
0,
"a failed poll fires nothing: the action must not have run"
);
Ok(())
}
#[test]
fn a_failed_sibling_does_not_publish_the_successful_sources_digest() -> TestResult {
let left_polls = Rc::new(Cell::new(0));
let right_polls = Rc::new(Cell::new(0));
let seen = Rc::new(RefCell::new(Vec::new()));
let left = Switched::new(1, false, Rc::clone(&left_polls));
let mid = Switched::new(9, true, Rc::clone(&right_polls));
let mid_fail = Rc::clone(&mid.fail);
let mut bot = EcsBot::builder("sibling")
.observe(left)
.on(|value: &u16| *value >= 1, Record(Rc::clone(&seen)))
.observe(mid)
.on(
|value: &u16| *value >= 9,
Record(Rc::new(RefCell::new(Vec::new()))),
)
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a failing sibling was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the sibling's own error, got {error:?}"
);
assert!(
seen.borrow().is_empty(),
"the failed fold must not commit A's payload"
);
assert_eq!(left_polls.get(), 1, "A was polled once on the first tick");
assert_eq!(right_polls.get(), 1, "and so was B");
mid_fail.set(false);
assert_eq!(bot.tick()?, 2, "both chains recover and both actions fire");
assert_eq!(
*seen.borrow(),
vec![1],
"A's uncommitted payload is committed on recovery, with 1 as its input"
);
assert_eq!(
left_polls.get(),
2,
"A is polled again: its digest must not have been published early"
);
assert_eq!(right_polls.get(), 2, "B recovers on a poll of its own");
Ok(())
}
#[test]
fn a_committed_source_that_moves_during_a_sibling_failure_is_not_lost() -> TestResult {
let left_polls = Rc::new(Cell::new(0));
let right_polls = Rc::new(Cell::new(0));
let seen = Rc::new(RefCell::new(Vec::new()));
let left = Switched::new(0, false, Rc::clone(&left_polls));
let left_value = Rc::clone(&left.value);
let mid = Switched::new(9, false, Rc::clone(&right_polls));
let mid_fail = Rc::clone(&mid.fail);
let mut bot = EcsBot::builder("recovery")
.observe(left)
.on(|value: &u16| *value >= 1, Record(Rc::clone(&seen)))
.observe(mid)
.on(|_: &u16| true, Record(Rc::new(RefCell::new(Vec::new()))))
.with_effects(test_effects()?)
.build(&net_grants())?;
assert_eq!(bot.tick()?, 1, "the first tick commits A=0 and B=9");
assert_eq!(
*seen.borrow(),
Vec::<u16>::new(),
"A=0 does not meet its condition"
);
assert_eq!(left_polls.get(), 1);
left_value.set(1);
mid_fail.set(true);
let error = refused(bot.tick(), |fired| {
format!("a failing sibling was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the sibling's own error, got {error:?}"
);
assert!(
seen.borrow().is_empty(),
"the failed fold must not commit A=1"
);
mid_fail.set(false);
assert_eq!(
bot.tick()?,
1,
"recovery commits A's move; B is unchanged and fires nothing"
);
assert_eq!(
*seen.borrow(),
vec![1],
"the transition to 1 is retained, not the older committed 0"
);
assert_eq!(
left_polls.get(),
3,
"A is polled on the failing tick and again on recovery"
);
Ok(())
}
#[test]
fn several_failed_siblings_do_not_suppress_the_sources_that_succeeded() -> TestResult {
let polls_left = Rc::new(Cell::new(0));
let polls_mid = Rc::new(Cell::new(0));
let polls_right = Rc::new(Cell::new(0));
let seen_left = Rc::new(RefCell::new(Vec::new()));
let seen_right = Rc::new(RefCell::new(Vec::new()));
let left = Switched::new(1, false, Rc::clone(&polls_left));
let mid = Switched::new(2, true, Rc::clone(&polls_mid));
let right = Switched::new(3, false, Rc::clone(&polls_right));
let far = Switched::new(4, true, Rc::new(Cell::new(0)));
let mid_fail = Rc::clone(&mid.fail);
let far_fail = Rc::clone(&far.fail);
let mut bot = EcsBot::builder("multi-fail")
.observe(left)
.on(|value: &u16| *value >= 1, Record(Rc::clone(&seen_left)))
.observe(mid)
.on(|_: &u16| true, Record(Rc::new(RefCell::new(Vec::new()))))
.observe(right)
.on(|value: &u16| *value >= 3, Record(Rc::clone(&seen_right)))
.observe(far)
.on(|_: &u16| true, Record(Rc::new(RefCell::new(Vec::new()))))
.with_effects(test_effects()?)
.build(&net_grants())?;
if let Ok(fired) = bot.tick() {
return Err(format!("failed siblings were reported as {fired} fired").into());
}
assert!(
seen_left.borrow().is_empty(),
"nothing commits on a failed fold"
);
assert!(
seen_right.borrow().is_empty(),
"nothing commits on a failed fold"
);
assert_eq!(polls_left.get(), 1);
assert_eq!(polls_right.get(), 1);
mid_fail.set(false);
far_fail.set(false);
assert_eq!(bot.tick()?, 4, "every chain recovers and fires");
assert_eq!(*seen_left.borrow(), vec![1], "A delivers its payload");
assert_eq!(*seen_right.borrow(), vec![3], "C delivers its payload");
assert_eq!(
polls_left.get(),
2,
"A was re-polled after the aborted fold"
);
assert_eq!(
polls_right.get(),
2,
"C was re-polled after the aborted fold"
);
Ok(())
}
#[test]
fn cancelling_a_poll_leaves_the_next_observation_usable() -> TestResult {
let polls = Rc::new(Cell::new(0));
let seen = Rc::new(RefCell::new(Vec::new()));
let yielded = Rc::new(Cell::new(false));
let source = Yielding::new(200, Rc::clone(&yielded), Rc::clone(&polls));
let mut bot = EcsBot::builder("cancel")
.observe(source)
.on(|value: &u16| *value >= 200, Record(Rc::clone(&seen)))
.with_effects(test_effects()?)
.build(&net_grants())?;
{
let mut tick = Box::pin(bot.tick_async());
let waker = Waker::noop();
let mut cx = Context::from_waker(waker);
match tick.as_mut().poll(&mut cx) {
Poll::Ready(_) => {
return Err("the yielding source resolved before it was cancelled".into());
}
Poll::Pending => {}
}
drop(tick);
}
assert!(seen.borrow().is_empty(), "a cancelled tick commits nothing");
assert_eq!(bot.tick()?, 1, "the recovery tick completes and fires");
assert_eq!(*seen.borrow(), vec![200], "with the payload as its input");
assert_eq!(
polls.get(),
2,
"the source is polled once per tick across the cancel"
);
assert_eq!(bot.tick()?, 0, "and the next tick compares the same value");
assert_eq!(polls.get(), 3, "the next tick polls after a cancellation");
Ok(())
}
#[test]
fn a_held_transition_keeps_its_input_across_a_sibling_failure() -> TestResult {
let logs = Rc::new(RefCell::new(Vec::new()));
let mid_polls = Rc::new(Cell::new(0));
let mid = Switched::new(9, false, Rc::clone(&mid_polls));
let mid_fail = Rc::clone(&mid.fail);
let left = Holds::new(200);
let mut bot = EcsBot::builder("held-newer")
.observe(left)
.on(
|value: &u16| *value >= 200,
Flaky {
remaining: Rc::new(Cell::new(1)),
log: Rc::clone(&logs),
kind: FlakyKind::Indeterminate,
},
)
.observe(mid)
.on(|_: &u16| true, Record(Rc::new(RefCell::new(Vec::new()))))
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("an indeterminate effect was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectIndeterminate { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(identities(&bot), vec![(0, 0)], "the held entry is named");
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::OutcomeUnknown { attempts: 1, .. })
),
"the first tick holds the attempt: {:?}",
first_hold(&bot)
);
assert!(logs.borrow().is_empty(), "nothing was recorded yet");
mid_fail.set(true);
let error = refused(bot.tick(), |fired| {
format!("a failing sibling was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the sibling's own error, got {error:?}"
);
assert_eq!(
identities(&bot),
vec![(0, 0)],
"the held entry survives an aborted fold"
);
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::OutcomeUnknown { attempts: 1, .. })
),
"and is still held on the same attempt: {:?}",
first_hold(&bot)
);
mid_fail.set(false);
match bot.tick() {
Ok(_) => {}
Err(error) => assert!(
matches!(error, BotError::PendingTransition { .. }),
"the held entry is still outstanding, got {error:?}"
),
}
assert!(
logs.borrow().is_empty(),
"a held effect is not attempted across a recovery without evidence"
);
let held = bot
.pending()
.into_iter()
.find(|work| identity(work) == (0, 0))
.ok_or("the held entry disappeared from the report")?;
bot.resolve_effect(&held_key(&held)?, EffectEvidence::NotApplied)?;
assert_eq!(bot.tick()?, 1, "the entry runs once evidence authorises it");
assert_eq!(
*logs.borrow(),
vec![200],
"the retained transition's immutable input is the one it was opened under"
);
Ok(())
}
#[test]
fn a_failure_mid_chain_leaves_the_untouched_entries_pending() -> TestResult {
let first = Rc::new(RefCell::new(Vec::new()));
let second = Rc::new(RefCell::new(Vec::new()));
let third = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("partial")
.observe(Script::new(vec![200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&first)))
.on(
|value: &u16| *value >= 200,
Flaky {
remaining: Rc::new(Cell::new(1)),
log: Rc::clone(&second),
kind: FlakyKind::Refused,
},
)
.on(|value: &u16| *value >= 200, Record(Rc::clone(&third)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refused action was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(
*first.borrow(),
vec![200],
"the action before the failure ran and its effect is not rolled back"
);
assert!(
second.borrow().is_empty(),
"the refused action recorded nothing"
);
assert!(
third.borrow().is_empty(),
"the entry after the failure was not attempted on the failing tick"
);
assert_eq!(
identities(&bot),
vec![(0, 1), (0, 2)],
"the refused entry and the entry behind it are both still outstanding"
);
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::Failed {
attempts: 1,
budget: 3,
..
})
),
"the refused entry is held as a definite failure with the attempt counted: {:?}",
first_hold(&bot)
);
assert_eq!(
bot.tick()?,
2,
"the refused entry is retried and the entry behind it runs, on a still source"
);
assert_eq!(
*first.borrow(),
vec![200],
"the acknowledged first effect is not replayed to reach the third entry"
);
assert_eq!(
*second.borrow(),
vec![200],
"the retried action ran exactly once, on the second tick"
);
assert_eq!(
*third.borrow(),
vec![200],
"the effect the old code lost ran once the entry ahead of it settled"
);
assert!(
bot.pending().is_empty(),
"a clean tick means nothing is left holding the transition: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn a_failure_before_the_first_effect_still_runs_it_once() -> TestResult {
let attempts = Rc::new(Cell::new(0));
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("refusing")
.observe(Script::new(vec![200, 200, 200, 200, 0, 0, 0]))
.on(
|value: &u16| *value >= 200,
Refuses::undelivered(Rc::clone(&attempts)),
)
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
for round in 1..=2 {
let error = refused(bot.tick(), |fired| {
format!("round {round} reported {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"round {round}: expected the action's own error, got {error:?}"
);
}
assert_eq!(attempts.get(), 2, "two ticks, two attempts");
assert!(
log.borrow().is_empty(),
"the second entry is not run ahead of the entry holding it back"
);
assert_eq!(
identities(&bot),
vec![(0, 0), (0, 1)],
"both entries are outstanding while the first keeps failing"
);
let error = refused(bot.tick(), |fired| {
format!("a spent budget was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"the abandonment does not replace the error that caused it: {error:?}"
);
assert_eq!(attempts.get(), 3, "exactly the declared budget was spent");
assert_eq!(
first_hold(&bot),
Some(TransitionHold::Abandoned {
reason: AbandonReason::AttemptsExhausted { attempts: 3 },
cause: "test::refuses: refused".into(),
}),
"the given-up-on entry is named after the tick that gave up on it"
);
assert!(
log.borrow().is_empty(),
"the successor is not run ahead of the entry that was just decided"
);
let error = refused(bot.tick(), |fired| {
format!(
"a tick with an abandoned prerequisite reported {fired} fired and ran {:?}",
log.borrow()
)
})?;
assert!(
matches!(error, BotError::PendingTransition { .. }),
"expected the unresolved chain to be reported, got {error:?}"
);
assert!(
log.borrow().is_empty(),
"the unattempted work does not run past the entry that was given up on"
);
let error = refused(bot.tick(), |fired| {
format!(
"an abandoned entry was reported as {fired} fired while pending() still \
names it: {:?}",
bot.pending()
)
})?;
assert!(
matches!(error, BotError::PendingTransition { .. }),
"expected the unresolved chain to be reported, got {error:?}"
);
assert!(
identities(&bot) == vec![(0, 0), (0, 1)],
"the abandoned entry and the entry it blocks are what is left to report: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn one_attempt_is_a_policy_a_caller_can_declare() -> TestResult {
let attempts = Rc::new(Cell::new(0));
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("impatient")
.with_retry_policy(RetryPolicy::ONE_ATTEMPT)
.observe(Script::new(vec![200, 200, 200]))
.on(
|value: &u16| *value >= 200,
Refuses::undelivered(Rc::clone(&attempts)),
)
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refusal was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(attempts.get(), 1, "the declared budget is one attempt");
assert!(
log.borrow().is_empty(),
"the entry behind the abandoned one is not run in the tick that abandoned it"
);
let error = refused(bot.tick(), |fired| {
format!(
"a tick with an abandoned prerequisite reported {fired} fired and ran {:?}",
log.borrow()
)
})?;
assert!(
matches!(error, BotError::PendingTransition { .. }),
"expected the unresolved chain to be reported, got {error:?}"
);
assert_eq!(attempts.get(), 1, "the abandoned entry is not retried");
assert!(
log.borrow().is_empty(),
"the successor of an abandoned entry is not attempted while it stands abandoned"
);
assert_eq!(
identities(&bot),
vec![(0, 0), (0, 1)],
"the abandonment and the entry it blocks are both still reported: {:?}",
bot.pending()
);
let blocked = bot.pending().into_iter().next().ok_or("held")?;
bot.resolve_effect(&held_key(&blocked)?, EffectEvidence::NotApplied)?;
let error = refused(bot.tick(), |fired| {
format!("a refusal was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the retry's own error, got {error:?}"
);
assert_eq!(attempts.get(), 2, "the revived entry attempted again");
assert!(
log.borrow().is_empty(),
"and the successor is behind the barrier again, because the retry failed too"
);
Ok(())
}
#[test]
fn a_failure_does_not_stop_a_later_chain() -> TestResult {
let attempts = Rc::new(Cell::new(0));
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("two-chains")
.observe(Script::new(vec![200, 200]))
.on(
|value: &u16| *value >= 200,
Refuses::undelivered(Rc::clone(&attempts)),
)
.observe(Script::new(vec![200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refused action was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(
*log.borrow(),
vec![200],
"the chain after the failure ran: a failure is contained to its own chain"
);
assert_eq!(
identities(&bot),
vec![(0, 0)],
"only the refusing chain has work outstanding: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn an_indeterminate_effect_is_held_until_evidence_says_what_happened() -> TestResult {
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = indeterminate_once_bot("indeterminate", &log)?;
let error = refused(bot.tick(), |fired| {
format!("an indeterminate effect was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectIndeterminate { .. }),
"expected the action's own error, got {error:?}"
);
assert!(
log.borrow().is_empty(),
"an indeterminate effect is not an effect: nothing was recorded"
);
let error = refused(bot.tick(), |fired| {
format!("a held effect was reported as {fired} fired")
})?;
let BotError::PendingTransition { work, outstanding } = error else {
return Err(format!("expected PendingTransition, got {error:?}").into());
};
assert_eq!(identity(&work), (0, 0), "the held entry is named");
assert!(
matches!(
work.hold(),
TransitionHold::OutcomeUnknown { attempts: 1, .. }
),
"the hold carries the attempt that could not be settled: {:?}",
work.hold()
);
assert_eq!(outstanding, 1, "one entry is still open");
assert!(
bot.pending()
.iter()
.all(|work| !matches!(work.hold(), TransitionHold::Abandoned { .. })),
"nothing was given up on: {:?}",
bot.pending()
);
assert!(
log.borrow().is_empty(),
"a held effect is not attempted again without evidence"
);
let held = bot
.pending()
.into_iter()
.next()
.ok_or("the held entry disappeared from the report")?;
bot.resolve_effect(&held_key(&held)?, EffectEvidence::NotApplied)?;
assert_eq!(bot.tick()?, 1, "the entry runs once evidence authorises it");
assert_eq!(*log.borrow(), vec![200], "the effect happened, once");
assert!(
bot.pending().is_empty(),
"and nothing is left holding the transition: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn evidence_that_the_effect_happened_records_it_without_replaying_it() -> TestResult {
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = indeterminate_once_bot("acknowledged", &log)?;
let error = refused(bot.tick(), |fired| {
format!("an indeterminate effect was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectIndeterminate { .. }),
"expected the action's own error, got {error:?}"
);
let held = bot
.pending()
.into_iter()
.next()
.ok_or("the held entry is not reported")?;
bot.resolve_effect(&held_key(&held)?, EffectEvidence::Applied)?;
assert_eq!(
bot.tick()?,
0,
"an acknowledged effect is never attempted again"
);
assert!(
log.borrow().is_empty(),
"the action that may already have run was not run a second time"
);
assert_eq!(bot.revisions(), vec![1], "and the transition is resolved");
assert!(
bot.pending().is_empty(),
"nothing is left pending: {:?}",
bot.pending()
);
bot.resolve_effect(&held_key(&held)?, EffectEvidence::Applied)?;
match bot.resolve_effect(&held_key(&held)?, EffectEvidence::NotApplied) {
Ok(()) => Err("contradicting evidence was accepted".into()),
Err(error) => {
assert!(
matches!(
error,
BotError::EvidenceContradicted {
settled: EffectEvidence::Applied,
submitted: EffectEvidence::NotApplied,
..
}
),
"expected EvidenceContradicted, got {error:?}"
);
Ok(())
}
}
}
#[test]
fn a_source_that_moves_while_work_is_outstanding_neither_loses_nor_duplicates_it() -> TestResult
{
let first = Rc::new(RefCell::new(Vec::new()));
let second = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("moving")
.observe(Script::new(vec![200, 503, 200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&first)))
.on(
|value: &u16| *value >= 200,
Flaky {
remaining: Rc::new(Cell::new(1)),
log: Rc::clone(&second),
kind: FlakyKind::Refused,
},
)
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refusal was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(*first.borrow(), vec![200], "the first entry ran once");
assert!(second.borrow().is_empty(), "the second entry refused");
assert_eq!(bot.tick()?, 1, "only the outstanding entry runs");
assert_eq!(
*first.borrow(),
vec![200],
"an effect the transition already acknowledged is not replayed for a later movement"
);
assert_eq!(
*second.borrow(),
vec![200],
"the outstanding entry reads the payload its transition was opened under, \
not the one that arrived since"
);
assert!(
bot.pending().is_empty(),
"the transition is fully resolved: {:?}",
bot.pending()
);
assert_eq!(
bot.tick()?,
0,
"a source that has returned to the value the transition was acknowledged \
for is not new work"
);
assert_eq!(
*first.borrow(),
vec![200],
"and the acknowledged effect is not replayed to make it look like one"
);
assert_eq!(*second.borrow(), vec![200], "nor is the second entry");
assert_eq!(
bot.tick()?,
0,
"a still source with no outstanding work fires nothing"
);
assert_eq!(
bot.revisions(),
vec![3],
"three polls moved the value: 200, 503, 200"
);
Ok(())
}
#[test]
fn an_attempt_whose_outcome_was_never_recorded_is_held_not_replayed() -> TestResult {
let log = Rc::new(RefCell::new(Vec::new()));
let journal = FaultJournal::new();
let refuse_outcome = journal.refuse_outcome();
refuse_outcome.set(true);
let mut bot = EcsBot::builder("interrupted")
.observe(Script::new(vec![200, 200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects_with(Box::new(journal))?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refused outcome append reported {fired} fired")
})?;
assert!(
matches!(error, BotError::EffectUnrecorded { .. }),
"expected the outcome append's own failure, got {error:?}"
);
assert_eq!(*log.borrow(), vec![200], "the action ran once");
refuse_outcome.set(false);
{
let revision = *bot.revisions().first().ok_or("the source has a revision")?;
let input = bot
.world
.non_send::<Ledger>()
.transitions
.first()
.and_then(Option::as_ref)
.map_or_else(
|| AdmittedInput {
identity: derive_input_stamp(1),
event: false,
},
|held| AdmittedInput {
identity: held.input,
event: held.event,
},
);
let mut transition = Transition::opened(input, revision, 1, Some(Erased::new(200_u16)));
*transition
.entries
.first_mut()
.ok_or("the chain declares no entries")? = EntryState::Unrecorded;
*transition
.attempts
.first_mut()
.ok_or("the chain declares no entries")? = AttemptRecord {
begun: Some(AttemptId::FIRST),
settled: None,
};
bot.world.non_send_mut::<Ledger>().put(0, Some(transition));
}
assert!(
matches!(first_hold(&bot), Some(TransitionHold::Unrecorded)),
"the interrupted attempt is reported as unrecorded: {:?}",
first_hold(&bot)
);
let error = refused(bot.tick(), |fired| {
format!("an unrecorded attempt fired {fired}")
})?;
let BotError::PendingTransition { work, .. } = error else {
return Err(format!("expected PendingTransition, got {error:?}").into());
};
assert!(
matches!(work.hold(), TransitionHold::Unrecorded),
"the tick reports it as held, not as handled: {:?}",
work.hold()
);
assert_eq!(
*log.borrow(),
vec![200],
"the action was not run a second time: the first attempt may have taken effect"
);
let held = bot
.pending()
.into_iter()
.next()
.ok_or("the held entry is not reported")?;
bot.resolve_effect(&held_key(&held)?, EffectEvidence::NotApplied)?;
assert_eq!(bot.tick()?, 1, "and then it runs");
assert_eq!(*log.borrow(), vec![200, 200], "exactly once more");
Ok(())
}
#[test]
fn an_action_identity_is_framed_and_portable() -> TestResult {
let ab_c = derive_action_id("ab", 0, 0, "c");
let a_bc = derive_action_id("a", 0, 0, "bc");
assert_ne!(
ab_c, a_bc,
"length framing keeps \"ab\"+\"c\" apart from \"a\"+\"bc\""
);
assert_eq!(
derive_action_id("bot", 2, 3, "dom"),
derive_action_id("bot", 2, 3, "dom")
);
assert_ne!(
derive_action_id("bot", 2, 3, "dom"),
derive_action_id("bot", 2, 3, "dom2"),
"a different domain is a different action"
);
assert_ne!(
derive_action_id("bot", 2, 3, "dom"),
derive_action_id("bot", 3, 2, "dom"),
"chain and entry do not commute"
);
assert_ne!(
derive_action_id("bot", 0, 0, "dom"),
derive_action_id("bo", 0, 0, "tdom"),
"length framing keeps \"bot\"+\"dom\" apart from \"bo\"+\"tdom\""
);
Ok(())
}
#[test]
fn an_unattempted_false_condition_still_skips_and_lets_its_successor_run() -> TestResult {
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("legitimate-skip")
.observe(Script::new(vec![200]))
.on(|value: &u16| *value < 200, Record(Rc::clone(&log)))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
assert_eq!(bot.tick()?, 1, "the successor ran behind a skipped entry");
assert_eq!(
*log.borrow(),
vec![200],
"exactly one action, the successor's: {:?}",
*log.borrow()
);
assert_eq!(bot.revisions(), vec![1], "the transition resolved");
assert!(
bot.pending().is_empty(),
"a skipped entry owes nothing: {:?}",
bot.pending()
);
Ok(())
}
#[test]
fn a_condition_that_cannot_be_evaluated_holds_the_chain_without_claiming_anything() -> TestResult
{
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("mismatched")
.observe(Script::new(vec![200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
bot.world.non_send_mut::<Chains>().0[0].entries.insert(
0,
typed_entry::<_, _, String>(|value: &String| value.len() >= 3, Record(Rc::clone(&log))),
);
let error = refused(bot.tick(), |fired| {
format!("an evaluate error was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EvaluateError { .. }),
"expected the condition's own error, got {error:?}"
);
assert!(log.borrow().is_empty(), "nothing was attempted");
assert_eq!(
identities(&bot),
vec![(0, 0), (0, 1)],
"the entry whose condition failed is still outstanding, and so is the entry \
behind it: {:?}",
bot.pending()
);
assert_eq!(
first_hold(&bot),
Some(TransitionHold::NotStarted),
"a condition that cannot be evaluated is not an attempt, so nothing is abandoned"
);
let error = refused(bot.tick(), |fired| {
format!("the second tick reported {fired} fired")
})?;
assert!(
matches!(error, BotError::EvaluateError { .. }),
"the condition is evaluated again, not silently given up on: {error:?}"
);
assert!(log.borrow().is_empty(), "and still nothing was attempted");
Ok(())
}
#[test]
fn a_condition_failure_stops_the_walk_after_the_effects_it_cleared() -> TestResult {
let log = Rc::new(RefCell::new(Vec::new()));
let after = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("condition-failure")
.observe(Script::new(vec![200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.on(RefusesToEvaluate, Count(Rc::clone(&after)))
.on(|value: &u16| *value >= 200, Count(Rc::clone(&after)))
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a refused condition was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::EvaluateError { .. }),
"expected the condition's own error, got {error:?}"
);
assert_eq!(
*log.borrow(),
vec![200],
"the effect the walk cleared before the failing condition ran"
);
assert_eq!(
after.get(),
0,
"the entries after the failing condition were not attempted"
);
Ok(())
}
#[test]
fn an_ambiguous_schedule_is_refused_at_build() -> TestResult {
let mut world = World::new();
let mut ambiguous = Schedule::default();
ambiguous.set_build_settings(ScheduleBuildSettings {
ambiguity_detection: LogLevel::Error,
..Default::default()
});
ambiguous.add_systems((observe_fold, fire_plan));
match validate(&mut ambiguous, &mut world) {
Ok(()) => return Err("an ambiguous schedule was accepted at build".into()),
Err(error) => assert!(
matches!(error, BotError::DomainError { .. }),
"a schedule-build failure must surface as a typed BotError, got {error:?}"
),
}
let mut ordered = Schedule::default();
ordered.set_build_settings(ScheduleBuildSettings {
ambiguity_detection: LogLevel::Error,
..Default::default()
});
ordered.add_systems((observe_fold, fire_plan).chain());
validate(&mut ordered, &mut world)?;
Ok(())
}
struct CountsU32(Rc<Cell<usize>>);
impl Execute for CountsU32 {
type Input = u32;
type Output = ();
fn required_caps(&self) -> &[Cap] {
&[]
}
fn effect_lifetime(&self) -> EffectLifetime {
EffectLifetime::Local
}
async fn execute_action(&self, call: (Auth, &u32)) -> Result<(), BotError> {
call.0.check(&[])?;
self.0.set(self.0.get().saturating_add(1));
Ok(())
}
fn domain_id(&self) -> &str {
"test::counts_u32"
}
}
#[test]
fn a_mispairing_behind_the_erasure_is_a_witness_miss() -> TestResult {
let ran = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("mispairing")
.observe(Script::new(vec![200]))
.on(
|value: &u16| *value >= 200,
Refuses::undelivered(Rc::clone(&ran)),
)
.with_effects(test_effects()?)
.build(&net_grants())?;
bot.world.non_send_mut::<Polled>().0 = vec![Ok(Some(Erased::new(300u32)))];
bot.schedule.run(&mut bot.world);
match bot.world.resource::<TickError>().0 {
Some(BotError::TypeMismatch {
site,
chain,
expected,
observed,
}) => {
assert_eq!(site, "observe_fold rendezvous", "the site is greppable");
assert_eq!(chain, Some(0), "the report names which chain disagreed");
assert_eq!(expected, "u16", "the chain's claim");
assert_eq!(observed, "u32", "what actually arrived");
}
ref other => {
return Err(format!("expected a witness miss, got {other:?}").into());
}
}
assert!(
bot.world
.non_send::<Observed>()
.0
.iter()
.all(Option::is_none),
"the fold commits nothing on a mismatch, so the previous tick's state survives"
);
assert_eq!(ran.get(), 0, "the action was never reached");
Ok(())
}
#[test]
fn a_mismatched_pairing_costs_one_attempt_not_the_retry_budget() -> TestResult {
let ran = Rc::new(Cell::new(0));
let log = Rc::new(RefCell::new(Vec::new()));
let mut bot = EcsBot::builder("mismatched")
.observe(Script::new(vec![200, 200, 200, 200]))
.on(|value: &u16| *value >= 200, Record(Rc::clone(&log)))
.with_effects(test_effects()?)
.build(&net_grants())?;
bot.world.non_send_mut::<Chains>().0[0].entries[0] =
typed_entry::<_, _, u16>(|value: &u16| *value >= 200, CountsU32(Rc::clone(&ran)));
let error = refused(bot.tick(), |fired| {
format!("a wiring mismatch was reported as {fired} fired")
})?;
assert!(
matches!(
error,
BotError::TypeMismatch {
site: "spec::typed_entry",
expected: "u32",
observed: "u16",
..
}
),
"a wiring defect is a type mismatch naming both types, not a domain failure: {error:?}"
);
assert_eq!(
ran.get(),
0,
"the downcast is checked before the action runs, so the mismatch was not a side effect"
);
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::Abandoned {
reason: AbandonReason::Terminal,
..
})
),
"one attempt, and the disposition is terminal rather than `AttemptsExhausted`: {:?}",
first_hold(&bot)
);
let error = refused(bot.tick(), |fired| {
format!("a retried wiring mismatch reported {fired} fired")
})?;
assert!(
matches!(error, BotError::PendingTransition { .. }),
"the abandoned entry is the barrier the next tick reports, got {error:?}"
);
Ok(())
}
#[test]
fn a_permanent_refusal_does_not_burn_the_retry_budget() -> TestResult {
let attempts = Rc::new(Cell::new(0));
let mut bot = EcsBot::builder("permanent")
.observe(Script::new(vec![200, 200, 200, 200]))
.on(
|value: &u16| *value >= 200,
Refuses::permanently(Rc::clone(&attempts)),
)
.with_effects(test_effects()?)
.build(&net_grants())?;
let error = refused(bot.tick(), |fired| {
format!("a permanent refusal was reported as {fired} fired")
})?;
assert!(
matches!(error, BotError::DomainError { .. }),
"expected the action's own error, got {error:?}"
);
assert_eq!(attempts.get(), 1, "a permanent refusal is attempted once");
assert!(
matches!(
first_hold(&bot),
Some(TransitionHold::Abandoned {
reason: AbandonReason::Terminal,
..
})
),
"the disposition is terminal, not `AttemptsExhausted`: {:?}",
first_hold(&bot)
);
Ok(())
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct PrefixProbe;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct PrefixProbeTail;
const PREFIX_PROBE_SUFFIX: &[u8] = b"Tail";
impl InputIdentity for PrefixProbe {
const SCHEMA_ID: &'static [u8] = b"lgwks.bot.schema.v1.test-prefix-probe";
fn write_identity(&self, hasher: &mut Hasher) {
hasher.update(PREFIX_PROBE_SUFFIX);
hasher.update(b"payload");
}
}
impl InputIdentity for PrefixProbeTail {
const SCHEMA_ID: &'static [u8] = b"lgwks.bot.schema.v1.test-prefix-probe-tail";
fn write_identity(&self, hasher: &mut Hasher) {
hasher.update(b"payload");
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct PrefixProbeNextSchema;
impl InputIdentity for PrefixProbeNextSchema {
const SCHEMA_ID: &'static [u8] = b"lgwks.bot.schema.v2.test-prefix-probe";
fn write_identity(&self, hasher: &mut Hasher) {
hasher.update(PREFIX_PROBE_SUFFIX);
hasher.update(b"payload");
}
}
struct Unpolled<T>(PhantomData<fn() -> T>);
impl<T> Observe for Unpolled<T> {
type Output = T;
fn required_caps(&self) -> &[Cap] {
&[]
}
async fn poll(&self, _call: (Auth, ())) -> Result<T, BotError> {
domain_refusal("test::unpolled", "identity fixtures are never polled")
}
fn domain_id(&self) -> &str {
"test::unpolled"
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct AltTick(u64);
impl InputIdentity for AltTick {
const SCHEMA_ID: &'static [u8] = b"lgwks.bot.schema.v1.test-alt-tick";
fn write_identity(&self, hasher: &mut Hasher) {
hasher.update(&self.0.to_le_bytes());
}
}
#[test]
fn two_output_types_sharing_a_name_prefix_hash_two_identities() {
let plain = identify_output::<Unpolled<PrefixProbe>>(&PrefixProbe);
let tail = identify_output::<Unpolled<PrefixProbeTail>>(&PrefixProbeTail);
assert_ne!(
plain.identity, tail.identity,
"the declared schema ids must separate the pair their name and \
byte streams would collide"
);
}
#[test]
fn a_schema_bump_moves_the_identity_of_identical_bytes() {
let current = identify_output::<Unpolled<PrefixProbe>>(&PrefixProbe);
let next = identify_output::<Unpolled<PrefixProbeNextSchema>>(&PrefixProbeNextSchema);
assert_ne!(current.identity, next.identity);
}
#[test]
fn the_same_output_hashes_one_identity_every_time() {
let first = identify_output::<Script>(&7_u16);
let second = identify_output::<Script>(&7_u16);
assert_eq!(first.identity, second.identity);
assert_eq!(first.event, second.event);
}
#[test]
fn a_downcast_miss_still_names_the_binding_it_refused() {
let hit = identify_output::<Script>(&7_u16);
let miss = identify_output::<Script>(&1_000_u32);
assert_ne!(
hit.identity, miss.identity,
"a missed downcast must not read as an admission of the value"
);
}
#[test]
fn the_same_event_twice_and_two_events_hash_as_documented() {
let redelivery = identify_output::<Unpolled<EventId<u64>>>(&EventId::new(5_u64, 7_u64));
let again = identify_output::<Unpolled<EventId<u64>>>(&EventId::new(5_u64, 7_u64));
assert_eq!(
redelivery.identity, again.identity,
"the same event id over the same payload is one event"
);
assert!(redelivery.event, "EventId names an event");
let other_id = identify_output::<Unpolled<EventId<u64>>>(&EventId::new(6_u64, 7_u64));
assert_ne!(redelivery.identity, other_id.identity);
let other_payload =
identify_output::<Unpolled<EventId<AltTick>>>(&EventId::new(5_u64, AltTick(7_u64)));
assert_ne!(redelivery.identity, other_payload.identity);
}
#[test]
fn identity_known_answer_vectors_hold() {
let hex = |identity: [u8; 16]| {
identity
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>()
};
let u16_seven = hex(identify_output::<Script>(&7_u16).identity);
assert_eq!(
u16_seven, "878e8add50e816abdbb0901de9ed7788",
"the u16 vector moved: the v2 identity scheme changed"
);
let str_value = identify_output::<Unpolled<String>>(&String::from("settle")).identity;
assert_eq!(
hex(str_value),
"174d75182e8e31c49374da2f22e0b353",
"the string vector moved: the v2 identity scheme changed"
);
let event = identify_output::<Unpolled<EventId<u64>>>(&EventId::new(9_u64, true)).identity;
assert_eq!(
hex(event),
"f4d5d8a77b9fed69dd1474cf0620348a",
"the event-id vector moved: the v2 identity scheme changed"
);
let miss = identify_output::<Script>(&1_000_u32).identity;
assert_eq!(
hex(miss),
"f2e18c137ba6e62e78f672c22c02bc66",
"the refused-binding vector moved: the v2 identity scheme changed"
);
}
#[test]
fn applied_membership_answers_at_scale_where_a_scan_would_be_quadratic() -> TestResult {
const COUNT: u64 = 50_000;
fn action_for(index: u64) -> Result<ActionId, Box<dyn std::error::Error>> {
Ok(ActionId::from_hex(&format!(
"{:032x}",
index.wrapping_add(1)
))?)
}
fn coordinates(index: u64) -> Result<(usize, usize), Box<dyn std::error::Error>> {
let chain = index.checked_rem(16).ok_or("sixteen chains is not zero")?;
let entry = index.checked_div(16).ok_or("sixteen chains is not zero")?;
Ok((usize::try_from(chain)?, usize::try_from(entry)?))
}
let scope = test_effects()?;
let mut effects = Effects::new(scope, JournalPosition::genesis());
let identity = effects.identity();
let flow = identity.flow();
let run = identity.run();
let environment = identity.environment();
let epoch = EnvironmentEpoch::from_decimal("1")?;
let attempt = AttemptId::from_decimal("1")?;
let input = [0u8; 16];
for index in 0..COUNT {
let (chain, entry) = coordinates(index)?;
let key = crate::effect::EffectIdentity::new(run, environment, flow).key(
action_for(index)?,
attempt,
derive_action_digest(flow, chain, entry, &input),
epoch,
);
effects.note_applied(key);
}
let (chain, entry) = coordinates(0)?;
effects.note_applied(
crate::effect::EffectIdentity::new(run, environment, flow).key(
action_for(0)?,
attempt,
derive_action_digest(flow, chain, entry, &input),
epoch,
),
);
assert_eq!(
effects.applied_events.len(),
usize::try_from(COUNT)?,
"a duplicate settlement changed the event-identity index"
);
assert_eq!(
effects.applied_latest.len(),
usize::try_from(COUNT)?,
"a duplicate settlement changed the latest-key index"
);
for index in 0..COUNT {
let (chain, entry) = coordinates(index)?;
let action = action_for(index)?;
assert!(
effects.applied_in(action, chain, entry, input, true),
"the event identity for obligation {index} is absent from the index"
);
assert!(
effects.applied_in(action, chain, entry, input, false),
"the contract identity for obligation {index} is absent from the index"
);
}
let absent = action_for(COUNT)?;
assert!(
!effects.applied_in(absent, 0, 0, input, true)
&& !effects.applied_in(absent, 0, 0, input, false),
"an unapplied obligation must not resolve as done"
);
Ok(())
}
struct NoSettlementRoom {
inner: MemoryJournal,
appends: Rc<Cell<usize>>,
}
impl EffectJournal for NoSettlementRoom {
fn durability(&self) -> DurabilityPromise {
DurabilityPromise::ProcessCrash
}
fn tail(&self) -> JournalPosition {
self.inner.tail()
}
fn committed(&self) -> Result<Vec<EffectEvent>, JournalError> {
EffectJournal::committed(&self.inner)
}
fn compare_and_append(
&mut self,
expected_tail: JournalPosition,
event: &EffectEvent,
) -> Result<DurableAck, JournalError> {
self.appends.set(self.appends.get().saturating_add(1));
self.inner.compare_and_append(expected_tail, event)
}
fn reserve_handoff_capacity(&self, rungs: u64) -> Result<(), JournalError> {
Err(JournalError::CapacityExceeded {
resource: crate::journal::JournalLimitKind::Events,
limit: 0,
requested: rungs,
})
}
}
#[test]
fn an_external_handoff_reserves_settlement_capacity_before_any_rung() -> TestResult {
let appends = Rc::new(Cell::new(0));
let journal = NoSettlementRoom {
inner: MemoryJournal::new(),
appends: Rc::clone(&appends),
};
let scope = test_effects_with(Box::new(journal))?;
let mut effects = Effects::new(scope, JournalPosition::genesis());
let key = crate::effect::EffectIdentity::new(
RunId::from_hex(TEST_RUN)?,
EnvironmentId::from_hex(TEST_ENV)?,
FlowRevision::from_tagged("blake3_256", TEST_FLOW)?,
)
.key(
ActionId::from_hex("1112131415161718191a1b1c1d1e1f20")?,
AttemptId::from_decimal("1")?,
ActionDigest::from_tagged(
"blake3_256",
"f0f1f2f3f4f5f6f7f8f9fafbfcfdfeffe0e1e2e3e4e5e6e7e8e9eaebecedeeef",
)?,
EnvironmentEpoch::from_decimal("1")?,
);
match lgwks_std::task::block_on(effects.prepare(key, EffectLifetime::External)) {
Err(DispatchError::Journal(JournalError::CapacityExceeded { .. })) => {}
Err(other) => {
return Err(format!("expected a capacity refusal, got {other:?}").into());
}
Ok(_) => {
return Err("a handoff with no settlement room must be refused".into());
}
}
assert_eq!(
appends.get(),
0,
"the handoff must be refused before the first rung is written"
);
Ok(())
}
}