use std::sync::Arc;
use aion_core::{ContentType, Payload, WorkflowId};
use aion_store::workloop::WorkloopStore;
use beamr::native::ProcessContext;
use beamr::term::Term;
use beamr::term::binary_ref::BinaryRef;
use beamr::term::heap_borrow::HeapBorrow;
use crate::runtime::nif_result_term::{NifRefusal, error_result_term, ok_result_term};
use crate::runtime::nif_state::EngineNifState;
use crate::workloop::close::IterationCloseContext;
use crate::workloop::iteration::WorkloopIterationClose;
pub(crate) const CLOSE_ITERATION_NIF: &str = "close_iteration";
pub(crate) struct WorkloopNifBridge {
close: IterationCloseContext,
tokio_handle: tokio::runtime::Handle,
}
impl WorkloopNifBridge {
pub(crate) const fn new(
close: IterationCloseContext,
tokio_handle: tokio::runtime::Handle,
) -> Self {
Self {
close,
tokio_handle,
}
}
fn workloop_store(&self) -> &Arc<dyn WorkloopStore> {
&self.close.workloop_store
}
}
pub(crate) fn install_workloop_nif_bridge(state: &EngineNifState, bridge: Arc<WorkloopNifBridge>) {
match state.workloop_bridge.write() {
Ok(mut slot) => *slot = Some(bridge),
Err(poisoned) => *poisoned.into_inner() = Some(bridge),
}
}
pub(crate) fn release_workloop_nif_bridge(state: &EngineNifState) {
match state.workloop_bridge.write() {
Ok(mut slot) => *slot = None,
Err(poisoned) => *poisoned.into_inner() = None,
}
}
fn workloop_bridge(state: &EngineNifState) -> Option<Arc<WorkloopNifBridge>> {
match state.workloop_bridge.read() {
Ok(slot) => slot.clone(),
Err(poisoned) => poisoned.into_inner().clone(),
}
}
pub(crate) fn is_registered_workloop(
state: &EngineNifState,
workflow_id: &WorkflowId,
) -> Result<bool, String> {
let Some(bridge) = workloop_bridge(state) else {
return Ok(false);
};
bridge
.tokio_handle
.block_on(bridge.workloop_store().get_workloop(workflow_id))
.map(|record| record.is_some())
.map_err(|error| format!("workloop_registration_probe:{error}"))
}
pub(crate) fn workflow_path_refusal(workflow_id: &WorkflowId) -> String {
format!(
"continue_as_new refused: workflow {workflow_id} is a REGISTERED WORKLOOP, and the \
workflow continue-as-new path would record a bare WorkflowContinuedAsNew and let the \
process-exit monitor start a resident successor process — defeating the park a loop \
relies on between fires (R13.3), and closing the iteration without its routes, health \
samples, or invariant current-state records. A workloop's `route start` must compile to \
aion_flow_ffi:{CLOSE_ITERATION_NIF}/3 (carry, routes, invariant_states). This is a \
compiled-output defect, not a runtime condition: nothing is recorded and the run stays \
live"
)
}
pub(super) fn close_iteration_impl(args: &[Term], ctx: &mut ProcessContext) -> Result<Term, Term> {
match run_close_iteration(args, ctx) {
Ok(term) => Ok(term),
Err(refusal) => refusal.into_nif_result(),
}
}
fn run_close_iteration(args: &[Term], ctx: &mut ProcessContext) -> Result<Term, NifRefusal> {
if args.len() != 3 {
return Err(close_refusal(
ctx,
&format!("expected 3 arguments, got {}", args.len()),
));
}
let carry_text = decode_string_arg(args[0], ctx.borrow_terms())
.map_err(|error| close_refusal(ctx, &format!("carry:{error}")))?;
let routes_text = decode_string_arg(args[1], ctx.borrow_terms())
.map_err(|error| close_refusal(ctx, &format!("routes:{error}")))?;
let states_text = decode_string_arg(args[2], ctx.borrow_terms())
.map_err(|error| close_refusal(ctx, &format!("invariant_states:{error}")))?;
let routes = decode_routes(&routes_text).map_err(|error| close_refusal(ctx, &error))?;
let invariant_states =
decode_invariant_states(&states_text).map_err(|error| close_refusal(ctx, &error))?;
let state = crate::runtime::nif_state::engine_nif_state(ctx)
.map_err(|message| close_refusal(ctx, &message))?;
let pid = ctx
.pid()
.ok_or_else(|| close_refusal(ctx, "missing_caller_pid"))?;
crate::runtime::nif_query_pump::ensure_not_servicing_query(&state, pid, CLOSE_ITERATION_NIF)
.map_err(|message| close_refusal(ctx, &message))?;
let Some(bridge) = workloop_bridge(&state) else {
return Err(close_refusal(
ctx,
"no workloop service is configured on this engine, so no iteration can be closed \
(EngineBuilder::with_workloop_service)",
));
};
let runtime = crate::runtime::nif_activity::runtime_context(&state)
.map_err(|error| close_refusal(ctx, &error.error_reason()))?;
let handle = crate::runtime::nif_context::NifContext::workflow_handle_for_pid(
pid,
runtime.registry.as_ref(),
runtime.runtime.signal_delivery(),
)
.map_err(|error| close_refusal(ctx, &error.error_reason()))?;
let loop_id = handle.workflow_id().clone();
let close = WorkloopIterationClose {
routes,
carry: Payload::new(ContentType::Json, carry_text.into_bytes()),
invariant_states,
};
match bridge
.tokio_handle
.block_on(crate::workloop::close::close_iteration(
&bridge.close,
&loop_id,
close,
)) {
Ok(next_run_id) => {
if let Err(error) = runtime.runtime.cancel_pid(pid) {
tracing::error!(
workflow_id = %loop_id,
next_run_id = %next_run_id,
error = %error,
"workloop iteration closed durably but the closing generation's process \
could not be ended; it is still runnable on a generation that has already \
continued, and every durable NIF it goes on to call writes into a closed \
history"
);
return Err(close_refusal(
ctx,
&format!("iteration closed but process termination failed: {error}"),
));
}
ok_result_term(ctx, b"iteration_closed").map_err(NifRefusal::Unbuildable)
}
Err(error) => {
tracing::warn!(
workflow_id = %loop_id,
error = %error,
"workloop iteration close refused; nothing was recorded and the generation \
stays live, so the dead-man switch counts this as a missed window"
);
Err(close_refusal(ctx, &format!("close_refused:{error}")))
}
}
}
fn decode_routes(text: &str) -> Result<Vec<String>, String> {
let value: serde_json::Value =
serde_json::from_str(text).map_err(|error| format!("routes_not_json:{error}"))?;
let array = value
.as_array()
.ok_or_else(|| String::from("routes_not_an_array"))?;
if array.is_empty() {
return Err(String::from(
"routes_empty: an iteration close must name the routes it took; an empty list \
confirms no invariant and would red every declared invariant on this tick",
));
}
array
.iter()
.map(|entry| {
let name = entry
.as_str()
.ok_or_else(|| String::from("route_not_a_string"))?;
if name.is_empty() {
return Err(String::from("route_name_empty"));
}
Ok(name.to_owned())
})
.collect()
}
fn decode_invariant_states(text: &str) -> Result<Vec<(String, Payload)>, String> {
let value: serde_json::Value =
serde_json::from_str(text).map_err(|error| format!("invariant_states_not_json:{error}"))?;
let object = value
.as_object()
.ok_or_else(|| String::from("invariant_states_not_an_object"))?;
object
.iter()
.map(|(name, entry)| {
if name.is_empty() {
return Err(String::from("invariant_name_empty"));
}
let bytes = serde_json::to_vec(entry)
.map_err(|error| format!("invariant_state_unencodable:{error}"))?;
Ok((name.clone(), Payload::new(ContentType::Json, bytes)))
})
.collect()
}
fn close_refusal(ctx: &mut ProcessContext, message: &str) -> NifRefusal {
NifRefusal::reported(error_result_term(
ctx,
&format!("{CLOSE_ITERATION_NIF}:{message}"),
))
}
fn decode_string_arg(term: Term, heap: HeapBorrow<'_>) -> Result<String, String> {
let bin = BinaryRef::new(term).ok_or_else(|| "argument is not a binary".to_owned())?;
String::from_utf8(bin.as_bytes(heap).to_vec())
.map_err(|_| "argument is not valid UTF-8".to_owned())
}
#[cfg(test)]
mod tests {
use super::{decode_invariant_states, decode_routes};
#[test]
fn routes_decode_in_the_order_taken() -> Result<(), Box<dyn std::error::Error>> {
assert_eq!(
decode_routes(r#"["sweep","dispatch","start"]"#)?,
vec![
String::from("sweep"),
String::from("dispatch"),
String::from("start")
]
);
Ok(())
}
#[test]
fn an_empty_route_list_is_refused() -> Result<(), Box<dyn std::error::Error>> {
let error = decode_routes("[]")
.err()
.ok_or("empty routes were accepted")?;
assert!(error.starts_with("routes_empty"), "error: {error}");
Ok(())
}
#[test]
fn malformed_route_arguments_are_typed_refusals() -> Result<(), Box<dyn std::error::Error>> {
assert!(
decode_routes("not json")
.err()
.ok_or("non-JSON accepted")?
.starts_with("routes_not_json")
);
assert_eq!(
decode_routes(r#"{"a":1}"#).err().ok_or("object accepted")?,
"routes_not_an_array"
);
assert_eq!(
decode_routes("[1]").err().ok_or("number accepted")?,
"route_not_a_string"
);
assert_eq!(
decode_routes(r#"[""]"#)
.err()
.ok_or("empty name accepted")?,
"route_name_empty"
);
Ok(())
}
#[test]
fn invariant_states_carry_their_values_as_opaque_json() -> Result<(), Box<dyn std::error::Error>>
{
let states = decode_invariant_states(r#"{"serving":{"queued":3},"drained":true}"#)?;
let mut names: Vec<&str> = states.iter().map(|(name, _)| name.as_str()).collect();
names.sort_unstable();
assert_eq!(names, vec!["drained", "serving"]);
let serving = states
.iter()
.find(|(name, _)| name == "serving")
.ok_or("serving state missing")?;
assert_eq!(serving.1.bytes(), br#"{"queued":3}"#);
Ok(())
}
#[test]
fn no_invariant_states_is_accepted() -> Result<(), Box<dyn std::error::Error>> {
assert!(decode_invariant_states("{}")?.is_empty());
Ok(())
}
#[test]
fn malformed_invariant_states_are_typed_refusals() -> Result<(), Box<dyn std::error::Error>> {
assert!(
decode_invariant_states("[]")
.err()
.ok_or("array accepted")?
.contains("not_an_object")
);
assert!(
decode_invariant_states("nope")
.err()
.ok_or("non-JSON accepted")?
.starts_with("invariant_states_not_json")
);
Ok(())
}
}