use std::collections::{BTreeMap, HashMap};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use sha2::{Digest, Sha256};
use thiserror::Error;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema)]
pub struct VectorClock(BTreeMap<String, u64>);
impl VectorClock {
pub fn new() -> Self {
Self::default()
}
pub fn get(&self, agent: &str) -> u64 {
self.0.get(agent).copied().unwrap_or(0)
}
pub fn tick(&mut self, agent: &str) -> u64 {
let counter = self.0.entry(agent.to_string()).or_insert(0);
*counter += 1;
*counter
}
pub fn merge(&mut self, other: &VectorClock) {
for (agent, &count) in &other.0 {
let entry = self.0.entry(agent.clone()).or_insert(0);
if count > *entry {
*entry = count;
}
}
}
pub fn le(&self, other: &VectorClock) -> bool {
self.0
.iter()
.all(|(agent, &count)| other.get(agent) >= count)
}
pub fn concurrent(&self, other: &VectorClock) -> bool {
!self.le(other) && !other.le(self)
}
pub fn digest(&self) -> String {
let bytes =
serde_json::to_vec(&self.0).expect("BTreeMap<String, u64> serialization cannot fail");
hex_encode(&Sha256::digest(&bytes))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema)]
pub struct ReadSetEntry {
pub domain: String,
pub seq: u64,
}
impl ReadSetEntry {
pub fn new(domain: impl Into<String>, seq: u64) -> Self {
Self {
domain: domain.into(),
seq,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, schemars::JsonSchema)]
pub struct EffectRequest {
pub intent: String,
#[serde(default)]
pub parent: Option<String>,
pub domain: String,
pub tool: String,
#[serde(default)]
pub args: Value,
#[serde(default)]
pub read_set: Vec<ReadSetEntry>,
pub agent: String,
#[serde(default)]
pub known_clock: VectorClock,
}
#[derive(Debug, Clone, Copy)]
pub struct FenceConfig {
pub lease_ttl: Duration,
pub result_ttl: Duration,
}
impl Default for FenceConfig {
fn default() -> Self {
Self {
lease_ttl: Duration::from_secs(60),
result_ttl: Duration::from_secs(24 * 60 * 60),
}
}
}
#[derive(Debug)]
enum IntentState {
InFlight { lease_expires_at: Instant },
Done {
cert: Arc<EffectCert>,
recorded_at: Instant,
},
Failed {
reason: Arc<str>,
recorded_at: Instant,
},
}
#[derive(Debug)]
struct Inner {
config: FenceConfig,
domains: Mutex<HashMap<String, Arc<AtomicU64>>>,
intents: Mutex<HashMap<String, IntentState>>,
}
#[derive(Debug, Clone)]
pub struct EffectFence {
inner: Arc<Inner>,
}
impl Default for EffectFence {
fn default() -> Self {
Self::new()
}
}
impl EffectFence {
pub fn new() -> Self {
Self::with_config(FenceConfig::default())
}
pub fn with_config(config: FenceConfig) -> Self {
Self {
inner: Arc::new(Inner {
config,
domains: Mutex::new(HashMap::new()),
intents: Mutex::new(HashMap::new()),
}),
}
}
pub fn current(&self, domain: &str) -> u64 {
let domains = self.inner.domains.lock().expect("fence mutex poisoned");
domains
.get(domain)
.map(|counter| counter.load(Ordering::SeqCst))
.unwrap_or(0)
}
pub fn reserve(&self, domain: &str, expected_seq: u64) -> Result<u64, u64> {
let counter = {
let mut domains = self.inner.domains.lock().expect("fence mutex poisoned");
domains
.entry(domain.to_string())
.or_insert_with(|| Arc::new(AtomicU64::new(0)))
.clone()
};
let next_seq = expected_seq + 1;
counter
.compare_exchange(expected_seq, next_seq, Ordering::SeqCst, Ordering::SeqCst)
.map(|_previous| next_seq)
}
pub fn domain_count(&self) -> usize {
self.inner
.domains
.lock()
.expect("fence mutex poisoned")
.len()
}
pub fn intent_count(&self) -> usize {
self.inner
.intents
.lock()
.expect("fence mutex poisoned")
.len()
}
pub fn clear_intent(&self, intent: &str) -> bool {
self.inner
.intents
.lock()
.expect("fence mutex poisoned")
.remove(intent)
.is_some()
}
pub fn evict_domain(&self, domain: &str) -> bool {
self.inner
.domains
.lock()
.expect("fence mutex poisoned")
.remove(domain)
.is_some()
}
pub fn sweep(&self) {
let now = Instant::now();
let result_ttl = self.inner.config.result_ttl;
let mut intents = self.inner.intents.lock().expect("fence mutex poisoned");
intents.retain(|_, state| match state {
IntentState::InFlight { lease_expires_at } => now < *lease_expires_at,
IntentState::Done { recorded_at, .. } | IntentState::Failed { recorded_at, .. } => {
now.saturating_duration_since(*recorded_at) < result_ttl
}
});
}
fn record_outcome(&self, intent: &str, state: IntentState) {
self.inner
.intents
.lock()
.expect("fence mutex poisoned")
.insert(intent.to_string(), state);
}
fn release_intent(&self, intent: &str) {
self.inner
.intents
.lock()
.expect("fence mutex poisoned")
.remove(intent);
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema)]
pub struct EffectCert {
pub intent: String,
pub parent: Option<String>,
pub domain: String,
pub seq: u64,
pub tool: String,
pub args: Value,
pub result: Value,
pub vector_clock: VectorClock,
pub read_set: Vec<ReadSetEntry>,
pub agent: String,
pub hash: String,
}
#[derive(Serialize)]
struct CertPreimage<'a> {
intent: &'a str,
parent: &'a Option<String>,
domain: &'a str,
seq: u64,
tool: &'a str,
args: &'a Value,
result: &'a Value,
vector_clock: &'a VectorClock,
read_set: &'a [ReadSetEntry],
agent: &'a str,
}
impl EffectCert {
fn new(prepared: PreparedEffect, result: Value) -> Self {
let hash = Self::compute_hash(
&prepared.intent,
&prepared.parent,
&prepared.domain,
prepared.seq,
&prepared.tool,
&prepared.args,
&result,
&prepared.vector_clock,
&prepared.read_set,
&prepared.agent,
);
Self {
intent: prepared.intent,
parent: prepared.parent,
domain: prepared.domain,
seq: prepared.seq,
tool: prepared.tool,
args: prepared.args,
result,
vector_clock: prepared.vector_clock,
read_set: prepared.read_set,
agent: prepared.agent,
hash,
}
}
#[allow(clippy::too_many_arguments)]
fn compute_hash(
intent: &str,
parent: &Option<String>,
domain: &str,
seq: u64,
tool: &str,
args: &Value,
result: &Value,
vector_clock: &VectorClock,
read_set: &[ReadSetEntry],
agent: &str,
) -> String {
let preimage = CertPreimage {
intent,
parent,
domain,
seq,
tool,
args,
result,
vector_clock,
read_set,
agent,
};
let bytes = serde_json::to_vec(&preimage).expect("CertPreimage serialization cannot fail");
hex_encode(&Sha256::digest(&bytes))
}
pub fn verify(&self) -> bool {
let expected = Self::compute_hash(
&self.intent,
&self.parent,
&self.domain,
self.seq,
&self.tool,
&self.args,
&self.result,
&self.vector_clock,
&self.read_set,
&self.agent,
);
expected == self.hash
}
}
#[derive(Debug, Clone, Serialize, Deserialize, schemars::JsonSchema)]
pub struct PreparedEffect {
pub intent: String,
pub parent: Option<String>,
pub domain: String,
pub seq: u64,
pub tool: String,
pub args: Value,
pub vector_clock: VectorClock,
pub read_set: Vec<ReadSetEntry>,
pub agent: String,
}
#[derive(Debug)]
pub enum Admission {
Fresh(PreparedEffect),
Replay(Arc<EffectCert>),
}
#[derive(Debug, Clone, PartialEq, Eq, Error)]
pub enum FenceError {
#[error(
"read set entry for domain `{domain}` is stale: expected seq {expected}, actual {actual}"
)]
ReadSetStale {
domain: String,
expected: u64,
actual: u64,
},
#[error(
"domain `{domain}` race: expected to reserve seq {expected}, but the domain was already at {actual}"
)]
DomainRace {
domain: String,
expected: u64,
actual: u64,
},
#[error("intent `{intent}` is already in flight: another attempt is executing it right now")]
IntentInFlight { intent: String },
#[error(
"intent `{intent}` previously failed ({reason}); outcome is unknown -- reconcile with \
the downstream system, then clear_intent to allow a retry"
)]
IntentFailed { intent: String, reason: String },
}
pub fn prepare_effect_fence(
fence: &EffectFence,
req: EffectRequest,
) -> Result<Admission, FenceError> {
let now = Instant::now();
let config = fence.inner.config;
{
let mut intents = fence.inner.intents.lock().expect("fence mutex poisoned");
match intents.get(&req.intent) {
Some(IntentState::Done { cert, recorded_at })
if now.saturating_duration_since(*recorded_at) < config.result_ttl =>
{
return Ok(Admission::Replay(cert.clone()));
}
Some(IntentState::Failed {
reason,
recorded_at,
}) if now.saturating_duration_since(*recorded_at) < config.result_ttl => {
return Err(FenceError::IntentFailed {
intent: req.intent.clone(),
reason: reason.to_string(),
});
}
Some(IntentState::InFlight { lease_expires_at }) if now < *lease_expires_at => {
return Err(FenceError::IntentInFlight {
intent: req.intent.clone(),
});
}
_ => {}
}
intents.insert(
req.intent.clone(),
IntentState::InFlight {
lease_expires_at: now + config.lease_ttl,
},
);
}
for entry in &req.read_set {
let actual = fence.current(&entry.domain);
if actual != entry.seq {
fence.release_intent(&req.intent);
return Err(FenceError::ReadSetStale {
domain: entry.domain.clone(),
expected: entry.seq,
actual,
});
}
}
let expected_seq = fence.current(&req.domain);
let seq = match fence.reserve(&req.domain, expected_seq) {
Ok(seq) => seq,
Err(actual) => {
fence.release_intent(&req.intent);
return Err(FenceError::DomainRace {
domain: req.domain.clone(),
expected: expected_seq,
actual,
});
}
};
let mut vector_clock = req.known_clock;
vector_clock.tick(&req.agent);
Ok(Admission::Fresh(PreparedEffect {
intent: req.intent,
parent: req.parent,
domain: req.domain,
seq,
tool: req.tool,
args: req.args,
vector_clock,
read_set: req.read_set,
agent: req.agent,
}))
}
pub fn commit_effect_cert(
fence: &EffectFence,
prepared: PreparedEffect,
result: Value,
) -> Result<EffectCert, FenceError> {
for entry in &prepared.read_set {
let actual = fence.current(&entry.domain);
if actual != entry.seq {
fence.record_outcome(
&prepared.intent,
IntentState::Failed {
reason: Arc::from(format!(
"commit rejected: read-set entry for `{}` went stale (expected {}, \
actual {}) while the effect was running",
entry.domain, entry.seq, actual
)),
recorded_at: Instant::now(),
},
);
return Err(FenceError::ReadSetStale {
domain: entry.domain.clone(),
expected: entry.seq,
actual,
});
}
}
let intent = prepared.intent.clone();
let cert = EffectCert::new(prepared, result);
fence.record_outcome(
&intent,
IntentState::Done {
cert: Arc::new(cert.clone()),
recorded_at: Instant::now(),
},
);
Ok(cert)
}
pub fn abort_effect(fence: &EffectFence, prepared: PreparedEffect, reason: impl Into<String>) {
fence.record_outcome(
&prepared.intent,
IntentState::Failed {
reason: Arc::from(reason.into()),
recorded_at: Instant::now(),
},
);
}
fn hex_encode(bytes: &[u8]) -> String {
use std::fmt::Write;
let mut out = String::with_capacity(bytes.len() * 2);
for byte in bytes {
write!(out, "{byte:02x}").expect("writing to a String cannot fail");
}
out
}
#[cfg(test)]
mod tests {
use super::*;
fn req(intent: &str, domain: &str, agent: &str) -> EffectRequest {
EffectRequest {
intent: intent.to_string(),
parent: None,
domain: domain.to_string(),
tool: "charge_card".to_string(),
args: serde_json::json!({"amount_cents": 1999}),
read_set: vec![],
agent: agent.to_string(),
known_clock: VectorClock::new(),
}
}
fn fresh(admission: Admission) -> PreparedEffect {
match admission {
Admission::Fresh(prepared) => prepared,
Admission::Replay(cert) => panic!("expected Fresh, got Replay of {}", cert.hash),
}
}
#[test]
fn vector_clock_tick_and_get() {
let mut vc = VectorClock::new();
assert_eq!(vc.get("agent-a"), 0);
assert_eq!(vc.tick("agent-a"), 1);
assert_eq!(vc.tick("agent-a"), 2);
assert_eq!(vc.get("agent-a"), 2);
assert_eq!(vc.get("agent-b"), 0);
}
#[test]
fn vector_clock_merge_takes_elementwise_max() {
let mut a = VectorClock::new();
a.tick("x"); a.tick("x"); let mut b = VectorClock::new();
b.tick("x"); b.tick("y"); b.tick("y");
a.merge(&b);
assert_eq!(a.get("x"), 2);
assert_eq!(a.get("y"), 2);
}
#[test]
fn vector_clock_partial_order_and_concurrency() {
let mut a = VectorClock::new();
a.tick("x");
let mut b = a.clone();
b.tick("x");
assert!(a.le(&b));
assert!(!b.le(&a));
assert!(!a.concurrent(&b));
let mut c = VectorClock::new();
c.tick("y"); assert!(!a.le(&c));
assert!(!c.le(&a));
assert!(a.concurrent(&c));
}
#[test]
fn vector_clock_digest_is_stable_and_order_independent() {
let mut a = VectorClock::new();
a.tick("agent-a");
a.tick("agent-b");
let mut b = VectorClock::new();
b.tick("agent-b");
b.tick("agent-a");
assert_eq!(a, b);
assert_eq!(a.digest(), b.digest());
}
#[test]
fn domain_reserve_advances_and_rejects_stale_expectation() {
let fence = EffectFence::new();
assert_eq!(fence.current("d"), 0);
assert_eq!(fence.reserve("d", 0), Ok(1));
assert_eq!(fence.current("d"), 1);
assert_eq!(fence.reserve("d", 0), Err(1));
assert_eq!(fence.reserve("d", 1), Ok(2));
assert_eq!(fence.current("d"), 2);
}
#[test]
fn asking_about_a_domain_does_not_create_tracking_state() {
let fence = EffectFence::new();
assert_eq!(fence.current("ghost:1"), 0);
assert_eq!(fence.current("ghost:2"), 0);
assert_eq!(
fence.domain_count(),
0,
"current() must be a pure read -- queries must not grow the domain map"
);
fence.reserve("real:1", 0).unwrap();
assert_eq!(fence.domain_count(), 1);
assert!(fence.evict_domain("real:1"));
assert_eq!(fence.domain_count(), 0);
}
#[test]
fn prepare_then_commit_happy_path_produces_a_verifiable_cert() {
let fence = EffectFence::new();
let prepared = fresh(
prepare_effect_fence(&fence, req("charge:order-1", "order:1", "agent-a"))
.expect("first attempt on a free intent+domain must succeed"),
);
assert_eq!(prepared.seq, 1);
assert_eq!(prepared.vector_clock.get("agent-a"), 1);
let cert = commit_effect_cert(&fence, prepared, serde_json::json!({"charge_id": "ch_123"}))
.expect("commit with an untouched read set must succeed");
assert!(cert.verify());
assert_eq!(cert.intent, "charge:order-1");
assert_eq!(cert.domain, "order:1");
assert_eq!(cert.seq, 1);
assert!(cert.parent.is_none());
}
#[test]
fn late_duplicate_of_a_done_intent_is_replayed_not_rerun() {
let fence = EffectFence::new();
let prepared = fresh(
prepare_effect_fence(&fence, req("charge:order-2", "order:2", "agent-a")).unwrap(),
);
let cert =
commit_effect_cert(&fence, prepared, serde_json::json!({"charge_id": "ch_1"})).unwrap();
assert_eq!(fence.current("order:2"), 1);
let replayed =
match prepare_effect_fence(&fence, req("charge:order-2", "order:2", "agent-b"))
.expect("duplicate of a done intent is not an error")
{
Admission::Replay(cert) => cert,
Admission::Fresh(_) => panic!("duplicate intent must NOT be admitted to run"),
};
assert_eq!(replayed.hash, cert.hash);
assert_eq!(
fence.current("order:2"),
1,
"a replayed duplicate must not consume a new domain sequence"
);
}
#[test]
fn duplicate_of_an_in_flight_intent_is_rejected() {
let fence = EffectFence::new();
let _prepared = fresh(
prepare_effect_fence(&fence, req("charge:order-3", "order:3", "agent-a")).unwrap(),
);
let err = prepare_effect_fence(&fence, req("charge:order-3", "order:3", "agent-b"))
.expect_err("an in-flight intent must reject duplicates");
assert_eq!(
err,
FenceError::IntentInFlight {
intent: "charge:order-3".to_string()
}
);
}
#[test]
fn aborted_intent_is_fenced_until_explicitly_cleared() {
let fence = EffectFence::new();
let prepared = fresh(
prepare_effect_fence(&fence, req("charge:order-4", "order:4", "agent-a")).unwrap(),
);
abort_effect(&fence, prepared, "card declined");
match prepare_effect_fence(&fence, req("charge:order-4", "order:4", "agent-a")) {
Err(FenceError::IntentFailed { intent, reason }) => {
assert_eq!(intent, "charge:order-4");
assert!(reason.contains("card declined"));
}
other => panic!("expected IntentFailed, got {other:?}"),
}
assert!(fence.clear_intent("charge:order-4"));
let retried = fresh(
prepare_effect_fence(&fence, req("charge:order-4", "order:4", "agent-a"))
.expect("cleared intent must be runnable again"),
);
assert_eq!(retried.seq, 2, "the aborted attempt's seq 1 stays burned");
}
#[test]
fn expired_lease_can_be_taken_over_by_a_later_attempt() {
let fence = EffectFence::with_config(FenceConfig {
lease_ttl: Duration::from_millis(5),
..FenceConfig::default()
});
let _abandoned = fresh(
prepare_effect_fence(&fence, req("charge:order-5", "order:5", "agent-a")).unwrap(),
);
std::thread::sleep(Duration::from_millis(20));
let takeover = fresh(
prepare_effect_fence(&fence, req("charge:order-5", "order:5", "agent-b"))
.expect("an expired lease must be claimable"),
);
assert_eq!(takeover.seq, 2);
}
#[test]
fn rejected_admission_releases_the_intent_lease() {
let fence = EffectFence::new();
fence.reserve("inventory:sku-9", 0).unwrap();
let mut r = req("charge:order-6", "order:6", "agent-a");
r.read_set = vec![ReadSetEntry::new("inventory:sku-9", 0)];
let err = prepare_effect_fence(&fence, r).expect_err("stale read set must reject");
assert!(matches!(err, FenceError::ReadSetStale { .. }));
let ok = fresh(
prepare_effect_fence(&fence, {
let mut r = req("charge:order-6", "order:6", "agent-a");
r.read_set = vec![ReadSetEntry::new("inventory:sku-9", 1)];
r
})
.expect("a corrected retry of a rejected (never-run) attempt must be admitted"),
);
assert_eq!(ok.intent, "charge:order-6");
}
#[test]
fn commit_rejection_marks_the_intent_failed_not_runnable() {
let fence = EffectFence::new();
fence.reserve("inventory:sku-1", 0).unwrap();
let mut r = req("reserve:order-7", "order:7", "agent-a");
r.tool = "reserve_inventory".to_string();
r.read_set = vec![ReadSetEntry::new("inventory:sku-1", 1)];
let prepared = fresh(prepare_effect_fence(&fence, r).expect("read set matches"));
fence.reserve("inventory:sku-1", 1).unwrap();
let err = commit_effect_cert(&fence, prepared, serde_json::json!({}))
.expect_err("dependency moved -- commit must be rejected");
assert!(matches!(err, FenceError::ReadSetStale { .. }));
match prepare_effect_fence(&fence, req("reserve:order-7", "order:7", "agent-a")) {
Err(FenceError::IntentFailed { .. }) => {}
other => panic!("expected IntentFailed after a rejected commit, got {other:?}"),
}
}
#[test]
fn sweep_evicts_expired_outcomes() {
let fence = EffectFence::with_config(FenceConfig {
result_ttl: Duration::from_millis(5),
..FenceConfig::default()
});
let prepared = fresh(
prepare_effect_fence(&fence, req("charge:order-8", "order:8", "agent-a")).unwrap(),
);
commit_effect_cert(&fence, prepared, serde_json::json!({})).unwrap();
assert_eq!(fence.intent_count(), 1);
std::thread::sleep(Duration::from_millis(20));
fence.sweep();
assert_eq!(fence.intent_count(), 0, "expired outcome must be evicted");
let rerun = fresh(
prepare_effect_fence(&fence, req("charge:order-8", "order:8", "agent-a")).unwrap(),
);
assert_eq!(rerun.seq, 2);
}
#[test]
fn different_intents_on_the_same_domain_queue_up_sequentially() {
let fence = EffectFence::new();
let first = fresh(
prepare_effect_fence(&fence, req("charge:june", "account:9", "agent-a")).unwrap(),
);
commit_effect_cert(&fence, first, serde_json::json!({})).unwrap();
let second = fresh(
prepare_effect_fence(&fence, req("charge:july", "account:9", "agent-b"))
.expect("a different intent is a different action and may proceed"),
);
assert_eq!(second.seq, 2);
}
#[test]
fn cert_chain_records_parent_for_causal_lineage() {
let fence = EffectFence::new();
let first = commit_effect_cert(
&fence,
fresh(
prepare_effect_fence(&fence, req("charge:order-10", "order:10", "agent-a"))
.unwrap(),
),
serde_json::json!({"charge_id": "ch_1"}),
)
.unwrap();
let mut receipt_req = req("receipt:order-10", "order:10", "agent-b");
receipt_req.tool = "send_receipt".to_string();
receipt_req.parent = Some(first.hash.clone());
receipt_req.known_clock = first.vector_clock.clone();
let second = commit_effect_cert(
&fence,
fresh(prepare_effect_fence(&fence, receipt_req).unwrap()),
serde_json::json!({"receipt_id": "rc_1"}),
)
.unwrap();
assert_eq!(second.parent, Some(first.hash));
assert!(second.vector_clock.get("agent-a") >= 1);
assert!(second.vector_clock.get("agent-b") >= 1);
assert!(second.verify());
}
}