diff --git a/crates/agentflare-backend/src/item.rs b/crates/agentflare-backend/src/item.rs index 6c7abb9..9f4c715 100644 --- a/crates/agentflare-backend/src/item.rs +++ b/crates/agentflare-backend/src/item.rs @@ -233,6 +233,18 @@ pub fn list_by_project(conn: &Connection, project_id: &str) -> Result> Ok(rows.collect::>()?) } +pub fn list_by_label(conn: &Connection, project_id: &str, label_id: &str) -> Result> { + let mut stmt = conn.prepare( + "SELECT items.id, items.project_id, items.state_id, items.name, items.description, items.priority, items.parent_id, items.assignee_agent, items.sequence_id, items.sort_order, items.started_at, items.completed_at, items.archived_at, items.external_source, items.external_id, items.metadata, items.created_at, items.updated_at, items.deleted_at + FROM items + INNER JOIN item_labels ON item_labels.item_id = items.id + WHERE item_labels.label_id = ?1 AND items.project_id = ?2 AND items.deleted_at IS NULL + ORDER BY items.sort_order", + )?; + let rows = stmt.query_map(rusqlite::params![label_id, project_id], row_to_item)?; + Ok(rows.collect::>()?) +} + /// List non-deleted items assigned to an agent (excludes completed/cancelled). pub fn list_by_assignee_agent( conn: &Connection, @@ -1900,4 +1912,71 @@ mod tests { .unwrap(); assert_eq!(updated.assignee_agent.as_deref(), Some("claude-code")); } + + #[test] + fn list_by_label_returns_only_items_carrying_that_label() { + let conn = db::open_in_memory().unwrap(); + let (pid, sid) = seed_project(&conn, "label"); + let ws_id = crate::project::get(&conn, &pid).unwrap().workspace_id; + let label = crate::label::create( + &conn, + crate::label::CreateLabel { + project_id: Some(pid.clone()), + workspace_id: ws_id, + name: "ready-for-work".into(), + color: None, + parent_id: None, + sort_order: None, + external_source: None, + external_id: None, + }, + ) + .unwrap(); + + let labeled = create( + &conn, + CreateItem { + project_id: pid.clone(), + state_id: sid.clone(), + name: "Labeled".into(), + description: None, + priority: None, + parent_id: None, + assignee_agent: None, + sort_order: None, + external_source: None, + external_id: None, + metadata: None, + label_ids: vec![], + assignee_ids: vec![], + dependency_ids: vec![], + }, + ) + .unwrap(); + create( + &conn, + CreateItem { + project_id: pid.clone(), + state_id: sid, + name: "Unlabeled".into(), + description: None, + priority: None, + parent_id: None, + assignee_agent: None, + sort_order: None, + external_source: None, + external_id: None, + metadata: None, + label_ids: vec![], + assignee_ids: vec![], + dependency_ids: vec![], + }, + ) + .unwrap(); + add_label(&conn, &labeled.id, &label.id).unwrap(); + + let found = list_by_label(&conn, &pid, &label.id).unwrap(); + assert_eq!(found.len(), 1); + assert_eq!(found[0].id, labeled.id); + } } diff --git a/src/dashboard/server.rs b/src/dashboard/server.rs index f99be45..5508eb8 100644 --- a/src/dashboard/server.rs +++ b/src/dashboard/server.rs @@ -142,6 +142,7 @@ const COST_REFRESH: std::time::Duration = std::time::Duration::from_secs(30); /// nothing else ever calls `Queue::cleanup`. const JOB_CLEANUP_INTERVAL: std::time::Duration = std::time::Duration::from_secs(3600); const JOB_RETENTION_SECS: i64 = 7 * 24 * 3600; +const SUPERVISOR_DISCOVERY_INTERVAL: std::time::Duration = std::time::Duration::from_secs(12); /// Runs for the lifetime of the process (like `snapshot_broadcaster`'s /// producer task): wakes on `JOB_CLEANUP_INTERVAL`, deletes finished jobs @@ -165,6 +166,39 @@ fn spawn_job_cleanup(queue: agentflare_jobs::Queue) { }); } +/// Runs for the lifetime of the process: wakes on `SUPERVISOR_DISCOVERY_INTERVAL`, +/// lists items labeled `ready-for-work` and dispatches an `agentflare work` job +/// for each confirmed-autonomous assignee. The blocking SQLite work happens off +/// the async worker threads via `spawn_blocking`, same as `spawn_job_cleanup`. +fn spawn_supervisor_discovery( + queue: agentflare_jobs::Queue, + mcp: std::sync::Arc, + interval: std::time::Duration, +) { + tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + loop { + ticker.tick().await; + let queue = queue.clone(); + let mcp = mcp.clone(); + let result = tokio::task::spawn_blocking(move || { + crate::supervisor::run_discovery_tick(&mcp, &queue) + }) + .await; + match result { + Ok(summary) if summary.dispatched > 0 || summary.skipped > 0 => { + eprintln!( + "agentflare-supervisor: dispatched {}, skipped {}", + summary.dispatched, summary.skipped + ); + } + Ok(_) => {} + Err(e) => eprintln!("agentflare-supervisor: tick task panicked: {e}"), + } + } + }); +} + /// Single shared broadcast of the live `{ claims, cost_today }` snapshot. Every /// `/events` client subscribes to this one channel, so there is no per-client /// work — the producer below runs at most one refresh cycle per `TICK`, @@ -557,6 +591,11 @@ pub async fn run(host: &str, port: u16, open: bool, yes_expose: bool) { let mut worker_pool = agentflare_jobs::WorkerPool::new(queue.clone()); worker_pool.start(2); spawn_job_cleanup(queue.clone()); + spawn_supervisor_discovery( + queue.clone(), + std::sync::Arc::new(crate::mcp_server::AgentflareMcp::default()), + SUPERVISOR_DISCOVERY_INTERVAL, + ); let listener = tokio::net::TcpListener::bind((host, port)) .await @@ -1044,4 +1083,14 @@ mod tests { "expected dashboard to serve on a local bind without --yes-expose" ); } + + #[tokio::test] + async fn spawn_supervisor_discovery_runs_without_panicking_on_an_empty_project() { + let dir = tempfile::tempdir().unwrap().keep(); + let queue = agentflare_jobs::Queue::open_memory(dir.join("logs")).unwrap(); + let mcp = std::sync::Arc::new(crate::mcp_server::AgentflareMcp::for_test_memory()); + + spawn_supervisor_discovery(queue, mcp, std::time::Duration::from_millis(20)); + tokio::time::sleep(std::time::Duration::from_millis(80)).await; + } } diff --git a/src/main.rs b/src/main.rs index e22475e..065a795 100644 --- a/src/main.rs +++ b/src/main.rs @@ -57,6 +57,7 @@ mod skill_detect; mod skill_proactive; mod state; mod store; +mod supervisor; mod tool_install; mod ui; mod uninstall; diff --git a/src/mcp_server.rs b/src/mcp_server.rs index 2cc27a9..38a65d0 100644 --- a/src/mcp_server.rs +++ b/src/mcp_server.rs @@ -82,7 +82,7 @@ pub struct AgentflareMcp { /// its own source of truth, so nothing to refresh. backend_db: std::sync::Mutex>, /// Tests inject a temp path here so they never touch the shared backend.db. - backend_db_override: Option, + pub(crate) backend_db_override: Option, /// Tests inject a temp file path here so project-link resolution never /// reads/writes this actual repo's `.agentflare/project.json`. backend_project_link_override: Option, @@ -580,6 +580,17 @@ impl AgentflareMcp { } } + /// Create an in-memory backend for tests that don't need a real repo/disk. + /// The other fields (skills, gateway, store, etc.) stay defaulted so they + /// lazily open `:memory:` SQLite connections or no-ops — no I/O to ~/. + #[cfg(test)] + pub(crate) fn for_test_memory() -> Self { + Self { + backend_db_override: Some(std::path::PathBuf::from(":memory:")), + ..Default::default() + } + } + /// Pure walk-up so the non-git fallback path is unit-testable without /// touching process-global state: neither this process's real cwd nor /// `crate::paths::home()` (which itself reads the `AGENTFLARE_HOME_OVERRIDE` diff --git a/src/mcp_server/item.rs b/src/mcp_server/item.rs index 7db5ef3..d1685b4 100644 --- a/src/mcp_server/item.rs +++ b/src/mcp_server/item.rs @@ -610,7 +610,7 @@ impl AgentflareMcp { })? } - pub(super) fn item_add_label(&self, req: ItemRequest) -> Result { + pub(crate) fn item_add_label(&self, req: ItemRequest) -> Result { let raw = req .id .ok_or_else(|| ErrorData::invalid_params("id is required for add_label", None))?; @@ -634,7 +634,7 @@ impl AgentflareMcp { })? } - pub(super) fn item_remove_label(&self, req: ItemRequest) -> Result { + pub(crate) fn item_remove_label(&self, req: ItemRequest) -> Result { let raw = req .id .ok_or_else(|| ErrorData::invalid_params("id is required for remove_label", None))?; diff --git a/src/supervisor.rs b/src/supervisor.rs new file mode 100644 index 0000000..b569892 --- /dev/null +++ b/src/supervisor.rs @@ -0,0 +1,294 @@ +//! Background discovery loop: finds items labeled `ready-for-work` and +//! dispatches an `agentflare work` job for each one whose assignee is a +//! confirmed-autonomous agent (skips the rest with a comment). + +use crate::mcp_server::AgentflareMcp; +use crate::mcp_server::types::{CommentRequest, ItemRequest}; + +const READY_LABEL: &str = "ready-for-work"; +const DISPATCHED_LABEL: &str = "dispatched"; +const NEEDS_MANUAL_LABEL: &str = "needs-manual-dispatch"; + +/// Returns the matching `Agent` only if `agent_registry::autonomous_args` +/// confirms it has a headless permission-bypass flag — the same gate +/// `agentflare work` itself uses (`src/cli/work.rs`'s `run_work`). +pub(crate) fn resolve_confirmed_agent(assignee: &str) -> Option { + let agent = agent_registry::REGISTRY + .iter() + .find(|s| s.id.as_str() == assignee) + .map(|s| s.id)?; + agent_registry::autonomous_args(agent).map(|_| agent) +} + +pub(crate) struct DiscoveryTickResult { + pub dispatched: usize, + pub skipped: usize, +} + +/// One pass: list items labeled `ready-for-work`, dispatch a job for each +/// one with a confirmed-autonomous assignee, skip (+ comment + relabel) the +/// rest. Ends after enqueueing — it does not watch job completion, since +/// `agentflare work` itself reports outcome back onto the item. +pub(crate) fn run_discovery_tick( + mcp: &AgentflareMcp, + queue: &agentflare_jobs::Queue, +) -> DiscoveryTickResult { + let mut result = DiscoveryTickResult { + dispatched: 0, + skipped: 0, + }; + + let fetched = mcp.with_backend_db(|conn| { + let project = mcp.resolve_project(conn).ok()?; + let labels = agentflare_backend::label::list_by_project(conn, &project.id).ok()?; + let mut label_id_by_name = std::collections::HashMap::new(); + for l in &labels { + label_id_by_name.insert(l.name.clone(), l.id.clone()); + } + let ready_id = label_id_by_name.get(READY_LABEL)?.clone(); + let items = agentflare_backend::item::list_by_label(conn, &project.id, &ready_id).ok()?; + Some((items, label_id_by_name)) + }); + + let Ok(Some((items, label_id_by_name))) = fetched else { + return result; + }; + let Some(ready_id) = label_id_by_name.get(READY_LABEL).cloned() else { + return result; + }; + + for item in items { + match item + .assignee_agent + .as_deref() + .and_then(resolve_confirmed_agent) + { + None => { + skip_item(mcp, &item, &label_id_by_name, &ready_id); + result.skipped += 1; + } + Some(agent) => { + if dispatch_item(mcp, queue, &item, agent, &label_id_by_name, &ready_id) { + result.dispatched += 1; + } + } + } + } + result +} + +fn skip_item( + mcp: &AgentflareMcp, + item: &agentflare_backend::item::Item, + label_id_by_name: &std::collections::HashMap, + ready_id: &str, +) { + let reason = match &item.assignee_agent { + None => "no assignee_agent set — cannot auto-dispatch".to_string(), + Some(a) => format!("assignee '{a}' is not a confirmed-autonomous agent"), + }; + let _ = mcp.comment_impl(CommentRequest { + action: "create".into(), + item_id: Some(item.id.clone()), + body: Some(format!( + "## supervisor — skipped\n\n{reason}. Run `agentflare work` manually." + )), + ..Default::default() + }); + let _ = mcp.item_remove_label(ItemRequest { + action: "remove_label".into(), + id: Some(item.id.clone()), + label_id: Some(ready_id.to_string()), + ..Default::default() + }); + if let Some(needs_manual_id) = label_id_by_name.get(NEEDS_MANUAL_LABEL) { + let _ = mcp.item_add_label(ItemRequest { + action: "add_label".into(), + id: Some(item.id.clone()), + label_id: Some(needs_manual_id.clone()), + ..Default::default() + }); + } +} + +fn dispatch_item( + mcp: &AgentflareMcp, + queue: &agentflare_jobs::Queue, + item: &agentflare_backend::item::Item, + agent: agent_registry::Agent, + label_id_by_name: &std::collections::HashMap, + ready_id: &str, +) -> bool { + let command = std::env::current_exe() + .map(|p| p.to_string_lossy().to_string()) + .unwrap_or_else(|_| "agentflare".to_string()); + let job = agentflare_jobs::AgentJob::new(command).args([ + "work".to_string(), + item.id.clone(), + "--agent".to_string(), + agent.as_str().to_string(), + ]); + let Ok(info) = queue.enqueue(&job) else { + return false; + }; + + let _ = mcp.item_remove_label(ItemRequest { + action: "remove_label".into(), + id: Some(item.id.clone()), + label_id: Some(ready_id.to_string()), + ..Default::default() + }); + if let Some(dispatched_id) = label_id_by_name.get(DISPATCHED_LABEL) { + let _ = mcp.item_add_label(ItemRequest { + action: "add_label".into(), + id: Some(item.id.clone()), + label_id: Some(dispatched_id.clone()), + ..Default::default() + }); + } + let _ = mcp.comment_impl(CommentRequest { + action: "create".into(), + item_id: Some(item.id.clone()), + body: Some(format!("## supervisor — dispatched\n\njob: {}", info.id)), + ..Default::default() + }); + true +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn resolve_confirmed_agent_accepts_claude_code() { + assert_eq!( + resolve_confirmed_agent("claude-code"), + Some(agent_registry::Agent::ClaudeCode) + ); + } + + #[test] + fn resolve_confirmed_agent_rejects_opencode() { + assert_eq!(resolve_confirmed_agent("opencode"), None); + } + + #[test] + fn resolve_confirmed_agent_rejects_unknown_agent_string() { + assert_eq!(resolve_confirmed_agent("not-a-real-agent"), None); + } + + fn test_mcp() -> AgentflareMcp { + AgentflareMcp::for_test_memory() + } + + fn test_queue() -> agentflare_jobs::Queue { + let dir = tempfile::tempdir().unwrap().keep(); + agentflare_jobs::Queue::open_memory(dir.join("logs")).unwrap() + } + + fn seed_ready_item(mcp: &AgentflareMcp, assignee: Option<&str>) -> String { + mcp.with_backend_db(|conn| { + let project = mcp.resolve_project(conn).unwrap(); + for name in ["ready-for-work", "dispatched", "needs-manual-dispatch"] { + agentflare_backend::label::create( + conn, + agentflare_backend::label::CreateLabel { + project_id: Some(project.id.clone()), + workspace_id: project.workspace_id.clone(), + name: name.into(), + color: None, + parent_id: None, + sort_order: None, + external_source: None, + external_id: None, + }, + ) + .unwrap(); + } + let states = agentflare_backend::state::list_by_project(conn, &project.id).unwrap(); + let state_id = states.iter().find(|s| s.is_default).unwrap().id.clone(); + let item = agentflare_backend::item::create( + conn, + agentflare_backend::item::CreateItem { + project_id: project.id.clone(), + state_id, + name: "Do the thing".into(), + description: Some("do it well".into()), + priority: None, + parent_id: None, + assignee_agent: assignee.map(str::to_string), + sort_order: None, + external_source: None, + external_id: None, + metadata: None, + label_ids: vec![], + assignee_ids: vec![], + dependency_ids: vec![], + }, + ) + .unwrap(); + let labels = agentflare_backend::label::list_by_project(conn, &project.id).unwrap(); + let ready_id = &labels + .iter() + .find(|l| l.name == "ready-for-work") + .unwrap() + .id; + agentflare_backend::item::add_label(conn, &item.id, ready_id).unwrap(); + item.id + }) + .unwrap() + } + + fn labels_contain_name(mcp: &AgentflareMcp, label_ids: &[String], name: &str) -> bool { + mcp.with_backend_db(|conn| { + let project = mcp.resolve_project(conn).unwrap(); + let all = agentflare_backend::label::list_by_project(conn, &project.id).unwrap(); + let target = all.iter().find(|l| l.name == name).unwrap(); + label_ids.contains(&target.id) + }) + .unwrap() + } + + #[test] + fn confirmed_agent_gets_dispatched_and_relabeled() { + let mcp = test_mcp(); + let queue = test_queue(); + let item_id = seed_ready_item(&mcp, Some("claude-code")); + + let result = run_discovery_tick(&mcp, &queue); + + assert_eq!(result.dispatched, 1); + assert_eq!(result.skipped, 0); + + let jobs = queue.list(None).unwrap(); + assert_eq!(jobs.len(), 1); + assert!(jobs[0].args.contains(&"work".to_string())); + assert!(jobs[0].args.contains(&item_id)); + assert!(jobs[0].args.contains(&"claude-code".to_string())); + + let labels = mcp + .with_backend_db(|conn| agentflare_backend::item::list_labels(conn, &item_id).unwrap()) + .unwrap(); + assert!(!labels_contain_name(&mcp, &labels, "ready-for-work")); + assert!(labels_contain_name(&mcp, &labels, "dispatched")); + } + + #[test] + fn unconfirmed_agent_gets_skipped_not_dispatched() { + let mcp = test_mcp(); + let queue = test_queue(); + let item_id = seed_ready_item(&mcp, Some("opencode")); + + let result = run_discovery_tick(&mcp, &queue); + + assert_eq!(result.dispatched, 0); + assert_eq!(result.skipped, 1); + assert!(queue.list(None).unwrap().is_empty()); + + let labels = mcp + .with_backend_db(|conn| agentflare_backend::item::list_labels(conn, &item_id).unwrap()) + .unwrap(); + assert!(!labels_contain_name(&mcp, &labels, "ready-for-work")); + assert!(labels_contain_name(&mcp, &labels, "needs-manual-dispatch")); + } +}