use chrono::{DateTime, Datelike, Duration, Timelike, Utc};
use axon_frontend::cron::{cron_expr, CronSchedule};
use axon_frontend::ir_nodes::{IRDaemon, IRFlowNode, IRListenStep, IRRun, IRWindow};
use crate::window::{decide, WindowAction};
pub trait Clock: Send + Sync {
fn now(&self) -> DateTime<Utc>;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct SystemClock;
impl Clock for SystemClock {
fn now(&self) -> DateTime<Utc> {
Utc::now()
}
}
pub fn next_fire_after(schedule: &CronSchedule, after: DateTime<Utc>) -> Option<DateTime<Utc>> {
let mut t = (after + Duration::minutes(1))
.with_second(0)
.and_then(|t| t.with_nanosecond(0))?;
let limit = after + Duration::days(366);
while t <= limit {
if matches_minute(schedule, t) {
return Some(t);
}
t += Duration::minutes(1);
}
None
}
pub fn next_fires(schedule: &CronSchedule, from: DateTime<Utc>, n: usize) -> Vec<DateTime<Utc>> {
let mut out = Vec::with_capacity(n);
let mut cursor = from;
for _ in 0..n {
match next_fire_after(schedule, cursor) {
Some(t) => {
out.push(t);
cursor = t;
}
None => break,
}
}
out
}
fn matches_minute(s: &CronSchedule, t: DateTime<Utc>) -> bool {
let dow = t.weekday().num_days_from_sunday(); s.minute.contains(&t.minute())
&& s.hour.contains(&t.hour())
&& s.day_of_month.contains(&t.day())
&& s.month.contains(&t.month())
&& s.day_of_week.contains(&dow)
}
pub struct CronListener<'a> {
pub schedule: CronSchedule,
pub body: &'a [IRFlowNode],
pub channel: String,
}
pub fn cron_listeners(daemon: &IRDaemon) -> Vec<CronListener<'_>> {
daemon
.listeners
.iter()
.filter_map(|l: &IRListenStep| {
let expr = cron_expr(&l.channel)?;
let schedule = CronSchedule::parse(expr).ok()?;
Some(CronListener {
schedule,
body: &l.body,
channel: l.channel.clone(),
})
})
.collect()
}
pub struct EventListener<'a> {
pub channel: String,
pub alias: String,
pub body: &'a [IRFlowNode],
}
pub fn event_listeners(daemon: &IRDaemon) -> Vec<EventListener<'_>> {
daemon
.listeners
.iter()
.filter(|l: &&IRListenStep| cron_expr(&l.channel).is_none())
.map(|l| EventListener {
channel: l.channel.clone(),
alias: l.event_alias.clone(),
body: &l.body,
})
.collect()
}
pub fn bound_window<'a>(
ir: &'a axon_frontend::ir_nodes::IRProgram,
daemon: &IRDaemon,
) -> Option<&'a IRWindow> {
if daemon.window_ref.is_empty() {
return None;
}
ir.windows.iter().find(|w| w.name == daemon.window_ref)
}
pub fn run_invocations(body: &[IRFlowNode]) -> Vec<&IRRun> {
body.iter()
.filter_map(|n| match n {
IRFlowNode::Run(r) => Some(r),
_ => None,
})
.collect()
}
pub fn execute_listener_body(
ir: &axon_frontend::ir_nodes::IRProgram,
body: &[IRFlowNode],
backend: &str,
source_file: &str,
budget: Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>,
) -> Vec<(String, Result<crate::runner::ServerRunnerMetrics, String>)> {
let empty = std::collections::HashMap::new();
run_invocations(body)
.into_iter()
.map(|run| {
let result = crate::runner::execute_server_flow(
ir,
&run.flow_name,
backend,
&crate::tenant_context::current_tenant_id(),
source_file,
None,
None,
&empty,
&empty,
None,
None,
None,
budget.clone(),
None,
None,
None,
None, None, None, None, );
(run.flow_name.clone(), result)
})
.collect()
}
#[allow(clippy::type_complexity)]
pub fn deliver_typed_event(
ir: &axon_frontend::ir_nodes::IRProgram,
channel: &str,
payload: &serde_json::Value,
backend: &str,
source_file: &str,
budget: Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>,
) -> Vec<(String, String, Result<crate::runner::ServerRunnerMetrics, String>)> {
let empty = std::collections::HashMap::new();
let mut out = Vec::new();
for daemon in &ir.daemons {
for listener in event_listeners(daemon) {
if listener.channel != channel {
continue;
}
for run in run_invocations(listener.body) {
let result = crate::runner::execute_server_flow(
ir,
&run.flow_name,
backend,
&crate::tenant_context::current_tenant_id(),
source_file,
None,
Some(payload),
&empty,
&empty,
None,
None,
None,
budget.clone(),
None, None,
None,
None, None, None, None, );
out.push((daemon.name.clone(), run.flow_name.clone(), result));
}
}
}
out
}
pub const MAX_DELIVERY_ATTEMPTS: u32 = 5;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DeliveryOutcome {
Acked { attempts: u32 },
DeadLettered { attempts: u32, last_error: String },
Dropped { error: String },
}
pub fn deliver_with_retry(
qos: &str,
max_attempts: u32,
mut run: impl FnMut() -> std::result::Result<(), String>,
) -> DeliveryOutcome {
let cap = max_attempts.max(1);
if qos == "at_most_once" {
return match run() {
Ok(()) => DeliveryOutcome::Acked { attempts: 1 },
Err(e) => DeliveryOutcome::Dropped { error: e },
};
}
let mut last_error = String::new();
for attempt in 1..=cap {
match run() {
Ok(()) => return DeliveryOutcome::Acked { attempts: attempt },
Err(e) => last_error = e,
}
}
DeliveryOutcome::DeadLettered { attempts: cap, last_error }
}
pub fn fan_out_count(qos: &str, matching: usize) -> usize {
if qos == "broadcast" {
matching
} else {
matching.min(1)
}
}
#[allow(clippy::type_complexity)]
pub fn deliver_typed_event_reliable(
ir: &axon_frontend::ir_nodes::IRProgram,
channel: &str,
payload: &serde_json::Value,
backend: &str,
source_file: &str,
budget: Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>,
) -> Vec<(String, DeliveryOutcome)> {
let qos = ir
.channels
.iter()
.find(|c| c.name == channel)
.map(|c| c.qos.clone())
.unwrap_or_else(|| "at_least_once".to_string());
let mut matching: Vec<(String, EventListener)> = Vec::new();
for daemon in &ir.daemons {
for listener in event_listeners(daemon) {
if listener.channel == channel {
matching.push((daemon.name.clone(), listener));
}
}
}
let take = fan_out_count(&qos, matching.len());
let empty = std::collections::HashMap::new();
let mut out = Vec::new();
for (daemon_name, listener) in matching.into_iter().take(take) {
{
let outcome = deliver_with_retry(&qos, MAX_DELIVERY_ATTEMPTS, || {
for run in run_invocations(listener.body) {
let r = crate::runner::execute_server_flow(
ir,
&run.flow_name,
backend,
&crate::tenant_context::current_tenant_id(),
source_file,
None,
Some(payload),
&empty,
&empty,
None,
None,
None,
budget.clone(),
None, None,
None,
None, None, None, None, );
if let Err(e) = r {
return Err(format!("flow '{}': {e}", run.flow_name));
}
}
Ok(())
});
if let DeliveryOutcome::DeadLettered { attempts, last_error } = &outcome {
eprintln!(
"daemon '{daemon_name}' listener on '{channel}' DEAD-LETTERED after \
{attempts} attempts: {last_error}"
);
}
out.push((daemon_name, outcome));
}
}
out
}
pub async fn run_event_listeners(
ir: std::sync::Arc<axon_frontend::ir_nodes::IRProgram>,
bus: std::sync::Arc<crate::runtime::channels::TypedEventBus>,
backend: String,
cancel: crate::cancel_token::CancellationFlag,
replay_log: Option<std::sync::Arc<dyn crate::replay_token::ReplayLog>>,
) {
let mut channels: Vec<String> = ir
.daemons
.iter()
.flat_map(|d| event_listeners(d).into_iter().map(|l| l.channel))
.collect();
channels.sort();
channels.dedup();
let mut tasks = Vec::new();
for channel in channels {
let ir = ir.clone();
let bus = bus.clone();
let cancel = cancel.clone();
let backend = backend.clone();
let replay_log = replay_log.clone();
tasks.push(tokio::spawn(async move {
loop {
tokio::select! {
_ = cancel.cancelled() => return,
received = bus.receive(&channel) => {
let Ok(event) = received else { return }; let payload = match event.payload {
crate::runtime::channels::TypedPayload::Scalar(v) => v,
crate::runtime::channels::TypedPayload::Handle(_) => continue,
};
let _ = deliver_and_record(
&ir, &channel, &payload, &backend, "<daemon-event>", None,
replay_log.as_deref(),
)
.await;
}
}
}
}));
}
for t in tasks {
let _ = t.await;
}
}
pub fn mint_channel_event_token(
effect_name: &str,
flow_id: &str,
payload: &serde_json::Value,
outputs: serde_json::Value,
) -> crate::replay_token::ReplayToken {
crate::replay_token::ReplayTokenBuilder::new()
.effect_name(effect_name)
.inputs(serde_json::json!({ "flow_id": flow_id, "payload": payload }))
.outputs(outputs)
.model_version("axon.builtin.channel.v1")
.mint()
}
pub async fn deliver_and_record(
ir: &axon_frontend::ir_nodes::IRProgram,
channel: &str,
payload: &serde_json::Value,
backend: &str,
source_file: &str,
budget: Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>,
log: Option<&dyn crate::replay_token::ReplayLog>,
) -> Vec<(String, DeliveryOutcome)> {
let outcomes =
deliver_typed_event_reliable(ir, channel, payload, backend, source_file, budget);
if let Some(log) = log {
for (daemon_name, outcome) in &outcomes {
let summary = match outcome {
DeliveryOutcome::Acked { attempts } => {
serde_json::json!({ "outcome": "acked", "attempts": attempts })
}
DeliveryOutcome::DeadLettered { attempts, last_error } => serde_json::json!({
"outcome": "dead_lettered", "attempts": attempts, "error": last_error
}),
DeliveryOutcome::Dropped { error } => {
serde_json::json!({ "outcome": "dropped", "error": error })
}
};
let token = mint_channel_event_token(
&format!("deliver:{channel}"),
daemon_name,
payload,
summary,
);
if let Err(e) = log.append(token).await {
eprintln!("deliver token append failed for '{channel}': {e}");
}
}
}
outcomes
}
#[allow(clippy::type_complexity)]
pub fn drain_outbox(
ir: &axon_frontend::ir_nodes::IRProgram,
outbox: &dyn crate::event_outbox::EventOutbox,
channel: &str,
backend: &str,
source_file: &str,
budget: Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>,
) -> Vec<(u64, String, DeliveryOutcome)> {
let mut out = Vec::new();
for entry in outbox.unprocessed(channel) {
let outcomes = deliver_typed_event_reliable(
ir,
channel,
&entry.payload,
backend,
source_file,
budget.clone(),
);
outbox.mark_processed(channel, entry.offset);
for (daemon, outcome) in outcomes {
out.push((entry.offset, daemon, outcome));
}
}
out
}
pub async fn run_daemon(
ir: std::sync::Arc<axon_frontend::ir_nodes::IRProgram>,
daemon_name: String,
backend: String,
clock: std::sync::Arc<dyn Clock>,
cancel: crate::cancel_token::CancellationFlag,
) {
type SharedBudget = Option<std::sync::Arc<std::sync::Mutex<crate::runtime::budget_kernel::BudgetGate>>>;
let (listeners, window, budget): (
Vec<(CronSchedule, Vec<IRFlowNode>, String)>,
Option<IRWindow>,
SharedBudget,
) = {
let Some(daemon) = ir.daemons.iter().find(|d| d.name == daemon_name) else {
eprintln!("run_daemon: daemon '{daemon_name}' not in IR — nothing to drive");
return;
};
let window = bound_window(&ir, daemon).cloned();
let budget = daemon.budget.as_ref().map(|b| {
std::sync::Arc::new(std::sync::Mutex::new(
crate::runtime::budget_kernel::BudgetGate::from_ir(b, &daemon_name, clock.now()),
))
});
let listeners = cron_listeners(daemon)
.into_iter()
.map(|l| (l.schedule, l.body.to_vec(), l.channel))
.collect();
(listeners, window, budget)
};
let mut tasks = Vec::new();
for (schedule, body, channel) in listeners {
let ir = ir.clone();
let clock = clock.clone();
let cancel = cancel.clone();
let daemon_name = daemon_name.clone();
let backend = backend.clone();
let window = window.clone();
let budget = budget.clone();
tasks.push(tokio::spawn(async move {
loop {
let now = clock.now();
let Some(next) = next_fire_after(&schedule, now) else {
eprintln!(
"daemon '{daemon_name}' listener '{channel}': schedule never \
fires within the horizon — stopping this listener"
);
return;
};
let wait = (next - now).to_std().unwrap_or(std::time::Duration::ZERO);
tokio::select! {
_ = cancel.cancelled() => return,
_ = tokio::time::sleep(wait) => {
if let Some(w) = &window {
match decide(next, w) {
WindowAction::Fire => {} WindowAction::Skip => {
eprintln!(
"daemon '{daemon_name}' tick at {next} is OUTSIDE \
window '{}' (on_outside: skip) — dropped",
w.name
);
continue;
}
WindowAction::Warn => {
eprintln!(
"daemon '{daemon_name}' tick at {next} is OUTSIDE \
window '{}' (on_outside: warn) — firing anyway",
w.name
);
}
WindowAction::Defer { open_at } => {
eprintln!(
"daemon '{daemon_name}' tick at {next} is OUTSIDE \
window '{}' (on_outside: defer) — OSS degrades defer to \
skip (next opening {open_at:?}); the enterprise \
defer-ledger fires it once when the window opens",
w.name
);
continue;
}
}
}
let results =
execute_listener_body(&ir, &body, &backend, "<daemon>", budget.clone());
for (flow, res) in results {
if let Err(e) = res {
eprintln!(
"daemon '{daemon_name}' tick → flow '{flow}' failed: {e}"
);
}
}
}
}
}
}));
}
for t in tasks {
let _ = t.await;
}
}
#[cfg(test)]
mod tests {
use super::*;
use chrono::TimeZone;
fn at(y: i32, mo: u32, d: u32, h: u32, mi: u32) -> DateTime<Utc> {
Utc.with_ymd_and_hms(y, mo, d, h, mi, 0).unwrap()
}
#[test]
fn every_five_minutes_next_fire() {
let s = CronSchedule::parse("*/5 * * * *").unwrap();
assert_eq!(next_fire_after(&s, at(2026, 6, 26, 10, 2)), Some(at(2026, 6, 26, 10, 5)));
assert_eq!(next_fire_after(&s, at(2026, 6, 26, 10, 5)), Some(at(2026, 6, 26, 10, 10)));
assert_eq!(next_fire_after(&s, at(2026, 6, 26, 10, 57)), Some(at(2026, 6, 26, 11, 0)));
}
#[test]
fn daily_at_specific_time_rolls_to_next_day() {
let s = CronSchedule::parse("30 9 * * *").unwrap(); assert_eq!(next_fire_after(&s, at(2026, 6, 26, 9, 0)), Some(at(2026, 6, 26, 9, 30)));
assert_eq!(next_fire_after(&s, at(2026, 6, 26, 9, 30)), Some(at(2026, 6, 27, 9, 30)));
}
#[test]
fn weekday_business_hours_skips_weekend() {
let s = CronSchedule::parse("0 9 * * 1-5").unwrap();
let fired = next_fire_after(&s, at(2026, 6, 26, 10, 0)).unwrap();
assert_eq!(fired, at(2026, 6, 29, 9, 0));
assert_eq!(fired.weekday().num_days_from_sunday(), 1, "Monday");
}
#[test]
fn next_fires_sequence() {
let s = CronSchedule::parse("*/15 * * * *").unwrap();
let fires = next_fires(&s, at(2026, 6, 26, 8, 0), 3);
assert_eq!(
fires,
vec![at(2026, 6, 26, 8, 15), at(2026, 6, 26, 8, 30), at(2026, 6, 26, 8, 45)]
);
}
#[test]
fn impossible_schedule_yields_none() {
let s = CronSchedule::parse("0 0 30 2 *").unwrap();
assert_eq!(next_fire_after(&s, at(2026, 1, 1, 0, 0)), None);
}
fn ir_with_daemon(src: &str) -> axon_frontend::ir_nodes::IRProgram {
let tokens = axon_frontend::lexer::Lexer::new(src, "d.axon").tokenize().unwrap();
let program = axon_frontend::parser::Parser::new(tokens).parse().unwrap();
axon_frontend::ir_generator::IRGenerator::new().generate(&program)
}
#[test]
fn extracts_cron_listeners_and_invocations() {
let ir = ir_with_daemon(
"flow HibernateSession() -> Unit { step S { ask: \"x\" output: Unit } }\n\
daemon Cleaner {\n\
goal: \"clean\"\n\
listen \"cron:*/5 * * * *\" as tick { run HibernateSession() }\n\
listen \"user_events\" as e { run HibernateSession() }\n\
}",
);
let daemon = ir.daemons.iter().find(|d| d.name == "Cleaner").unwrap();
let crons = cron_listeners(daemon);
assert_eq!(crons.len(), 1);
assert_eq!(crons[0].channel, "cron:*/5 * * * *");
let invs = run_invocations(crons[0].body);
assert_eq!(invs.len(), 1);
assert_eq!(invs[0].flow_name, "HibernateSession");
}
#[test]
fn event_listeners_are_the_complement_of_cron() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once persistence: persistent_axonstore }\n\
flow Learn() -> Unit { probe p }\n\
daemon Mixed {\n\
requires: [flow.execute]\n\
listen \"cron:*/5 * * * *\" as tick { run Learn() }\n\
listen HibCh as ev { run Learn() }\n\
}",
);
let daemon = ir.daemons.iter().find(|d| d.name == "Mixed").unwrap();
assert_eq!(cron_listeners(daemon).len(), 1);
let evs = event_listeners(daemon);
assert_eq!(evs.len(), 1);
assert_eq!(evs[0].channel, "HibCh");
assert_eq!(evs[0].alias, "ev");
}
#[tokio::test]
async fn emit_reaches_a_daemon_listener_in_process() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String session_id: String }\n\
channel HibCh { message: Hib qos: at_least_once persistence: persistent_axonstore }\n\
flow Learn(tenant_id: String, session_id: String) -> String { return \"learned\" }\n\
daemon IntentLearner {\n\
requires: [flow.execute]\n\
listen HibCh as ev { run Learn() }\n\
}",
);
let bus = crate::runtime::channels::TypedEventBus::from_ir_program(&ir);
let payload = serde_json::json!({ "tenant_id": "acme", "session_id": "s1" });
bus.emit(
"HibCh",
crate::runtime::channels::TypedPayload::Scalar(payload.clone()),
)
.await
.expect("emit to a registered channel succeeds");
let event = bus.receive("HibCh").await.expect("the emitted event is queued");
let received = match event.payload {
crate::runtime::channels::TypedPayload::Scalar(v) => v,
other => panic!("expected scalar payload, got {other:?}"),
};
assert_eq!(received["tenant_id"], "acme");
let results = deliver_typed_event(&ir, "HibCh", &received, "stub", "<test>", None);
assert_eq!(results.len(), 1, "exactly one listener fired");
let (daemon, flow, res) = &results[0];
assert_eq!(daemon, "IntentLearner");
assert_eq!(flow, "Learn");
assert!(
res.is_ok(),
"the consumer flow ran to completion (err: {:?})",
res.as_ref().err()
);
}
#[test]
fn deliver_with_retry_acks_on_a_later_attempt() {
let mut n = 0;
let outcome = deliver_with_retry("at_least_once", 5, || {
n += 1;
if n < 3 { Err(format!("transient {n}")) } else { Ok(()) }
});
assert_eq!(outcome, DeliveryOutcome::Acked { attempts: 3 });
}
#[test]
fn deliver_with_retry_dead_letters_after_the_ceiling() {
let mut n = 0;
let outcome = deliver_with_retry("at_least_once", MAX_DELIVERY_ATTEMPTS, || {
n += 1;
Err(format!("perma-fail {n}"))
});
match outcome {
DeliveryOutcome::DeadLettered { attempts, last_error } => {
assert_eq!(attempts, MAX_DELIVERY_ATTEMPTS);
assert!(last_error.contains("perma-fail"));
}
other => panic!("expected DeadLettered, got {other:?}"),
}
assert_eq!(n, MAX_DELIVERY_ATTEMPTS, "tried exactly the ceiling, no more");
}
#[test]
fn deliver_with_retry_at_most_once_drops_without_retry() {
let mut n = 0;
let outcome = deliver_with_retry("at_most_once", 5, || {
n += 1;
Err("boom".to_string())
});
assert!(matches!(outcome, DeliveryOutcome::Dropped { .. }));
assert_eq!(n, 1, "at_most_once never retries");
}
#[tokio::test]
async fn reliable_delivery_acks_a_passing_listener() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once }\n\
flow Learn(tenant_id: String) -> String { return \"ok\" }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Learn() } }",
);
let payload = serde_json::json!({ "tenant_id": "acme" });
let out = deliver_typed_event_reliable(&ir, "HibCh", &payload, "stub", "<test>", None);
assert_eq!(out.len(), 1);
assert_eq!(out[0].0, "D");
assert!(matches!(out[0].1, DeliveryOutcome::Acked { .. }), "{:?}", out[0].1);
}
#[tokio::test]
async fn reliable_delivery_dead_letters_a_failing_listener() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Ghost() } }",
);
let out = deliver_typed_event_reliable(
&ir, "HibCh", &serde_json::json!({}), "stub", "<test>", None,
);
assert_eq!(out.len(), 1);
match &out[0].1 {
DeliveryOutcome::DeadLettered { attempts, .. } => {
assert_eq!(*attempts, MAX_DELIVERY_ATTEMPTS)
}
other => panic!("expected DeadLettered, got {other:?}"),
}
}
#[tokio::test]
async fn reliable_delivery_at_most_once_drops_a_failing_listener() {
let ir = ir_with_daemon(
"type T { x: String }\n\
channel C { message: T qos: at_most_once }\n\
daemon D { requires: [flow.execute] listen C as ev { run Ghost() } }",
);
let out = deliver_typed_event_reliable(&ir, "C", &serde_json::json!({}), "stub", "<test>", None);
assert_eq!(out.len(), 1);
assert!(matches!(out[0].1, DeliveryOutcome::Dropped { .. }), "{:?}", out[0].1);
}
#[tokio::test]
async fn outbox_event_redelivers_after_the_consumer_was_down() {
use crate::event_outbox::{EventOutbox, InMemoryEventOutbox};
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once persistence: persistent_axonstore }\n\
flow Learn(tenant_id: String) -> String { return \"ok\" }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Learn() } }",
);
let outbox = InMemoryEventOutbox::new();
outbox.append("HibCh", serde_json::json!({ "tenant_id": "acme" }));
assert_eq!(outbox.unprocessed("HibCh").len(), 1, "queued, awaiting delivery");
let outcomes = drain_outbox(&ir, &outbox, "HibCh", "stub", "<test>", None);
assert_eq!(outcomes.len(), 1);
assert!(matches!(outcomes[0].2, DeliveryOutcome::Acked { .. }), "{:?}", outcomes[0].2);
assert!(outbox.unprocessed("HibCh").is_empty(), "acked → not redelivered");
assert!(drain_outbox(&ir, &outbox, "HibCh", "stub", "<test>", None).is_empty());
}
#[tokio::test]
async fn outbox_dead_letters_a_failing_listener_then_stops_redelivering() {
use crate::event_outbox::{EventOutbox, InMemoryEventOutbox};
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once persistence: persistent_axonstore }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Ghost() } }",
);
let outbox = InMemoryEventOutbox::new();
outbox.append("HibCh", serde_json::json!({}));
let outcomes = drain_outbox(&ir, &outbox, "HibCh", "stub", "<test>", None);
assert!(matches!(outcomes[0].2, DeliveryOutcome::DeadLettered { .. }));
assert!(outbox.unprocessed("HibCh").is_empty(), "dead-lettered entries are acked");
}
#[test]
fn fan_out_count_broadcast_is_all_else_one() {
assert_eq!(fan_out_count("broadcast", 3), 3, "broadcast → every listener");
for single in ["at_least_once", "at_most_once", "exactly_once", "queue"] {
assert_eq!(fan_out_count(single, 3), 1, "{single} → single-consumer");
}
assert_eq!(fan_out_count("broadcast", 0), 0);
assert_eq!(fan_out_count("at_least_once", 0), 0);
}
#[tokio::test]
async fn broadcast_fans_out_to_every_listener() {
let ir = ir_with_daemon(
"type E { x: String }\n\
channel BroadcastCh { message: E qos: broadcast }\n\
flow A() -> String { return \"a\" }\n\
flow B() -> String { return \"b\" }\n\
daemon DA { requires: [flow.execute] listen BroadcastCh as ev { run A() } }\n\
daemon DB { requires: [flow.execute] listen BroadcastCh as ev { run B() } }",
);
let out = deliver_typed_event_reliable(
&ir, "BroadcastCh", &serde_json::json!({}), "stub", "<test>", None,
);
assert_eq!(out.len(), 2, "broadcast → BOTH daemons fire: {out:?}");
assert!(out.iter().all(|(_, o)| matches!(o, DeliveryOutcome::Acked { .. })));
}
#[tokio::test]
async fn single_consumer_qos_delivers_to_one_listener_only() {
let ir = ir_with_daemon(
"type E { x: String }\n\
channel QCh { message: E qos: at_least_once }\n\
flow A() -> String { return \"a\" }\n\
daemon DA { requires: [flow.execute] listen QCh as ev { run A() } }\n\
daemon DB { requires: [flow.execute] listen QCh as ev { run A() } }",
);
let out =
deliver_typed_event_reliable(&ir, "QCh", &serde_json::json!({}), "stub", "<test>", None);
assert_eq!(out.len(), 1, "single-consumer → exactly one fired: {out:?}");
}
#[test]
fn channel_event_token_captures_the_causal_receipt() {
let tok = mint_channel_event_token(
"deliver:HibCh",
"IntentLearner",
&serde_json::json!({ "tenant_id": "acme" }),
serde_json::json!({ "outcome": "acked", "attempts": 1 }),
);
assert_eq!(tok.effect_name, "deliver:HibCh");
assert_eq!(tok.model_version, "axon.builtin.channel.v1");
assert_eq!(tok.inputs["flow_id"], "IntentLearner");
assert_eq!(tok.inputs["payload"]["tenant_id"], "acme");
assert_eq!(tok.outputs["outcome"], "acked");
assert!(!tok.token_hash_hex.is_empty(), "the receipt is hashed");
}
#[tokio::test]
async fn deliver_and_record_logs_a_deliver_token_per_outcome() {
use crate::replay_token::{InMemoryReplayLog, ReplayLog};
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib qos: at_least_once }\n\
flow Learn(tenant_id: String) -> String { return \"ok\" }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Learn() } }",
);
let log = InMemoryReplayLog::new();
let payload = serde_json::json!({ "tenant_id": "acme" });
let outcomes = deliver_and_record(
&ir, "HibCh", &payload, "stub", "<test>", None, Some(&log),
)
.await;
assert_eq!(outcomes.len(), 1);
assert_eq!(log.len(), 1, "one deliver token recorded");
let tokens = log.tokens_for_flow("D").await.unwrap();
assert_eq!(tokens.len(), 1);
assert_eq!(tokens[0].effect_name, "deliver:HibCh");
}
#[tokio::test]
async fn deliver_and_record_without_a_log_is_pure_delivery() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib }\n\
flow Learn(tenant_id: String) -> String { return \"ok\" }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Learn() } }",
);
let out = deliver_and_record(
&ir, "HibCh", &serde_json::json!({ "tenant_id": "x" }), "stub", "<test>", None, None,
)
.await;
assert_eq!(out.len(), 1, "delivery still happens with no replay sink");
}
#[test]
fn deliver_to_a_channel_with_no_listener_is_an_empty_fan_out() {
let ir = ir_with_daemon(
"type Hib { tenant_id: String }\n\
channel HibCh { message: Hib }\n\
flow Learn() -> String { return \"ok\" }\n\
daemon D { requires: [flow.execute] listen HibCh as ev { run Learn() } }",
);
let out = deliver_typed_event(&ir, "OtherCh", &serde_json::json!({}), "stub", "<test>", None);
assert!(out.is_empty());
}
#[test]
fn bound_window_resolves_the_daemon_guard() {
let ir = ir_with_daemon(
"flow Send() -> Unit { step S { ask: \"x\" output: Unit } }\n\
window BusinessHours {\n\
timezone: \"America/Bogota\"\n\
allow: [ { days: Mon..Fri, hours: 9..18 } ]\n\
on_outside: skip\n\
}\n\
daemon Scheduler {\n\
window: BusinessHours\n\
requires: [flow.execute]\n\
listen \"cron:*/5 * * * *\" as tick { run Send() }\n\
}",
);
let daemon = ir.daemons.iter().find(|d| d.name == "Scheduler").unwrap();
assert_eq!(daemon.window_ref, "BusinessHours");
let w = bound_window(&ir, daemon).expect("window resolves");
assert_eq!(w.name, "BusinessHours");
assert_eq!(w.timezone, "America/Bogota");
assert_eq!(w.on_outside, "skip");
let unguarded = ir_with_daemon(
"flow Send() -> Unit { step S { ask: \"x\" output: Unit } }\n\
daemon Plain {\n\
listen \"cron:*/5 * * * *\" as tick { run Send() }\n\
}",
);
let d2 = unguarded.daemons.iter().find(|d| d.name == "Plain").unwrap();
assert!(bound_window(&unguarded, d2).is_none());
}
#[test]
fn budget_gate_builds_from_a_parsed_daemon_and_enforces() {
let ir = ir_with_daemon(
"tool TelnyxCall { provider: http timeout: 5s }\n\
flow SendBatch() -> Unit { step S { ask: \"x\" output: Unit } }\n\
daemon OutboundScheduler {\n\
requires: [flow.execute]\n\
budget {\n\
max: 1 per hour on Tool(TelnyxCall)\n\
on_exhausted: block\n\
}\n\
listen \"cron:*/5 * * * *\" as t { run SendBatch() }\n\
}",
);
let daemon = ir.daemons.iter().find(|d| d.name == "OutboundScheduler").unwrap();
let budget = daemon.budget.as_ref().expect("budget lowered onto the daemon");
let now: chrono::DateTime<chrono::Utc> = "2026-06-29T00:00:00Z".parse().unwrap();
let mut gate = crate::runtime::budget_kernel::BudgetGate::from_ir(budget, &daemon.name, now);
use crate::runtime::budget_kernel::GateDecision;
assert_eq!(gate.gate("TelnyxCall", now), GateDecision::Allow);
match gate.gate("TelnyxCall", now) {
GateDecision::Deny { on_exhausted, .. } => assert_eq!(on_exhausted, "block"),
other => panic!("expected Deny, got {other:?}"),
}
assert_eq!(gate.gate("SomeOtherTool", now), GateDecision::Allow);
}
}