use std::sync::Arc;
use std::sync::atomic::AtomicU32;
use serde_json::json;
use crate::client::GatewayClient;
use crate::debug::DebugCapture;
use crate::lua::{LuaFanoutResult, LuaProgram, SectionVm, ToolBindings};
use crate::model::ModelBindings;
use crate::observe::{Observation, Observer, detail};
use crate::parser::Section;
use crate::store::StoreRef;
use crate::tools::SharedTools;
use crate::{Error, Result, cancel, subst};
use super::proxies::ProxyObserver;
pub(crate) struct ArmPayload {
pub(crate) worker: Section,
pub(crate) item_text: String,
pub(crate) index: usize,
pub(crate) store: StoreRef,
pub(crate) client: Option<GatewayClient>,
pub(crate) args: String,
pub(crate) execution: String,
pub(crate) when: String,
pub(crate) last_reply: Option<String>,
pub(crate) shared: Option<LuaProgram>,
pub(crate) bindings: ToolBindings,
pub(crate) models: ModelBindings,
pub(crate) analysis: crate::execute::ToolAnalysis,
pub(crate) shared_tools: SharedTools,
pub(crate) max_tool_iterations: usize,
pub(crate) lua_memory_bytes: usize,
pub(crate) lua_log_events: u32,
pub(crate) parent_id: usize,
pub(crate) section_count: usize,
pub(crate) turns: Arc<AtomicU32>,
pub(crate) observer: Arc<ProxyObserver>,
pub(crate) debug: Option<Arc<dyn DebugCapture>>,
pub(crate) cancel: Option<cancel::CancelHandle>,
}
pub(crate) struct ArmFinalizer {
observer: Arc<ProxyObserver>,
execution: String,
section: String,
finished: bool,
}
impl ArmFinalizer {
pub(crate) fn new(observer: Arc<ProxyObserver>, execution: String, section: String) -> Self {
Self {
observer,
execution,
section,
finished: false,
}
}
pub(crate) fn finish(&mut self, event: Observation) {
self.finished = true;
(self.observer.as_ref() as &dyn Observer).observe(&self.execution, &self.section, event);
}
}
impl Drop for ArmFinalizer {
fn drop(&mut self) {
if !self.finished {
(self.observer.as_ref() as &dyn Observer).observe(
&self.execution,
&self.section,
detail::FANOUT_ARM_CANCELLED,
);
}
}
}
#[expect(
clippy::too_many_lines,
reason = "the arm body is one cohesive linear sequence of fallible steps"
)]
pub(crate) async fn run_one_arm(payload: ArmPayload) -> Result<(usize, LuaFanoutResult)> {
let ArmPayload {
worker,
item_text,
index,
store,
client,
args,
execution,
when,
last_reply,
shared,
bindings,
models,
analysis,
shared_tools,
max_tool_iterations,
lua_memory_bytes,
lua_log_events,
parent_id,
section_count,
turns,
observer,
debug,
cancel,
} = payload;
let taskid = (index + 1).to_string();
let observer_arc = observer;
let observer = observer_arc.as_ref() as &dyn Observer;
observer.observe(&execution, &worker.name, detail::FANOUT_ARM_STARTED);
let mut finalizer = ArmFinalizer::new(
Arc::clone(&observer_arc),
execution.clone(),
worker.name.clone(),
);
let mut vm = match SectionVm::new_for_section(
shared.as_ref(),
&bindings,
&models,
&execution,
observer,
&worker.name,
) {
Ok(vm) => vm,
Err(error) => {
finalizer.finish(detail::FANOUT_ARM_FAILED);
return Err(error);
}
};
let body = async {
vm.apply_lua_limits(lua_memory_bytes, lua_log_events)?;
let now = crate::execute::now_rfc3339_checked()?;
let sys = json!({
"when": when,
"now": now,
"id": parent_id,
"taskid": taskid,
"section_name": worker.name,
"execution": execution,
"section_count": section_count,
});
vm.inject_host(&args, &sys, &store, last_reply.as_deref())?;
vm.set_global_string("item", &item_text)?;
if let Some(program) = worker.prologue()
&& let Some(value) = vm.run_prologue(program, observer, &worker.name)?
{
return Ok((
LuaFanoutResult::success(&item_text, value),
detail::FANOUT_ARM_SUCCEEDED,
));
}
let scopes = vm.close_scopes(observer, &worker.name)?;
let scope = scopes.tools;
let counts = Some(vm.install_tool_call_counts(&scope)?);
let sys = if let Some(model_binding) = scopes.model.as_ref() {
let current = vm.current_sys(&sys)?;
let enriched = crate::lua::enrich_sys_model(¤t, model_binding);
vm.re_seal_sys(&enriched)?;
enriched
} else {
sys
};
let var = vm.var()?;
let prose = subst::substitute(
worker.prose(),
&args,
last_reply.as_deref(),
Some(&item_text),
&var,
&sys,
)?;
let mut arm_reply: Option<String> = None;
if !prose.trim().is_empty() {
let Some(model_binding) = scopes.model else {
return Err(Error::ModelRequired {
section: worker.name.clone(),
});
};
let completion_options = model_binding.completion_options();
let registry = shared_tools.registry();
let (schemas, dispatch) = crate::execute::prepare_effective_scope(
&analysis,
&scope,
®istry,
&execution,
observer,
&worker.name,
)?;
if let Some(client) = client.as_ref() {
let global_aliases = Some(&analysis.alias_to_id);
let debug_ref = debug.as_deref();
match crate::execute::run_tool_loop(
client,
&schemas,
&dispatch,
®istry,
prose,
max_tool_iterations,
crate::execute::SectionProgress {
execution: &execution,
observer,
section: &worker.name,
turns: turns.as_ref(),
debug: debug_ref,
completion_options: &completion_options,
},
counts.as_ref(),
global_aliases,
)
.await
{
Ok((text, finish_reason)) => {
let current = vm.current_sys(&sys)?;
let enriched = crate::lua::enrich_sys_reply_finish_reason(
¤t,
finish_reason.as_deref(),
);
vm.re_seal_sys(&enriched)?;
vm.bind_reply(&text, observer, &worker.name)?;
arm_reply = Some(text);
}
Err(Error::ToolLoopExhausted) => {
let stub = format!(
"## {item_text}\n\nUNKNOWN\n\n(section incomplete: tool loop exhausted)"
);
return Ok((
LuaFanoutResult::exhausted_stub(&item_text, stub),
detail::FANOUT_ARM_EXHAUSTED,
));
}
Err(error) => return Err(error),
}
}
}
let epilog_return = if let Some(program) = worker.epilog() {
vm.run_epilog(program, observer, &worker.name)?
} else {
None
};
let text = epilog_return.or(arm_reply).unwrap_or_default();
Ok((
LuaFanoutResult::success(item_text.clone(), text),
detail::FANOUT_ARM_SUCCEEDED,
))
};
let outcome: Result<(LuaFanoutResult, Observation)> = match cancel {
Some(handle) => cancel::scope(handle, body).await,
None => body.await,
};
vm.teardown(observer, &worker.name);
match outcome {
Ok((result, event)) => {
finalizer.finish(event);
Ok((index, result))
}
Err(error) => {
finalizer.finish(detail::FANOUT_ARM_FAILED);
Err(error)
}
}
}