diff --git a/Cargo.lock b/Cargo.lock index 7778e359a0..d95172195d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3866,7 +3866,7 @@ dependencies = [ [[package]] name = "openfang-api" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "axum", @@ -3902,7 +3902,7 @@ dependencies = [ [[package]] name = "openfang-channels" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "axum", @@ -3933,7 +3933,7 @@ dependencies = [ [[package]] name = "openfang-cli" -version = "0.2.6" +version = "0.2.7" dependencies = [ "clap", "clap_complete", @@ -3960,7 +3960,7 @@ dependencies = [ [[package]] name = "openfang-desktop" -version = "0.2.6" +version = "0.2.7" dependencies = [ "axum", "open", @@ -3986,7 +3986,7 @@ dependencies = [ [[package]] name = "openfang-extensions" -version = "0.2.6" +version = "0.2.7" dependencies = [ "aes-gcm", "argon2", @@ -4014,7 +4014,7 @@ dependencies = [ [[package]] name = "openfang-hands" -version = "0.2.6" +version = "0.2.7" dependencies = [ "chrono", "dashmap", @@ -4031,7 +4031,7 @@ dependencies = [ [[package]] name = "openfang-kernel" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "chrono", @@ -4067,7 +4067,7 @@ dependencies = [ [[package]] name = "openfang-memory" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "chrono", @@ -4086,7 +4086,7 @@ dependencies = [ [[package]] name = "openfang-migrate" -version = "0.2.6" +version = "0.2.7" dependencies = [ "chrono", "dirs 6.0.0", @@ -4105,7 +4105,7 @@ dependencies = [ [[package]] name = "openfang-runtime" -version = "0.2.6" +version = "0.2.7" dependencies = [ "anyhow", "async-trait", @@ -4137,7 +4137,7 @@ dependencies = [ [[package]] name = "openfang-skills" -version = "0.2.6" +version = "0.2.7" dependencies = [ "chrono", "hex", @@ -4159,7 +4159,7 @@ dependencies = [ [[package]] name = "openfang-types" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "chrono", @@ -4178,7 +4178,7 @@ dependencies = [ [[package]] name = "openfang-wire" -version = "0.2.6" +version = "0.2.7" dependencies = [ "async-trait", "chrono", @@ -8790,7 +8790,7 @@ checksum = "b9cc00251562a284751c9973bace760d86c0276c471b4be569fe6b068ee97a56" [[package]] name = "xtask" -version = "0.2.6" +version = "0.2.7" [[package]] name = "yoke" diff --git a/crates/openfang-api/src/lib.rs b/crates/openfang-api/src/lib.rs index 856d443212..033841ecf8 100644 --- a/crates/openfang-api/src/lib.rs +++ b/crates/openfang-api/src/lib.rs @@ -4,6 +4,7 @@ //! The kernel runs in-process; the CLI connects over HTTP. pub mod channel_bridge; +pub mod loops; pub mod middleware; pub mod openai_compat; pub mod rate_limiter; diff --git a/crates/openfang-api/src/loops.rs b/crates/openfang-api/src/loops.rs new file mode 100644 index 0000000000..9abc53a2ae --- /dev/null +++ b/crates/openfang-api/src/loops.rs @@ -0,0 +1,199 @@ +//! LoopRegistry — persistent store for Loop records. +//! +//! A Loop ties together: +//! - a long-running agent (spawned from a manifest) +//! - a binding (cron job OR event trigger) that drives it +//! - metadata about safety limits and run history location +//! +//! Persistence mirrors the CronScheduler pattern (kernel.rs:820): JSON file under +//! `config.home_dir`, atomically rewritten on every mutation. + +use serde::{Deserialize, Serialize}; +use std::io; +use std::path::PathBuf; +use std::sync::RwLock; + +// --------------------------------------------------------------------------- +// Data types +// --------------------------------------------------------------------------- + +/// Safety guardrails embedded in each loop definition. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LoopSafety { + /// File-path prefixes the loop agent must not write to. + pub denylist_paths: Vec, + /// Maximum consecutive tool-call iterations per run. + pub max_attempts: u32, + /// Max PRs the loop may open in a calendar day. + pub daily_pr_cap: u32, + /// Approximate LLM token ceiling per run (input + output combined). + pub token_budget: u64, +} + +impl Default for LoopSafety { + fn default() -> Self { + Self { + denylist_paths: vec![], + max_attempts: 20, + daily_pr_cap: 5, + token_budget: 100_000, + } + } +} + +/// What kind of binding drives the loop. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum TriggerKind { + /// Repeating schedule — e.g. "every 3600 secs" or cron expr. + Every, + Cron, + /// File-system / process event via kernel trigger subsystem. + Event, +} + +/// Persistent record for one loop. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LoopRecord { + /// Stable UUID string assigned at creation time. + pub id: String, + /// Human-readable name / pattern tag (e.g. "nightly-review"). + pub pattern_id: String, + /// Workspace root path the loop operates on. + pub workspace: String, + /// Arbitrary level label (e.g. "low", "medium", "high") for steering. + pub level: String, + /// Discriminant for the binding type. + pub trigger_kind: TriggerKind, + /// For every: number of seconds as string; for cron: the expr; for event: JSON pattern. + pub trigger_param: String, + /// Safety configuration. + pub safety: LoopSafety, + /// Raw TOML manifest used to spawn the agent. + pub manifest_toml: String, + /// Message sent to the agent on each scheduled fire. + pub turn_prompt: String, + /// UUID string of the live agent (set after successful spawn). + pub agent_id: Option, + /// UUID string of the cron job or trigger id that drives the agent. + pub binding_id: Option, + /// Whether the loop is currently active. + pub enabled: bool, + /// RFC 3339 creation timestamp. + pub created_at: String, +} + +// --------------------------------------------------------------------------- +// Registry +// --------------------------------------------------------------------------- + +/// In-process registry backed by a JSON file on disk. +/// +/// All public methods are synchronous because mutations are cheap local +/// JSON serialisation; no async I/O needed. The RwLock is an `std` lock so +/// handlers calling `state.loop_registry.*` do not need to be async. +pub struct LoopRegistry { + /// Path to `/loops.json`. + path: PathBuf, + records: RwLock>, +} + +impl LoopRegistry { + /// Load from disk, or start with an empty list if the file is absent. + pub fn load(path: PathBuf) -> Self { + let records = if path.exists() { + std::fs::read_to_string(&path) + .ok() + .and_then(|s| serde_json::from_str::>(&s).ok()) + .unwrap_or_default() + } else { + Vec::new() + }; + Self { + path, + records: RwLock::new(records), + } + } + + /// Write current records to disk atomically (write to `.tmp`, then rename). + pub fn persist(&self) -> io::Result<()> { + let records = self + .records + .read() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "lock poisoned"))?; + let json = serde_json::to_vec_pretty(&*records) + .map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?; + // Atomic write: temp file alongside the real path, then rename. + let tmp_path = self.path.with_extension("json.tmp"); + std::fs::write(&tmp_path, &json)?; + std::fs::rename(&tmp_path, &self.path)?; + Ok(()) + } + + /// Insert a new record. Persists immediately. + pub fn insert(&self, record: LoopRecord) -> io::Result<()> { + { + let mut records = self + .records + .write() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "lock poisoned"))?; + records.push(record); + } + self.persist() + } + + /// Remove by id. Returns true if found. Persists if found. + pub fn remove(&self, id: &str) -> io::Result { + let found = { + let mut records = self + .records + .write() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "lock poisoned"))?; + let before = records.len(); + records.retain(|r| r.id != id); + records.len() < before + }; + if found { + self.persist()?; + } + Ok(found) + } + + /// Get a clone of the record matching `id`. + pub fn get(&self, id: &str) -> Option { + self.records + .read() + .ok()? + .iter() + .find(|r| r.id == id) + .cloned() + } + + /// Replace the record with matching id. Persists if found. + pub fn update(&self, id: &str, updated: LoopRecord) -> io::Result { + let found = { + let mut records = self + .records + .write() + .map_err(|_| io::Error::new(io::ErrorKind::Other, "lock poisoned"))?; + if let Some(slot) = records.iter_mut().find(|r| r.id == id) { + *slot = updated; + true + } else { + false + } + }; + if found { + self.persist()?; + } + Ok(found) + } + + /// Return a snapshot of all records. + pub fn list(&self) -> Vec { + self.records + .read() + .map(|g| g.clone()) + .unwrap_or_default() + } +} diff --git a/crates/openfang-api/src/routes.rs b/crates/openfang-api/src/routes.rs index 4c7161deaf..054b6ec30b 100644 --- a/crates/openfang-api/src/routes.rs +++ b/crates/openfang-api/src/routes.rs @@ -1,5 +1,6 @@ //! Route handlers for the OpenFang API. +use crate::loops::{LoopRegistry, LoopRecord, LoopSafety, TriggerKind}; use crate::types::*; use axum::extract::{Path, Query, State}; use axum::http::StatusCode; @@ -33,6 +34,9 @@ pub struct AppState { pub channels_config: tokio::sync::RwLock, /// Notify handle to trigger graceful HTTP server shutdown from the API. pub shutdown_notify: Arc, + /// Persistent registry of Loop records (agent + binding pairs). + /// Mirrors CronScheduler persistence: JSON file under config.home_dir (kernel.rs:820). + pub loop_registry: Arc, } /// POST /api/agents — Spawn a new agent. @@ -9664,3 +9668,811 @@ pub async fn comms_task( ), } } + +// --------------------------------------------------------------------------- +// Loop endpoints — /api/loops +// --------------------------------------------------------------------------- + +/// Request body for POST /api/loops. +#[derive(serde::Deserialize)] +pub struct CreateLoopRequest { + /// Human-readable name / pattern tag. + pub pattern_id: String, + /// Workspace root the loop operates on. + pub workspace: String, + /// Steering level label (e.g. "low", "medium", "high"). + pub level: String, + /// "every", "cron", or "event". + pub trigger_kind: String, + /// Seconds (every), cron expr (cron), or JSON-encoded TriggerPattern (event). + pub trigger_param: String, + /// Safety limits; all fields are optional with sensible defaults. + #[serde(default)] + pub safety: Option, + /// Raw TOML agent manifest. + pub manifest_toml: String, + /// Message sent on each scheduled fire / event match. + pub turn_prompt: String, +} + +/// Request body for PUT /api/loops/{id}/level. +#[derive(serde::Deserialize)] +pub struct SetLoopLevelRequest { + pub level: String, + pub manifest_toml: String, + pub turn_prompt: String, +} + +/// Request body for PUT /api/loops/{id} — full loop update. +/// +/// workspace and pattern_id are immutable; all other fields may be replaced. +#[derive(serde::Deserialize)] +pub struct UpdateLoopRequest { + pub trigger_kind: String, + pub trigger_param: String, + #[serde(default)] + pub safety: Option, + pub level: String, + pub manifest_toml: String, + pub turn_prompt: String, +} + +/// Sanitize a client-supplied pattern id into a valid cron job name fragment. +/// +/// `CronJob::validate` only permits alphanumeric characters plus space, hyphen, +/// and underscore (openfang-types/scheduler.rs). Any other character (e.g. `:`, +/// `/`) would make `cron_create` fail at runtime, so we fold them to `-`. +fn sanitize_cron_name(pattern_id: &str) -> String { + let cleaned: String = pattern_id + .chars() + .map(|c| { + if c.is_alphanumeric() || c == ' ' || c == '-' || c == '_' { + c + } else { + '-' + } + }) + .collect(); + if cleaned.is_empty() { + "unnamed".to_string() + } else { + cleaned + } +} + +/// POST /api/loops — create a new loop (spawn agent + bind schedule/event). +pub async fn create_loop( + State(state): State>, + Json(req): Json, +) -> impl IntoResponse { + // Parse manifest_toml → AgentManifest (same as spawn_agent handler) + let manifest: AgentManifest = match toml::from_str(&req.manifest_toml) { + Ok(m) => m, + Err(e) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Invalid manifest TOML: {e}")})), + ); + } + }; + + // Spawn the agent (kernel.rs:973) + let agent_id = match state.kernel.spawn_agent(manifest) { + Ok(id) => id, + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Failed to spawn agent: {e}")})), + ); + } + }; + let agent_id_str = agent_id.to_string(); + + // Determine trigger_kind + let trigger_kind_parsed = match req.trigger_kind.as_str() { + "every" => TriggerKind::Every, + "cron" => TriggerKind::Cron, + "event" => TriggerKind::Event, + other => { + // Agent was spawned; kill it before returning error to avoid orphan. + let _ = state.kernel.kill_agent(agent_id); + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Unknown trigger_kind '{other}', expected every/cron/event")})), + ); + } + }; + + // Bind: create cron job or event trigger + let binding_id: Option = match trigger_kind_parsed { + TriggerKind::Every | TriggerKind::Cron => { + // Build job_json matching the shape cron_create expects (kernel.rs:4749). + let schedule = if trigger_kind_parsed == TriggerKind::Every { + let every_secs: u64 = match req.trigger_param.parse() { + Ok(n) => n, + Err(_) => { + let _ = state.kernel.kill_agent(agent_id); + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "trigger_param must be a positive integer (seconds) for trigger_kind 'every'"})), + ); + } + }; + serde_json::json!({"kind": "every", "every_secs": every_secs}) + } else { + serde_json::json!({"kind": "cron", "expr": req.trigger_param}) + }; + + let job_json = serde_json::json!({ + "agent_id": agent_id_str, + "name": format!("loop-{}", sanitize_cron_name(&req.pattern_id)), + "schedule": schedule, + "action": { + "kind": "agent_turn", + "message": req.turn_prompt, + "timeout_secs": 300 + }, + "delivery": {"kind": "none"}, + "one_shot": false + }); + + match state.kernel.cron_create(&agent_id_str, job_json).await { + Ok(result_str) => { + // cron_create returns JSON string: {"job_id": "...", "status": "created"} (kernel.rs:4801) + let parsed: serde_json::Value = + serde_json::from_str(&result_str).unwrap_or_default(); + parsed["job_id"].as_str().map(String::from) + } + Err(e) => { + let _ = state.kernel.kill_agent(agent_id); + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Failed to create cron binding: {e}")})), + ); + } + } + } + TriggerKind::Event => { + // Parse TriggerPattern from trigger_param JSON (create_trigger handler: routes.rs:696) + let pattern: TriggerPattern = match serde_json::from_str(&req.trigger_param) { + Ok(p) => p, + Err(e) => { + let _ = state.kernel.kill_agent(agent_id); + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Invalid TriggerPattern JSON: {e}")})), + ); + } + }; + match state + .kernel + .register_trigger(agent_id, pattern, req.turn_prompt.clone(), 0) + { + Ok(trigger_id) => Some(trigger_id.to_string()), + Err(e) => { + let _ = state.kernel.kill_agent(agent_id); + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Failed to register trigger: {e}")})), + ); + } + } + } + }; + + let loop_id = uuid::Uuid::new_v4().to_string(); + let record = LoopRecord { + id: loop_id.clone(), + pattern_id: req.pattern_id, + workspace: req.workspace, + level: req.level, + trigger_kind: trigger_kind_parsed, + trigger_param: req.trigger_param, + safety: req.safety.unwrap_or_default(), + manifest_toml: req.manifest_toml, + turn_prompt: req.turn_prompt, + agent_id: Some(agent_id_str.clone()), + binding_id: binding_id.clone(), + enabled: true, + created_at: chrono::Utc::now().to_rfc3339(), + }; + + if let Err(e) = state.loop_registry.insert(record) { + tracing::warn!("Failed to persist loop record {loop_id}: {e}"); + } + + ( + StatusCode::CREATED, + Json(serde_json::json!({ + "id": loop_id, + "agent_id": agent_id_str, + "binding_id": binding_id, + })), + ) +} + +/// GET /api/loops — list all loop records. +pub async fn list_loops(State(state): State>) -> impl IntoResponse { + (StatusCode::OK, Json(serde_json::json!(state.loop_registry.list()))) +} + +/// DELETE /api/loops/{id} — remove binding, kill agent, delete record. +pub async fn delete_loop( + State(state): State>, + Path(loop_id): Path, +) -> impl IntoResponse { + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + // Remove binding first + if let Some(ref binding_id) = record.binding_id { + match record.trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + // Remove cron job (delete_cron_job pattern: routes.rs:8644) + if let Ok(uuid) = uuid::Uuid::parse_str(binding_id) { + let job_id = openfang_types::scheduler::CronJobId(uuid); + if state.kernel.cron_scheduler.remove_job(job_id).is_ok() { + let _ = state.kernel.cron_scheduler.persist(); + } + } + } + TriggerKind::Event => { + // Remove event trigger (kernel.rs:3002) + if let Ok(trigger_uuid) = uuid::Uuid::parse_str(binding_id) { + let trigger_id = openfang_kernel::triggers::TriggerId(trigger_uuid); + state.kernel.remove_trigger(trigger_id); + } + } + } + } + + // Kill the agent if still alive (kernel.rs:2645) + if let Some(ref agent_id_str) = record.agent_id { + if let Ok(agent_id) = agent_id_str.parse::() { + if let Err(e) = state.kernel.kill_agent(agent_id) { + tracing::warn!("delete_loop: kill_agent {agent_id_str} failed: {e}"); + } + } + } + + match state.loop_registry.remove(&loop_id) { + Ok(_) => (StatusCode::OK, Json(serde_json::json!({"status": "deleted"}))), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Registry persist failed: {e}")})), + ), + } +} + +/// POST /api/loops/{id}/run — send turn_prompt to the loop agent immediately. +pub async fn run_loop( + State(state): State>, + Path(loop_id): Path, +) -> impl IntoResponse { + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + let agent_id_str = match record.agent_id { + Some(ref id) => id.clone(), + None => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "Loop has no agent (not fully initialised)"})), + ); + } + }; + + let agent_id: AgentId = match agent_id_str.parse() { + Ok(id) => id, + Err(_) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": "Stored agent_id is not a valid UUID"})), + ); + } + }; + + // send_message pattern (routes.rs:263): pass Arc as KernelHandle + let kernel_handle: Arc = + state.kernel.clone() as Arc; + match state + .kernel + .send_message_with_handle(agent_id, &record.turn_prompt, Some(kernel_handle)) + .await + { + Ok(result) => ( + StatusCode::OK, + Json(serde_json::json!({ + "response": result.response, + "input_tokens": result.total_usage.input_tokens, + "output_tokens": result.total_usage.output_tokens, + "iterations": result.iterations, + })), + ), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Run failed: {e}")})), + ), + } +} + +/// PUT /api/loops/{id}/level — replace agent + rebind with new manifest/prompt/level. +pub async fn set_loop_level( + State(state): State>, + Path(loop_id): Path, + Json(req): Json, +) -> impl IntoResponse { + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + // Parse new manifest + let manifest: AgentManifest = match toml::from_str(&req.manifest_toml) { + Ok(m) => m, + Err(e) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Invalid manifest TOML: {e}")})), + ); + } + }; + + // Kill old agent + remove old binding + if let Some(ref binding_id) = record.binding_id { + match record.trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + if let Ok(uuid) = uuid::Uuid::parse_str(binding_id) { + let job_id = openfang_types::scheduler::CronJobId(uuid); + if state.kernel.cron_scheduler.remove_job(job_id).is_ok() { + let _ = state.kernel.cron_scheduler.persist(); + } + } + } + TriggerKind::Event => { + if let Ok(trigger_uuid) = uuid::Uuid::parse_str(binding_id) { + let trigger_id = openfang_kernel::triggers::TriggerId(trigger_uuid); + state.kernel.remove_trigger(trigger_id); + } + } + } + } + if let Some(ref old_agent_id_str) = record.agent_id { + if let Ok(old_agent_id) = old_agent_id_str.parse::() { + if let Err(e) = state.kernel.kill_agent(old_agent_id) { + tracing::warn!("set_loop_level: kill old agent {old_agent_id_str} failed: {e}"); + } + } + } + + // Spawn new agent + let new_agent_id = match state.kernel.spawn_agent(manifest) { + Ok(id) => id, + Err(e) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Failed to spawn new agent: {e}")})), + ); + } + }; + let new_agent_id_str = new_agent_id.to_string(); + + // Rebind using same trigger_kind + trigger_param from original record + let new_binding_id: Option = match record.trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + let schedule = if record.trigger_kind == TriggerKind::Every { + let every_secs: u64 = record.trigger_param.parse().unwrap_or(3600); + serde_json::json!({"kind": "every", "every_secs": every_secs}) + } else { + serde_json::json!({"kind": "cron", "expr": record.trigger_param}) + }; + let job_json = serde_json::json!({ + "agent_id": new_agent_id_str, + "name": format!("loop-{}", sanitize_cron_name(&record.pattern_id)), + "schedule": schedule, + "action": { + "kind": "agent_turn", + "message": req.turn_prompt, + "timeout_secs": 300 + }, + "delivery": {"kind": "none"}, + "one_shot": false + }); + match state.kernel.cron_create(&new_agent_id_str, job_json).await { + Ok(result_str) => { + let parsed: serde_json::Value = + serde_json::from_str(&result_str).unwrap_or_default(); + parsed["job_id"].as_str().map(String::from) + } + Err(e) => { + tracing::warn!("set_loop_level: cron rebind failed: {e}"); + None + } + } + } + TriggerKind::Event => { + match serde_json::from_str::(&record.trigger_param) { + Ok(pattern) => { + match state + .kernel + .register_trigger(new_agent_id, pattern, req.turn_prompt.clone(), 0) + { + Ok(tid) => Some(tid.to_string()), + Err(e) => { + tracing::warn!("set_loop_level: trigger rebind failed: {e}"); + None + } + } + } + Err(e) => { + tracing::warn!("set_loop_level: cannot parse stored trigger_param: {e}"); + None + } + } + } + }; + + let mut updated = record.clone(); + updated.level = req.level; + updated.manifest_toml = req.manifest_toml; + updated.turn_prompt = req.turn_prompt; + updated.agent_id = Some(new_agent_id_str.clone()); + updated.binding_id = new_binding_id.clone(); + + if let Err(e) = state.loop_registry.update(&loop_id, updated) { + tracing::warn!("set_loop_level: persist failed: {e}"); + } + + ( + StatusCode::OK, + Json(serde_json::json!({ + "id": loop_id, + "agent_id": new_agent_id_str, + "binding_id": new_binding_id, + })), + ) +} + +/// GET /api/loops/{id}/runs — read the loop run log from the workspace. +/// +/// Log lives at `/.antidev/loops//loop-run-log.md`. +pub async fn loop_runs( + State(state): State>, + Path(loop_id): Path, +) -> impl IntoResponse { + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + let log_path = std::path::PathBuf::from(&record.workspace) + .join(".antidev") + .join("loops") + .join(&loop_id) + .join("loop-run-log.md"); + + let content = std::fs::read_to_string(&log_path).unwrap_or_else(|_| String::from("no runs")); + + ( + StatusCode::OK, + Json(serde_json::json!({"log": content, "path": log_path.to_string_lossy()})), + ) +} + +/// POST /api/loops/{id}/enable — enable or disable the binding that drives a loop. +/// +/// Mirrors toggle_cron_job (routes.rs:8673) for cron bindings and the +/// set_trigger_enabled pattern (routes.rs:4554) for event bindings. +pub async fn enable_loop( + State(state): State>, + Path(loop_id): Path, + Json(body): Json, +) -> impl IntoResponse { + let enabled = match body.get("enabled").and_then(|v| v.as_bool()) { + Some(b) => b, + None => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "'enabled' boolean field is required"})), + ); + } + }; + + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + let binding_id = match record.binding_id { + Some(ref id) => id.clone(), + None => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "Loop has no binding (not fully initialised)"})), + ); + } + }; + + // Toggle the underlying scheduler/trigger binding. + match record.trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + // set_enabled mirrors toggle_cron_job (routes.rs:8682). + let uuid = match uuid::Uuid::parse_str(&binding_id) { + Ok(u) => u, + Err(_) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": "Stored binding_id is not a valid UUID"})), + ); + } + }; + let job_id = openfang_types::scheduler::CronJobId(uuid); + if let Err(e) = state.kernel.cron_scheduler.set_enabled(job_id, enabled) { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": format!("Cron job not found: {e}")})), + ); + } + // Persist scheduler state; non-fatal on failure (mirrors delete_loop, routes.rs:9908). + let _ = state.kernel.cron_scheduler.persist(); + } + TriggerKind::Event => { + // set_trigger_enabled mirrors routes.rs:4554. + let uuid = match uuid::Uuid::parse_str(&binding_id) { + Ok(u) => u, + Err(_) => { + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": "Stored binding_id is not a valid UUID"})), + ); + } + }; + let trigger_id = openfang_kernel::triggers::TriggerId(uuid); + if !state.kernel.set_trigger_enabled(trigger_id, enabled) { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Event trigger not found"})), + ); + } + } + } + + // Persist the updated enabled flag in the loop record. + let mut updated = record.clone(); + updated.enabled = enabled; + if let Err(e) = state.loop_registry.update(&loop_id, updated) { + tracing::warn!("enable_loop: registry persist failed: {e}"); + } + + ( + StatusCode::OK, + Json(serde_json::json!({"id": loop_id, "enabled": enabled})), + ) +} + +/// PUT /api/loops/{id} — replace trigger, safety, level, manifest, and prompt. +/// +/// Kills the old agent and removes the old binding, then spawns a fresh agent +/// and creates a new binding. workspace and pattern_id are immutable. +/// +/// Mirrors set_loop_level (routes.rs:10000) but extends it to also replace +/// trigger_kind/trigger_param and safety, not just the manifest/level/prompt. +pub async fn update_loop( + State(state): State>, + Path(loop_id): Path, + Json(req): Json, +) -> impl IntoResponse { + let record = match state.loop_registry.get(&loop_id) { + Some(r) => r, + None => { + return ( + StatusCode::NOT_FOUND, + Json(serde_json::json!({"error": "Loop not found"})), + ); + } + }; + + // Parse new manifest before touching anything that can't be undone. + let manifest: AgentManifest = match toml::from_str(&req.manifest_toml) { + Ok(m) => m, + Err(e) => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Invalid manifest TOML: {e}")})), + ); + } + }; + + // Validate new trigger_kind before mutating state. + let new_trigger_kind = match req.trigger_kind.as_str() { + "every" => TriggerKind::Every, + "cron" => TriggerKind::Cron, + "event" => TriggerKind::Event, + other => { + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": format!("Unknown trigger_kind '{other}', expected every/cron/event")})), + ); + } + }; + + // Remove OLD binding (mirrors set_loop_level, routes.rs:10027). + if let Some(ref binding_id) = record.binding_id { + match record.trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + if let Ok(uuid) = uuid::Uuid::parse_str(binding_id) { + let job_id = openfang_types::scheduler::CronJobId(uuid); + if state.kernel.cron_scheduler.remove_job(job_id).is_ok() { + let _ = state.kernel.cron_scheduler.persist(); + } + } + } + TriggerKind::Event => { + if let Ok(trigger_uuid) = uuid::Uuid::parse_str(binding_id) { + let trigger_id = openfang_kernel::triggers::TriggerId(trigger_uuid); + state.kernel.remove_trigger(trigger_id); + } + } + } + } + + // Kill OLD agent (mirrors set_loop_level, routes.rs:10045). + if let Some(ref old_agent_id_str) = record.agent_id { + if let Ok(old_agent_id) = old_agent_id_str.parse::() { + if let Err(e) = state.kernel.kill_agent(old_agent_id) { + tracing::warn!("update_loop: kill old agent {old_agent_id_str} failed: {e}"); + } + } + } + + // Spawn NEW agent. On failure: record agent_id=None so it doesn't point + // at the killed agent, then return 500. + let new_agent_id = match state.kernel.spawn_agent(manifest) { + Ok(id) => id, + Err(e) => { + // Old agent is dead; clear agent_id to avoid stale reference. + let mut cleared = record.clone(); + cleared.agent_id = None; + cleared.binding_id = None; + if let Err(persist_err) = state.loop_registry.update(&loop_id, cleared) { + tracing::warn!("update_loop: could not clear stale agent_id: {persist_err}"); + } + return ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({"error": format!("Failed to spawn new agent: {e}")})), + ); + } + }; + let new_agent_id_str = new_agent_id.to_string(); + + // Create NEW binding with the new trigger_kind/param (mirrors create_loop binding + // block, routes.rs:9773, using sanitize_cron_name for cron names). + let new_binding_id: Option = match new_trigger_kind { + TriggerKind::Every | TriggerKind::Cron => { + let schedule = if new_trigger_kind == TriggerKind::Every { + let every_secs: u64 = match req.trigger_param.parse() { + Ok(n) => n, + Err(_) => { + // New agent spawned but binding failed; record what we have. + let mut partial = record.clone(); + partial.agent_id = Some(new_agent_id_str.clone()); + partial.binding_id = None; + partial.trigger_kind = new_trigger_kind; + partial.trigger_param = req.trigger_param.clone(); + partial.safety = req.safety.unwrap_or(record.safety.clone()); + partial.level = req.level; + partial.manifest_toml = req.manifest_toml; + partial.turn_prompt = req.turn_prompt; + if let Err(e) = state.loop_registry.update(&loop_id, partial) { + tracing::warn!("update_loop: persist after bad trigger_param: {e}"); + } + return ( + StatusCode::BAD_REQUEST, + Json(serde_json::json!({"error": "trigger_param must be a positive integer (seconds) for trigger_kind 'every'"})), + ); + } + }; + serde_json::json!({"kind": "every", "every_secs": every_secs}) + } else { + serde_json::json!({"kind": "cron", "expr": req.trigger_param}) + }; + + let job_json = serde_json::json!({ + "agent_id": new_agent_id_str, + "name": format!("loop-{}", sanitize_cron_name(&record.pattern_id)), + "schedule": schedule, + "action": { + "kind": "agent_turn", + "message": req.turn_prompt, + "timeout_secs": 300 + }, + "delivery": {"kind": "none"}, + "one_shot": false + }); + match state.kernel.cron_create(&new_agent_id_str, job_json).await { + Ok(result_str) => { + let parsed: serde_json::Value = + serde_json::from_str(&result_str).unwrap_or_default(); + parsed["job_id"].as_str().map(String::from) + } + Err(e) => { + tracing::warn!("update_loop: cron rebind failed: {e}"); + None + } + } + } + TriggerKind::Event => { + match serde_json::from_str::(&req.trigger_param) { + Ok(pattern) => { + match state + .kernel + .register_trigger(new_agent_id, pattern, req.turn_prompt.clone(), 0) + { + Ok(tid) => Some(tid.to_string()), + Err(e) => { + tracing::warn!("update_loop: trigger rebind failed: {e}"); + None + } + } + } + Err(e) => { + tracing::warn!("update_loop: cannot parse trigger_param as TriggerPattern: {e}"); + None + } + } + } + }; + + // Persist updated record; keep id, pattern_id, workspace, created_at immutable. + let mut updated = record.clone(); + updated.trigger_kind = new_trigger_kind; + updated.trigger_param = req.trigger_param; + updated.safety = req.safety.unwrap_or(record.safety); + updated.level = req.level; + updated.manifest_toml = req.manifest_toml; + updated.turn_prompt = req.turn_prompt; + updated.agent_id = Some(new_agent_id_str.clone()); + updated.binding_id = new_binding_id.clone(); + + if let Err(e) = state.loop_registry.update(&loop_id, updated) { + tracing::warn!("update_loop: registry persist failed: {e}"); + } + + ( + StatusCode::OK, + Json(serde_json::json!({ + "id": loop_id, + "agent_id": new_agent_id_str, + "binding_id": new_binding_id, + })), + ) +} diff --git a/crates/openfang-api/src/server.rs b/crates/openfang-api/src/server.rs index 3df41cbfbd..3a01cfaacf 100644 --- a/crates/openfang-api/src/server.rs +++ b/crates/openfang-api/src/server.rs @@ -1,6 +1,7 @@ //! OpenFang daemon server — boots the kernel and serves the HTTP API. use crate::channel_bridge; +use crate::loops::LoopRegistry; use crate::middleware; use crate::rate_limiter; use crate::routes::{self, AppState}; @@ -42,6 +43,10 @@ pub async fn build_router( let bridge = channel_bridge::start_channel_bridge(kernel.clone()).await; let channels_config = kernel.config.channels.clone(); + // Load loop registry from /loops.json (same pattern as CronScheduler: kernel.rs:820) + let loops_path = kernel.config.home_dir.join("loops.json"); + let loop_registry = Arc::new(LoopRegistry::load(loops_path)); + let state = Arc::new(AppState { kernel: kernel.clone(), started_at: Instant::now(), @@ -49,6 +54,7 @@ pub async fn build_router( bridge_manager: tokio::sync::Mutex::new(bridge), channels_config: tokio::sync::RwLock::new(channels_config), shutdown_notify: Arc::new(tokio::sync::Notify::new()), + loop_registry, }); // CORS: allow localhost origins by default. If API key is set, the API @@ -411,6 +417,31 @@ pub async fn build_router( "/api/comms/task", axum::routing::post(routes::comms_task), ) + // Loop endpoints + .route( + "/api/loops", + axum::routing::get(routes::list_loops).post(routes::create_loop), + ) + .route( + "/api/loops/{id}", + axum::routing::delete(routes::delete_loop).put(routes::update_loop), + ) + .route( + "/api/loops/{id}/enable", + axum::routing::post(routes::enable_loop), + ) + .route( + "/api/loops/{id}/run", + axum::routing::post(routes::run_loop), + ) + .route( + "/api/loops/{id}/level", + axum::routing::put(routes::set_loop_level), + ) + .route( + "/api/loops/{id}/runs", + axum::routing::get(routes::loop_runs), + ) // Tools endpoint .route("/api/tools", axum::routing::get(routes::list_tools)) // Config endpoints diff --git a/crates/openfang-api/static/css/components.css b/crates/openfang-api/static/css/components.css index b13cc1d945..6fa6e018e9 100644 --- a/crates/openfang-api/static/css/components.css +++ b/crates/openfang-api/static/css/components.css @@ -1101,14 +1101,15 @@ mark.search-highlight { /* Theme switcher — 3-mode pill (Light / System / Dark) */ .theme-switcher { display: inline-flex; + flex-shrink: 0; border-radius: var(--radius-sm); border: 1px solid var(--border); overflow: hidden; } .theme-opt { cursor: pointer; - padding: 4px 8px; - font-size: 14px; + padding: 3px 6px; + font-size: 12px; background: none; border: none; color: var(--text-muted); diff --git a/crates/openfang-api/static/css/layout.css b/crates/openfang-api/static/css/layout.css index 4b4163192f..d587a833ee 100644 --- a/crates/openfang-api/static/css/layout.css +++ b/crates/openfang-api/static/css/layout.css @@ -29,24 +29,37 @@ .sidebar.collapsed .nav-item { justify-content: center; padding: 12px 0; } .sidebar-header { - padding: 16px; + padding: 14px 12px; border-bottom: 1px solid var(--border); display: flex; align-items: center; justify-content: space-between; - min-height: 60px; + gap: 8px; + min-height: 56px; +} + +.sidebar-header-text { + min-width: 0; + flex: 1 1 auto; } .sidebar-logo { display: flex; align-items: center; - gap: 10px; + gap: 8px; + min-width: 0; +} + +.sidebar-logo > div { + min-width: 0; + overflow: hidden; } .sidebar-logo img { - width: 28px; - height: 28px; - opacity: 0.8; + width: 22px; + height: 22px; + flex-shrink: 0; + opacity: 0.85; transition: opacity 0.2s, transform 0.2s; } @@ -56,11 +69,15 @@ } .sidebar-header h1 { - font-size: 14px; + font-size: 11px; font-weight: 700; color: var(--accent); - letter-spacing: 3px; + letter-spacing: 1.5px; font-family: var(--font-mono); + white-space: nowrap; + overflow: hidden; + text-overflow: ellipsis; + margin: 0; } .sidebar-header .version { diff --git a/crates/openfang-api/static/css/theme.css b/crates/openfang-api/static/css/theme.css index 73a9f6ec7d..d5543e5e5f 100644 --- a/crates/openfang-api/static/css/theme.css +++ b/crates/openfang-api/static/css/theme.css @@ -20,12 +20,12 @@ --text-dim: #6B6560; --text-muted: #9A958F; - /* Brand — Orange accent */ - --accent: #FF5C00; - --accent-light: #FF7A2E; - --accent-dim: #E05200; - --accent-glow: rgba(255, 92, 0, 0.1); - --accent-subtle: rgba(255, 92, 0, 0.05); + /* Brand — AntiDev Accent (Indigo) */ + --accent: #4F46E5; + --accent-light: #6366F1; + --accent-dim: #4338CA; + --accent-glow: rgba(79, 70, 229, 0.1); + --accent-subtle: rgba(79, 70, 229, 0.05); /* Status colors */ --success: #22C55E; @@ -70,7 +70,7 @@ --shadow-lg: 0 12px 28px rgba(0,0,0,0.08), 0 4px 10px rgba(0,0,0,0.05); --shadow-xl: 0 20px 40px rgba(0,0,0,0.1), 0 8px 16px rgba(0,0,0,0.06); --shadow-glow: 0 0 40px rgba(0,0,0,0.05); - --shadow-accent: 0 4px 16px rgba(255, 92, 0, 0.12); + --shadow-accent: 0 4px 16px rgba(79, 70, 229, 0.12); --shadow-inset: inset 0 1px 0 rgba(255,255,255,0.5); /* Typography — dual font system */ @@ -101,11 +101,11 @@ --text-secondary: #C4C0BC; --text-dim: #8A8380; --text-muted: #5C5754; - --accent: #FF5C00; - --accent-light: #FF7A2E; - --accent-dim: #E05200; - --accent-glow: rgba(255, 92, 0, 0.15); - --accent-subtle: rgba(255, 92, 0, 0.08); + --accent: #6366F1; + --accent-light: #818CF8; + --accent-dim: #4F46E5; + --accent-glow: rgba(99, 102, 241, 0.15); + --accent-subtle: rgba(99, 102, 241, 0.08); --success: #4ADE80; --success-dim: #22C55E; --success-subtle: rgba(74, 222, 128, 0.1); @@ -132,7 +132,7 @@ --shadow-lg: 0 12px 28px rgba(0,0,0,0.35), 0 4px 10px rgba(0,0,0,0.3); --shadow-xl: 0 20px 40px rgba(0,0,0,0.4), 0 8px 16px rgba(0,0,0,0.3); --shadow-glow: 0 0 80px rgba(0,0,0,0.6); - --shadow-accent: 0 4px 16px rgba(255, 92, 0, 0.2); + --shadow-accent: 0 4px 16px rgba(99, 102, 241, 0.2); --shadow-inset: inset 0 1px 0 rgba(255,255,255,0.03); } diff --git a/crates/openfang-api/static/favicon.ico b/crates/openfang-api/static/favicon.ico index aea6f126c2..7c7868519f 100644 Binary files a/crates/openfang-api/static/favicon.ico and b/crates/openfang-api/static/favicon.ico differ diff --git a/crates/openfang-api/static/index_body.html b/crates/openfang-api/static/index_body.html index e349446291..71998bee7f 100644 --- a/crates/openfang-api/static/index_body.html +++ b/crates/openfang-api/static/index_body.html @@ -16,9 +16,9 @@

API Key Required