use crate::flow_dispatcher::{DispatchCtx, DispatchError, NodeOutcome};
use crate::flow_execution_event::{now_ms, FlowExecutionEvent};
use crate::ir_nodes::{
IRConsensusBlock, IRDeliberateBlock, IRDiscover, IREmit, IRMutateStep,
IRPersistStep, IRPublish, IRPurgeStep, IRQuant, IRRetrieveStep, IRTransactBlock, IRYield,
};
use crate::store::audit_chain::StoreMutationKind;
use crate::store::capability;
use crate::store::epistemic;
use crate::store::row_stream;
use crate::store::filter::SqlValue;
use crate::store::postgres_backend::{PostgresStoreBackend, StoreError};
use crate::store::registry::StoreHandle;
pub fn emit_to_channel(channel_ref: &str, value: &str, ctx: &mut DispatchCtx) -> String {
let key = format!("__channel_{channel_ref}");
let existing = ctx.let_bindings.get(&key).cloned().unwrap_or_default();
let updated = if existing.is_empty() {
value.to_string()
} else {
format!("{existing}\n{value}")
};
ctx.let_bindings.insert(key, updated);
value.to_string()
}
pub fn publish_capability(channel_ref: &str, shield_ref: &str, ctx: &mut DispatchCtx) -> String {
let key = format!("__pub_{channel_ref}");
ctx.let_bindings.insert(key, shield_ref.to_string());
format!("published {channel_ref} with {shield_ref}")
}
pub fn discover_capability(capability_ref: &str, ctx: &DispatchCtx) -> String {
let key = format!("__pub_{capability_ref}");
ctx.let_bindings.get(&key).cloned().unwrap_or_default()
}
pub fn persist_to_store(store_name: &str, ctx: &mut DispatchCtx) -> usize {
let prefix = format!("__store_{store_name}_");
let user_bindings: Vec<(String, String)> = ctx
.let_bindings
.iter()
.filter(|(k, _)| !k.starts_with("__"))
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
let count = user_bindings.len();
for (k, v) in user_bindings {
ctx.let_bindings.insert(format!("{prefix}{k}"), v);
}
count
}
pub fn retrieve_from_store(
store_name: &str,
where_expr: &str,
ctx: &DispatchCtx,
) -> String {
let key = format!("__store_{store_name}_{where_expr}");
ctx.let_bindings.get(&key).cloned().unwrap_or_default()
}
pub fn mutate_store(store_name: &str, where_expr: &str, ctx: &mut DispatchCtx) -> u64 {
let key = format!("__store_{store_name}_{where_expr}");
if !ctx.let_bindings.contains_key(&key) {
return 0;
}
let new_value = ctx
.let_bindings
.get(where_expr)
.cloned()
.unwrap_or_else(|| where_expr.to_string());
ctx.let_bindings.insert(key, new_value);
1
}
pub fn purge_from_store(
store_name: &str,
where_expr: &str,
ctx: &mut DispatchCtx,
) -> u64 {
let key = format!("__store_{store_name}_{where_expr}");
if ctx.let_bindings.remove(&key).is_some() {
1
} else {
0
}
}
fn resolve_pg_backend(
ctx: &DispatchCtx,
store_name: &str,
) -> Result<Option<(PostgresStoreBackend, Option<f64>)>, StoreError> {
let Some(registry) = ctx.store_registry.as_ref() else {
return Ok(None);
};
match registry.resolve(store_name)? {
StoreHandle::InMemory => Ok(None),
StoreHandle::Postgres(backend) => {
let floor =
registry.spec(store_name).and_then(|s| s.confidence_floor);
Ok(Some((backend, floor)))
}
}
}
fn sql_row_from_bindings(ctx: &DispatchCtx) -> Vec<(String, SqlValue)> {
let mut row: Vec<(String, SqlValue)> = ctx
.let_bindings
.iter()
.filter(|(k, _)| !k.starts_with("__"))
.map(|(k, v)| (k.clone(), SqlValue::Text(v.clone())))
.collect();
row.sort_by(|a, b| a.0.cmp(&b.0));
row
}
fn store_row(fields: &[(String, String)], ctx: &DispatchCtx) -> Vec<(String, SqlValue)> {
if fields.is_empty() {
return sql_row_from_bindings(ctx);
}
fields
.iter()
.map(|(col, expr)| {
(
col.clone(),
SqlValue::Text(crate::exec_context::interpolate_vars(
expr,
&ctx.let_bindings,
)),
)
})
.collect()
}
fn unresolved_reference(value: &str) -> Option<String> {
let bytes = value.as_bytes();
let mut i = 0;
while i + 2 < bytes.len() {
if bytes[i] == b'$' && bytes[i + 1] == b'{' {
let start = i + 2;
if bytes[start].is_ascii_alphabetic() || bytes[start] == b'_' {
if let Some(close) = value[start..].find('}') {
let inner = &value[start..start + close];
if inner
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '.')
{
return Some(inner.to_string());
}
}
}
}
i += 1;
}
None
}
fn reject_unresolved_row(
store_name: &str,
op: &str,
row: &[(String, SqlValue)],
) -> Result<(), DispatchError> {
for (col, val) in row {
if let SqlValue::Text(s) = val {
if let Some(reference) = unresolved_reference(s) {
return Err(DispatchError::BackendError {
name: "axonstore".to_string(),
message: format!(
"{op} into `{store_name}`: column `{col}` carries an \
UNRESOLVED reference `${{{reference}}}` after \
interpolation — it did not resolve to a binding (check \
the loop variable / step output). Refusing to write the \
literal `${{{reference}}}` to the database."
),
});
}
}
}
Ok(())
}
fn sql_dispatch_error(e: StoreError) -> DispatchError {
DispatchError::BackendError {
name: "axonstore".to_string(),
message: e.to_string(),
}
}
fn enforce_store_capability(
ctx: &DispatchCtx,
store_name: &str,
) -> Result<(), DispatchError> {
let Some(held) = ctx.held_capabilities.as_ref() else {
return Ok(());
};
let required = ctx
.store_registry
.as_ref()
.and_then(|r| r.spec(store_name))
.map(|s| s.capability.as_str())
.unwrap_or("");
capability::check_store_capability(store_name, required, held).map_err(
|denied| DispatchError::BackendError {
name: "axonstore.capability".to_string(),
message: denied.to_string(),
},
)
}
fn record_store_mutation(
ctx: &DispatchCtx,
kind: StoreMutationKind,
store: &str,
summary: &str,
) {
let mut chain = ctx
.audit_chain
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
chain.record(kind, store, summary);
}
pub async fn run_emit(
node: &IREmit,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.channel_ref.is_empty() {
"Emit".to_string()
} else {
node.channel_ref.clone()
};
emit_step_start(ctx, &step_name, step_index, "emit")?;
let resolved_value = ctx
.let_bindings
.get(&node.value_ref)
.cloned()
.unwrap_or_else(|| node.value_ref.clone());
let bus_handle = ctx
.event_bus
.as_ref()
.and_then(|b| b.get_handle(&node.channel_ref).ok());
let is_durable = bus_handle
.as_ref()
.map_or(false, |h| h.persistence == "persistent_axonstore");
let emitted = if let (true, Some(outbox)) = (is_durable, ctx.event_outbox.clone()) {
let payload = serde_json::from_str::<serde_json::Value>(&resolved_value)
.unwrap_or_else(|_| serde_json::Value::String(resolved_value.clone()));
outbox.append(&node.channel_ref, payload);
resolved_value.clone()
} else if let Some(bus) = ctx.event_bus.clone() {
if bus.get_handle(&node.channel_ref).is_ok() {
let payload = serde_json::from_str::<serde_json::Value>(&resolved_value)
.unwrap_or_else(|_| serde_json::Value::String(resolved_value.clone()));
bus.emit(
&node.channel_ref,
crate::runtime::channels::TypedPayload::Scalar(payload),
)
.await
.map_err(|e| DispatchError::BackendError {
name: node.channel_ref.clone(),
message: format!("emit to channel '{}' failed: {e}", node.channel_ref),
})?;
resolved_value.clone()
} else {
emit_to_channel(&node.channel_ref, &resolved_value, ctx)
}
} else {
emit_to_channel(&node.channel_ref, &resolved_value, ctx)
};
if let Some(log) = ctx.replay_log.clone() {
let payload = serde_json::from_str::<serde_json::Value>(&emitted)
.unwrap_or_else(|_| serde_json::Value::String(emitted.clone()));
let token = crate::daemon::mint_channel_event_token(
&format!("emit:{}", node.channel_ref),
&ctx.flow_name,
&payload,
serde_json::Value::Null,
);
if let Err(e) = log.append(token).await {
eprintln!("§74.e emit token append failed for '{}': {e}", node.channel_ref);
}
}
emit_step_complete(ctx, &step_name, step_index, &emitted, 0)?;
Ok(NodeOutcome::Completed {
output: emitted,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_publish(
node: &IRPublish,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.channel_ref.is_empty() {
"Publish".to_string()
} else {
node.channel_ref.clone()
};
emit_step_start(ctx, &step_name, step_index, "publish")?;
let output = publish_capability(&node.channel_ref, &node.shield_ref, ctx);
emit_step_complete(ctx, &step_name, step_index, &output, 0)?;
Ok(NodeOutcome::Completed {
output,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_discover(
node: &IRDiscover,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.alias.is_empty() {
"Discover".to_string()
} else {
node.alias.clone()
};
emit_step_start(ctx, &step_name, step_index, "discover")?;
let discovered = discover_capability(&node.capability_ref, ctx);
if !node.alias.is_empty() {
ctx.let_bindings.insert(node.alias.clone(), discovered.clone());
}
emit_step_complete(ctx, &step_name, step_index, &discovered, 0)?;
Ok(NodeOutcome::Completed {
output: discovered,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_persist(
node: &IRPersistStep,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.store_name.is_empty() {
"Persist".to_string()
} else {
node.store_name.clone()
};
enforce_store_capability(ctx, &node.store_name)?;
emit_step_start(ctx, &step_name, step_index, "persist")?;
let output = match resolve_pg_backend(ctx, &node.store_name) {
Ok(Some((backend, floor))) => {
let row = store_row(&node.fields, ctx);
reject_unresolved_row(&node.store_name, "persist", &row)?;
epistemic::enforce_persist_floor(&row, floor, &node.store_name)
.map_err(|e| sql_dispatch_error(StoreError::from(e)))?;
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> = {
ctx.pinned_conns.lock().unwrap().remove(&node.store_name)
};
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
pin = Some(p);
}
}
let n = {
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
backend
.insert(&mut store_conn, &node.store_name, &row)
.await
.map_err(sql_dispatch_error)?
};
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(node.store_name.clone(), p);
}
ctx.record_store_rows(crate::flow_dispatcher::StoreRowKind::Persisted, n as u64);
format!("persisted {n} row(s) to `{}`", node.store_name)
}
Ok(None) => {
let count = persist_to_store(&node.store_name, ctx);
ctx.record_store_rows(
crate::flow_dispatcher::StoreRowKind::Persisted,
count as u64,
);
format!("persisted {count} entries to `{}`", node.store_name)
}
Err(e) => return Err(sql_dispatch_error(e)),
};
record_store_mutation(ctx, StoreMutationKind::Persist, &node.store_name, &output);
emit_step_complete(ctx, &step_name, step_index, &output, 0)?;
Ok(NodeOutcome::Completed {
output,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_retrieve(
node: &IRRetrieveStep,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.alias.is_empty() {
"Retrieve".to_string()
} else {
node.alias.clone()
};
enforce_store_capability(ctx, &node.store_name)?;
emit_step_start(ctx, &step_name, step_index, "retrieve")?;
let value = match resolve_pg_backend(ctx, &node.store_name) {
Ok(Some((backend, floor))) => {
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> = {
ctx.pinned_conns.lock().unwrap().remove(&node.store_name)
};
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
crate::store::pin_observability::emit_pin_acquire(
&node.store_name,
&ctx.flow_name,
"",
"lazy",
if ctx.branch_path.is_empty() {
None
} else {
Some(ctx.branch_path.len())
},
);
pin = Some(p);
}
}
let stream_outcome_result = {
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
row_stream::stream_retrieve(
&backend,
&mut store_conn,
&node.store_name,
&node.where_expr,
&node.order_by,
&node.limit_expr,
row_stream::DEFAULT_RETRIEVE_POLICY,
row_stream::DEFAULT_MAX_ROWS,
&ctx.cancel,
&ctx.let_bindings,
)
.await
};
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(node.store_name.clone(), p);
}
let stream_outcome = stream_outcome_result
.map_err(sql_dispatch_error)?;
ctx.record_store_rows(
crate::flow_dispatcher::StoreRowKind::Retrieved,
stream_outcome.rows.len() as u64,
);
let metadata = row_stream::stream_metadata(
row_stream::DEFAULT_RETRIEVE_POLICY,
&stream_outcome,
);
let floored = epistemic::enforce_retrieve_floor(
epistemic::mark_retrieved(stream_outcome.rows),
floor,
);
let mut envelope = epistemic::retrieve_envelope(&floored, floor);
envelope["stream"] = metadata;
serde_json::to_string(&envelope).unwrap_or_else(|_| "{}".to_string())
}
Ok(None) => retrieve_from_store(&node.store_name, &node.where_expr, ctx),
Err(e) => return Err(sql_dispatch_error(e)),
};
if !node.alias.is_empty() {
ctx.let_bindings.insert(node.alias.clone(), value.clone());
}
emit_step_complete(ctx, &step_name, step_index, &value, 0)?;
Ok(NodeOutcome::Completed {
output: value,
tokens_emitted: 0,
step_index,
})
}
pub async fn read_all_store_rows(
ctx: &mut DispatchCtx,
store_name: &str,
where_expr: &str,
) -> Result<Option<Vec<crate::store::postgres_backend::StoreRow>>, DispatchError> {
match resolve_pg_backend(ctx, store_name) {
Ok(Some((backend, _floor))) => {
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> =
{ ctx.pinned_conns.lock().unwrap().remove(store_name) };
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
pin = Some(p);
}
}
let outcome = {
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
row_stream::stream_retrieve(
&backend,
&mut store_conn,
store_name,
where_expr,
"",
"",
row_stream::DEFAULT_RETRIEVE_POLICY,
row_stream::DEFAULT_MAX_ROWS,
&ctx.cancel,
&ctx.let_bindings,
)
.await
};
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(store_name.to_string(), p);
}
let outcome = outcome.map_err(sql_dispatch_error)?;
Ok(Some(outcome.rows))
}
Ok(None) => Ok(None),
Err(e) => Err(sql_dispatch_error(e)),
}
}
pub fn extract_corpus_rows(
doc_rows: &[crate::store::postgres_backend::StoreRow],
edge_rows: &[crate::store::postgres_backend::StoreRow],
src: &crate::ir_nodes::IRCorpusStoreSource,
) -> (Vec<(String, String)>, Vec<(String, String, String, f64)>) {
let col = |row: &crate::store::postgres_backend::StoreRow, name: &str| {
row.columns
.iter()
.find(|(c, _)| c == name)
.map(|(_, v)| v.clone())
};
let as_str = |v: &serde_json::Value| -> String {
match v {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
}
};
let docs = doc_rows
.iter()
.filter_map(|r| {
let id = col(r, &src.doc_id)?;
let title = col(r, &src.doc_title)?;
Some((as_str(&id), as_str(&title)))
})
.collect();
let edges = edge_rows
.iter()
.filter_map(|r| {
let from = as_str(&col(r, &src.edge_from)?);
let to = as_str(&col(r, &src.edge_to)?);
let etype = as_str(&col(r, &src.edge_type)?);
let weight = col(r, &src.edge_weight)?.as_f64().unwrap_or(0.0);
Some((from, to, etype, weight))
})
.collect();
(docs, edges)
}
pub fn plan_edge_reinforcements(
corpus: &crate::mdn::Corpus,
selected: &[crate::mdn::DocId],
docs: &[(String, String)],
score: f64,
mean_score: f64,
eta: f64,
) -> Vec<(String, String, String, f64)> {
let delta = eta * (score - mean_score);
if delta == 0.0 {
return Vec::new();
}
let mut out = Vec::new();
for pair in selected.windows(2) {
let (a, b) = (pair[0], pair[1]);
for e in corpus.edges() {
if e.from == a && e.to == b {
if let (Some(fa), Some(tb)) = (docs.get(a as usize), docs.get(b as usize)) {
out.push((fa.0.clone(), tb.0.clone(), e.etype.slug().to_string(), delta));
}
}
}
}
out
}
pub async fn persist_reinforcements(
ctx: &mut DispatchCtx,
edge_store: &str,
weight_col: &str,
from_col: &str,
to_col: &str,
etype_col: &str,
plan: &[(String, String, String, f64)],
epsilon: f64,
) -> Result<(), DispatchError> {
if plan.is_empty() {
return Ok(());
}
let Ok(Some((backend, _floor))) = resolve_pg_backend(ctx, edge_store) else {
return Ok(());
};
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> =
{ ctx.pinned_conns.lock().unwrap().remove(edge_store) };
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
pin = Some(p);
}
}
{
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
for (from_id, to_id, etype, delta) in plan {
let _ = backend
.reinforce(
&mut store_conn,
edge_store,
weight_col,
from_col,
to_col,
etype_col,
&crate::store::filter::SqlValue::Text(from_id.clone()),
&crate::store::filter::SqlValue::Text(to_id.clone()),
&crate::store::filter::SqlValue::Text(etype.clone()),
*delta,
epsilon,
)
.await;
}
}
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(edge_store.to_string(), p);
}
Ok(())
}
pub async fn run_mutate(
node: &IRMutateStep,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.store_name.is_empty() {
"Mutate".to_string()
} else {
node.store_name.clone()
};
enforce_store_capability(ctx, &node.store_name)?;
emit_step_start(ctx, &step_name, step_index, "mutate")?;
let output = match resolve_pg_backend(ctx, &node.store_name) {
Ok(Some((backend, _floor))) => {
let row = store_row(&node.fields, ctx);
reject_unresolved_row(&node.store_name, "mutate", &row)?;
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> = {
ctx.pinned_conns.lock().unwrap().remove(&node.store_name)
};
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
crate::store::pin_observability::emit_pin_acquire(
&node.store_name,
&ctx.flow_name,
"",
"lazy",
if ctx.branch_path.is_empty() {
None
} else {
Some(ctx.branch_path.len())
},
);
pin = Some(p);
}
}
let n = {
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
backend
.mutate(&mut store_conn, &node.store_name, &node.where_expr, &row, &ctx.let_bindings)
.await
.map_err(sql_dispatch_error)?
};
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(node.store_name.clone(), p);
}
ctx.record_store_rows(crate::flow_dispatcher::StoreRowKind::Mutated, n as u64);
format!("mutated {n} row(s) in `{}`", node.store_name)
}
Ok(None) => {
let count = mutate_store(&node.store_name, &node.where_expr, ctx);
ctx.record_store_rows(
crate::flow_dispatcher::StoreRowKind::Mutated,
count as u64,
);
format!("mutated {count} entries in `{}`", node.store_name)
}
Err(e) => return Err(sql_dispatch_error(e)),
};
record_store_mutation(ctx, StoreMutationKind::Mutate, &node.store_name, &output);
emit_step_complete(ctx, &step_name, step_index, &output, 0)?;
Ok(NodeOutcome::Completed {
output,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_purge(
node: &IRPurgeStep,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
let step_name = if node.store_name.is_empty() {
"Purge".to_string()
} else {
node.store_name.clone()
};
enforce_store_capability(ctx, &node.store_name)?;
emit_step_start(ctx, &step_name, step_index, "purge")?;
let output = match resolve_pg_backend(ctx, &node.store_name) {
Ok(Some((backend, _floor))) => {
let mut pin: Option<sqlx::pool::PoolConnection<sqlx::Postgres>> = {
ctx.pinned_conns.lock().unwrap().remove(&node.store_name)
};
if pin.is_none() {
if let Ok(p) = backend.acquire_pin().await {
crate::store::pin_observability::emit_pin_acquire(
&node.store_name,
&ctx.flow_name,
"",
"lazy",
if ctx.branch_path.is_empty() {
None
} else {
Some(ctx.branch_path.len())
},
);
pin = Some(p);
}
}
let n = {
let mut store_conn = match &mut pin {
Some(p) => crate::store::store_conn::StoreConn::Pinned(p),
None => crate::store::store_conn::StoreConn::Pool(backend.pool()),
};
backend
.purge(&mut store_conn, &node.store_name, &node.where_expr, &ctx.let_bindings)
.await
.map_err(sql_dispatch_error)?
};
if let Some(p) = pin {
ctx.pinned_conns
.lock()
.unwrap()
.insert(node.store_name.clone(), p);
}
ctx.record_store_rows(crate::flow_dispatcher::StoreRowKind::Purged, n as u64);
format!("purged {n} row(s) from `{}`", node.store_name)
}
Ok(None) => {
let count = purge_from_store(&node.store_name, &node.where_expr, ctx);
ctx.record_store_rows(
crate::flow_dispatcher::StoreRowKind::Purged,
count as u64,
);
format!("purged {count} entries from `{}`", node.store_name)
}
Err(e) => return Err(sql_dispatch_error(e)),
};
record_store_mutation(ctx, StoreMutationKind::Purge, &node.store_name, &output);
emit_step_complete(ctx, &step_name, step_index, &output, 0)?;
Ok(NodeOutcome::Completed {
output,
tokens_emitted: 0,
step_index,
})
}
pub async fn run_transact(
_node: &IRTransactBlock,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
emit_step_start(ctx, "Transact", step_index, "transact")?;
ctx.let_bindings
.insert("__txn_active".to_string(), "true".to_string());
emit_step_complete(ctx, "Transact", step_index, "", 0)?;
Ok(NodeOutcome::Completed {
output: String::new(),
tokens_emitted: 0,
step_index,
})
}
pub async fn run_quant(
node: &IRQuant,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
emit_step_start(ctx, "Quant", step_index, "quant")?;
ctx.let_bindings
.insert("__quant_backend".to_string(), node.effect.clone());
emit_step_complete(ctx, "Quant", step_index, "", 0)?;
Ok(NodeOutcome::Completed {
output: String::new(),
tokens_emitted: 0,
step_index,
})
}
pub async fn run_yield(
node: &IRYield,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
emit_step_start(ctx, "Yield", step_index, "yield")?;
emit_step_complete(ctx, "Yield", step_index, "", 0)?;
ctx.let_bindings
.insert("__quant_yield".to_string(), node.value_expr.clone());
Ok(NodeOutcome::Completed {
output: String::new(),
tokens_emitted: 0,
step_index,
})
}
pub async fn run_deliberate(
_node: &IRDeliberateBlock,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
emit_step_start(ctx, "Deliberate", step_index, "deliberate")?;
emit_step_complete(ctx, "Deliberate", step_index, "", 0)?;
Ok(NodeOutcome::Completed {
output: String::new(),
tokens_emitted: 0,
step_index,
})
}
pub async fn run_consensus(
_node: &IRConsensusBlock,
ctx: &mut DispatchCtx,
) -> Result<NodeOutcome, DispatchError> {
if ctx.cancel.is_cancelled() {
return Err(DispatchError::UpstreamCancelled);
}
let step_index = ctx.step_counter;
ctx.step_counter += 1;
emit_step_start(ctx, "Consensus", step_index, "consensus")?;
emit_step_complete(ctx, "Consensus", step_index, "", 0)?;
Ok(NodeOutcome::Completed {
output: String::new(),
tokens_emitted: 0,
step_index,
})
}
fn emit_step_start(
ctx: &mut DispatchCtx,
step_name: &str,
step_index: usize,
step_type: &str,
) -> Result<(), DispatchError> {
ctx.tx
.send(FlowExecutionEvent::StepStart {
step_name: step_name.to_string(),
step_index,
step_type: step_type.to_string(),
branch_path: ctx.branch_path_string(),
timestamp_ms: now_ms(),
})
.map_err(|_| DispatchError::ChannelClosed)
}
fn emit_step_complete(
ctx: &mut DispatchCtx,
step_name: &str,
step_index: usize,
full_output: &str,
tokens_output: u64,
) -> Result<(), DispatchError> {
ctx.tx
.send(FlowExecutionEvent::StepComplete {
step_name: step_name.to_string(),
step_index,
success: true,
full_output: full_output.to_string(),
tokens_input: 0,
tokens_output,
branch_path: ctx.branch_path_string(),
timestamp_ms: now_ms(),
})
.map_err(|_| DispatchError::ChannelClosed)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cancel_token::CancellationFlag;
use crate::ir_nodes::*;
use tokio::sync::mpsc;
#[test]
fn unresolved_reference_flags_a_surviving_dollar_brace_identifier() {
assert_eq!(unresolved_reference("${e.to_id}").as_deref(), Some("e.to_id"));
assert_eq!(unresolved_reference("${missing}").as_deref(), Some("missing"));
assert_eq!(unresolved_reference("11111111-1111-1111-1111-111111111111"), None);
assert_eq!(unresolved_reference("plain text"), None);
assert_eq!(unresolved_reference("cost ${100}"), None); assert_eq!(unresolved_reference(""), None);
}
#[test]
fn reject_unresolved_row_errors_instead_of_writing_the_literal() {
let ok = vec![("id".to_string(), SqlValue::Text("abc".to_string()))];
assert!(reject_unresolved_row("s", "persist", &ok).is_ok());
let bad = vec![("to_id".to_string(), SqlValue::Text("${e.to_id}".to_string()))];
let err = reject_unresolved_row("ltm_edges", "persist", &bad).unwrap_err();
match err {
DispatchError::BackendError { message, .. } => {
assert!(message.contains("UNRESOLVED"), "names the failure: {message}");
assert!(message.contains("e.to_id"), "quotes the reference: {message}");
assert!(message.contains("ltm_edges"), "names the store: {message}");
}
other => panic!("expected BackendError, got {other:?}"),
}
}
fn mk_store_row(pairs: &[(&str, serde_json::Value)]) -> crate::store::postgres_backend::StoreRow {
crate::store::postgres_backend::StoreRow {
columns: pairs.iter().map(|(k, v)| (k.to_string(), v.clone())).collect(),
}
}
#[test]
fn extract_corpus_rows_projects_the_mapped_columns() {
use serde_json::json;
let src = IRCorpusStoreSource {
doc_store: "LtmSummaries".into(),
doc_id: "id".into(),
doc_title: "summary".into(),
edge_store: "LtmEdges".into(),
edge_from: "from_id".into(),
edge_to: "to_id".into(),
edge_type: "etype".into(),
edge_weight: "weight".into(),
};
let doc_rows = vec![
mk_store_row(&[("id", json!("uuid-a")), ("summary", json!("A")), ("noise", json!(1))]),
mk_store_row(&[("id", json!("uuid-b")), ("summary", json!("B"))]),
];
let edge_rows = vec![mk_store_row(&[
("from_id", json!("uuid-b")),
("to_id", json!("uuid-a")),
("etype", json!("cite")),
("weight", json!(0.9)),
])];
let (docs, edges) = extract_corpus_rows(&doc_rows, &edge_rows, &src);
assert_eq!(docs, vec![("uuid-a".into(), "A".into()), ("uuid-b".into(), "B".into())]);
assert_eq!(edges, vec![("uuid-b".into(), "uuid-a".into(), "cite".into(), 0.9)]);
}
#[test]
fn plan_edge_reinforcements_zero_for_one_outcome_nonzero_with_variance() {
let docs = vec![("id-a".to_string(), "A".to_string()), ("id-b".to_string(), "B".to_string())];
let edges = vec![("id-a".to_string(), "id-b".to_string(), "cite".to_string(), 0.5)];
let corpus = crate::mdn::Corpus::from_rows(&docs, &edges).unwrap();
let selected = vec![0u32, 1u32];
let p0 = plan_edge_reinforcements(&corpus, &selected, &docs, 0.8, 0.8, 0.1);
assert!(p0.is_empty(), "a single outcome reinforces nothing (relative signal)");
let p1 = plan_edge_reinforcements(&corpus, &selected, &docs, 0.9, 0.5, 0.1);
assert_eq!(p1.len(), 1);
assert_eq!(p1[0].0, "id-a");
assert_eq!(p1[0].1, "id-b");
assert_eq!(p1[0].2, "cite");
assert!((p1[0].3 - 0.04).abs() < 1e-9, "Δ = η(s−s̄): got {}", p1[0].3);
}
#[test]
fn extract_corpus_rows_drops_rows_missing_a_mapped_column() {
use serde_json::json;
let src = IRCorpusStoreSource {
doc_store: "S".into(),
doc_id: "id".into(),
doc_title: "summary".into(),
edge_store: "E".into(),
edge_from: "f".into(),
edge_to: "t".into(),
edge_type: "ty".into(),
edge_weight: "w".into(),
};
let doc_rows = vec![
mk_store_row(&[("id", json!("a")), ("summary", json!("A"))]),
mk_store_row(&[("id", json!("b"))]),
];
let (docs, _edges) = extract_corpus_rows(&doc_rows, &[], &src);
assert_eq!(docs.len(), 1, "a row missing a mapped column is dropped");
assert_eq!(docs[0].0, "a");
}
fn fresh_ctx() -> (
DispatchCtx,
mpsc::UnboundedReceiver<FlowExecutionEvent>,
) {
let (tx, rx) = mpsc::unbounded_channel();
let ctx = DispatchCtx::new(
"TestFlow",
"stub",
"",
CancellationFlag::new(),
tx,
);
(ctx, rx)
}
#[test]
fn emit_to_channel_appends_to_buffer() {
let (mut ctx, _rx) = fresh_ctx();
emit_to_channel("c1", "v1", &mut ctx);
emit_to_channel("c1", "v2", &mut ctx);
let buffer = ctx.let_bindings.get("__channel_c1").unwrap();
assert_eq!(buffer, "v1\nv2");
}
#[test]
fn publish_then_discover_round_trip() {
let (mut ctx, _rx) = fresh_ctx();
publish_capability("user_inbox", "shield_pii", &mut ctx);
assert_eq!(discover_capability("user_inbox", &ctx), "shield_pii");
}
#[test]
fn discover_missing_returns_empty() {
let (ctx, _rx) = fresh_ctx();
assert_eq!(discover_capability("never_set", &ctx), "");
}
#[test]
fn persist_snapshots_user_bindings() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("name".into(), "alice".into());
ctx.let_bindings.insert("age".into(), "30".into());
ctx.let_bindings
.insert("__internal".into(), "should_not_be_snapshotted".into());
let count = persist_to_store("users", &mut ctx);
assert_eq!(count, 2);
assert_eq!(ctx.let_bindings.get("__store_users_name").unwrap(), "alice");
assert_eq!(ctx.let_bindings.get("__store_users_age").unwrap(), "30");
assert!(!ctx.let_bindings.contains_key("__store_users___internal"));
}
#[test]
fn retrieve_from_store_returns_persisted_value() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("city".into(), "Bogota".into());
persist_to_store("locations", &mut ctx);
assert_eq!(retrieve_from_store("locations", "city", &ctx), "Bogota");
}
#[test]
fn mutate_store_updates_existing_entry() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("counter".into(), "1".into());
persist_to_store("metrics", &mut ctx);
ctx.let_bindings.insert("counter".into(), "2".into());
let count = mutate_store("metrics", "counter", &mut ctx);
assert_eq!(count, 1);
assert_eq!(
ctx.let_bindings.get("__store_metrics_counter").unwrap(),
"2"
);
}
#[test]
fn mutate_missing_entry_returns_zero() {
let (mut ctx, _rx) = fresh_ctx();
assert_eq!(mutate_store("empty", "k", &mut ctx), 0);
}
#[test]
fn purge_removes_entry() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("key".into(), "value".into());
persist_to_store("s", &mut ctx);
let count = purge_from_store("s", "key", &mut ctx);
assert_eq!(count, 1);
assert!(!ctx.let_bindings.contains_key("__store_s_key"));
}
#[test]
fn purge_missing_returns_zero() {
let (mut ctx, _rx) = fresh_ctx();
assert_eq!(purge_from_store("s", "absent", &mut ctx), 0);
}
#[tokio::test]
async fn run_emit_appends_to_channel_buffer() {
let (mut ctx, mut rx) = fresh_ctx();
ctx.let_bindings.insert("payload".into(), "hello".into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "out_channel".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
let outcome = run_emit(&node, &mut ctx).await.unwrap();
match outcome {
NodeOutcome::Completed { output, .. } => assert_eq!(output, "hello"),
other => panic!("expected Completed, got {other:?}"),
}
assert_eq!(
ctx.let_bindings.get("__channel_out_channel").unwrap(),
"hello"
);
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "emit");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn run_emit_routes_to_the_event_bus_when_attached() {
use crate::runtime::channels::{TypedChannelHandle, TypedEventBus, TypedPayload};
let bus = std::sync::Arc::new(TypedEventBus::new());
bus.register(TypedChannelHandle::new("HibCh", "Hib"));
let (mut ctx, _rx) = fresh_ctx();
ctx = ctx.with_event_bus(bus.clone());
ctx.let_bindings
.insert("payload".into(), r#"{"tenant_id":"acme"}"#.into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "HibCh".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
run_emit(&node, &mut ctx).await.unwrap();
assert!(
!ctx.let_bindings.contains_key("__channel_HibCh"),
"a bus-routed emit must not also append to the legacy buffer"
);
let event = bus.receive("HibCh").await.expect("the event is on the bus");
match event.payload {
TypedPayload::Scalar(v) => assert_eq!(v["tenant_id"], "acme"),
other => panic!("expected the scalar payload, got {other:?}"),
}
}
#[tokio::test]
async fn run_emit_appends_to_the_outbox_for_a_persistent_channel() {
use crate::event_outbox::{EventOutbox, InMemoryEventOutbox};
use crate::runtime::channels::{TypedChannelHandle, TypedEventBus};
let bus = std::sync::Arc::new(TypedEventBus::new());
let mut h = TypedChannelHandle::new("HibCh", "Hib");
h.persistence = "persistent_axonstore".into();
bus.register(h);
let outbox = std::sync::Arc::new(InMemoryEventOutbox::new());
let (mut ctx, _rx) = fresh_ctx();
ctx = ctx.with_event_bus(bus).with_event_outbox(outbox.clone());
ctx.let_bindings
.insert("payload".into(), r#"{"tenant_id":"acme"}"#.into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "HibCh".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
run_emit(&node, &mut ctx).await.unwrap();
assert_eq!(outbox.pending_total(), 1, "the event is durably queued");
assert!(!ctx.let_bindings.contains_key("__channel_HibCh"));
let tail = outbox.unprocessed("HibCh");
assert_eq!(tail[0].payload["tenant_id"], "acme");
}
#[tokio::test]
async fn run_emit_records_a_replay_token_when_a_log_is_attached() {
use crate::replay_token::{InMemoryReplayLog, ReplayLog};
use crate::runtime::channels::{TypedChannelHandle, TypedEventBus};
let bus = std::sync::Arc::new(TypedEventBus::new());
bus.register(TypedChannelHandle::new("HibCh", "Hib"));
let log = std::sync::Arc::new(InMemoryReplayLog::new());
let (mut ctx, _rx) = fresh_ctx();
ctx = ctx.with_event_bus(bus).with_replay_log(log.clone());
ctx.let_bindings
.insert("payload".into(), r#"{"tenant_id":"acme"}"#.into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "HibCh".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
run_emit(&node, &mut ctx).await.unwrap();
assert_eq!(log.len(), 1, "the emit recorded one replay token");
let tokens = log.tokens_for_flow(&ctx.flow_name).await.unwrap();
assert_eq!(tokens.len(), 1);
assert_eq!(tokens[0].effect_name, "emit:HibCh");
assert_eq!(tokens[0].inputs["payload"]["tenant_id"], "acme");
}
#[tokio::test]
async fn run_emit_ephemeral_channel_uses_the_bus_not_the_outbox() {
use crate::event_outbox::{EventOutbox, InMemoryEventOutbox};
use crate::runtime::channels::{TypedChannelHandle, TypedEventBus};
let bus = std::sync::Arc::new(TypedEventBus::new());
bus.register(TypedChannelHandle::new("Tick", "T")); let outbox = std::sync::Arc::new(InMemoryEventOutbox::new());
let (mut ctx, _rx) = fresh_ctx();
ctx = ctx.with_event_bus(bus.clone()).with_event_outbox(outbox.clone());
ctx.let_bindings.insert("payload".into(), "t".into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "Tick".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
run_emit(&node, &mut ctx).await.unwrap();
assert_eq!(outbox.pending_total(), 0, "ephemeral does not touch the outbox");
assert!(bus.receive("Tick").await.is_ok(), "ephemeral went to the bus");
}
#[tokio::test]
async fn run_emit_falls_back_to_buffer_for_an_unregistered_channel() {
use crate::runtime::channels::TypedEventBus;
let bus = std::sync::Arc::new(TypedEventBus::new());
let (mut ctx, _rx) = fresh_ctx();
ctx = ctx.with_event_bus(bus);
ctx.let_bindings.insert("payload".into(), "hello".into());
let node = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "unregistered".into(),
value_ref: "payload".into(),
value_is_channel: false,
};
run_emit(&node, &mut ctx).await.unwrap();
assert_eq!(ctx.let_bindings.get("__channel_unregistered").unwrap(), "hello");
}
#[tokio::test]
async fn run_publish_records_capability() {
let (mut ctx, mut rx) = fresh_ctx();
let node = IRPublish {
node_type: "publish",
source_line: 0,
source_column: 0,
channel_ref: "secure_chan".into(),
shield_ref: "hipaa".into(),
};
run_publish(&node, &mut ctx).await.unwrap();
assert_eq!(
ctx.let_bindings.get("__pub_secure_chan").unwrap(),
"hipaa"
);
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "publish");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn run_discover_binds_under_alias() {
let (mut ctx, mut rx) = fresh_ctx();
publish_capability("secure_chan", "hipaa", &mut ctx);
let node = IRDiscover {
node_type: "discover",
source_line: 0,
source_column: 0,
capability_ref: "secure_chan".into(),
alias: "found".into(),
};
run_discover(&node, &mut ctx).await.unwrap();
assert_eq!(ctx.let_bindings.get("found").unwrap(), "hipaa");
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "discover");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn run_persist_then_retrieve_round_trip() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("id".into(), "42".into());
ctx.let_bindings.insert("name".into(), "test".into());
let persist = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "entities".into(),
};
run_persist(&persist, &mut ctx).await.unwrap();
let retrieve = IRRetrieveStep {
node_type: "retrieve",
source_line: 0,
source_column: 0,
store_name: "entities".into(),
where_expr: "id".into(),
alias: "retrieved_id".into(),
order_by: String::new(),
limit_expr: String::new(),
};
run_retrieve(&retrieve, &mut ctx).await.unwrap();
assert_eq!(ctx.let_bindings.get("retrieved_id").unwrap(), "42");
}
#[tokio::test]
async fn store_ops_fold_observable_row_counts_into_the_ctx() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("id".into(), "1".into());
let persist = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "entities".into(),
};
assert_eq!(
*ctx.store_row_counts.lock().unwrap(),
crate::flow_dispatcher::StoreRowCounts::default()
);
run_persist(&persist, &mut ctx).await.unwrap();
let after_one = ctx.store_row_counts.lock().unwrap().persisted;
assert!(after_one >= 1, "a persist must count >=1 row, got {after_one}");
run_persist(&persist, &mut ctx).await.unwrap();
let after_two = ctx.store_row_counts.lock().unwrap().persisted;
assert_eq!(after_two, after_one * 2, "counts accumulate across ops");
let c = *ctx.store_row_counts.lock().unwrap();
assert_eq!((c.retrieved, c.mutated, c.purged), (0, 0, 0));
}
#[test]
fn store_row_scopes_to_the_declared_field_block() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("message".into(), "hello".into());
ctx.let_bindings.insert("tenant_id".into(), "acme".into());
ctx.let_bindings
.insert("channel_kind".into(), "whatsapp".into());
let node = IRPersistStep {
node_type: "persist",
source_line: 0,
source_column: 0,
store_name: "chat_history".into(),
fields: vec![
("sender".into(), "user".into()),
("content".into(), "${message}".into()),
("tenant_id".into(), "${tenant_id}".into()),
],
};
let row = store_row(&node.fields, &ctx);
assert_eq!(
row,
vec![
("sender".to_string(), SqlValue::Text("user".into())),
("content".to_string(), SqlValue::Text("hello".into())),
("tenant_id".to_string(), SqlValue::Text("acme".into())),
]
);
assert!(!row
.iter()
.any(|(c, _)| c == "channel_kind" || c == "message"));
}
#[test]
fn store_row_without_a_block_falls_back_to_user_bindings() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("a".into(), "1".into());
ctx.let_bindings.insert("b".into(), "2".into());
let node = IRPersistStep {
node_type: "persist",
source_line: 0,
source_column: 0,
store_name: "s".into(),
fields: Vec::new(),
};
assert_eq!(store_row(&node.fields, &ctx), sql_row_from_bindings(&ctx));
}
#[test]
fn store_row_for_a_mutate_node_scopes_to_its_set_block() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("new_balance".into(), "500".into());
ctx.let_bindings.insert("tenant_id".into(), "acme".into());
let node = IRMutateStep {
node_type: "mutate",
source_line: 0,
source_column: 0,
store_name: "accounts".into(),
where_expr: "id = 1".into(),
fields: vec![
("balance".into(), "${new_balance}".into()),
("status".into(), "active".into()),
],
};
let row = store_row(&node.fields, &ctx);
assert_eq!(
row,
vec![
("balance".to_string(), SqlValue::Text("500".into())),
("status".to_string(), SqlValue::Text("active".into())),
]
);
assert!(!row.iter().any(|(c, _)| c == "tenant_id"));
}
#[tokio::test]
async fn run_mutate_updates_existing() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("counter".into(), "1".into());
let persist = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "stats".into(),
};
run_persist(&persist, &mut ctx).await.unwrap();
ctx.let_bindings.insert("counter".into(), "2".into());
let mutate = IRMutateStep {
node_type: "mutate",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "stats".into(),
where_expr: "counter".into(),
};
let outcome = run_mutate(&mutate, &mut ctx).await.unwrap();
match outcome {
NodeOutcome::Completed { output, .. } => {
assert!(output.contains("mutated 1 entries"));
}
other => panic!("expected Completed, got {other:?}"),
}
assert_eq!(
ctx.let_bindings.get("__store_stats_counter").unwrap(),
"2"
);
}
#[tokio::test]
async fn run_purge_removes_persisted_entry() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("tmp".into(), "data".into());
run_persist(
&IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "scratch".into(),
},
&mut ctx,
)
.await
.unwrap();
let outcome = run_purge(
&IRPurgeStep {
node_type: "purge",
source_line: 0,
source_column: 0,
store_name: "scratch".into(),
where_expr: "tmp".into(),
},
&mut ctx,
)
.await
.unwrap();
match outcome {
NodeOutcome::Completed { output, .. } => {
assert!(output.contains("purged 1 entries"));
}
other => panic!("expected Completed, got {other:?}"),
}
assert!(!ctx.let_bindings.contains_key("__store_scratch_tmp"));
}
#[tokio::test]
async fn run_transact_sets_active_marker() {
let (mut ctx, mut rx) = fresh_ctx();
run_transact(
&IRTransactBlock {
node_type: "transact",
source_line: 0,
source_column: 0,
},
&mut ctx,
)
.await
.unwrap();
assert_eq!(ctx.let_bindings.get("__txn_active").unwrap(), "true");
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "transact");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn run_deliberate_canonical_wire_shape() {
let (mut ctx, mut rx) = fresh_ctx();
run_deliberate(
&IRDeliberateBlock {
node_type: "deliberate",
source_line: 0,
source_column: 0,
},
&mut ctx,
)
.await
.unwrap();
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "deliberate");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn run_consensus_canonical_wire_shape() {
let (mut ctx, mut rx) = fresh_ctx();
run_consensus(
&IRConsensusBlock {
node_type: "consensus",
source_line: 0,
source_column: 0,
},
&mut ctx,
)
.await
.unwrap();
let first = rx.try_recv().unwrap();
match first {
FlowExecutionEvent::StepStart { step_type, .. } => {
assert_eq!(step_type, "consensus");
}
e => panic!("expected StepStart, got {e:?}"),
}
}
#[tokio::test]
async fn every_handler_short_circuits_on_cancel() {
let cancel = CancellationFlag::new();
cancel.cancel();
let (tx, _rx) = mpsc::unbounded_channel();
let mut ctx = DispatchCtx::new("F", "stub", "", cancel, tx);
let emit = IREmit {
node_type: "emit",
source_line: 0,
source_column: 0,
channel_ref: "c".into(),
value_ref: "v".into(),
value_is_channel: false,
};
assert!(matches!(run_emit(&emit, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let publish = IRPublish {
node_type: "publish",
source_line: 0,
source_column: 0,
channel_ref: "c".into(),
shield_ref: "s".into(),
};
assert!(matches!(run_publish(&publish, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let discover = IRDiscover {
node_type: "discover",
source_line: 0,
source_column: 0,
capability_ref: "c".into(),
alias: "a".into(),
};
assert!(matches!(run_discover(&discover, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let persist = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "s".into(),
};
assert!(matches!(run_persist(&persist, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let retrieve = IRRetrieveStep {
node_type: "retrieve",
source_line: 0,
source_column: 0,
store_name: "s".into(),
where_expr: "w".into(),
alias: "a".into(),
order_by: String::new(),
limit_expr: String::new(),
};
assert!(matches!(run_retrieve(&retrieve, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let mutate = IRMutateStep {
node_type: "mutate",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "s".into(),
where_expr: "w".into(),
};
assert!(matches!(run_mutate(&mutate, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let purge = IRPurgeStep {
node_type: "purge",
source_line: 0,
source_column: 0,
store_name: "s".into(),
where_expr: "w".into(),
};
assert!(matches!(run_purge(&purge, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let transact = IRTransactBlock {
node_type: "transact",
source_line: 0,
source_column: 0,
};
assert!(matches!(run_transact(&transact, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let deliberate = IRDeliberateBlock {
node_type: "deliberate",
source_line: 0,
source_column: 0,
};
assert!(matches!(run_deliberate(&deliberate, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
let consensus = IRConsensusBlock {
node_type: "consensus",
source_line: 0,
source_column: 0,
};
assert!(matches!(run_consensus(&consensus, &mut ctx).await, Err(DispatchError::UpstreamCancelled)));
}
fn axonstore(name: &str, backend: &str, connection: &str) -> IRAxonStore {
IRAxonStore {
node_type: "axonstore",
source_line: 0,
source_column: 0,
name: name.to_string(),
backend: backend.to_string(),
connection: connection.to_string(),
confidence_floor: None,
isolation: String::new(),
on_breach: String::new(),
capability: String::new(),
column_schema: None,
}
}
fn ctx_with_registry(
specs: &[IRAxonStore],
) -> (DispatchCtx, mpsc::UnboundedReceiver<FlowExecutionEvent>) {
let (tx, rx) = mpsc::unbounded_channel();
let registry = crate::store::registry::StoreRegistry::build(specs).unwrap();
let ctx = DispatchCtx::new("TestFlow", "stub", "", CancellationFlag::new(), tx)
.with_store_registry(std::sync::Arc::new(registry));
(ctx, rx)
}
#[test]
fn resolve_pg_backend_no_registry_is_kv() {
let (ctx, _rx) = fresh_ctx();
assert!(resolve_pg_backend(&ctx, "anything").unwrap().is_none());
}
#[test]
fn resolve_pg_backend_in_memory_store_is_kv() {
let (ctx, _rx) = ctx_with_registry(&[axonstore("cache", "in_memory", "")]);
assert!(resolve_pg_backend(&ctx, "cache").unwrap().is_none());
assert!(resolve_pg_backend(&ctx, "undeclared").unwrap().is_none());
}
#[test]
fn resolve_pg_backend_missing_env_var_errors_not_kv_fallback() {
let (ctx, _rx) = ctx_with_registry(&[axonstore(
"tenants",
"postgresql",
"env:AXON_NONEXISTENT_VAR_FASE35F",
)]);
assert!(matches!(
resolve_pg_backend(&ctx, "tenants"),
Err(StoreError::MissingEnvVar { .. })
));
}
#[tokio::test]
async fn run_retrieve_postgresql_missing_env_surfaces_backend_error() {
let (mut ctx, _rx) = ctx_with_registry(&[axonstore(
"tenants",
"postgresql",
"env:AXON_NONEXISTENT_VAR_FASE35F",
)]);
let node = IRRetrieveStep {
node_type: "retrieve",
source_line: 0,
source_column: 0,
store_name: "tenants".into(),
where_expr: "id = 1".into(),
alias: "found".into(),
order_by: String::new(),
limit_expr: String::new(),
};
assert!(matches!(
run_retrieve(&node, &mut ctx).await,
Err(DispatchError::BackendError { .. })
));
}
#[tokio::test]
async fn run_persist_postgresql_malformed_dsn_surfaces_backend_error() {
let (mut ctx, _rx) =
ctx_with_registry(&[axonstore("events", "postgresql", "not a dsn")]);
ctx.let_bindings.insert("kind".into(), "login".into());
let node = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "events".into(),
};
assert!(matches!(
run_persist(&node, &mut ctx).await,
Err(DispatchError::BackendError { .. })
));
}
#[tokio::test]
async fn run_persist_in_memory_store_keeps_byte_identical_kv_path() {
let (mut ctx, _rx) = ctx_with_registry(&[axonstore("cache", "in_memory", "")]);
ctx.let_bindings.insert("k".into(), "v".into());
let node = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "cache".into(),
};
match run_persist(&node, &mut ctx).await.unwrap() {
NodeOutcome::Completed { output, .. } => {
assert!(output.contains("entries"), "KV path output shape");
}
other => panic!("expected Completed, got {other:?}"),
}
assert_eq!(ctx.let_bindings.get("__store_cache_k").unwrap(), "v");
}
#[tokio::test]
async fn run_persist_below_confidence_floor_is_blocked() {
let mut store =
axonstore("ledger", "postgresql", "postgresql://u:p@localhost:5432/db");
store.confidence_floor = Some(0.8);
let (mut ctx, _rx) = ctx_with_registry(&[store]);
ctx.let_bindings.insert("amount".into(), "100".into()); let node = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "ledger".into(),
};
assert!(matches!(
run_persist(&node, &mut ctx).await,
Err(DispatchError::BackendError { .. })
));
}
fn gated_kv(name: &str, capability: &str) -> IRAxonStore {
let mut s = axonstore(name, "in_memory", "");
s.capability = capability.to_string();
s
}
fn ctx_with_caps(
specs: &[IRAxonStore],
held: Vec<String>,
) -> (DispatchCtx, mpsc::UnboundedReceiver<FlowExecutionEvent>) {
let (tx, rx) = mpsc::unbounded_channel();
let registry = crate::store::registry::StoreRegistry::build(specs).unwrap();
let ctx = DispatchCtx::new("F", "stub", "", CancellationFlag::new(), tx)
.with_store_registry(std::sync::Arc::new(registry))
.with_held_capabilities(held);
(ctx, rx)
}
fn retrieve_node(store: &str) -> IRRetrieveStep {
IRRetrieveStep {
node_type: "retrieve",
source_line: 0,
source_column: 0,
store_name: store.to_string(),
where_expr: "k".to_string(),
alias: "v".to_string(),
order_by: String::new(),
limit_expr: String::new(),
}
}
#[tokio::test]
async fn retrieve_denied_when_capability_not_held() {
let (mut ctx, _rx) = ctx_with_caps(
&[gated_kv("tenants", "tenant.read")],
vec!["audit.write".to_string()],
);
assert!(matches!(
run_retrieve(&retrieve_node("tenants"), &mut ctx).await,
Err(DispatchError::BackendError { .. })
));
}
#[tokio::test]
async fn retrieve_allowed_when_capability_held() {
let (mut ctx, _rx) = ctx_with_caps(
&[gated_kv("tenants", "tenant.read")],
vec!["tenant.read".to_string()],
);
assert!(run_retrieve(&retrieve_node("tenants"), &mut ctx).await.is_ok());
}
#[tokio::test]
async fn persist_into_gated_store_denied_without_capability() {
let (mut ctx, _rx) =
ctx_with_caps(&[gated_kv("ledger", "ledger.write")], vec![]);
let node = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "ledger".into(),
};
assert!(matches!(
run_persist(&node, &mut ctx).await,
Err(DispatchError::BackendError { .. })
));
}
#[tokio::test]
async fn ungated_store_needs_no_capability() {
let (mut ctx, _rx) =
ctx_with_caps(&[axonstore("cache", "in_memory", "")], vec![]);
assert!(run_retrieve(&retrieve_node("cache"), &mut ctx).await.is_ok());
}
#[tokio::test]
async fn no_capability_context_skips_the_runtime_recheck() {
let (mut ctx, _rx) = ctx_with_registry(&[gated_kv("tenants", "tenant.read")]);
assert!(ctx.held_capabilities.is_none());
assert!(run_retrieve(&retrieve_node("tenants"), &mut ctx).await.is_ok());
}
#[tokio::test]
async fn persist_appends_a_delta_to_the_audit_chain() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("k".into(), "v".into());
let node = IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "s".into(),
};
run_persist(&node, &mut ctx).await.unwrap();
let chain = ctx.audit_chain.lock().unwrap();
assert_eq!(chain.len(), 1);
assert_eq!(
chain.verify(),
crate::store::audit_chain::ChainVerdict::Intact
);
}
#[tokio::test]
async fn retrieve_does_not_append_an_audit_delta() {
let (mut ctx, _rx) = fresh_ctx();
run_retrieve(&retrieve_node("s"), &mut ctx).await.unwrap();
assert!(ctx.audit_chain.lock().unwrap().is_empty());
}
#[tokio::test]
async fn persist_mutate_purge_chain_into_one_verifiable_history() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("k".into(), "v".into());
run_persist(
&IRPersistStep {
node_type: "persist",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "s".into(),
},
&mut ctx,
)
.await
.unwrap();
run_mutate(
&IRMutateStep {
node_type: "mutate",
fields: Vec::new(),
source_line: 0,
source_column: 0,
store_name: "s".into(),
where_expr: "k".into(),
},
&mut ctx,
)
.await
.unwrap();
run_purge(
&IRPurgeStep {
node_type: "purge",
source_line: 0,
source_column: 0,
store_name: "s".into(),
where_expr: "k".into(),
},
&mut ctx,
)
.await
.unwrap();
let chain = ctx.audit_chain.lock().unwrap();
assert_eq!(chain.len(), 3, "three mutations → three chained deltas");
assert_eq!(
chain.verify(),
crate::store::audit_chain::ChainVerdict::Intact
);
}
#[test]
fn sql_row_from_bindings_excludes_namespace_keys_and_sorts() {
let (mut ctx, _rx) = fresh_ctx();
ctx.let_bindings.insert("name".into(), "Alice".into());
ctx.let_bindings.insert("id".into(), "7".into());
ctx.let_bindings
.insert("__store_internal".into(), "bookkeeping".into());
let row = sql_row_from_bindings(&ctx);
assert_eq!(
row,
vec![
("id".to_string(), SqlValue::Text("7".to_string())),
("name".to_string(), SqlValue::Text("Alice".to_string())),
]
);
}
}