diff --git a/src/core/jsonrpc.rs b/src/core/jsonrpc.rs index 9531cfe53..91aade58e 100644 --- a/src/core/jsonrpc.rs +++ b/src/core/jsonrpc.rs @@ -2322,6 +2322,14 @@ fn register_domain_subscribers( // originating chat as an idle-gated, batched, system-injected turn. crate::openhuman::agent_orchestration::background_delivery::register_background_delivery(); + // Run-ledger finalizer: detached `spawn_async_subagent` runs outlive + // their parent turn, so their terminal `AgentProgress` never reaches the + // per-turn progress bridge that settles the ledger. This global-bus + // subscriber settles `agent_runs` from `DomainEvent::Subagent{Completed, + // Failed}` (always fired from the detached task), preventing rows from + // leaking as perpetual `running` timeline entries on thread reopen. + crate::openhuman::agent_orchestration::run_ledger_finalize::register_run_ledger_finalize_subscriber(&config); + // MCP clients lifecycle subscriber: logs McpServer{Installed,Connected, // Disconnected} + McpClientToolExecuted for observability. The boot-time // spawn of installed servers (boot::spawn_installed_servers) runs later @@ -2392,6 +2400,18 @@ pub async fn bootstrap_core_runtime(host_kind: crate::core::types::HostKind) { } } + // --- Run-ledger recovery ------------------------------------------- + // Detached sub-agent runs (`spawn_async_subagent`) from a previous process + // are gone with that process. Any `agent_runs` row still marked `running` + // at boot is orphaned — its driver died without firing a terminal event, so + // the finalizer never settled it. Stamp such rows `interrupted` so they stop + // rendering as perpetual "running" timeline entries on thread reopen. + match crate::openhuman::session_db::run_ledger::interrupt_orphaned_agent_runs(&cfg) { + Ok(0) => {} + Ok(count) => log::info!("[runtime] settled {count} orphaned agent run(s) on startup"), + Err(err) => log::warn!("[runtime] failed to settle orphaned agent runs: {err}"), + } + // --- Cost dashboard tracker --- // Activates the previously-dormant CostTracker so the dashboard RPC // surface (`openhuman.cost_get_dashboard`) and `record_provider_usage` diff --git a/src/openhuman/agent_orchestration/mod.rs b/src/openhuman/agent_orchestration/mod.rs index 3444abe7f..cf43dffe6 100644 --- a/src/openhuman/agent_orchestration/mod.rs +++ b/src/openhuman/agent_orchestration/mod.rs @@ -11,6 +11,7 @@ pub mod background_delivery; pub mod command_center; mod ops; pub(crate) mod parent_context; +pub mod run_ledger_finalize; pub mod running_subagents; pub mod subagent_control; pub mod subagent_sessions; diff --git a/src/openhuman/agent_orchestration/run_ledger_finalize.rs b/src/openhuman/agent_orchestration/run_ledger_finalize.rs new file mode 100644 index 000000000..c8bbe98f5 --- /dev/null +++ b/src/openhuman/agent_orchestration/run_ledger_finalize.rs @@ -0,0 +1,117 @@ +//! Global-bus finalizer that settles subagent rows in the run ledger. +//! +//! # Why this exists +//! +//! Every spawn path (`spawn_subagent`, `spawn_async_subagent`, +//! `spawn_parallel_agents`, `continue_subagent`, `dispatch`) creates an +//! `agent_runs` row at `status = "running"` and, on completion, fires **both**: +//! * `publish_global(DomainEvent::SubagentCompleted/Failed)` — the global bus, +//! * `progress_sink.send(AgentProgress::SubagentCompleted/Failed)` — the +//! *spawning turn's* progress channel. +//! +//! The run-ledger's terminal transition used to live **only** in the per-turn +//! [`progress_bridge`](crate::openhuman::channels::providers::web::progress_bridge), +//! which consumes that progress channel. That works for synchronous subagents +//! (they finish inside the turn, sink still alive) but **leaks** for detached +//! `spawn_async_subagent` runs: they outlive the parent turn, so when they +//! finish the progress sink is already dropped and `let _ = tx.send(...)` fails +//! silently. The ledger row stays `running` forever, and every thread reopen +//! re-renders it as a perpetual "Tinyplace Agent" timeline row. +//! +//! This subscriber closes that gap by finalizing the ledger from the **global +//! bus**, which always fires from the detached task regardless of the parent +//! turn's lifecycle. It is idempotent with the progress-bridge path: a run that +//! was already settled there is simply re-stamped to the same terminal status. + +use std::sync::Arc; + +use async_trait::async_trait; + +use crate::core::event_bus::{subscribe_global, DomainEvent, EventHandler}; +use crate::openhuman::config::Config; +use crate::openhuman::session_db::run_ledger::{transition_agent_run_status, AgentRunStatus}; + +const LOG_PREFIX: &str = "[run_ledger][finalize]"; + +/// Subscribes to terminal subagent [`DomainEvent`]s and transitions the +/// matching `agent_runs` row to a terminal status. Holds its own [`Config`] +/// clone so it can reach the session DB from the global-bus context (handlers +/// receive no config). +struct RunLedgerFinalizeSubscriber { + config: Config, +} + +#[async_trait] +impl EventHandler for RunLedgerFinalizeSubscriber { + fn name(&self) -> &str { + "agent_orchestration::run_ledger_finalize" + } + + async fn handle(&self, event: &DomainEvent) { + let (task_id, status, error) = match event { + DomainEvent::SubagentCompleted { task_id, .. } => { + (task_id.clone(), AgentRunStatus::Completed, None) + } + DomainEvent::SubagentFailed { task_id, error, .. } => { + (task_id.clone(), AgentRunStatus::Failed, Some(error.clone())) + } + _ => return, + }; + + // Single fast UPDATE, but keep it off the async runtime to honour the + // EventHandler "must not block" contract. + let config = self.config.clone(); + let completed_at = chrono::Utc::now(); + let result = tokio::task::spawn_blocking(move || { + transition_agent_run_status( + &config, + &task_id, + status, + error.as_deref(), + Some(completed_at), + ) + .map(|run| (task_id, run)) + }) + .await; + + match result { + Ok(Ok((task_id, Some(_)))) => { + log::debug!( + "{LOG_PREFIX} settled run id={task_id} status={}", + status.as_str() + ); + } + Ok(Ok((task_id, None))) => { + // No row matched — the spawn upsert never reached the ledger, or + // the run was deleted. Nothing to settle. + log::debug!("{LOG_PREFIX} no run row to settle id={task_id}"); + } + Ok(Err(err)) => { + log::warn!("{LOG_PREFIX} failed to settle run: {err}"); + } + Err(join_err) => { + log::warn!("{LOG_PREFIX} settle task panicked: {join_err}"); + } + } + } +} + +/// Register the run-ledger finalizer on the global event bus. Leaks the +/// subscription handle so it lives for the whole process (its `Drop` would +/// cancel the subscriber). Called once from `register_domain_subscribers`. +pub fn register_run_ledger_finalize_subscriber(config: &Config) { + if let Some(handle) = subscribe_global(Arc::new(RunLedgerFinalizeSubscriber { + config: config.clone(), + })) { + std::mem::forget(handle); + log::info!("{LOG_PREFIX} run-ledger finalize subscriber registered"); + } else { + log::warn!( + "{LOG_PREFIX} failed to register run-ledger finalize subscriber — bus not initialized" + ); + } +} + +#[cfg(test)] +#[path = "run_ledger_finalize_tests.rs"] +mod tests; diff --git a/src/openhuman/agent_orchestration/run_ledger_finalize_tests.rs b/src/openhuman/agent_orchestration/run_ledger_finalize_tests.rs new file mode 100644 index 000000000..cdb5667bb --- /dev/null +++ b/src/openhuman/agent_orchestration/run_ledger_finalize_tests.rs @@ -0,0 +1,115 @@ +//! Tests for the global-bus run-ledger finalizer. + +use super::*; + +use serde_json::json; +use tempfile::TempDir; + +use crate::core::event_bus::EventHandler; +use crate::openhuman::session_db::run_ledger::{ + get_agent_run, upsert_agent_run, AgentRunKind, AgentRunStatus, AgentRunUpsert, +}; + +fn test_config(dir: &TempDir) -> Config { + let mut config = Config::default(); + config.workspace_dir = dir.path().to_path_buf(); + config.action_dir = dir.path().join("actions"); + config +} + +fn seed_running(config: &Config, id: &str) { + upsert_agent_run( + config, + AgentRunUpsert { + id: id.into(), + kind: AgentRunKind::Subagent, + parent_run_id: Some("parent-turn".into()), + parent_thread_id: Some("thread-1".into()), + agent_id: Some("tinyplace_agent".into()), + status: AgentRunStatus::Running, + prompt_ref: None, + worker_thread_id: None, + task_board_id: None, + task_card_id: None, + checkpoint_path: None, + checkpoint: None, + summary: None, + error: None, + metadata: json!({}), + started_at: None, + completed_at: None, + }, + ) + .unwrap(); +} + +#[tokio::test] +async fn settles_running_run_on_subagent_completed() { + let dir = TempDir::new().unwrap(); + let config = test_config(&dir); + seed_running(&config, "sub-1"); + + let sub = RunLedgerFinalizeSubscriber { + config: config.clone(), + }; + sub.handle(&DomainEvent::SubagentCompleted { + parent_session: "session-1".into(), + task_id: "sub-1".into(), + agent_id: "tinyplace_agent".into(), + elapsed_ms: 1234, + output_chars: 760, + iterations: 3, + }) + .await; + + let run = get_agent_run(&config, "sub-1") + .unwrap() + .expect("run present"); + assert_eq!(run.status, AgentRunStatus::Completed); + assert!(run.completed_at.is_some()); +} + +#[tokio::test] +async fn settles_running_run_on_subagent_failed_with_error() { + let dir = TempDir::new().unwrap(); + let config = test_config(&dir); + seed_running(&config, "sub-2"); + + let sub = RunLedgerFinalizeSubscriber { + config: config.clone(), + }; + sub.handle(&DomainEvent::SubagentFailed { + parent_session: "session-1".into(), + task_id: "sub-2".into(), + agent_id: "tinyplace_agent".into(), + error: "boom".into(), + }) + .await; + + let run = get_agent_run(&config, "sub-2") + .unwrap() + .expect("run present"); + assert_eq!(run.status, AgentRunStatus::Failed); + assert_eq!(run.error.as_deref(), Some("boom")); +} + +#[tokio::test] +async fn ignores_unrelated_events_and_missing_runs() { + let dir = TempDir::new().unwrap(); + let config = test_config(&dir); + + let sub = RunLedgerFinalizeSubscriber { + config: config.clone(), + }; + // Completion for a run that was never recorded — no-op, no panic. + sub.handle(&DomainEvent::SubagentCompleted { + parent_session: "session-1".into(), + task_id: "ghost".into(), + agent_id: "tinyplace_agent".into(), + elapsed_ms: 1, + output_chars: 0, + iterations: 1, + }) + .await; + assert!(get_agent_run(&config, "ghost").unwrap().is_none()); +} diff --git a/src/openhuman/session_db/run_ledger/mod.rs b/src/openhuman/session_db/run_ledger/mod.rs index 9dc05487f..310bbc36d 100644 --- a/src/openhuman/session_db/run_ledger/mod.rs +++ b/src/openhuman/session_db/run_ledger/mod.rs @@ -11,12 +11,12 @@ pub mod types; pub use ops::{ append_run_event, claim_agent_team_task, complete_agent_team_task, get_agent_run, - get_agent_team, get_agent_team_member, get_agent_team_task, get_workflow_run, list_agent_runs, - list_agent_team_members, list_agent_team_tasks, list_agent_teams, list_recent_run_events, - list_workflow_runs, mark_agent_team_member_idle, mark_agent_team_member_running, - release_agent_team_task, shutdown_agent_team_member, transition_agent_run_status, - upsert_agent_run, upsert_agent_team, upsert_agent_team_member, upsert_agent_team_task, - upsert_run_telemetry, upsert_workflow_run, + get_agent_team, get_agent_team_member, get_agent_team_task, get_workflow_run, + interrupt_orphaned_agent_runs, list_agent_runs, list_agent_team_members, list_agent_team_tasks, + list_agent_teams, list_recent_run_events, list_workflow_runs, mark_agent_team_member_idle, + mark_agent_team_member_running, release_agent_team_task, shutdown_agent_team_member, + transition_agent_run_status, upsert_agent_run, upsert_agent_team, upsert_agent_team_member, + upsert_agent_team_task, upsert_run_telemetry, upsert_workflow_run, }; pub use types::{ AgentRun, AgentRunKind, AgentRunListRequest, AgentRunListResponse, AgentRunStatus, diff --git a/src/openhuman/session_db/run_ledger/ops.rs b/src/openhuman/session_db/run_ledger/ops.rs index f4d6f3274..0b07c816c 100644 --- a/src/openhuman/session_db/run_ledger/ops.rs +++ b/src/openhuman/session_db/run_ledger/ops.rs @@ -274,6 +274,41 @@ pub fn transition_agent_run_status( }) } +/// Settle non-terminal `agent_runs` rows left behind by a previous process. +/// +/// A freshly-booted core has no in-flight subagents — any detached run task +/// from a prior process is gone with that process. So a row still marked +/// `running` (or `pending`) at startup is, by definition, orphaned: its driver +/// died without firing a terminal `DomainEvent::Subagent{Completed,Failed}`, so +/// the [`register_run_ledger_finalize_subscriber`] never settled it. Without +/// this sweep those rows render as perpetual "running" timeline entries on every +/// thread reopen. +/// +/// We stamp them `interrupted` (outcome unknown — mirrors the turn-state +/// `mark_all_interrupted` recovery) and set `completed_at`. `awaiting_user` / +/// `paused` are intentionally left untouched: those are resumable states a user +/// may still continue. +/// +/// [`register_run_ledger_finalize_subscriber`]: crate::openhuman::agent_orchestration::run_ledger_finalize::register_run_ledger_finalize_subscriber +pub fn interrupt_orphaned_agent_runs(config: &Config) -> Result { + let now = Utc::now(); + crate::openhuman::session_db::store::with_connection(config, |conn| { + init_run_ledger_schema(conn)?; + let rows_affected = conn + .execute( + "UPDATE agent_runs + SET status = ?1, completed_at = COALESCE(completed_at, ?2), updated_at = ?2 + WHERE status IN ('running', 'pending')", + params![AgentRunStatus::Interrupted.as_str(), now.to_rfc3339()], + ) + .context("interrupt orphaned agent runs")?; + if rows_affected > 0 { + log::info!("{LOG_PREFIX} interrupted {rows_affected} orphaned agent run(s) on startup"); + } + Ok(rows_affected) + }) +} + pub fn list_agent_runs( config: &Config, request: &AgentRunListRequest, @@ -1824,4 +1859,57 @@ mod tests { let teams = list_agent_teams(&config, &AgentTeamListRequest::default()).unwrap(); assert_eq!(teams.count, 1); } + + fn seed_run(config: &Config, id: &str, status: AgentRunStatus) { + upsert_agent_run( + config, + AgentRunUpsert { + id: id.into(), + kind: super::super::types::AgentRunKind::Subagent, + parent_run_id: None, + parent_thread_id: Some("thread-1".into()), + agent_id: Some("tinyplace_agent".into()), + status, + prompt_ref: None, + worker_thread_id: None, + task_board_id: None, + task_card_id: None, + checkpoint_path: None, + checkpoint: None, + summary: None, + error: None, + metadata: json!({}), + started_at: None, + completed_at: None, + }, + ) + .unwrap(); + } + + #[test] + fn interrupt_orphaned_runs_settles_only_non_terminal_inflight_rows() { + let dir = TempDir::new().unwrap(); + let config = test_config(&dir); + + seed_run(&config, "run-running", AgentRunStatus::Running); + seed_run(&config, "run-pending", AgentRunStatus::Pending); + seed_run(&config, "run-completed", AgentRunStatus::Completed); + seed_run(&config, "run-awaiting", AgentRunStatus::AwaitingUser); + + let settled = interrupt_orphaned_agent_runs(&config).unwrap(); + assert_eq!(settled, 2, "only running + pending are orphaned at boot"); + + let get = |id: &str| get_agent_run(&config, id).unwrap().expect("run present"); + // Orphaned in-flight rows become terminal `interrupted` with a completion time… + let running = get("run-running"); + assert_eq!(running.status, AgentRunStatus::Interrupted); + assert!(running.completed_at.is_some()); + assert_eq!(get("run-pending").status, AgentRunStatus::Interrupted); + // …already-terminal and resumable rows are untouched. + assert_eq!(get("run-completed").status, AgentRunStatus::Completed); + assert_eq!(get("run-awaiting").status, AgentRunStatus::AwaitingUser); + + // Idempotent: a second sweep finds nothing left to settle. + assert_eq!(interrupt_orphaned_agent_runs(&config).unwrap(), 0); + } }