use std::sync::{Arc, OnceLock, RwLock};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, Default)]
#[serde(rename_all = "snake_case")]
pub enum ProvenanceMode {
#[default]
Attached,
Cleared,
}
impl ProvenanceMode {
pub fn from_ir(s: &str) -> Self {
match s {
"cleared" => ProvenanceMode::Cleared,
_ => ProvenanceMode::Attached,
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct FieldProvenance {
pub level: String,
pub confidence: f64,
pub source: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DeliveredField {
pub name: String,
pub value: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub provenance: Option<FieldProvenance>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CanonicalOp {
pub kind: String,
pub idempotency_key: String,
pub fields: Vec<DeliveredField>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DeliveryRequest {
pub name: String,
pub target: String,
pub provenance: ProvenanceMode,
pub ops: Vec<CanonicalOp>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DeliveryReceipt {
pub kind: String,
pub idempotency_key: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub record_id: Option<String>,
pub created: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DeliveryError {
NoProviderConfigured,
Caller(String),
Provider(String),
Network(String),
TierDisabled,
}
impl std::fmt::Display for DeliveryError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
DeliveryError::NoProviderConfigured => write!(
f,
"deliver: no delivery transducer is configured — the OSS build has no engine \
(this is a typed refusal, never a fabricated receipt)"
),
DeliveryError::Caller(m) => write!(f, "deliver: malformed operation — {m}"),
DeliveryError::Provider(m) => write!(f, "deliver: the CRM rejected the write — {m}"),
DeliveryError::Network(m) => write!(f, "deliver: the CRM was unreachable — {m}"),
DeliveryError::TierDisabled => write!(
f,
"deliver: the tenant's `crm.deliver_enabled` tier is OFF (a human must enable it)"
),
}
}
}
pub trait DeliveryProvider: Send + Sync {
fn name(&self) -> &str;
fn deliver(&self, req: &DeliveryRequest) -> Result<Vec<DeliveryReceipt>, DeliveryError>;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct NoProvider;
impl DeliveryProvider for NoProvider {
fn name(&self) -> &str {
"none"
}
fn deliver(&self, _req: &DeliveryRequest) -> Result<Vec<DeliveryReceipt>, DeliveryError> {
Err(DeliveryError::NoProviderConfigured)
}
}
fn registry() -> &'static RwLock<Option<Arc<dyn DeliveryProvider>>> {
static REG: OnceLock<RwLock<Option<Arc<dyn DeliveryProvider>>>> = OnceLock::new();
REG.get_or_init(|| RwLock::new(None))
}
pub fn register_provider(provider: Arc<dyn DeliveryProvider>) {
*registry().write().expect("delivery registry poisoned") = Some(provider);
}
pub fn clear_provider() {
*registry().write().expect("delivery registry poisoned") = None;
}
pub fn active_provider() -> Option<Arc<dyn DeliveryProvider>> {
registry().read().expect("delivery registry poisoned").clone()
}
pub fn plan_delivery(
ir: &crate::ir_nodes::IRDeliver,
mut resolve: impl FnMut(&str) -> Option<String>,
) -> Result<DeliveryRequest, DeliveryError> {
let mut ops = Vec::with_capacity(ir.ops.len());
for op in &ir.ops {
let mut idempotency_key = None;
let mut fields = Vec::with_capacity(op.fields.len());
for f in &op.fields {
let value = match f.kind {
"ref" => resolve(&f.value).ok_or_else(|| {
DeliveryError::Caller(format!(
"operation `{}` binds unresolved flow value `{}`",
op.kind, f.value
))
})?,
_ => f.value.clone(),
};
if f.name == "key" {
idempotency_key = Some(value.clone());
}
fields.push(DeliveredField {
name: f.name.clone(),
value,
provenance: None,
});
}
let idempotency_key = idempotency_key.ok_or_else(|| {
DeliveryError::Caller(format!("operation `{}` resolved no idempotency `key`", op.kind))
})?;
ops.push(CanonicalOp {
kind: op.kind.clone(),
idempotency_key,
fields,
});
}
Ok(DeliveryRequest {
name: ir.name.clone(),
target: ir.target.clone(),
provenance: ProvenanceMode::from_ir(&ir.provenance),
ops,
})
}
pub fn run_delivery(req: &DeliveryRequest) -> Result<Vec<DeliveryReceipt>, DeliveryError> {
match active_provider() {
Some(p) => p.deliver(req),
None => Err(DeliveryError::NoProviderConfigured),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ir_nodes::{IRDeliver, IRDeliverOp, IRDocField};
static REG_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
fn ir_deliver(provenance: &str) -> IRDeliver {
IRDeliver {
node_type: "deliver",
source_line: 1,
source_column: 1,
name: "push_lead".into(),
target: "crm".into(),
provenance: provenance.into(),
secret: "crm_api_key".into(),
effect_row: vec!["web".into()],
epistemic_mode: String::new(),
ops: vec![IRDeliverOp {
kind: "upsert_contact".into(),
fields: vec![
IRDocField { name: "key".into(), kind: "ref", value: "resolved_email".into(), items: vec![] },
IRDocField { name: "email".into(), kind: "ref", value: "resolved_email".into(), items: vec![] },
IRDocField { name: "firstname".into(), kind: "text", value: "Ada".into(), items: vec![] },
],
}],
}
}
#[test]
fn plan_resolves_refs_and_extracts_idempotency_key() {
let ir = ir_deliver("attached");
let req = plan_delivery(&ir, |name| match name {
"resolved_email" => Some("ada@example.com".into()),
_ => None,
})
.expect("plan");
assert_eq!(req.provenance, ProvenanceMode::Attached);
assert_eq!(req.ops.len(), 1);
let op = &req.ops[0];
assert_eq!(op.idempotency_key, "ada@example.com");
assert_eq!(op.fields[1].value, "ada@example.com");
assert_eq!(op.fields[2].value, "Ada");
}
#[test]
fn plan_unresolved_ref_is_caller_blame_never_a_silent_skip() {
let ir = ir_deliver("attached");
let err = plan_delivery(&ir, |_| None).expect_err("must fail");
assert!(matches!(err, DeliveryError::Caller(_)));
}
#[test]
fn cleared_mode_lowers_from_ir() {
let ir = ir_deliver("cleared");
let req = plan_delivery(&ir, |_| Some("x".into())).expect("plan");
assert_eq!(req.provenance, ProvenanceMode::Cleared);
}
#[test]
fn no_provider_is_a_typed_refusal_never_a_fabricated_receipt() {
let _g = REG_LOCK.lock().unwrap();
clear_provider();
let ir = ir_deliver("attached");
let req = plan_delivery(&ir, |_| Some("ada@example.com".into())).expect("plan");
assert_eq!(run_delivery(&req), Err(DeliveryError::NoProviderConfigured));
}
#[test]
fn registered_provider_is_dispatched_and_idempotent() {
let _g = REG_LOCK.lock().unwrap();
struct Echo;
impl DeliveryProvider for Echo {
fn name(&self) -> &str {
"echo"
}
fn deliver(&self, req: &DeliveryRequest) -> Result<Vec<DeliveryReceipt>, DeliveryError> {
Ok(req
.ops
.iter()
.map(|o| DeliveryReceipt {
kind: o.kind.clone(),
idempotency_key: o.idempotency_key.clone(),
record_id: Some(format!("rec_{}", o.idempotency_key)),
created: true,
})
.collect())
}
}
register_provider(Arc::new(Echo));
let ir = ir_deliver("attached");
let req = plan_delivery(&ir, |_| Some("ada@example.com".into())).expect("plan");
let receipts = run_delivery(&req).expect("deliver");
assert_eq!(receipts.len(), 1);
assert_eq!(receipts[0].record_id.as_deref(), Some("rec_ada@example.com"));
clear_provider();
}
}