use std::collections::{HashMap, HashSet, VecDeque};
use std::time::Duration;
use bevy_ecs::entity::Entity;
use tokio::sync::broadcast;
use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
use crate::components::{
AgentMessage, AgentState, AgentStatus, AwaitingInteraction, ContextWindow, ParentRef,
SubAgentChildren, WaitReason,
};
use crate::interaction_hub::InteractionHub;
use crate::persistence::{RunMetadata, TokenTotals};
use crate::world::{AgentId, LaneSnapshot, PipelineWorld};
mod events;
pub use events::*;
mod types;
pub use types::*;
pub struct WorldHost {
world: PipelineWorld,
by_run_id: HashMap<String, AgentId>,
interactions: InteractionHub,
spawner: Option<Spawner>,
spawn_preprocessor: Option<SpawnPreprocessor>,
reloader: Option<Reloader>,
force_terminator: Option<ForceTerminator>,
reaper: Option<Reaper>,
events: broadcast::Sender<WorldEvent>,
emitted: HashMap<String, Emitted>,
emitted_interactions: HashSet<String>,
subagent_tx: UnboundedSender<SubAgentOp>,
subagent_rx: UnboundedReceiver<SubAgentOp>,
redrive: Duration,
dead_cycles: u32,
last_progress: Option<u64>,
relief_granted: usize,
healthy_cycles: u32,
dead_cycles_before_relief: u32,
finished: VecDeque<(i64, RunListEntry)>,
finished_retention_secs: u64,
parked: HashMap<String, RunListEntry>,
}
const HEALTHY_CYCLES_BEFORE_DECAY: u32 = 4;
const DEFAULT_REDRIVE_INTERVAL: Duration = Duration::from_secs(30);
pub const DEFAULT_DEAD_CYCLES_BEFORE_RELIEF: u32 = 10;
pub const DEFAULT_FINISHED_RETENTION_SECS: u64 = 300;
const MAX_RETAINED_FINISHED: usize = 256;
impl WorldHost {
pub fn new(world: PipelineWorld) -> Self {
Self::with_interactions(world, InteractionHub::new())
}
pub fn with_interactions(mut world: PipelineWorld, interactions: InteractionHub) -> Self {
let (events, _) = broadcast::channel(256);
world
.world_mut()
.insert_resource(WorldEventSink(events.clone()));
let (subagent_tx, subagent_rx) = tokio::sync::mpsc::unbounded_channel();
Self {
world,
by_run_id: HashMap::new(),
interactions,
spawner: None,
spawn_preprocessor: None,
reloader: None,
force_terminator: None,
reaper: None,
events,
emitted: HashMap::new(),
emitted_interactions: HashSet::new(),
parked: HashMap::new(),
subagent_tx,
subagent_rx,
redrive: DEFAULT_REDRIVE_INTERVAL,
dead_cycles: 0,
last_progress: None,
relief_granted: 0,
healthy_cycles: 0,
dead_cycles_before_relief: DEFAULT_DEAD_CYCLES_BEFORE_RELIEF,
finished: VecDeque::new(),
finished_retention_secs: DEFAULT_FINISHED_RETENTION_SECS,
}
}
}
mod emit;
mod health;
mod listing;
mod subagents;
impl WorldHost {
fn adopt_unregistered_runs(&mut self) {
let live: Vec<(String, Entity)> = self
.world
.world_mut()
.query::<(Entity, &RunMetadata)>()
.iter(self.world.world())
.map(|(entity, md)| (md.run_id.clone(), entity))
.collect();
for (run_id, entity) in live {
let agent = self.world.own_agent(entity);
if self.live_entity(&run_id) != Some(agent) {
self.by_run_id.insert(run_id, agent);
}
}
}
fn parkable(&self, entity: Entity, status: &AgentStatus) -> bool {
if !matches!(status, AgentStatus::Paused) {
return false;
}
if self.reloader.is_none() {
return false;
}
let world = self.world.world();
let paused_persisted = world
.get::<crate::pipeline::PersistWatermark>(entity)
.and_then(|w| w.persisted_status())
== Some(leviath_core::run_meta::RunStatus::Paused);
paused_persisted
&& world.get::<crate::components::ParentRef>(entity).is_none()
&& world.get::<SubAgentChildren>(entity).is_none()
&& world.get::<crate::fanout::FanOutWaiting>(entity).is_none()
&& world
.get::<crate::interaction_points::AwaitingInteractionPoint>(entity)
.is_none()
&& world.get::<AwaitingInteraction>(entity).is_none()
}
fn no_live_parent(&self, entity: Entity) -> bool {
let world = self.world.world();
match world.get::<crate::components::ParentRef>(entity) {
None => true,
Some(parent_ref) => match world.get::<AgentState>(parent_ref.parent_entity) {
None => true,
Some(state) => matches!(
state.status,
AgentStatus::Complete | AgentStatus::Error { .. } | AgentStatus::Cancelled
),
},
}
}
pub fn set_spawner(&mut self, spawner: Spawner) {
self.spawner = Some(spawner);
}
pub fn set_spawn_preprocessor(&mut self, pp: SpawnPreprocessor) {
self.spawn_preprocessor = Some(pp);
}
pub fn set_reloader(&mut self, reloader: Reloader) {
self.reloader = Some(reloader);
}
pub fn set_force_terminator(&mut self, force_terminator: ForceTerminator) {
self.force_terminator = Some(force_terminator);
}
pub fn set_reaper(&mut self, reaper: Reaper) {
self.reaper = Some(reaper);
}
fn resolve_or_reload(&mut self, run_id: &str) -> Option<AgentId> {
if let Some(entity) = self.live_entity(run_id) {
return Some(entity);
}
let entity = (self.reloader.as_mut()?)(&mut self.world, run_id)?;
self.by_run_id.insert(run_id.to_string(), entity);
self.parked.remove(run_id);
Some(entity)
}
pub fn interactions(&self) -> InteractionHub {
self.interactions.clone()
}
pub fn world_mut(&mut self) -> &mut PipelineWorld {
&mut self.world
}
pub fn register(&mut self, run_id: impl Into<String>, agent: AgentId) {
let run_id = run_id.into();
self.parked.remove(&run_id);
self.by_run_id.insert(run_id, agent);
}
fn live_entity(&self, run_id: &str) -> Option<AgentId> {
let agent = *self.by_run_id.get(run_id)?;
self.world
.world()
.get::<AgentState>(agent.entity())
.map(|_| agent)
}
pub fn handle(&mut self, op: ControlOp) {
match op {
ControlOp::Spawn { args, reply } => {
let result = match self.spawner.as_mut() {
Some(spawner) => {
match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
spawner(&mut self.world, &args)
})) {
Ok(Ok(entity)) => {
let agent = self.world.own_agent(entity);
self.by_run_id.insert(args.run_id.clone(), agent);
Ok(args.run_id.clone())
}
Ok(Err(e)) => Err(e),
Err(_) => Err("agent spawn panicked".to_string()),
}
}
None => Err("this daemon cannot spawn agents".to_string()),
};
if let Err(error) = &result {
tracing::error!(
run_id = %args.run_id,
blueprint = %args.blueprint_path,
workdir = %args.workdir,
error = %error,
"agent spawn failed"
);
}
let _ = reply.send(result);
}
ControlOp::Result { run_id, reply } => {
let output = self
.live_entity(&run_id)
.and_then(|agent| {
self.world
.world()
.get::<crate::persistence::FinalOutput>(agent.entity())
})
.map(|o| o.0.clone());
let _ = reply.send(output);
}
ControlOp::Status { run_id, reply } => {
let status = self
.live_entity(&run_id)
.and_then(|e| self.world.agent_status(e))
.or_else(|| self.parked.get(&run_id).map(|e| e.status.clone()))
.or_else(|| {
self.finished
.iter()
.find(|(_, e)| e.run_id == run_id)
.map(|(_, e)| e.status.clone())
});
let _ = reply.send(status);
}
ControlOp::Pause { run_id, reply } => {
let ok = self
.resolve_or_reload(&run_id)
.is_some_and(|e| self.world.pause(e));
let _ = reply.send(ok);
}
ControlOp::Resume { run_id, reply } => {
let ok = self
.resolve_or_reload(&run_id)
.is_some_and(|e| self.world.resume(e));
let _ = reply.send(ok);
}
ControlOp::Cancel { run_id, reply } => {
let ok = self.cancel_tree(&run_id)
|| self
.force_terminator
.as_mut()
.is_some_and(|terminate| terminate(&run_id));
let _ = reply.send(ok);
}
ControlOp::List { reply } => {
let _ = reply.send(RunListing {
runs: self.list(),
finished: self.finished(),
health: self.health(),
});
}
ControlOp::Message {
agent_id,
content,
target_region,
reply,
} => {
self.resolve_or_reload(&agent_id);
let ok = self
.world
.send_message(AgentMessage {
agent_id,
content,
target_region,
})
.is_ok();
let _ = reply.send(ok);
}
ControlOp::ListInteractions { reply } => {
let _ = reply.send(self.interactions.pending());
}
ControlOp::AnswerInteraction { response, reply } => {
let _ = reply.send(self.interactions.answer(response));
}
ControlOp::CancelInteraction { request_id, reply } => {
let _ = reply.send(self.interactions.cancel(&request_id));
}
ControlOp::Shutdown { reply } => {
let _ = reply.send(true);
self.world.shutdown();
}
}
}
pub async fn flush_and_stop(&mut self) {
self.world.flush_and_stop().await;
}
pub async fn serve(&mut self, mut control_rx: UnboundedReceiver<ControlOp>) {
let wake = self.world.wake_handle();
let shutdown = self.world.shutdown_handle();
let mut redrive =
tokio::time::interval_at(tokio::time::Instant::now() + self.redrive, self.redrive);
redrive.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
'serve: loop {
self.world.run_to_fixed_point();
self.emit_events();
tokio::select! {
_ = wake.notified() => {}
_ = shutdown.notified() => break 'serve,
_ = redrive.tick() => self.observe_redrive(),
op = control_rx.recv() => {
match op {
Some(op) => {
let pre = match &op {
ControlOp::Spawn { args, .. } => {
self.spawn_preprocessor.as_ref().map(|pp| pp(args))
}
_ => None,
};
if let Some(fut) = pre {
fut.await;
}
self.handle(op);
}
None => break 'serve, }
}
Some(sub) = self.subagent_rx.recv() => {
let pre = match &sub {
SubAgentOp::Spawn { args, .. } => {
self.spawn_preprocessor.as_ref().map(|pp| pp(args))
}
_ => None,
};
if let Some(fut) = pre {
fut.await;
}
self.handle_subagent(sub);
}
}
}
self.flush_and_stop().await;
}
}
#[cfg(test)]
#[path = "../host_tests.rs"]
mod tests;