diff --git a/README.md b/README.md index 4d9336a6..26a1c388 100644 --- a/README.md +++ b/README.md @@ -21,7 +21,7 @@ A modular daemon system for managing AI agents, notifications, interactive quest **agentd** is a suite of services and tools designed to orchestrate AI agents and provide intelligent, context-aware notifications and interactions. It consists of: - **agent** - Command-line interface for interacting with all services -- **agentd-core** - Core service providing user and organization management +- **agentd-core** - Core service providing user, organization, and project management - **agentd-orchestrator** - Agent lifecycle management, WebSocket SDK server, workflow scheduler, and tool policy enforcement - **agentd-notify** - Notification service with REST API and SQLite storage - **agentd-ask** - Interactive question service with tmux integration @@ -97,6 +97,7 @@ agent teardown .agentd/ # delete in reverse order - **Declarative templates** - `agent apply` / `agent teardown` for YAML-based agent and workflow management - **Agent management** - create, list, get, delete, attach, send-message, stream - **Workflow management** - create, list, get, update, delete, history, validate-template +- **Project management** - create, list, show, update, delete; add/remove agent and workflow associations - **Tool policies** - get-policy, set-policy, `--tool-policy` flag on create-agent - **Approval management** - list-approvals, approve, deny (for RequireApproval policy) - **Health monitoring** - `agent status` checks all services concurrently; per-service `health` commands @@ -362,12 +363,13 @@ For the complete configuration reference including all environment variables, da | Service | Dev Port | Prod Port | Description | |---------|----------|-----------|-------------| +| agentd-core | 17000 | 7000 | Core API (users, organizations, projects) | | agentd-ask | 17001 | 7001 | Interactive question service | | agentd-hook | 17002 | 7002 | Shell hook integration | | agentd-monitor | 17003 | 7003 | System monitoring | | agentd-notify | 17004 | 7004 | Notification service | | agentd-wrap | 17005 | 7005 | Tmux session management | -| agentd-orchestrator | 17006 | 7006 | Agent orchestration | +| agentd-orchestrator | 17006 | 7006 | Agent orchestration and project associations | | agentd-index | 17012 | 17012 | Semantic code search and indexing | | agentd-mcp | - | - | MCP server (stdio transport, no HTTP port) | @@ -419,6 +421,12 @@ For the complete configuration reference including all environment variables, da - ✅ Structured JSON logging (`AGENTD_LOG_FORMAT=json`) - ✅ GitHub Actions CI/CD pipeline +**Project Entity Migration (v0.15.0):** +- ✅ Project CRUD moved from orchestrator to core service (`agentd-core`) +- ✅ CLI, MCP tools, and UI repointed to core for project CRUD +- ✅ Orchestrator retains project association endpoints (add/remove agent and workflow) +- ⚠️ **Operators upgrading to v0.15.0** must run `agent admin backfill-projects` once after the upgrade to assign `organization_id` to existing project rows created before multi-tenancy was introduced. + **In Progress:** - 🔄 Hook service - 🔄 Monitor service diff --git a/crates/cli/src/commands/admin.rs b/crates/cli/src/commands/admin.rs index aa92da63..ddc9d295 100644 --- a/crates/cli/src/commands/admin.rs +++ b/crates/cli/src/commands/admin.rs @@ -8,6 +8,8 @@ //! //! - **backfill-tenant** — Assign an `organization_id` to all existing rows //! that have a NULL value, backfilling legacy unscoped data. +//! - **backfill-projects** — Copy project rows from orchestrator's database +//! into core's database, preserving UUIDs. //! //! # Usage //! @@ -17,14 +19,22 @@ //! //! # Dry run to preview affected row counts without modifying data //! agent admin backfill-tenant --org-id acme-corp --dry-run +//! +//! # Copy all orchestrator projects into core +//! agent admin backfill-projects +//! +//! # Dry-run scoped to one organization +//! agent admin backfill-projects --org-id acme-corp --dry-run //! ``` +use std::collections::HashSet; + use agentd_common::storage::{create_connection, get_db_path}; use agentd_install::migrate::DB_SERVICES; use anyhow::{Context, Result}; use clap::Subcommand; use colored::*; -use sea_orm::{ConnectionTrait, Statement, Value}; +use sea_orm::{ConnectionTrait, FromQueryResult, Statement, Value}; /// Admin subcommands. #[derive(Debug, Subcommand)] @@ -56,6 +66,35 @@ pub enum AdminCommand { #[arg(long)] dry_run: bool, }, + + /// Copy project rows from orchestrator's database into core's database. + /// + /// Phase 2 of the projects-to-core migration: after core's `projects` + /// table exists (Phase 1), this command seeds it with all existing rows + /// from orchestrator, preserving UUIDs so that foreign-key references in + /// other services continue to resolve. + /// + /// The command is **idempotent**: rows already present in core (matched by + /// `id`) are silently skipped, so it is safe to run more than once. + /// + /// **When to use:** + /// - After upgrading to a release that adds `projects` to core + /// - Before repointing consumers (communicate, knowledge, orchestrator) + /// from the orchestrator project API to the core project API + #[command(name = "backfill-projects")] + BackfillProjects { + /// Preview which rows would be inserted without writing anything. + #[arg(long)] + dry_run: bool, + + /// Scope the backfill to a specific organization. + /// + /// When provided, only projects with a matching `organization_id` + /// **or** a NULL `organization_id` (legacy data) are copied. + /// When omitted, all orchestrator projects are copied. + #[arg(long)] + org_id: Option, + }, } impl AdminCommand { @@ -65,6 +104,9 @@ impl AdminCommand { AdminCommand::BackfillTenant { org_id, dry_run } => { backfill_tenant(org_id, *dry_run, json).await } + AdminCommand::BackfillProjects { dry_run, org_id } => { + backfill_projects(org_id.as_deref(), *dry_run, json).await + } } } } @@ -82,7 +124,8 @@ struct ServiceTables { } const SCOPED_TABLES: &[ServiceTables] = &[ - ServiceTables { service: "orchestrator", tables: &["agents", "workflows", "projects"] }, + ServiceTables { service: "orchestrator", tables: &["agents", "workflows"] }, + ServiceTables { service: "core", tables: &["projects"] }, ServiceTables { service: "notify", tables: &["notifications"] }, ServiceTables { service: "communicate", tables: &["rooms"] }, ServiceTables { service: "memory", tables: &["memory_entries"] }, @@ -239,6 +282,266 @@ async fn backfill_tenant(org_id: &str, dry_run: bool, json: bool) -> Result<()> Ok(()) } +// --------------------------------------------------------------------------- +// backfill_projects +// --------------------------------------------------------------------------- + +/// A project row read from orchestrator's `projects` table. +#[derive(Debug, FromQueryResult)] +struct OrchestratorProjectRow { + id: String, + name: String, + description: Option, + organization_id: Option, + created_at: String, + updated_at: String, +} + +/// Per-project outcome reported by [`run_backfill_projects`]. +#[derive(serde::Serialize, Clone, PartialEq, Debug)] +#[serde(rename_all = "snake_case")] +pub enum ProjectBackfillStatus { + /// Row was inserted into core's `projects` table. + Inserted, + /// Row already existed in core (matched by `id`); skipped. + Skipped, + /// Dry-run mode: row would be inserted. + WouldInsert, + /// Insert failed (e.g. unique-name conflict with a different id). + Error, +} + +/// Per-project result returned by [`run_backfill_projects`]. +#[derive(serde::Serialize, Clone, Debug)] +pub struct ProjectBackfillResult { + pub id: String, + pub name: String, + pub status: ProjectBackfillStatus, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +/// Core logic for the backfill-projects command. +/// +/// Separated from [`backfill_projects`] so it can be exercised in tests with +/// in-memory SQLite connections instead of real service databases. +pub async fn run_backfill_projects( + orch_db: &sea_orm::DatabaseConnection, + core_db: &sea_orm::DatabaseConnection, + org_id: Option<&str>, + dry_run: bool, +) -> Result> { + // Read projects from orchestrator (with optional org scope). + let rows = if let Some(oid) = org_id { + OrchestratorProjectRow::find_by_statement(Statement::from_sql_and_values( + orch_db.get_database_backend(), + "SELECT id, name, description, organization_id, created_at, updated_at \ + FROM projects \ + WHERE organization_id = ? OR organization_id IS NULL", + vec![Value::from(oid.to_owned())], + )) + .all(orch_db) + .await + .context("failed to read orchestrator projects")? + } else { + OrchestratorProjectRow::find_by_statement(Statement::from_string( + orch_db.get_database_backend(), + "SELECT id, name, description, organization_id, created_at, updated_at \ + FROM projects", + )) + .all(orch_db) + .await + .context("failed to read orchestrator projects")? + }; + + // Load all project IDs already present in core so we can skip them cheaply. + #[derive(FromQueryResult)] + struct IdRow { + id: String, + } + + let existing: HashSet = IdRow::find_by_statement(Statement::from_string( + core_db.get_database_backend(), + "SELECT id FROM projects", + )) + .all(core_db) + .await + .context("failed to read existing project IDs from core")? + .into_iter() + .map(|r| r.id) + .collect(); + + let mut results = Vec::with_capacity(rows.len()); + + for row in rows { + // Skip rows already present in core (idempotency). + if existing.contains(&row.id) { + results.push(ProjectBackfillResult { + id: row.id, + name: row.name, + status: ProjectBackfillStatus::Skipped, + error: None, + }); + continue; + } + + if dry_run { + results.push(ProjectBackfillResult { + id: row.id, + name: row.name, + status: ProjectBackfillStatus::WouldInsert, + error: None, + }); + continue; + } + + let stmt = Statement::from_sql_and_values( + core_db.get_database_backend(), + "INSERT INTO projects \ + (id, name, description, organization_id, created_at, updated_at) \ + VALUES (?, ?, ?, ?, ?, ?)", + vec![ + Value::from(row.id.clone()), + Value::from(row.name.clone()), + Value::from(row.description.clone()), + Value::from(row.organization_id.clone()), + Value::from(row.created_at.clone()), + Value::from(row.updated_at.clone()), + ], + ); + + match core_db.execute(stmt).await { + Ok(_) => results.push(ProjectBackfillResult { + id: row.id, + name: row.name, + status: ProjectBackfillStatus::Inserted, + error: None, + }), + Err(e) => results.push(ProjectBackfillResult { + id: row.id, + name: row.name, + status: ProjectBackfillStatus::Error, + error: Some(e.to_string()), + }), + } + } + + Ok(results) +} + +/// Render backfill results to stdout in either JSON or human-readable form. +fn output_backfill_projects_results( + results: &[ProjectBackfillResult], + dry_run: bool, + org_id: Option<&str>, + json: bool, +) -> Result<()> { + if json { + println!("{}", serde_json::to_string_pretty(results)?); + return Ok(()); + } + + let mode_label = if dry_run { " (dry-run)" } else { "" }; + let scope_label = match org_id { + Some(oid) => format!(": org_id = {oid}"), + None => String::new(), + }; + println!("{}", format!("Backfill projects{scope_label}{mode_label}").blue().bold()); + println!("{}", "=".repeat(60).cyan()); + + let mut n_inserted = 0u64; + let mut n_skipped = 0u64; + let mut n_errors = 0u64; + let mut n_would_insert = 0u64; + + for r in results { + let (icon, verb) = match r.status { + ProjectBackfillStatus::Inserted => { + n_inserted += 1; + ("✅", "inserted") + } + ProjectBackfillStatus::Skipped => { + n_skipped += 1; + (" ", "skipped") + } + ProjectBackfillStatus::WouldInsert => { + n_would_insert += 1; + ("🔍", "would insert") + } + ProjectBackfillStatus::Error => { + n_errors += 1; + ("❌", "error") + } + }; + let colored_verb = match r.status { + ProjectBackfillStatus::Error => verb.red().bold(), + _ => verb.green().bold(), + }; + println!(" {} {:<40} {}", icon, r.name.bright_black(), colored_verb); + if let Some(err) = &r.error { + println!(" error: {}", err.red()); + } + } + + println!(); + let total = results.len(); + if dry_run { + println!( + "Total: {} rows found in orchestrator — {} would be inserted, {} already in core", + total.to_string().green().bold(), + n_would_insert.to_string().green().bold(), + n_skipped.to_string().bright_black() + ); + println!("{}", "Run without --dry-run to apply the backfill.".yellow()); + } else { + println!( + "Total: {} rows processed — {} inserted, {} skipped, {} errors", + total.to_string().green().bold(), + n_inserted.to_string().green().bold(), + n_skipped.to_string().bright_black(), + if n_errors > 0 { n_errors.to_string().red() } else { n_errors.to_string().green() } + ); + } + + Ok(()) +} + +/// Open both service databases and run the project backfill. +async fn backfill_projects(org_id: Option<&str>, dry_run: bool, json: bool) -> Result<()> { + // Reject an explicitly empty --org-id before touching any database. + if let Some(id) = org_id { + if id.is_empty() { + anyhow::bail!("--org-id must not be empty"); + } + } + + // Open orchestrator database. + let orch_path = get_db_path("agentd-orchestrator", "orchestrator.db") + .context("cannot resolve orchestrator database path")?; + if !orch_path.exists() { + anyhow::bail!( + "orchestrator database not found at {}; is the orchestrator service initialized?", + orch_path.display() + ); + } + let orch_db = + create_connection(&orch_path).await.context("failed to open orchestrator database")?; + + // Open core database. + let core_path = + get_db_path("agentd-core", "core.db").context("cannot resolve core database path")?; + if !core_path.exists() { + anyhow::bail!( + "core database not found at {}; is the core service initialized?", + core_path.display() + ); + } + let core_db = create_connection(&core_path).await.context("failed to open core database")?; + + let results = run_backfill_projects(&orch_db, &core_db, org_id, dry_run).await?; + output_backfill_projects_results(&results, dry_run, org_id, json) +} + // --------------------------------------------------------------------------- // Tests // --------------------------------------------------------------------------- @@ -372,4 +675,173 @@ mod tests { "expected an empty-org-id error, got: {err}" ); } + + // ----------------------------------------------------------------------- + // backfill_projects tests + // ----------------------------------------------------------------------- + + /// Create the `projects` schema in an in-memory SQLite connection. + async fn create_projects_table(db: &sea_orm::DatabaseConnection) { + db.execute_unprepared( + "CREATE TABLE IF NOT EXISTS projects (\ + id TEXT PRIMARY KEY, \ + name TEXT NOT NULL UNIQUE, \ + description TEXT, \ + organization_id TEXT, \ + created_at TEXT NOT NULL, \ + updated_at TEXT NOT NULL\ + )", + ) + .await + .unwrap(); + } + + /// Seed the `projects` table with rows of the form `(id, name, org_id)`. + /// + /// `org_id = None` → NULL in the database. + /// + /// Uses parameterized statements so that values containing special characters + /// (e.g. single quotes) are handled safely, consistent with production paths. + async fn seed_projects(db: &sea_orm::DatabaseConnection, rows: &[(&str, &str, Option<&str>)]) { + create_projects_table(db).await; + for (id, name, org_id) in rows { + let stmt = Statement::from_sql_and_values( + db.get_database_backend(), + "INSERT INTO projects \ + (id, name, description, organization_id, created_at, updated_at) \ + VALUES (?, ?, NULL, ?, '2024-01-01T00:00:00Z', '2024-01-01T00:00:00Z')", + vec![ + Value::from((*id).to_string()), + Value::from((*name).to_string()), + Value::from(org_id.map(|s| s.to_string())), + ], + ); + db.execute(stmt).await.unwrap(); + } + } + + #[tokio::test] + async fn backfill_projects_rejects_empty_org_id() { + // Guard fires before any database path is resolved. + let err = backfill_projects(Some(""), false, true).await.unwrap_err(); + assert!( + err.to_string().contains("must not be empty"), + "expected an empty-org-id error, got: {err}" + ); + } + + #[tokio::test] + async fn backfill_projects_inserts_all_rows_first_run() { + let (orch_db, _orch_tmp) = create_test_connection().await; + let (core_db, _core_tmp) = create_test_connection().await; + + seed_projects(&orch_db, &[("uuid-1", "Alpha", Some("org-a")), ("uuid-2", "Beta", None)]) + .await; + create_projects_table(&core_db).await; + + let results = run_backfill_projects(&orch_db, &core_db, None, false).await.unwrap(); + + assert_eq!(results.len(), 2, "should process both orchestrator rows"); + assert!( + results.iter().all(|r| r.status == ProjectBackfillStatus::Inserted), + "both rows should be inserted on first run" + ); + } + + #[tokio::test] + async fn backfill_projects_idempotent() { + let (orch_db, _orch_tmp) = create_test_connection().await; + let (core_db, _core_tmp) = create_test_connection().await; + + seed_projects( + &orch_db, + &[("uuid-1", "Alpha", Some("org-a")), ("uuid-2", "Beta", Some("org-b"))], + ) + .await; + create_projects_table(&core_db).await; + + // First run: both rows should be inserted. + let first = run_backfill_projects(&orch_db, &core_db, None, false).await.unwrap(); + assert!( + first.iter().all(|r| r.status == ProjectBackfillStatus::Inserted), + "first run should insert all rows" + ); + + // Second run: both rows already exist → all skipped. + let second = run_backfill_projects(&orch_db, &core_db, None, false).await.unwrap(); + assert_eq!(second.len(), 2, "second run should still process both rows"); + assert!( + second.iter().all(|r| r.status == ProjectBackfillStatus::Skipped), + "second run should skip all rows (already in core)" + ); + } + + #[tokio::test] + async fn backfill_projects_dry_run_does_not_insert() { + let (orch_db, _orch_tmp) = create_test_connection().await; + let (core_db, _core_tmp) = create_test_connection().await; + + seed_projects(&orch_db, &[("uuid-1", "Alpha", None)]).await; + create_projects_table(&core_db).await; + + let results = + run_backfill_projects(&orch_db, &core_db, None, /* dry_run */ true).await.unwrap(); + + assert_eq!(results.len(), 1); + assert_eq!(results[0].status, ProjectBackfillStatus::WouldInsert); + + // Core must remain empty. + assert_eq!(total_rows(&core_db, "projects").await, 0, "dry-run must not insert anything"); + } + + #[tokio::test] + async fn backfill_projects_org_filter_includes_null_rows() { + let (orch_db, _orch_tmp) = create_test_connection().await; + let (core_db, _core_tmp) = create_test_connection().await; + + seed_projects( + &orch_db, + &[ + ("uuid-1", "OrgA Project", Some("org-a")), + ("uuid-2", "OrgB Project", Some("org-b")), + ("uuid-3", "Legacy Project", None), + ], + ) + .await; + create_projects_table(&core_db).await; + + // Filter to org-a: should include OrgA Project AND Legacy Project (NULL). + let results = + run_backfill_projects(&orch_db, &core_db, Some("org-a"), false).await.unwrap(); + + assert_eq!(results.len(), 2, "org filter should return org-a + NULL rows"); + let names: Vec<&str> = results.iter().map(|r| r.name.as_str()).collect(); + assert!(names.contains(&"OrgA Project")); + assert!(names.contains(&"Legacy Project")); + assert!(!names.contains(&"OrgB Project"), "org-b project must not appear"); + } + + #[tokio::test] + async fn backfill_projects_name_collision_reports_error() { + let (orch_db, _orch_tmp) = create_test_connection().await; + let (core_db, _core_tmp) = create_test_connection().await; + + // Orchestrator has "SharedName" with id "uuid-orch". + seed_projects(&orch_db, &[("uuid-orch", "SharedName", None)]).await; + + // Core already has "SharedName" but under a *different* id. + // Because the id doesn't match, the skip-by-id guard won't fire, + // and the INSERT will hit the UNIQUE constraint on `name`. + seed_projects(&core_db, &[("uuid-core", "SharedName", None)]).await; + + let results = run_backfill_projects(&orch_db, &core_db, None, false).await.unwrap(); + + assert_eq!(results.len(), 1); + assert_eq!( + results[0].status, + ProjectBackfillStatus::Error, + "name collision should yield Error status, not a panic" + ); + assert!(results[0].error.is_some(), "error field must contain the database error message"); + } } diff --git a/crates/cli/src/commands/project.rs b/crates/cli/src/commands/project.rs index da9e3223..d7a618d9 100644 --- a/crates/cli/src/commands/project.rs +++ b/crates/cli/src/commands/project.rs @@ -292,8 +292,6 @@ async fn show_project(client: &OrchestratorClient, id_or_name: &str, json: bool) if let Some(desc) = &project.description { println!(" Description: {}", desc); } - println!(" Agents: {}", project.agent_count); - println!(" Workflows: {}", project.workflow_count); println!(" Created: {}", project.created_at.format("%Y-%m-%d %H:%M:%S UTC")); println!(" Updated: {}", project.updated_at.format("%Y-%m-%d %H:%M:%S UTC")); diff --git a/crates/cli/src/main.rs b/crates/cli/src/main.rs index 07514408..f53847e4 100644 --- a/crates/cli/src/main.rs +++ b/crates/cli/src/main.rs @@ -712,12 +712,16 @@ async fn main() -> Result<()> { } Commands::Project { command } => { let token = load_token_or_warn(); - let url = gateway_url("orchestrator"); - let mut client = OrchestratorClient::new(url); + let orch_url = gateway_url("orchestrator"); + // Project CRUD now lives in core; association operations stay on the + // orchestrator. Build a client that routes CRUD to core and keeps + // associations on the orchestrator base URL. + let core_api_url = format!("{}/api/v1", client::core_url()); + let mut orch_client = OrchestratorClient::new(orch_url).with_core_url(core_api_url); if let Some(t) = token { - client = client.with_token(t); + orch_client = orch_client.with_token(t); } - command.execute(&client, cli.json).await?; + command.execute(&orch_client, cli.json).await?; } Commands::Control => { agentd_tui::run_control().await?; diff --git a/crates/cli/tests/integration_test.rs b/crates/cli/tests/integration_test.rs index 06f6885f..f79fddfc 100644 --- a/crates/cli/tests/integration_test.rs +++ b/crates/cli/tests/integration_test.rs @@ -1701,3 +1701,211 @@ mod client_tests { mock2.assert_async().await; } } + +// --------------------------------------------------------------------------- +// Project command tests +// --------------------------------------------------------------------------- +// +// Project CRUD uses the core service URL; association operations (add/remove +// agent or workflow) use the orchestrator URL. Both URL paths are exercised +// below using two independent mockito servers. + +mod project_tests { + use mockito::Server; + use orchestrator::client::OrchestratorClient; + use serde_json::json; + use uuid::Uuid; + + fn project_json(id: &str, name: &str) -> serde_json::Value { + json!({ + "id": id, + "name": name, + "description": null, + "organization_id": null, + "created_at": "2026-01-01T00:00:00Z", + "updated_at": "2026-01-01T00:00:00Z" + }) + } + + fn paginated(items: Vec) -> serde_json::Value { + let total = items.len(); + json!({ "items": items, "total": total, "limit": 50, "offset": 0 }) + } + + /// `list_projects` routes to the core URL (`/projects`). + #[tokio::test] + async fn test_list_projects_routes_to_core() { + let mut core_server = Server::new_async().await; + let orch_server = Server::new_async().await; // should not be called + + let id = Uuid::new_v4().to_string(); + let mock = core_server + .mock("GET", "/projects") + .with_status(200) + .with_header("content-type", "application/json") + .with_body(paginated(vec![project_json(&id, "My Project")]).to_string()) + .create_async() + .await; + + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client.list_projects().await; + + assert!(result.is_ok(), "list_projects should succeed: {result:?}"); + assert_eq!(result.unwrap().items.len(), 1); + mock.assert_async().await; + } + + /// `create_project` routes to the core URL. + #[tokio::test] + async fn test_create_project_routes_to_core() { + let mut core_server = Server::new_async().await; + let orch_server = Server::new_async().await; + + let id = Uuid::new_v4().to_string(); + let mock = core_server + .mock("POST", "/projects") + .with_status(201) + .with_header("content-type", "application/json") + .with_body(project_json(&id, "New Project").to_string()) + .create_async() + .await; + + use orchestrator::types::CreateProjectRequest; + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client + .create_project(&CreateProjectRequest { + name: "New Project".to_string(), + description: None, + }) + .await; + + assert!(result.is_ok(), "create_project should succeed: {result:?}"); + assert_eq!(result.unwrap().name, "New Project"); + mock.assert_async().await; + } + + /// `add_agent_to_project` routes to the orchestrator URL. + #[tokio::test] + async fn test_add_agent_to_project_routes_to_orchestrator() { + let core_server = Server::new_async().await; + let mut orch_server = Server::new_async().await; + + let project_id = Uuid::new_v4(); + let agent_id = Uuid::new_v4(); + + let mock = orch_server + .mock("POST", format!("/projects/{project_id}/agents/{agent_id}").as_str()) + .with_status(204) + .create_async() + .await; + + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client.associate_project_agent(&project_id, &agent_id).await; + + assert!(result.is_ok(), "associate_project_agent should succeed: {result:?}"); + mock.assert_async().await; + } + + /// `remove_agent_from_project` routes to the orchestrator URL. + #[tokio::test] + async fn test_remove_agent_from_project_routes_to_orchestrator() { + let core_server = Server::new_async().await; + let mut orch_server = Server::new_async().await; + + let project_id = Uuid::new_v4(); + let agent_id = Uuid::new_v4(); + + let mock = orch_server + .mock("DELETE", format!("/projects/{project_id}/agents/{agent_id}").as_str()) + .with_status(204) + .create_async() + .await; + + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client.dissociate_project_agent(&project_id, &agent_id).await; + + assert!(result.is_ok(), "dissociate_project_agent should succeed: {result:?}"); + mock.assert_async().await; + } + + /// `delete_project` pre-checks associations on the orchestrator before + /// deleting from core. When both return empty associations (total=0), + /// the delete request is forwarded to core. + #[tokio::test] + async fn test_delete_project_checks_associations_then_deletes_from_core() { + let mut core_server = Server::new_async().await; + let mut orch_server = Server::new_async().await; + + let project_id = Uuid::new_v4(); + + // Orchestrator: no agents associated + let agents_mock = orch_server + .mock("GET", format!("/projects/{project_id}/agents").as_str()) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(json!({ "items": [], "total": 0, "limit": 50, "offset": 0 }).to_string()) + .create_async() + .await; + + // Orchestrator: no workflows associated + let workflows_mock = orch_server + .mock("GET", format!("/projects/{project_id}/workflows").as_str()) + .with_status(200) + .with_header("content-type", "application/json") + .with_body(json!({ "items": [], "total": 0, "limit": 50, "offset": 0 }).to_string()) + .create_async() + .await; + + // Core: delete succeeds + let delete_mock = core_server + .mock("DELETE", format!("/projects/{project_id}").as_str()) + .with_status(204) + .create_async() + .await; + + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client.delete_project(&project_id).await; + + assert!(result.is_ok(), "delete_project should succeed: {result:?}"); + agents_mock.assert_async().await; + workflows_mock.assert_async().await; + delete_mock.assert_async().await; + } + + /// `delete_project` returns an error when agents are still associated. + #[tokio::test] + async fn test_delete_project_blocked_when_agents_associated() { + let core_server = Server::new_async().await; + let mut orch_server = Server::new_async().await; + + let project_id = Uuid::new_v4(); + let agent_id = Uuid::new_v4(); + + // Orchestrator: one agent associated + let _agents_mock = orch_server + .mock("GET", format!("/projects/{project_id}/agents").as_str()) + .with_status(200) + .with_header("content-type", "application/json") + .with_body( + json!({ + "items": [{ "id": agent_id, "name": "busy-agent", "status": "idle" }], + "total": 1, + "limit": 50, + "offset": 0 + }) + .to_string(), + ) + .create_async() + .await; + + let client = OrchestratorClient::new(orch_server.url()).with_core_url(core_server.url()); + let result = client.delete_project(&project_id).await; + + assert!(result.is_err(), "delete_project should fail when agents are associated"); + let msg = result.unwrap_err().to_string(); + assert!( + msg.contains("agent") || msg.contains("associated"), + "error message should mention agents: {msg}" + ); + } +} diff --git a/crates/common/src/config.rs b/crates/common/src/config.rs index a7986f4d..442ffb50 100644 --- a/crates/common/src/config.rs +++ b/crates/common/src/config.rs @@ -455,6 +455,8 @@ impl Default for CorePamConfig { #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] #[serde(default)] pub struct McpConfig { + /// Core service URL (API gateway). Defaults to `"http://localhost:17000"`. + pub core_url: String, /// Orchestrator service URL. Defaults to `"http://localhost:17006"`. pub orchestrator_url: String, /// Notify service URL. Defaults to `"http://localhost:17004"`. @@ -478,6 +480,7 @@ pub struct McpConfig { impl Default for McpConfig { fn default() -> Self { Self { + core_url: "http://localhost:17000".to_string(), orchestrator_url: "http://localhost:17006".to_string(), notify_url: "http://localhost:17004".to_string(), ask_url: "http://localhost:17001".to_string(), @@ -731,6 +734,7 @@ impl ValidateConfig for MonitorConfig { impl ValidateConfig for McpConfig { fn validate(&self) -> Result<()> { for (name, url) in [ + ("mcp.core_url", self.core_url.as_str()), ("mcp.orchestrator_url", self.orchestrator_url.as_str()), ("mcp.notify_url", self.notify_url.as_str()), ("mcp.ask_url", self.ask_url.as_str()), @@ -1164,6 +1168,11 @@ fn merge(base: AgentdConfig, file: AgentdConfig) -> AgentdConfig { }, }, mcp: McpConfig { + core_url: pick( + &base.services.mcp.core_url, + &file.services.mcp.core_url, + &d.services.mcp.core_url, + ), orchestrator_url: pick( &base.services.mcp.orchestrator_url, &file.services.mcp.orchestrator_url, @@ -1368,6 +1377,9 @@ fn apply_env_overrides(cfg: &mut AgentdConfig) { } // ── MCP ─────────────────────────────────────────────────────────────── + if let Ok(v) = env::var("AGENTD_MCP_CORE_URL") { + cfg.services.mcp.core_url = v; + } if let Ok(v) = env::var("AGENTD_MCP_ORCHESTRATOR_URL") { cfg.services.mcp.orchestrator_url = v; } diff --git a/crates/core/src/api.rs b/crates/core/src/api.rs index c04355e9..c92189d5 100644 --- a/crates/core/src/api.rs +++ b/crates/core/src/api.rs @@ -22,6 +22,11 @@ //! - `GET /api/v1/organizations/{id}/members` — list members //! - `POST /api/v1/organizations/{id}/members` — add member (owners only) //! - `DELETE /api/v1/organizations/{id}/members/{uid}` — remove member (owners only) +//! - `POST /api/v1/projects` — create project +//! - `GET /api/v1/projects` — list projects (tenant-scoped via X-Tenant-ID) +//! - `GET /api/v1/projects/{id}` — get project +//! - `PUT /api/v1/projects/{id}` — update project +//! - `DELETE /api/v1/projects/{id}` — delete project //! - `GET /api/v1/health` — aggregate downstream health check //! - `ANY /api/v1/{service}/*` — proxy to downstream service @@ -29,6 +34,7 @@ pub mod admin; pub mod auth; pub mod gateway; pub mod organizations; +pub mod projects; pub mod users; use std::sync::Arc; @@ -84,6 +90,7 @@ pub fn create_router_with_proxy(state: AppState, proxy: ProxyConfig) -> Router { let api_v1 = Router::new() .nest("/users", users::v1_router()) .nest("/organizations", organizations::router()) + .nest("/projects", projects::router()) // Product-admin routes — superuser only, product-wide (not tenant-scoped) .nest("/admin", admin::router()) // Gateway routes — /api/v1/health and /api/v1/{service}/* path diff --git a/crates/core/src/api/projects.rs b/crates/core/src/api/projects.rs new file mode 100644 index 00000000..db2f186a --- /dev/null +++ b/crates/core/src/api/projects.rs @@ -0,0 +1,911 @@ +//! Project management endpoint handlers. +//! +//! # Endpoints +//! +//! | Method | Path | Auth | Description | +//! |--------|------------------------|------|-------------------------------------------| +//! | POST | `/api/v1/projects` | Yes | Create a project (org from tenant header) | +//! | GET | `/api/v1/projects` | Yes | List projects, optionally scoped to org | +//! | GET | `/api/v1/projects/{id}`| Yes | Get project by UUID | +//! | PUT | `/api/v1/projects/{id}`| Yes | Update project name and/or description | +//! | DELETE | `/api/v1/projects/{id}`| Yes | Delete project | + +use axum::{ + extract::{Path, Query, State}, + http::StatusCode, + response::IntoResponse, + routing::{delete, get, post, put}, + Json, Router, +}; +use serde::{Deserialize, Serialize}; + +use agentd_common::error::ApiError; +use agentd_common::tenant::OptionalTenantId; + +use crate::middleware::auth::AuthUser; + +use super::AppState; + +// --------------------------------------------------------------------------- +// Request / response types +// --------------------------------------------------------------------------- + +#[derive(Debug, Deserialize)] +pub struct CreateProjectRequest { + pub name: String, + pub description: Option, +} + +#[derive(Debug, Deserialize)] +pub struct UpdateProjectRequest { + pub name: Option, + pub description: Option, +} + +#[derive(Debug, Deserialize)] +pub struct ProjectListQuery { + pub limit: Option, + pub offset: Option, +} + +#[derive(Debug, Serialize)] +pub struct ProjectResponse { + pub id: String, + pub name: String, + pub description: Option, + pub organization_id: Option, + pub created_at: String, + pub updated_at: String, +} + +// --------------------------------------------------------------------------- +// Router +// --------------------------------------------------------------------------- + +pub fn router() -> Router { + Router::new() + .route("/", post(create_project_handler)) + .route("/", get(list_projects_handler)) + .route("/{id}", get(get_project_handler)) + .route("/{id}", put(update_project_handler)) + .route("/{id}", delete(delete_project_handler)) +} + +// --------------------------------------------------------------------------- +// Handlers +// --------------------------------------------------------------------------- + +/// `POST /api/v1/projects` +/// +/// Creates a new project. The `organization_id` is taken from the +/// `X-Tenant-ID` header when present (forwarded by the core gateway). +/// +/// Returns `201 Created` with the project payload. +async fn create_project_handler( + _auth: AuthUser, + OptionalTenantId(org_id): OptionalTenantId, + State(state): State, + Json(body): Json, +) -> Result { + if body.name.trim().is_empty() { + return Err(ApiError::InvalidInput("project name must not be empty".to_string())); + } + + let project = state + .storage + .projects() + .create(&body.name, body.description.as_deref(), org_id.as_deref()) + .await + .map_err(|e| { + if e.to_string().contains("UNIQUE") { + ApiError::Conflict(format!("a project named '{}' already exists", body.name)) + } else { + ApiError::Internal(e) + } + })?; + + Ok(( + StatusCode::CREATED, + Json(ProjectResponse { + id: project.id, + name: project.name, + description: project.description, + organization_id: project.organization_id, + created_at: project.created_at, + updated_at: project.updated_at, + }), + )) +} + +/// `GET /api/v1/projects` +/// +/// Lists projects. When an `X-Tenant-ID` header is present (forwarded by +/// the core gateway), only projects belonging to that organization are +/// returned (plus legacy NULL-org rows). Supports `limit` and `offset` +/// query parameters for pagination. +/// +/// Note: rows are fetched in full then sliced in memory. A follow-up +/// `list_org_paginated` storage method should be added to avoid loading +/// unbounded result sets for large deployments. +async fn list_projects_handler( + _auth: AuthUser, + OptionalTenantId(org_id): OptionalTenantId, + State(state): State, + Query(query): Query, +) -> Result { + let all = + state.storage.projects().list_org(org_id.as_deref()).await.map_err(ApiError::Internal)?; + + let total = all.len(); + let limit = query.limit.unwrap_or(50).min(200); + let offset = query.offset.unwrap_or(0); + + let items: Vec = all + .into_iter() + .skip(offset) + .take(limit) + .map(|p| ProjectResponse { + id: p.id, + name: p.name, + description: p.description, + organization_id: p.organization_id, + created_at: p.created_at, + updated_at: p.updated_at, + }) + .collect(); + + Ok(Json(serde_json::json!({ + "items": items, + "total": total, + "limit": limit, + "offset": offset, + }))) +} + +/// `GET /api/v1/projects/{id}` +/// +/// Returns project details for the given UUID. Returns `404` if not found. +async fn get_project_handler( + _auth: AuthUser, + State(state): State, + Path(id): Path, +) -> Result { + let project = state + .storage + .projects() + .get_by_id(&id) + .await + .map_err(ApiError::Internal)? + .ok_or(ApiError::NotFound)?; + + Ok(Json(ProjectResponse { + id: project.id, + name: project.name, + description: project.description, + organization_id: project.organization_id, + created_at: project.created_at, + updated_at: project.updated_at, + })) +} + +/// `PUT /api/v1/projects/{id}` +/// +/// Updates the project's `name` and/or `description`. Fields absent from +/// the body are left unchanged. Returns `400` if the supplied name is +/// whitespace-only, `404` if not found, `409` if the new name conflicts +/// with an existing project. +async fn update_project_handler( + _auth: AuthUser, + State(state): State, + Path(id): Path, + Json(body): Json, +) -> Result { + if let Some(ref name) = body.name { + if name.trim().is_empty() { + return Err(ApiError::InvalidInput("project name must not be empty".to_string())); + } + } + + let project = state + .storage + .projects() + .update(&id, body.name.as_deref(), body.description.as_deref()) + .await + .map_err(|e| { + let msg = e.to_string(); + if msg.contains("not found") { + ApiError::NotFound + } else if msg.contains("UNIQUE") { + ApiError::Conflict("a project with that name already exists".to_string()) + } else { + ApiError::Internal(e) + } + })?; + + Ok(Json(ProjectResponse { + id: project.id, + name: project.name, + description: project.description, + organization_id: project.organization_id, + created_at: project.created_at, + updated_at: project.updated_at, + })) +} + +/// `DELETE /api/v1/projects/{id}` +/// +/// Deletes the project. Core only checks its own constraints — there is no +/// cross-service delete-guard at this layer. Returns `204 No Content`. +async fn delete_project_handler( + _auth: AuthUser, + State(state): State, + Path(id): Path, +) -> Result { + let deleted = state.storage.projects().delete(&id).await.map_err(ApiError::Internal)?; + + if !deleted { + return Err(ApiError::NotFound); + } + + Ok(StatusCode::NO_CONTENT) +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +#[cfg(test)] +mod tests { + use crate::api::AppState; + use crate::storage::Storage; + use agentd_common::storage::create_test_connection; + use axum::{ + body::Body, + http::{header, Request, StatusCode}, + Router, + }; + use http_body_util::BodyExt; + use tower::ServiceExt; + + async fn test_app() -> (Router, tempfile::TempDir) { + let (conn, tmp) = create_test_connection().await; + let storage = Storage::new(conn).await.unwrap(); + let state = AppState::new(storage); + let app = crate::api::create_router(state); + (app, tmp) + } + + async fn body_json(response: axum::response::Response) -> serde_json::Value { + let bytes = response.into_body().collect().await.unwrap().to_bytes(); + serde_json::from_slice(&bytes).unwrap() + } + + /// Register a user and return `(token, user_id)`. + async fn register(app: &Router, username: &str, email: &str) -> (String, String) { + let payload = + serde_json::json!({ "username": username, "email": email, "password": "testpass" }); + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/auth/register") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let body = body_json(response).await; + let token = body["token"].as_str().unwrap().to_string(); + let user_id = body["user"]["id"].as_str().unwrap().to_string(); + (token, user_id) + } + + // ----------------------------------------------------------------------- + // Create project + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn test_create_project_returns_201() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "alice", "alice@example.com").await; + + let payload = serde_json::json!({ "name": "Test Project" }); + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::CREATED); + let body = body_json(response).await; + assert_eq!(body["name"], "Test Project"); + assert!(body["id"].as_str().is_some()); + } + + #[tokio::test] + async fn test_create_project_unauthenticated() { + let (app, _tmp) = test_app().await; + + let payload = serde_json::json!({ "name": "Unauth Project" }); + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn test_create_project_with_description() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "bob", "bob@example.com").await; + + let payload = serde_json::json!({ + "name": "Described Project", + "description": "A helpful description" + }); + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::CREATED); + let body = body_json(response).await; + assert_eq!(body["description"], "A helpful description"); + } + + #[tokio::test] + async fn test_create_project_empty_name_returns_400() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "carol", "carol@example.com").await; + + let payload = serde_json::json!({ "name": " " }); + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn test_create_duplicate_project_returns_409() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "dan", "dan@example.com").await; + + let payload = serde_json::json!({ "name": "Duplicate Project" }); + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::CONFLICT); + } + + #[tokio::test] + async fn test_create_project_sets_org_from_tenant_header() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "eve", "eve@example.com").await; + + let payload = serde_json::json!({ "name": "Tenant Project" }); + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .header("X-Tenant-ID", "org-abc") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::CREATED); + let body = body_json(response).await; + assert_eq!(body["organization_id"], "org-abc"); + } + + // ----------------------------------------------------------------------- + // Get project + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn test_get_project_returns_200() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "frank", "frank@example.com").await; + + let payload = serde_json::json!({ "name": "Fetchable Project" }); + let create_resp = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let created = body_json(create_resp).await; + let id = created["id"].as_str().unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri(format!("/api/v1/projects/{id}")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["name"], "Fetchable Project"); + } + + #[tokio::test] + async fn test_get_project_unauthenticated() { + let (app, _tmp) = test_app().await; + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects/some-id") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn test_get_missing_project_returns_404() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "grace", "grace@example.com").await; + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects/nonexistent-id") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + + // ----------------------------------------------------------------------- + // List projects + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn test_list_projects_unauthenticated() { + let (app, _tmp) = test_app().await; + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn test_list_projects_returns_all() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "henry", "henry@example.com").await; + + for name in &["Alpha", "Beta", "Gamma"] { + let payload = serde_json::json!({ "name": name }); + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + } + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["total"], 3); + assert_eq!(body["items"].as_array().unwrap().len(), 3); + } + + #[tokio::test] + async fn test_list_projects_with_pagination() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "ida", "ida@example.com").await; + + for i in 0..5u32 { + let payload = serde_json::json!({ "name": format!("Proj {i}") }); + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + } + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects?limit=2&offset=1") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["total"], 5); + assert_eq!(body["items"].as_array().unwrap().len(), 2); + assert_eq!(body["limit"], 2); + assert_eq!(body["offset"], 1); + } + + #[tokio::test] + async fn test_list_projects_with_tenant_header_scopes_results() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "tenant_tester", "tenant_tester@example.com").await; + + // Project for org-x + for &(name, org) in &[ + ("Org-X Project", Some("org-x")), + ("No-Org Project", None), + ("Org-Y Project", Some("org-y")), + ] { + let payload = serde_json::json!({ "name": name }); + let mut builder = Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json"); + if let Some(oid) = org { + builder = builder.header("X-Tenant-ID", oid); + } + app.clone() + .oneshot(builder.body(Body::from(payload.to_string())).unwrap()) + .await + .unwrap(); + } + + // List with X-Tenant-ID: org-x — should return org-x project + NULL-org project (2) + let response = app + .clone() + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header("X-Tenant-ID", "org-x") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["total"], 2, "org-x scoped list should include org-x + NULL rows"); + + // List without tenant header — should return all 3 + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["total"], 3, "unscoped list should return all projects"); + } + + // ----------------------------------------------------------------------- + // Update project + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn test_update_project_unauthenticated() { + let (app, _tmp) = test_app().await; + + let payload = serde_json::json!({ "name": "New Name" }); + let response = app + .oneshot( + Request::builder() + .method("PUT") + .uri("/api/v1/projects/some-id") + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn test_update_project_returns_200() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "jack", "jack@example.com").await; + + let payload = serde_json::json!({ "name": "Original Name" }); + let create_resp = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let created = body_json(create_resp).await; + let id = created["id"].as_str().unwrap(); + + let update_payload = serde_json::json!({ "name": "Updated Name" }); + let response = app + .oneshot( + Request::builder() + .method("PUT") + .uri(format!("/api/v1/projects/{id}")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(update_payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let body = body_json(response).await; + assert_eq!(body["name"], "Updated Name"); + } + + #[tokio::test] + async fn test_update_project_empty_name_returns_400() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "kate", "kate@example.com").await; + + let payload = serde_json::json!({ "name": "Valid Name" }); + let create_resp = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let created = body_json(create_resp).await; + let id = created["id"].as_str().unwrap(); + + let update_payload = serde_json::json!({ "name": " " }); + let response = app + .oneshot( + Request::builder() + .method("PUT") + .uri(format!("/api/v1/projects/{id}")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(update_payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::BAD_REQUEST); + } + + #[tokio::test] + async fn test_update_missing_project_returns_404() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "liam", "liam@example.com").await; + + let payload = serde_json::json!({ "name": "New Name" }); + let response = app + .oneshot( + Request::builder() + .method("PUT") + .uri("/api/v1/projects/nonexistent-id") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + + // ----------------------------------------------------------------------- + // Delete project + // ----------------------------------------------------------------------- + + #[tokio::test] + async fn test_delete_project_unauthenticated() { + let (app, _tmp) = test_app().await; + + let response = app + .oneshot( + Request::builder() + .method("DELETE") + .uri("/api/v1/projects/some-id") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::UNAUTHORIZED); + } + + #[tokio::test] + async fn test_delete_project_returns_204() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "mia", "mia@example.com").await; + + let payload = serde_json::json!({ "name": "To Delete" }); + let create_resp = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri("/api/v1/projects") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .header(header::CONTENT_TYPE, "application/json") + .body(Body::from(payload.to_string())) + .unwrap(), + ) + .await + .unwrap(); + let created = body_json(create_resp).await; + let id = created["id"].as_str().unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/api/v1/projects/{id}")) + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NO_CONTENT); + } + + #[tokio::test] + async fn test_delete_missing_project_returns_404() { + let (app, _tmp) = test_app().await; + let (token, _) = register(&app, "noah", "noah@example.com").await; + + let response = app + .oneshot( + Request::builder() + .method("DELETE") + .uri("/api/v1/projects/nonexistent-id") + .header(header::AUTHORIZATION, format!("Bearer {token}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } +} diff --git a/crates/core/src/entity/mod.rs b/crates/core/src/entity/mod.rs index 3bc7b312..8e411137 100644 --- a/crates/core/src/entity/mod.rs +++ b/crates/core/src/entity/mod.rs @@ -1,13 +1,15 @@ //! SeaORM entity definitions for the core crate. //! -//! Four entities map to the core service tables: +//! Five entities map to the core service tables: //! //! - [`user`] → `users` table //! - [`organization`] → `organizations` table //! - [`membership`] → `memberships` table //! - [`session`] → `sessions` table +//! - [`project`] → `projects` table pub mod membership; pub mod organization; +pub mod project; pub mod session; pub mod user; diff --git a/crates/core/src/entity/project.rs b/crates/core/src/entity/project.rs new file mode 100644 index 00000000..dcbd0b5f --- /dev/null +++ b/crates/core/src/entity/project.rs @@ -0,0 +1,29 @@ +//! SeaORM entity for the `projects` table. +//! +//! Projects provide a logical grouping layer below organizations: +//! `org > project > {agents, workflows, rooms, docs}`. + +use sea_orm::entity::prelude::*; + +/// SeaORM model for the `projects` table. +/// +/// The `name` column carries a unique index enforced by the migration. +/// `organization_id` is optional — rows with `NULL` represent legacy data +/// created before multi-tenancy was introduced. +#[derive(Clone, Debug, PartialEq, DeriveEntityModel)] +#[sea_orm(table_name = "projects")] +pub struct Model { + #[sea_orm(primary_key, auto_increment = false)] + pub id: String, + #[sea_orm(unique)] + pub name: String, + pub description: Option, + pub organization_id: Option, + pub created_at: String, + pub updated_at: String, +} + +#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] +pub enum Relation {} + +impl ActiveModelBehavior for ActiveModel {} diff --git a/crates/core/src/lib.rs b/crates/core/src/lib.rs index 6e948a80..e8520d06 100644 --- a/crates/core/src/lib.rs +++ b/crates/core/src/lib.rs @@ -12,6 +12,7 @@ pub mod middleware; pub mod migration; pub mod organization_storage; pub mod pam_auth; +pub mod project_storage; pub mod proxy; pub mod session_storage; pub mod storage; diff --git a/crates/core/src/migration/m20260616_000004_create_projects_table.rs b/crates/core/src/migration/m20260616_000004_create_projects_table.rs new file mode 100644 index 00000000..dde336f2 --- /dev/null +++ b/crates/core/src/migration/m20260616_000004_create_projects_table.rs @@ -0,0 +1,72 @@ +//! Migration 4: create the `projects` table in the core service. +//! +//! Projects sit below organizations in the tenant hierarchy: +//! `org > project > {agents, workflows, rooms, docs}`. + +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .create_table( + Table::create() + .table(Projects::Table) + .if_not_exists() + .col(ColumnDef::new(Projects::Id).string().not_null().primary_key()) + .col(ColumnDef::new(Projects::Name).string().not_null()) + .col(ColumnDef::new(Projects::Description).string().null()) + .col(ColumnDef::new(Projects::OrganizationId).string().null()) + .col(ColumnDef::new(Projects::CreatedAt).string().not_null()) + .col(ColumnDef::new(Projects::UpdatedAt).string().not_null()) + .to_owned(), + ) + .await?; + + // Unique index on name — project names must be globally unique. + manager + .create_index( + Index::create() + .name("idx_projects_name") + .table(Projects::Table) + .col(Projects::Name) + .unique() + .if_not_exists() + .to_owned(), + ) + .await?; + + // Non-unique index on organization_id for tenant-scoped listing. + manager + .create_index( + Index::create() + .name("idx_projects_organization_id") + .table(Projects::Table) + .col(Projects::OrganizationId) + .if_not_exists() + .to_owned(), + ) + .await?; + + Ok(()) + } + + async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager.drop_table(Table::drop().table(Projects::Table).to_owned()).await + } +} + +/// Iden enum for the `projects` table columns. +#[derive(DeriveIden)] +enum Projects { + Table, + Id, + Name, + Description, + OrganizationId, + CreatedAt, + UpdatedAt, +} diff --git a/crates/core/src/migration/mod.rs b/crates/core/src/migration/mod.rs index 4a31b80b..543b9f2a 100644 --- a/crates/core/src/migration/mod.rs +++ b/crates/core/src/migration/mod.rs @@ -14,6 +14,7 @@ pub use sea_orm_migration::prelude::*; mod m20260305_000001_create_core_tables; mod m20260408_000002_add_username_to_users; mod m20260614_000003_add_is_superuser_to_users; +mod m20260616_000004_create_projects_table; mod m20260618_000004_add_auth_provider_to_users; /// The migration runner — applies all known migrations in order. @@ -26,6 +27,7 @@ impl MigratorTrait for Migrator { Box::new(m20260305_000001_create_core_tables::Migration), Box::new(m20260408_000002_add_username_to_users::Migration), Box::new(m20260614_000003_add_is_superuser_to_users::Migration), + Box::new(m20260616_000004_create_projects_table::Migration), Box::new(m20260618_000004_add_auth_provider_to_users::Migration), ] } diff --git a/crates/core/src/project_storage.rs b/crates/core/src/project_storage.rs new file mode 100644 index 00000000..5d4304a7 --- /dev/null +++ b/crates/core/src/project_storage.rs @@ -0,0 +1,296 @@ +//! SeaORM-based storage for project records. +//! +//! [`ProjectStorage`] provides CRUD operations for the `projects` table. +//! It holds a [`DatabaseConnection`] shared with the parent +//! [`crate::storage::Storage`]. +//! +//! ## Tenant-scoped listing +//! +//! [`ProjectStorage::list_org`] includes rows with a `NULL` `organization_id` +//! alongside rows that match the requested org. This preserves visibility of +//! legacy data created before multi-tenancy was introduced; those rows should +//! be backfilled by the `backfill-projects` admin command before NULL inclusion +//! is removed. + +use anyhow::{anyhow, Result}; +use sea_orm::{ + ActiveModelTrait, ColumnTrait, Condition, DatabaseConnection, EntityTrait, PaginatorTrait, + QueryFilter, QueryOrder, Set, +}; +use uuid::Uuid; + +use crate::entity::project::{self, Column}; + +/// Storage operations for the `projects` table. +#[derive(Clone)] +pub struct ProjectStorage { + db: DatabaseConnection, +} + +impl ProjectStorage { + pub fn new(db: DatabaseConnection) -> Self { + Self { db } + } + + /// Insert a new project and return the created record. + pub async fn create( + &self, + name: &str, + description: Option<&str>, + organization_id: Option<&str>, + ) -> Result { + let now = chrono::Utc::now().to_rfc3339(); + let model = project::ActiveModel { + id: Set(Uuid::new_v4().to_string()), + name: Set(name.to_string()), + description: Set(description.map(|s| s.to_string())), + organization_id: Set(organization_id.map(|s| s.to_string())), + created_at: Set(now.clone()), + updated_at: Set(now), + } + .insert(&self.db) + .await?; + Ok(model) + } + + /// Return the project with the given id, or `None` if not found. + pub async fn get_by_id(&self, id: &str) -> Result> { + Ok(project::Entity::find_by_id(id).one(&self.db).await?) + } + + /// Return the project with the given name, or `None` if not found. + pub async fn get_by_name(&self, name: &str) -> Result> { + Ok(project::Entity::find().filter(Column::Name.eq(name)).one(&self.db).await?) + } + + /// Return all projects ordered by `created_at` descending. + pub async fn list(&self) -> Result> { + Ok(project::Entity::find().order_by_desc(Column::CreatedAt).all(&self.db).await?) + } + + /// Return projects filtered by organization, or all projects when `org_id` is `None`. + /// + /// When `org_id` is provided, rows with a matching `organization_id` **or** + /// a `NULL` `organization_id` are included. The NULL-row inclusion is a + /// deliberate tenant-transition aid: pre-migration data remains visible to + /// authenticated tenants until the `backfill-projects` admin command is run. + pub async fn list_org(&self, org_id: Option<&str>) -> Result> { + let query = project::Entity::find().order_by_desc(Column::CreatedAt); + let query = if let Some(oid) = org_id { + // Include legacy NULL rows so pre-migration data is still visible + // to authenticated tenants until backfill-tenant is run. + query.filter( + Condition::any() + .add(Column::OrganizationId.eq(oid)) + .add(Column::OrganizationId.is_null()), + ) + } else { + query + }; + Ok(query.all(&self.db).await?) + } + + /// Return a paginated list of all projects ordered by `created_at` descending. + /// + /// Intended for product-admin use — not tenant-scoped. + pub async fn list_paginated( + &self, + limit: u64, + offset: u64, + ) -> Result> { + let paginator = + project::Entity::find().order_by_desc(Column::CreatedAt).paginate(&self.db, limit); + let total = paginator.num_items().await?; + let page = offset.checked_div(limit).unwrap_or(0); + let items = paginator.fetch_page(page).await?; + Ok(agentd_common::types::PaginatedResponse { + items, + total: total as usize, + limit: limit as usize, + offset: offset as usize, + }) + } + + /// Update name and/or description for the given project. + /// + /// Returns an error if the project does not exist. + pub async fn update( + &self, + id: &str, + name: Option<&str>, + description: Option<&str>, + ) -> Result { + let proj = project::Entity::find_by_id(id) + .one(&self.db) + .await? + .ok_or_else(|| anyhow!("project not found: {}", id))?; + + let mut active: project::ActiveModel = proj.into(); + if let Some(name) = name { + active.name = Set(name.to_string()); + } + if let Some(description) = description { + active.description = Set(Some(description.to_string())); + } + active.updated_at = Set(chrono::Utc::now().to_rfc3339()); + Ok(active.update(&self.db).await?) + } + + /// Delete the project with the given id. + /// + /// Returns `true` if a row was deleted, `false` if not found. + pub async fn delete(&self, id: &str) -> Result { + let result = project::Entity::delete_by_id(id).exec(&self.db).await?; + Ok(result.rows_affected > 0) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::migration::Migrator; + use agentd_common::storage::create_test_connection; + use sea_orm_migration::MigratorTrait; + + async fn setup() -> (ProjectStorage, tempfile::TempDir) { + let (conn, tmp) = create_test_connection().await; + Migrator::up(&conn, None).await.unwrap(); + (ProjectStorage::new(conn), tmp) + } + + #[tokio::test] + async fn test_create_and_get_by_id() { + let (storage, _tmp) = setup().await; + let proj = storage.create("My Project", Some("A test project"), None).await.unwrap(); + assert_eq!(proj.name, "My Project"); + assert_eq!(proj.description.as_deref(), Some("A test project")); + assert!(proj.organization_id.is_none()); + + let found = storage.get_by_id(&proj.id).await.unwrap().unwrap(); + assert_eq!(found.id, proj.id); + assert_eq!(found.name, "My Project"); + } + + #[tokio::test] + async fn test_get_by_name() { + let (storage, _tmp) = setup().await; + storage.create("Named Project", None, None).await.unwrap(); + + let found = storage.get_by_name("Named Project").await.unwrap().unwrap(); + assert_eq!(found.name, "Named Project"); + + let missing = storage.get_by_name("nope").await.unwrap(); + assert!(missing.is_none()); + } + + #[tokio::test] + async fn test_get_by_id_not_found() { + let (storage, _tmp) = setup().await; + let missing = storage.get_by_id("nonexistent-id").await.unwrap(); + assert!(missing.is_none()); + } + + #[tokio::test] + async fn test_list_ordered_by_created_at_desc() { + let (storage, _tmp) = setup().await; + storage.create("Alpha", None, None).await.unwrap(); + storage.create("Beta", None, None).await.unwrap(); + + let projects = storage.list().await.unwrap(); + assert_eq!(projects.len(), 2); + // Beta was created after Alpha so it should come first (desc order) + assert_eq!(projects[0].name, "Beta"); + assert_eq!(projects[1].name, "Alpha"); + } + + #[tokio::test] + async fn test_list_org_filters_by_org_id() { + let (storage, _tmp) = setup().await; + let org_a = "org-a"; + let org_b = "org-b"; + storage.create("Proj A", None, Some(org_a)).await.unwrap(); + storage.create("Proj B", None, Some(org_b)).await.unwrap(); + // Legacy row with no org + storage.create("Legacy", None, None).await.unwrap(); + + // list_org(Some(org_a)) should return Proj A and Legacy (NULL) + let for_a = storage.list_org(Some(org_a)).await.unwrap(); + let names_a: Vec<&str> = for_a.iter().map(|p| p.name.as_str()).collect(); + assert!(names_a.contains(&"Proj A"), "expected Proj A in {:?}", names_a); + assert!(names_a.contains(&"Legacy"), "expected Legacy (NULL row) in {:?}", names_a); + assert!(!names_a.contains(&"Proj B"), "did not expect Proj B in {:?}", names_a); + } + + #[tokio::test] + async fn test_list_org_none_returns_all() { + let (storage, _tmp) = setup().await; + storage.create("P1", None, Some("org-x")).await.unwrap(); + storage.create("P2", None, None).await.unwrap(); + + let all = storage.list_org(None).await.unwrap(); + assert_eq!(all.len(), 2); + } + + #[tokio::test] + async fn test_list_paginated() { + let (storage, _tmp) = setup().await; + for i in 0..5u32 { + storage.create(&format!("Project {i}"), None, None).await.unwrap(); + } + + let page = storage.list_paginated(2, 0).await.unwrap(); + assert_eq!(page.items.len(), 2); + assert_eq!(page.total, 5); + assert_eq!(page.limit, 2); + assert_eq!(page.offset, 0); + + let page2 = storage.list_paginated(2, 2).await.unwrap(); + assert_eq!(page2.items.len(), 2); + assert_eq!(page2.offset, 2); + } + + #[tokio::test] + async fn test_update() { + let (storage, _tmp) = setup().await; + let proj = storage.create("Old Name", Some("old desc"), None).await.unwrap(); + + let updated = storage.update(&proj.id, Some("New Name"), Some("new desc")).await.unwrap(); + assert_eq!(updated.name, "New Name"); + assert_eq!(updated.description.as_deref(), Some("new desc")); + assert!(updated.updated_at >= proj.updated_at); + } + + #[tokio::test] + async fn test_update_not_found() { + let (storage, _tmp) = setup().await; + let result = storage.update("nonexistent-id", Some("x"), None).await; + assert!(result.is_err()); + } + + #[tokio::test] + async fn test_delete() { + let (storage, _tmp) = setup().await; + let proj = storage.create("To Delete", None, None).await.unwrap(); + + let deleted = storage.delete(&proj.id).await.unwrap(); + assert!(deleted); + + let found = storage.get_by_id(&proj.id).await.unwrap(); + assert!(found.is_none()); + } + + #[tokio::test] + async fn test_delete_not_found() { + let (storage, _tmp) = setup().await; + let result = storage.delete("nonexistent-id").await.unwrap(); + assert!(!result); + } + + #[tokio::test] + async fn test_unique_name_constraint() { + let (storage, _tmp) = setup().await; + storage.create("Duplicate", None, None).await.unwrap(); + let result = storage.create("Duplicate", None, None).await; + assert!(result.is_err(), "duplicate name should be rejected"); + } +} diff --git a/crates/core/src/storage.rs b/crates/core/src/storage.rs index 850a14e0..72f20fe3 100644 --- a/crates/core/src/storage.rs +++ b/crates/core/src/storage.rs @@ -22,6 +22,7 @@ use sea_orm_migration::MigratorTrait; use crate::membership_storage::MembershipStorage; use crate::migration::Migrator; use crate::organization_storage::OrganizationStorage; +use crate::project_storage::ProjectStorage; use crate::session_storage::SessionStorage; use crate::user_storage::UserStorage; @@ -54,6 +55,11 @@ impl Storage { OrganizationStorage::new(self.db.clone()) } + /// Returns a [`ProjectStorage`] instance sharing this connection. + pub fn projects(&self) -> ProjectStorage { + ProjectStorage::new(self.db.clone()) + } + /// Returns a [`MembershipStorage`] instance sharing this connection. pub fn memberships(&self) -> MembershipStorage { MembershipStorage::new(self.db.clone()) diff --git a/crates/mcp/src/client.rs b/crates/mcp/src/client.rs index 386aec3e..8326ab62 100644 --- a/crates/mcp/src/client.rs +++ b/crates/mcp/src/client.rs @@ -22,6 +22,11 @@ impl AgentdClient { Self { inner: Client::new(), config } } + /// Returns the base URL for the core service (API gateway). + pub fn core_url(&self) -> &str { + &self.config.core_url + } + /// Returns the base URL for the orchestrator service. pub fn orchestrator_url(&self) -> &str { &self.config.orchestrator_url diff --git a/crates/mcp/src/config.rs b/crates/mcp/src/config.rs index 6ffa57dc..9cf15458 100644 --- a/crates/mcp/src/config.rs +++ b/crates/mcp/src/config.rs @@ -11,6 +11,8 @@ use std::env; /// Configuration for connecting to agentd services. #[derive(Debug, Clone)] pub struct AgentdMcpConfig { + /// Core service URL (API gateway, default: `http://127.0.0.1:17000`) + pub core_url: String, /// Orchestrator service URL (default: `http://127.0.0.1:17006`) pub orchestrator_url: String, /// Communicate service URL (default: `http://127.0.0.1:17010`) @@ -41,6 +43,7 @@ impl AgentdMcpConfig { /// /// | Variable | Default | /// |---------------------------------|-----------------------------| + /// | `AGENTD_CORE_URL` | `http://127.0.0.1:17000` | /// | `AGENTD_ORCHESTRATOR_URL` | `http://127.0.0.1:17006` | /// | `AGENTD_COMMUNICATE_URL` | `http://127.0.0.1:17010` | /// | `AGENTD_MEMORY_URL` | `http://127.0.0.1:17008` | @@ -58,6 +61,7 @@ impl AgentdMcpConfig { let base = shared.services.mcp; Self { + core_url: env::var("AGENTD_CORE_URL").unwrap_or(base.core_url), orchestrator_url: env::var("AGENTD_ORCHESTRATOR_URL").unwrap_or(base.orchestrator_url), communicate_url: env::var("AGENTD_COMMUNICATE_URL").unwrap_or(base.communicate_url), memory_url: env::var("AGENTD_MEMORY_URL").unwrap_or(base.memory_url), @@ -80,6 +84,7 @@ impl AgentdMcpConfig { impl ValidateConfig for AgentdMcpConfig { fn validate(&self) -> Result<()> { let urls = [ + ("mcp.core_url", self.core_url.as_str()), ("mcp.orchestrator_url", self.orchestrator_url.as_str()), ("mcp.communicate_url", self.communicate_url.as_str()), ("mcp.memory_url", self.memory_url.as_str()), @@ -111,6 +116,7 @@ mod tests { fn test_defaults() { let _g = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner()); let vars = [ + "AGENTD_CORE_URL", "AGENTD_ORCHESTRATOR_URL", "AGENTD_COMMUNICATE_URL", "AGENTD_MEMORY_URL", @@ -133,6 +139,7 @@ mod tests { env::set_var("AGENTD_CONFIG", "/nonexistent/agentd-mcp-test-config.toml"); let config = AgentdMcpConfig::from_env(); + assert_eq!(config.core_url, "http://localhost:17000"); assert_eq!(config.orchestrator_url, "http://localhost:17006"); assert_eq!(config.communicate_url, "http://localhost:17010"); assert_eq!(config.memory_url, "http://localhost:17008"); @@ -168,6 +175,7 @@ mod tests { fn test_validate_default_passes() { let _g = ENV_LOCK.lock().unwrap_or_else(|e| e.into_inner()); let vars = [ + "AGENTD_CORE_URL", "AGENTD_ORCHESTRATOR_URL", "AGENTD_COMMUNICATE_URL", "AGENTD_MEMORY_URL", @@ -195,6 +203,7 @@ mod tests { #[test] fn test_validate_bad_url_fails() { let config = AgentdMcpConfig { + core_url: "http://127.0.0.1:17000".to_string(), orchestrator_url: "not-a-url".to_string(), communicate_url: "http://127.0.0.1:17010".to_string(), memory_url: "http://127.0.0.1:17008".to_string(), diff --git a/crates/mcp/src/tools/orchestrator_debug.rs b/crates/mcp/src/tools/orchestrator_debug.rs index 5f77c52b..a9f17057 100644 --- a/crates/mcp/src/tools/orchestrator_debug.rs +++ b/crates/mcp/src/tools/orchestrator_debug.rs @@ -90,8 +90,6 @@ struct ProjectDetail { description: Option, created_at: String, updated_at: String, - agent_count: u64, - workflow_count: u64, } #[derive(Debug, Deserialize)] @@ -302,13 +300,13 @@ pub async fn run_get_conversation_summary(client: &AgentdClient, agent_id: &str) } pub async fn run_list_projects(client: &AgentdClient, limit: Option) -> String { - let base = client.orchestrator_url(); + let base = client.core_url(); let limit_val = limit.unwrap_or(50).clamp(1, 200); - let url = format!("{base}/projects?limit={limit_val}"); + let url = format!("{base}/api/v1/projects?limit={limit_val}"); let resp = match client.inner.get(&url).send().await { Ok(r) => r, - Err(e) => return format!("Error: orchestrator unreachable at {base}: {e}"), + Err(e) => return format!("Error: core unreachable at {base}: {e}"), }; if !resp.status().is_success() { return format!("Error: HTTP {} listing projects", resp.status()); @@ -338,12 +336,12 @@ pub async fn run_list_projects(client: &AgentdClient, limit: Option) -> Str } pub async fn run_get_project(client: &AgentdClient, project_id: &str) -> String { - let base = client.orchestrator_url(); - let url = format!("{base}/projects/{project_id}"); + let base = client.core_url(); + let url = format!("{base}/api/v1/projects/{project_id}"); let resp = match client.inner.get(&url).send().await { Ok(r) => r, - Err(e) => return format!("Error: orchestrator unreachable at {base}: {e}"), + Err(e) => return format!("Error: core unreachable at {base}: {e}"), }; if resp.status() == reqwest::StatusCode::NOT_FOUND { return format!("Project `{project_id}` not found."); @@ -361,8 +359,6 @@ pub async fn run_get_project(client: &AgentdClient, project_id: &str) -> String if let Some(ref d) = p.description { out.push_str(&format!("- **Description**: {d}\n")); } - out.push_str(&format!("- **Agents**: {}\n", p.agent_count)); - out.push_str(&format!("- **Workflows**: {}\n", p.workflow_count)); out.push_str(&format!("- **Created**: {}\n", p.created_at)); out.push_str(&format!("- **Updated**: {}\n", p.updated_at)); out diff --git a/crates/mcp/tests/common/mod.rs b/crates/mcp/tests/common/mod.rs index 4f137c1a..8ecb2bf1 100644 --- a/crates/mcp/tests/common/mod.rs +++ b/crates/mcp/tests/common/mod.rs @@ -418,6 +418,7 @@ pub async fn mock_monitor_server() -> MockServer { /// Build an `AgentdClient` pointed at the given mock servers. pub fn test_client(orch_addr: &str, notify_addr: &str, monitor_addr: &str) -> AgentdClient { let config = Arc::new(AgentdMcpConfig { + core_url: "http://127.0.0.1:1".to_string(), // unused orchestrator_url: orch_addr.to_string(), communicate_url: "http://127.0.0.1:1".to_string(), // unused memory_url: "http://127.0.0.1:1".to_string(), diff --git a/crates/orchestrator/src/api.rs b/crates/orchestrator/src/api.rs index 46cd414b..46e7a657 100644 --- a/crates/orchestrator/src/api.rs +++ b/crates/orchestrator/src/api.rs @@ -97,9 +97,7 @@ pub fn create_router(state: ApiState) -> Router { .route("/approvals/{id}/deny", post(deny_tool)) .route("/debug/agents", get(debug_agents)) .route("/events/ask", post(ask_event_handler)) - // Project management - .route("/projects", get(list_projects).post(create_project)) - .route("/projects/{id}", get(get_project).put(update_project).delete(delete_project)) + // Project association endpoints (CRUD lives in core service) .route("/projects/{id}/agents", get(list_project_agents)) .route( "/projects/{id}/agents/{agent_id}", @@ -1307,8 +1305,10 @@ async fn ask_event_handler( } // --------------------------------------------------------------------------- -// Project management handlers +// Project association handlers // --------------------------------------------------------------------------- +// Project CRUD (create/read/update/delete) has been moved to the core service. +// These handlers manage the associations between projects and agents/workflows. #[derive(Deserialize)] struct ProjectListQuery { @@ -1316,131 +1316,11 @@ struct ProjectListQuery { offset: Option, } -async fn list_projects( - OptionalTenantId(org_id): OptionalTenantId, - State(state): State, - Query(query): Query, -) -> Result { - let ps = state.manager.project_storage(); - let all = ps.list_org(org_id.as_deref()).await.map_err(ApiError::Internal)?; - let total = all.len(); - let limit = clamp_limit(query.limit); - let offset = query.offset.unwrap_or(0); - let items: Vec = all.into_iter().skip(offset).take(limit).collect(); - Ok(Json(PaginatedResponse { items, total, limit, offset })) -} - -async fn create_project( - OptionalTenantId(org_id): OptionalTenantId, - State(state): State, - Json(req): Json, -) -> Result { - if req.name.trim().is_empty() { - return Err(ApiError::InvalidInput("project name must not be empty".to_string())); - } - let mut project = Project::new(req.name, req.description); - project.organization_id = org_id; - let ps = state.manager.project_storage(); - let created = ps.create(&project).await.map_err(|e| { - if e.to_string().contains("UNIQUE") { - ApiError::Conflict(format!("a project named '{}' already exists", project.name)) - } else { - ApiError::Internal(e) - } - })?; - Ok((StatusCode::CREATED, Json(created))) -} - -async fn get_project( - State(state): State, - Path(id): Path, -) -> Result { - let ps = state.manager.project_storage(); - let project = ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; - - let agent_count = - state.manager.agent_storage().list(None, Some(id)).await.map_err(ApiError::Internal)?.len(); - - let workflow_count = - state.scheduler.storage().list_workflows(Some(id)).await.map_err(ApiError::Internal)?.len(); - - Ok(Json(ProjectResponse { - id: project.id, - name: project.name, - description: project.description, - created_at: project.created_at, - updated_at: project.updated_at, - agent_count, - workflow_count, - })) -} - -async fn update_project( - State(state): State, - Path(id): Path, - Json(req): Json, -) -> Result { - let ps = state.manager.project_storage(); - let updated = ps.update(&id, &req).await.map_err(|e| { - let msg = e.to_string(); - if msg.contains("not found") { - ApiError::NotFound - } else if msg.contains("UNIQUE") { - ApiError::Conflict("a project with that name already exists".to_string()) - } else { - ApiError::Internal(e) - } - })?; - Ok(Json(updated)) -} - -async fn delete_project( - State(state): State, - Path(id): Path, -) -> Result { - let ps = state.manager.project_storage(); - - // Verify project exists. - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; - - // Reject if any agents are still associated. - let associated_agents = - state.manager.agent_storage().list(None, Some(id)).await.map_err(ApiError::Internal)?; - if !associated_agents.is_empty() { - return Err(ApiError::InvalidInput(format!( - "cannot delete project: {} agent(s) still associated", - associated_agents.len() - ))); - } - - // Reject if any workflows are still associated. - let associated_workflows = - state.scheduler.storage().list_workflows(Some(id)).await.map_err(ApiError::Internal)?; - if !associated_workflows.is_empty() { - return Err(ApiError::InvalidInput(format!( - "cannot delete project: {} workflow(s) still associated", - associated_workflows.len() - ))); - } - - ps.delete(&id).await.map_err(|e| { - if e.to_string().contains("not found") { - ApiError::NotFound - } else { - ApiError::Internal(e) - } - })?; - Ok(StatusCode::NO_CONTENT) -} - async fn list_project_agents( State(state): State, Path(id): Path, Query(query): Query, ) -> Result { - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; - let limit = clamp_limit(query.limit); let offset = query.offset.unwrap_or(0); @@ -1465,9 +1345,7 @@ async fn associate_project_agent( State(state): State, Path((id, agent_id)): Path<(Uuid, Uuid)>, ) -> Result { - // Verify project and agent exist. - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; + // Verify agent exists. state .manager .get_agent(&agent_id) @@ -1486,10 +1364,8 @@ async fn associate_project_agent( async fn dissociate_project_agent( State(state): State, - Path((id, agent_id)): Path<(Uuid, Uuid)>, + Path((_id, agent_id)): Path<(Uuid, Uuid)>, ) -> Result { - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; state .manager .get_agent(&agent_id) @@ -1511,9 +1387,6 @@ async fn list_project_workflows( Path(id): Path, Query(query): Query, ) -> Result { - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; - let limit = clamp_limit(query.limit); let offset = query.offset.unwrap_or(0); @@ -1532,8 +1405,7 @@ async fn associate_project_workflow( State(state): State, Path((id, workflow_id)): Path<(Uuid, Uuid)>, ) -> Result { - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; + // Verify workflow exists. state .scheduler .storage() @@ -1553,10 +1425,9 @@ async fn associate_project_workflow( async fn dissociate_project_workflow( State(state): State, - Path((id, workflow_id)): Path<(Uuid, Uuid)>, + Path((_id, workflow_id)): Path<(Uuid, Uuid)>, ) -> Result { - let ps = state.manager.project_storage(); - ps.get(&id).await.map_err(ApiError::Internal)?.ok_or(ApiError::NotFound)?; + // Verify workflow exists. state .scheduler .storage() diff --git a/crates/orchestrator/src/client.rs b/crates/orchestrator/src/client.rs index af4cea5e..5c6b1142 100644 --- a/crates/orchestrator/src/client.rs +++ b/crates/orchestrator/src/client.rs @@ -37,9 +37,9 @@ use crate::types::{ AddDirRequest, AddDirResponse, AgentResponse, AgentUsageStats, ApprovalActionRequest, ClearContextRequest, ClearContextResponse, ConversationEventResponse, ConversationHistoryQuery, ConversationHistoryResponse, ConversationSummary, CreateAgentRequest, CreateProjectRequest, - HealthResponse, PaginatedResponse, PendingApproval, Project, ProjectResponse, - SendMessageRequest, SendMessageResponse, SetModelRequest, ToolPolicy, UpdateAgentRequest, - UpdateAgentResponse, UpdateProjectRequest, + HealthResponse, PaginatedResponse, PendingApproval, Project, SendMessageRequest, + SendMessageResponse, SetModelRequest, ToolPolicy, UpdateAgentRequest, UpdateAgentResponse, + UpdateProjectRequest, }; /// Typed HTTP client for the orchestrator service. @@ -58,6 +58,12 @@ use crate::types::{ pub struct OrchestratorClient { client: reqwest::Client, base_url: String, + /// Optional base URL for the core service (used for project CRUD). + /// + /// When set, project CRUD operations (`list_projects`, `create_project`, + /// `get_project`, `update_project`, `delete_project`) target this URL + /// instead of `base_url`. Association operations remain on `base_url`. + core_base_url: Option, token: Option, } @@ -72,7 +78,12 @@ impl OrchestratorClient { /// let client = OrchestratorClient::new("http://localhost:7006"); /// ``` pub fn new(base_url: impl Into) -> Self { - Self { client: reqwest::Client::new(), base_url: base_url.into(), token: None } + Self { + client: reqwest::Client::new(), + base_url: base_url.into(), + core_base_url: None, + token: None, + } } /// Attach a bearer token to all requests made by this client. @@ -83,6 +94,26 @@ impl OrchestratorClient { self } + /// Set the core service base URL for project CRUD operations. + /// + /// Project CRUD (`list_projects`, `create_project`, `get_project`, + /// `update_project`, `delete_project`) will target `{core_base_url}/projects` + /// instead of the orchestrator. Association operations remain on the + /// orchestrator `base_url`. + /// + /// # Examples + /// + /// ```ignore + /// use orchestrator::client::OrchestratorClient; + /// + /// let client = OrchestratorClient::new("http://localhost:17006") + /// .with_core_url("http://localhost:17000/api/v1"); + /// ``` + pub fn with_core_url(mut self, core_base_url: impl Into) -> Self { + self.core_base_url = Some(core_base_url.into()); + self + } + /// The base URL this client targets (e.g. the core gateway /// `{core_url}/api/v1/orchestrator`). /// @@ -408,37 +439,82 @@ impl OrchestratorClient { self.delete_with_response(&format!("/queues/{}", queue_name)).await } - // -- Project management -- + // -- Project management (CRUD targets core; associations stay on orchestrator) -- /// List all projects. + /// + /// Targets the core service when a core base URL is configured via + /// [`with_core_url`]; otherwise falls back to the orchestrator base URL. pub async fn list_projects(&self) -> Result> { - self.get("/projects").await + self.get_core("/projects").await } /// Create a new project. + /// + /// Targets the core service when a core base URL is configured. pub async fn create_project(&self, req: &CreateProjectRequest) -> Result { - self.post("/projects", req).await + self.post_core("/projects", req).await } - /// Get a project by UUID, including agent and workflow counts. - pub async fn get_project(&self, id: &Uuid) -> Result { - self.get(&format!("/projects/{id}")).await + /// Get a project by UUID. + /// + /// Targets the core service when a core base URL is configured. + /// Returns the project without agent/workflow counts (core does not + /// compute those; query the orchestrator association endpoints separately + /// if counts are needed). + pub async fn get_project(&self, id: &Uuid) -> Result { + self.get_core(&format!("/projects/{id}")).await } /// Find a project by name (client-side search). + /// + /// Targets the core service when a core base URL is configured. pub async fn get_project_by_name(&self, name: &str) -> Result> { - let resp: PaginatedResponse = self.get("/projects?limit=500").await?; + let resp: PaginatedResponse = self.get_core("/projects?limit=500").await?; Ok(resp.items.into_iter().find(|p| p.name == name)) } /// Update a project's name and/or description. + /// + /// Targets the core service when a core base URL is configured. pub async fn update_project(&self, id: &Uuid, req: &UpdateProjectRequest) -> Result { - self.put(&format!("/projects/{id}"), req).await + self.put_core(&format!("/projects/{id}"), req).await } - /// Delete a project (fails if agents or workflows are still associated). + /// Delete a project. + /// + /// Checks the orchestrator for active agent and workflow associations + /// before sending the DELETE to core, preventing orphaned references + /// (core does not enforce cross-service constraints at this layer). + /// Returns an error if any associations remain. + /// + /// Targets the core service when a core base URL is configured. pub async fn delete_project(&self, id: &Uuid) -> Result<()> { - self.delete(&format!("/projects/{id}")).await + // Guard: verify no agent associations remain on the orchestrator. + let agents = self + .list_project_agents(id) + .await + .context("Failed to check project agent associations before deletion")?; + if agents.total > 0 { + anyhow::bail!( + "cannot delete project {id}: {} agent(s) still associated \ + (dissociate them first with `project remove-agent`)", + agents.total + ); + } + // Guard: verify no workflow associations remain on the orchestrator. + let workflows = self + .list_project_workflows(id) + .await + .context("Failed to check project workflow associations before deletion")?; + if workflows.total > 0 { + anyhow::bail!( + "cannot delete project {id}: {} workflow(s) still associated \ + (dissociate them first with `project remove-workflow`)", + workflows.total + ); + } + self.delete_core(&format!("/projects/{id}")).await } /// List agents associated with a project. @@ -554,49 +630,126 @@ impl OrchestratorClient { } // -- Private HTTP helpers -- + // + // `url_for` is the single URL-computation function used by all helpers. + // Pass `use_core = true` to route to `core_base_url` (falling back to + // `base_url` when unset); `use_core = false` always uses `base_url`. + // + // The `*_core` wrappers are thin aliases that set `use_core = true` so + // call sites in the project-CRUD methods remain readable without repeating + // the flag everywhere. + + fn url_for(&self, path: &str, use_core: bool) -> String { + let base = if use_core { + self.core_base_url.as_deref().unwrap_or(&self.base_url) + } else { + &self.base_url + }; + format!("{base}{path}") + } + + // -- Core-routing wrappers (use_core = true) -- + + async fn get_core(&self, path: &str) -> Result { + self.get_url(self.url_for(path, true)).await + } + + async fn post_core( + &self, + path: &str, + body: &T, + ) -> Result { + self.post_url(self.url_for(path, true), body).await + } + + async fn put_core(&self, path: &str, body: &T) -> Result { + self.put_url(self.url_for(path, true), body).await + } + + async fn delete_core(&self, path: &str) -> Result<()> { + self.delete_url(self.url_for(path, true)).await + } + + // -- Orchestrator-routing wrappers (use_core = false) -- async fn get(&self, path: &str) -> Result { - let url = format!("{}{}", self.base_url, path); - let mut req = self.client.get(&url); + self.get_url(self.url_for(path, false)).await + } + + async fn post(&self, path: &str, body: &T) -> Result { + self.post_url(self.url_for(path, false), body).await + } + + async fn put(&self, path: &str, body: &T) -> Result { + self.put_url(self.url_for(path, false), body).await + } + + async fn patch(&self, path: &str, body: &T) -> Result { + let url = self.url_for(path, false); + let mut req = self.client.patch(&url).json(body); if let Some(t) = &self.token { req = req.bearer_auth(t); } - let response = req.send().await.context(format!("Failed to GET {url}"))?; + let response = req.send().await.context(format!("Failed to PATCH {url}"))?; Self::handle_response(response).await } - async fn post(&self, path: &str, body: &T) -> Result { - let url = format!("{}{}", self.base_url, path); - let mut req = self.client.post(&url).json(body); + async fn delete(&self, path: &str) -> Result<()> { + self.delete_url(self.url_for(path, false)).await + } + + async fn delete_with_body( + &self, + path: &str, + body: &T, + ) -> Result { + let url = self.url_for(path, false); + let mut req = self.client.delete(&url).json(body); if let Some(t) = &self.token { req = req.bearer_auth(t); } - let response = req.send().await.context(format!("Failed to POST {url}"))?; + let response = req.send().await.context(format!("Failed to DELETE {url}"))?; Self::handle_response(response).await } - async fn put(&self, path: &str, body: &T) -> Result { - let url = format!("{}{}", self.base_url, path); - let mut req = self.client.put(&url).json(body); + async fn delete_with_response(&self, path: &str) -> Result { + self.delete_with_response_url(self.url_for(path, false)).await + } + + // -- URL-based implementations (shared by both routing variants) -- + + async fn get_url(&self, url: String) -> Result { + let mut req = self.client.get(&url); if let Some(t) = &self.token { req = req.bearer_auth(t); } - let response = req.send().await.context(format!("Failed to PUT {url}"))?; + let response = req.send().await.context(format!("Failed to GET {url}"))?; Self::handle_response(response).await } - async fn patch(&self, path: &str, body: &T) -> Result { - let url = format!("{}{}", self.base_url, path); - let mut req = self.client.patch(&url).json(body); + async fn post_url( + &self, + url: String, + body: &T, + ) -> Result { + let mut req = self.client.post(&url).json(body); if let Some(t) = &self.token { req = req.bearer_auth(t); } - let response = req.send().await.context(format!("Failed to PATCH {url}"))?; + let response = req.send().await.context(format!("Failed to POST {url}"))?; Self::handle_response(response).await } - async fn delete(&self, path: &str) -> Result<()> { - let url = format!("{}{}", self.base_url, path); + async fn put_url(&self, url: String, body: &T) -> Result { + let mut req = self.client.put(&url).json(body); + if let Some(t) = &self.token { + req = req.bearer_auth(t); + } + let response = req.send().await.context(format!("Failed to PUT {url}"))?; + Self::handle_response(response).await + } + + async fn delete_url(&self, url: String) -> Result<()> { let mut req = self.client.delete(&url); if let Some(t) = &self.token { req = req.bearer_auth(t); @@ -611,22 +764,7 @@ impl OrchestratorClient { } } - async fn delete_with_body( - &self, - path: &str, - body: &T, - ) -> Result { - let url = format!("{}{}", self.base_url, path); - let mut req = self.client.delete(&url).json(body); - if let Some(t) = &self.token { - req = req.bearer_auth(t); - } - let response = req.send().await.context(format!("Failed to DELETE {url}"))?; - Self::handle_response(response).await - } - - async fn delete_with_response(&self, path: &str) -> Result { - let url = format!("{}{}", self.base_url, path); + async fn delete_with_response_url(&self, url: String) -> Result { let mut req = self.client.delete(&url); if let Some(t) = &self.token { req = req.bearer_auth(t); @@ -650,7 +788,21 @@ impl OrchestratorClient { mod tests { // reqwest::Client::new() triggers macOS system-configuration TLS // initialisation which panics when called from non-main test threads. - // These tests verify URL string handling without constructing the client. + // Tests that only need to verify URL string handling use a lightweight + // helper that skips reqwest construction. + + use super::*; + + /// Build a minimal `OrchestratorClient` for URL-computation tests only. + /// Does NOT construct a live reqwest client; safe to call from test threads. + fn url_only_client(base: &str, core: Option<&str>) -> OrchestratorClient { + OrchestratorClient { + client: reqwest::Client::new(), + base_url: base.to_string(), + core_base_url: core.map(|s| s.to_string()), + token: None, + } + } #[test] fn test_base_url_string_conversion() { @@ -670,4 +822,32 @@ mod tests { let url: String = String::from("http://localhost:7006"); assert_eq!(url, "http://localhost:7006"); } + + // -- url_for tests -- + + #[test] + fn test_url_for_uses_base_when_core_not_set() { + let c = url_only_client("http://orch", None); + // use_core = true falls back to base_url when core_base_url is None + assert_eq!(c.url_for("/projects", true), "http://orch/projects"); + // use_core = false always uses base_url + assert_eq!(c.url_for("/agents", false), "http://orch/agents"); + } + + #[test] + fn test_url_for_uses_core_base_when_set() { + let c = url_only_client("http://orch", Some("http://core/api/v1")); + // use_core = true targets core_base_url + assert_eq!(c.url_for("/projects", true), "http://core/api/v1/projects"); + // use_core = false still targets base_url (orchestrator) + assert_eq!(c.url_for("/agents", false), "http://orch/agents"); + } + + #[test] + fn test_with_core_url_builder_sets_core_base() { + // Verify the public builder wires up core_base_url correctly via url_for. + let c = url_only_client("http://orch", None); + let c = OrchestratorClient { core_base_url: Some("http://core/api/v1".into()), ..c }; + assert_eq!(c.url_for("/projects", true), "http://core/api/v1/projects"); + } } diff --git a/crates/orchestrator/src/entity/mod.rs b/crates/orchestrator/src/entity/mod.rs index a00cdd78..ecb690c3 100644 --- a/crates/orchestrator/src/entity/mod.rs +++ b/crates/orchestrator/src/entity/mod.rs @@ -3,7 +3,6 @@ pub mod agent; pub mod conversation_event; pub mod dispatch; -pub mod project; pub mod task_queue; pub mod usage_session; pub mod workflow; diff --git a/crates/orchestrator/src/manager.rs b/crates/orchestrator/src/manager.rs index 950dfcb3..efac9c5b 100644 --- a/crates/orchestrator/src/manager.rs +++ b/crates/orchestrator/src/manager.rs @@ -1,5 +1,5 @@ use crate::scheduler::events::SystemEvent; -use crate::storage::{AgentStorage, ProjectStorage}; +use crate::storage::AgentStorage; use crate::types::{ Agent, AgentConfig, AgentStatus, AgentUsageStats, ClearContextResponse, RetentionConfig, }; @@ -86,11 +86,6 @@ impl AgentManager { &self.registry } - /// Returns a [`ProjectStorage`] backed by the same database connection. - pub fn project_storage(&self) -> ProjectStorage { - ProjectStorage::from_db(self.storage.db().clone()) - } - /// Returns the underlying [`AgentStorage`] for direct access. pub fn agent_storage(&self) -> &AgentStorage { &self.storage diff --git a/crates/orchestrator/src/storage.rs b/crates/orchestrator/src/storage.rs index d23a937f..54657aed 100644 --- a/crates/orchestrator/src/storage.rs +++ b/crates/orchestrator/src/storage.rs @@ -6,13 +6,11 @@ use crate::{ entity::agent as agent_entity, entity::conversation_event as conv_entity, - entity::project as project_entity, entity::usage_session as session_entity, migration::Migrator, types::{ Agent, AgentConfig, AgentStatus, AgentUsageStats, ConversationEvent, ConversationQuery, - ConversationSummary, Project, SessionUsage, ToolPolicy, UpdateProjectRequest, - UsageSnapshot, + ConversationSummary, SessionUsage, ToolPolicy, UsageSnapshot, }, }; use anyhow::Result; @@ -1052,124 +1050,6 @@ impl AgentStorage { } } -// --------------------------------------------------------------------------- -// ProjectStorage -// --------------------------------------------------------------------------- - -/// Persistent storage backend for project records using SeaORM + SQLite. -/// -/// Shares the same [`DatabaseConnection`] as [`AgentStorage`] — construct via -/// [`ProjectStorage::from_db`] using the connection returned by -/// [`AgentStorage::db()`]. -#[derive(Clone)] -pub struct ProjectStorage { - db: DatabaseConnection, -} - -impl ProjectStorage { - /// Wrap an existing database connection (shared with `AgentStorage`). - pub fn from_db(db: DatabaseConnection) -> Self { - Self { db } - } - - /// Insert a new project and return it. - pub async fn create(&self, project: &Project) -> Result { - let model = project_entity::ActiveModel { - id: Set(project.id.to_string()), - name: Set(project.name.clone()), - description: Set(project.description.clone()), - created_at: Set(project.created_at.to_rfc3339()), - updated_at: Set(project.updated_at.to_rfc3339()), - organization_id: Set(project.organization_id.clone()), - }; - project_entity::Entity::insert(model).exec(&self.db).await?; - Ok(project.clone()) - } - - /// Retrieve a project by UUID. Returns `None` if not found. - pub async fn get(&self, id: &Uuid) -> Result> { - let model = project_entity::Entity::find_by_id(id.to_string()).one(&self.db).await?; - model.map(model_to_project).transpose() - } - - /// Retrieve a project by its unique name. Returns `None` if not found. - #[allow(dead_code)] - pub async fn get_by_name(&self, name: &str) -> Result> { - let model = project_entity::Entity::find() - .filter(project_entity::Column::Name.eq(name)) - .one(&self.db) - .await?; - model.map(model_to_project).transpose() - } - - /// List all projects ordered by creation time (newest first). - #[allow(dead_code)] - pub async fn list(&self) -> Result> { - self.list_org(None).await - } - - /// Like [`list`] but also filters by `org_id` when provided. - pub async fn list_org(&self, org_id: Option<&str>) -> Result> { - let mut query = - project_entity::Entity::find().order_by(project_entity::Column::CreatedAt, Order::Desc); - if let Some(oid) = org_id { - // Include legacy NULL rows so pre-migration data is still visible - // to authenticated tenants until backfill-tenant is run. - query = query.filter( - Condition::any() - .add(project_entity::Column::OrganizationId.eq(oid)) - .add(project_entity::Column::OrganizationId.is_null()), - ); - } - let models = query.all(&self.db).await?; - models.into_iter().map(model_to_project).collect() - } - - /// Apply an update request to the project identified by `id`. - /// - /// Only fields present in `req` are changed; `updated_at` is always - /// refreshed. Returns the updated project or errors if not found. - pub async fn update(&self, id: &Uuid, req: &UpdateProjectRequest) -> Result { - use sea_orm::sea_query::Expr; - - let now = Utc::now().to_rfc3339(); - - let mut update = project_entity::Entity::update_many() - .col_expr(project_entity::Column::UpdatedAt, Expr::value(&now)) - .filter(project_entity::Column::Id.eq(id.to_string())); - - if let Some(ref name) = req.name { - update = update.col_expr(project_entity::Column::Name, Expr::value(name.clone())); - } - if let Some(ref desc) = req.description { - update = - update.col_expr(project_entity::Column::Description, Expr::value(desc.clone())); - } - - let result = update.exec(&self.db).await?; - if result.rows_affected == 0 { - anyhow::bail!("Project not found"); - } - - self.get(id).await?.ok_or_else(|| anyhow::anyhow!("Project not found after update")) - } - - /// Delete a project by UUID. Errors if not found. - /// - /// Callers should verify no agents or workflows reference this project - /// before calling (enforced via FK in migration #828). - pub async fn delete(&self, id: &Uuid) -> Result<()> { - let result = project_entity::Entity::delete_many() - .filter(project_entity::Column::Id.eq(id.to_string())) - .exec(&self.db) - .await?; - if result.rows_affected == 0 { - anyhow::bail!("Project not found"); - } - Ok(()) - } -} - // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- @@ -1249,18 +1129,6 @@ fn model_to_session_usage(model: &session_entity::Model) -> Result }) } -/// Convert a raw [`project_entity::Model`] into the domain [`Project`] type. -fn model_to_project(model: project_entity::Model) -> Result { - Ok(Project { - id: Uuid::parse_str(&model.id)?, - name: model.name, - description: model.description, - created_at: DateTime::parse_from_rfc3339(&model.created_at)?.with_timezone(&Utc), - updated_at: DateTime::parse_from_rfc3339(&model.updated_at)?.with_timezone(&Utc), - organization_id: model.organization_id, - }) -} - #[allow(dead_code)] fn model_to_conversation_event(model: conv_entity::Model) -> Result { let metadata = model @@ -1818,162 +1686,6 @@ mod tests { assert_eq!(retrieved.config.docker_image, None); assert_eq!(retrieved.config.resource_limits, None); } - - // ----------------------------------------------------------------------- - // ProjectStorage tests - // ----------------------------------------------------------------------- - - fn project_storage(agent_storage: &AgentStorage) -> ProjectStorage { - ProjectStorage::from_db(agent_storage.db().clone()) - } - - #[tokio::test] - async fn test_project_create_and_get() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let project = Project::new("alpha".to_string(), Some("first project".to_string())); - let id = project.id; - - ps.create(&project).await.unwrap(); - - let retrieved = ps.get(&id).await.unwrap().unwrap(); - assert_eq!(retrieved.id, id); - assert_eq!(retrieved.name, "alpha"); - assert_eq!(retrieved.description, Some("first project".to_string())); - } - - #[tokio::test] - async fn test_project_get_returns_none_for_unknown_id() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let result = ps.get(&Uuid::new_v4()).await.unwrap(); - assert!(result.is_none()); - } - - #[tokio::test] - async fn test_project_get_by_name() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let project = Project::new("beta".to_string(), None); - ps.create(&project).await.unwrap(); - - let found = ps.get_by_name("beta").await.unwrap().unwrap(); - assert_eq!(found.id, project.id); - assert_eq!(found.description, None); - } - - #[tokio::test] - async fn test_project_get_by_name_missing() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - assert!(ps.get_by_name("nonexistent").await.unwrap().is_none()); - } - - #[tokio::test] - async fn test_project_list() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - ps.create(&Project::new("p1".to_string(), None)).await.unwrap(); - ps.create(&Project::new("p2".to_string(), None)).await.unwrap(); - ps.create(&Project::new("p3".to_string(), None)).await.unwrap(); - - let projects = ps.list().await.unwrap(); - assert_eq!(projects.len(), 3); - // Newest first - let names: Vec<&str> = projects.iter().map(|p| p.name.as_str()).collect(); - assert!(names.contains(&"p1")); - assert!(names.contains(&"p2")); - assert!(names.contains(&"p3")); - } - - #[tokio::test] - async fn test_project_list_empty() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - assert!(ps.list().await.unwrap().is_empty()); - } - - #[tokio::test] - async fn test_project_update_name() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let project = Project::new("old-name".to_string(), None); - ps.create(&project).await.unwrap(); - - let req = UpdateProjectRequest { name: Some("new-name".to_string()), description: None }; - let updated = ps.update(&project.id, &req).await.unwrap(); - - assert_eq!(updated.name, "new-name"); - assert_eq!(updated.description, None); - assert!(updated.updated_at >= project.updated_at); - } - - #[tokio::test] - async fn test_project_update_description() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let project = Project::new("proj".to_string(), None); - ps.create(&project).await.unwrap(); - - let req = - UpdateProjectRequest { name: None, description: Some("added description".to_string()) }; - let updated = ps.update(&project.id, &req).await.unwrap(); - - assert_eq!(updated.name, "proj"); - assert_eq!(updated.description, Some("added description".to_string())); - } - - #[tokio::test] - async fn test_project_update_not_found() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let req = UpdateProjectRequest { name: Some("x".to_string()), description: None }; - let result = ps.update(&Uuid::new_v4(), &req).await; - assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("not found")); - } - - #[tokio::test] - async fn test_project_delete() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let project = Project::new("to-delete".to_string(), None); - ps.create(&project).await.unwrap(); - assert!(ps.get(&project.id).await.unwrap().is_some()); - - ps.delete(&project.id).await.unwrap(); - assert!(ps.get(&project.id).await.unwrap().is_none()); - } - - #[tokio::test] - async fn test_project_delete_not_found() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - let result = ps.delete(&Uuid::new_v4()).await; - assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("not found")); - } - - #[tokio::test] - async fn test_project_name_unique() { - let (storage, _tmp) = create_test_storage().await; - let ps = project_storage(&storage); - - ps.create(&Project::new("unique".to_string(), None)).await.unwrap(); - let result = ps.create(&Project::new("unique".to_string(), None)).await; - assert!(result.is_err()); - } - // ----------------------------------------------------------------------- // Conversation event tests // ----------------------------------------------------------------------- diff --git a/crates/orchestrator/src/types.rs b/crates/orchestrator/src/types.rs index 7037c5c4..b0b4dc08 100644 --- a/crates/orchestrator/src/types.rs +++ b/crates/orchestrator/src/types.rs @@ -514,9 +514,15 @@ fn default_shell() -> String { // --------------------------------------------------------------------------- // Project types +// +// NOTE: Project CRUD has moved to the core service (epic #1306). These types +// are kept temporarily so that `client.rs` and the CLI can still compile while +// issue #1311 (repoint CLI/client to core) is in progress. Remove them once +// #1311 lands. // --------------------------------------------------------------------------- /// A project groups agents, workflows, and rooms under a named boundary. +#[allow(dead_code)] #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct Project { pub id: Uuid, @@ -529,6 +535,7 @@ pub struct Project { pub organization_id: Option, } +#[allow(dead_code)] impl Project { pub fn new(name: String, description: Option) -> Self { let now = Utc::now(); @@ -544,6 +551,7 @@ impl Project { } /// Detailed project response including associated resource counts. +#[allow(dead_code)] #[derive(Debug, Clone, Serialize, Deserialize)] pub struct ProjectResponse { pub id: Uuid, @@ -557,6 +565,7 @@ pub struct ProjectResponse { } /// Request body for POST /projects. +#[allow(dead_code)] #[derive(Debug, Clone, Serialize, Deserialize)] pub struct CreateProjectRequest { pub name: String, @@ -569,6 +578,7 @@ pub struct CreateProjectRequest { /// Fields that are `None` are left unchanged in the database. /// Pass `description: Some(None)` is not supported via this struct — /// to clear a description set it to `Some("")` or omit the field. +#[allow(dead_code)] #[derive(Debug, Clone, Serialize, Deserialize)] pub struct UpdateProjectRequest { #[serde(skip_serializing_if = "Option::is_none")] diff --git a/crates/orchestrator/tests/project_associations_http.rs b/crates/orchestrator/tests/project_associations_http.rs new file mode 100644 index 00000000..d995a070 --- /dev/null +++ b/crates/orchestrator/tests/project_associations_http.rs @@ -0,0 +1,414 @@ +//! Integration tests for project association endpoints. +//! +//! After project CRUD was moved to the core service (#1313), the orchestrator +//! retains four association endpoints: +//! +//! - `POST /projects/{id}/agents/{agent_id}` — associate agent with project +//! - `DELETE /projects/{id}/agents/{agent_id}` — dissociate agent from project +//! - `POST /projects/{id}/workflows/{wf_id}` — associate workflow with project +//! - `DELETE /projects/{id}/workflows/{wf_id}` — dissociate workflow from project +//! +//! The orchestrator verifies that the agent / workflow exists (returns 404 if +//! not) but does **not** cross-check with the core service to verify the +//! project ID — that validation is the caller's responsibility. + +use async_trait::async_trait; +use axum::{ + body::Body, + http::{Request, StatusCode}, +}; +use chrono::Utc; +use communicate::client::CommunicateClient; +use orchestrator::{ + api::{create_router, ApiState}, + manager::AgentManager, + scheduler::{ + storage::SchedulerStorage, + types::{TriggerConfig, WorkflowConfig}, + Scheduler, + }, + storage::AgentStorage, + types::{Agent, AgentConfig, AgentStatus}, + websocket::ConnectionRegistry, +}; +use std::sync::Arc; +use tempfile::TempDir; +use tower::ServiceExt; +use uuid::Uuid; +use wrap::{ + backend::{SessionConfig, SessionExitInfo, SessionHealth}, + types::BackendType, + ExecutionBackend, +}; + +// --------------------------------------------------------------------------- +// No-op backend +// --------------------------------------------------------------------------- + +struct NullBackend; + +#[async_trait] +impl ExecutionBackend for NullBackend { + async fn create_session(&self, _config: &SessionConfig) -> anyhow::Result<()> { + Ok(()) + } + async fn launch_agent(&self, _config: &SessionConfig) -> anyhow::Result<()> { + Ok(()) + } + async fn session_exists(&self, _session_name: &str) -> anyhow::Result { + Ok(false) + } + async fn kill_session(&self, _session_name: &str) -> anyhow::Result<()> { + Ok(()) + } + async fn send_command(&self, _session_name: &str, _command: &str) -> anyhow::Result<()> { + Ok(()) + } + async fn list_sessions(&self) -> anyhow::Result> { + Ok(vec![]) + } + fn prefix(&self) -> &str { + "test" + } + async fn session_health(&self, _session_name: &str) -> anyhow::Result { + Ok(SessionHealth::Unknown) + } + async fn session_exit_info( + &self, + _session_name: &str, + ) -> anyhow::Result> { + Ok(None) + } +} + +// --------------------------------------------------------------------------- +// Test app builder +// --------------------------------------------------------------------------- + +async fn build_app() -> (axum::Router, Arc, Arc, TempDir) { + let temp_dir = TempDir::new().unwrap(); + let db_path = temp_dir.path().join("test.db"); + + let storage = Arc::new(AgentStorage::with_path(&db_path).await.unwrap()); + let scheduler_storage = SchedulerStorage::new(storage.db().clone()); + let registry = ConnectionRegistry::new(); + let scheduler = Arc::new(Scheduler::new(scheduler_storage, registry.clone())); + let manager = Arc::new( + AgentManager::new( + storage.clone(), + Arc::new(NullBackend), + registry.clone(), + "ws://localhost:7006".to_string(), + ) + .with_mcp_config_dir(temp_dir.path().join("mcp")), + ); + let communicate = CommunicateClient::new("http://localhost:17010"); + + let state = ApiState { + manager, + registry, + scheduler: scheduler.clone(), + communicate, + backend_type: BackendType::Tmux, + }; + + (create_router(state), storage, scheduler, temp_dir) +} + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +/// Insert an agent with the given status directly into storage. +async fn insert_agent(storage: &AgentStorage, name: &str) -> Agent { + let config: AgentConfig = + serde_json::from_value(serde_json::json!({ "working_dir": "/tmp" })).unwrap(); + let mut agent = Agent::new(name.to_string(), config); + agent.status = AgentStatus::Pending; + storage.add(&agent).await.unwrap(); + agent +} + +/// Insert a workflow with a manual trigger directly into scheduler storage. +async fn insert_workflow(scheduler: &Scheduler, name: &str, agent_id: Uuid) -> WorkflowConfig { + let now = Utc::now(); + let config = WorkflowConfig { + id: Uuid::new_v4(), + name: name.to_string(), + agent_id, + trigger_config: TriggerConfig::Manual {}, + prompt_template: "Task: {{title}}".to_string(), + poll_interval_secs: 60, + enabled: true, + tool_policy: Default::default(), + created_at: now, + updated_at: now, + project_id: None, + organization_id: None, + }; + scheduler.storage().add_workflow(&config).await.unwrap(); + config +} + +// --------------------------------------------------------------------------- +// Agent association tests +// --------------------------------------------------------------------------- + +/// `POST /projects/{project_id}/agents/{agent_id}` with a valid agent returns 204. +#[tokio::test] +async fn test_associate_agent_with_project_returns_204() { + let (app, storage, _scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-assoc-agent-1").await; + let project_id = Uuid::new_v4(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/agents/{}", agent.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NO_CONTENT); +} + +/// `DELETE /projects/{project_id}/agents/{agent_id}` with a valid agent returns 204. +#[tokio::test] +async fn test_dissociate_agent_from_project_returns_204() { + let (app, storage, _scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-assoc-agent-2").await; + let project_id = Uuid::new_v4(); + + // Associate first, then dissociate. + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/agents/{}", agent.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/projects/{project_id}/agents/{}", agent.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NO_CONTENT); +} + +/// `POST /projects/{project_id}/agents/{agent_id}` with an unknown agent_id returns 404. +#[tokio::test] +async fn test_associate_agent_nonexistent_agent_returns_404() { + let (app, _storage, _scheduler, _tmp) = build_app().await; + let project_id = Uuid::new_v4(); + let missing_agent_id = Uuid::new_v4(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/agents/{missing_agent_id}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); +} + +/// Associating an agent with a non-existent project UUID still returns 204 +/// because the orchestrator does not validate project IDs against the core +/// service — that check is the caller's responsibility. +#[tokio::test] +async fn test_associate_agent_nonexistent_project_still_succeeds() { + let (app, storage, _scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-assoc-agent-3").await; + let nonexistent_project_id = Uuid::new_v4(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{nonexistent_project_id}/agents/{}", agent.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + // 204 — orchestrator only validates that the agent exists, not the project. + assert_eq!(response.status(), StatusCode::NO_CONTENT); +} + +// --------------------------------------------------------------------------- +// Workflow association tests +// --------------------------------------------------------------------------- + +/// `POST /projects/{project_id}/workflows/{wf_id}` with a valid workflow returns 204. +#[tokio::test] +async fn test_associate_workflow_with_project_returns_204() { + let (app, storage, scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-wf-assoc-agent-1").await; + let workflow = insert_workflow(&scheduler, "orch-wf-1", agent.id).await; + let project_id = Uuid::new_v4(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/workflows/{}", workflow.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NO_CONTENT); +} + +/// `DELETE /projects/{project_id}/workflows/{wf_id}` with a valid workflow returns 204. +#[tokio::test] +async fn test_dissociate_workflow_from_project_returns_204() { + let (app, storage, scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-wf-assoc-agent-2").await; + let workflow = insert_workflow(&scheduler, "orch-wf-2", agent.id).await; + let project_id = Uuid::new_v4(); + + // Associate first. + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/workflows/{}", workflow.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("DELETE") + .uri(format!("/projects/{project_id}/workflows/{}", workflow.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NO_CONTENT); +} + +/// `POST /projects/{project_id}/workflows/{wf_id}` with an unknown workflow returns 404. +#[tokio::test] +async fn test_associate_workflow_nonexistent_workflow_returns_404() { + let (app, _storage, _scheduler, _tmp) = build_app().await; + let project_id = Uuid::new_v4(); + let missing_wf_id = Uuid::new_v4(); + + let response = app + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/workflows/{missing_wf_id}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::NOT_FOUND); +} + +// --------------------------------------------------------------------------- +// List association tests +// --------------------------------------------------------------------------- + +/// After associating an agent, `GET /projects/{id}/agents` returns that agent. +#[tokio::test] +async fn test_list_project_agents_returns_associated_agents() { + let (app, storage, _scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-list-agent").await; + let project_id = Uuid::new_v4(); + + // Associate the agent. + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/agents/{}", agent.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri(format!("/projects/{project_id}/agents")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await.unwrap(); + let body: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + assert_eq!(body["total"], 1); + assert_eq!(body["items"][0]["id"], agent.id.to_string()); +} + +/// After associating a workflow, `GET /projects/{id}/workflows` returns that workflow. +#[tokio::test] +async fn test_list_project_workflows_returns_associated_workflows() { + let (app, storage, scheduler, _tmp) = build_app().await; + let agent = insert_agent(&storage, "orch-list-wf-agent").await; + let workflow = insert_workflow(&scheduler, "orch-list-wf", agent.id).await; + let project_id = Uuid::new_v4(); + + // Associate the workflow. + app.clone() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/projects/{project_id}/workflows/{}", workflow.id)) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + let response = app + .oneshot( + Request::builder() + .method("GET") + .uri(format!("/projects/{project_id}/workflows")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX).await.unwrap(); + let body: serde_json::Value = serde_json::from_slice(&bytes).unwrap(); + assert_eq!(body["total"], 1); + assert_eq!(body["items"][0]["id"], workflow.id.to_string()); +} diff --git a/docs/storage.md b/docs/storage.md index 0f2abea3..1a516749 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -852,6 +852,28 @@ let scheduler_storage = SchedulerStorage::new(agent_storage.db().clone()); --- +## Service-to-Table Ownership + +Each agentd service owns its own SQLite database. The table below maps key +entities to their canonical service after the v0.15.0 project migration: + +| Entity | Owning service | Database | Notes | +|--------|---------------|----------|-------| +| `users` | agentd-core | `agentd-core/core.db` | | +| `organizations` | agentd-core | `agentd-core/core.db` | | +| `projects` | agentd-core | `agentd-core/core.db` | Moved from orchestrator in v0.15.0 | +| `agents` | agentd-orchestrator | `agentd/agent.db` | `project_id` FK points to core `projects.id` | +| `workflows` | agentd-orchestrator | `agentd/agent.db` | `project_id` FK points to core `projects.id` | +| `notifications` | agentd-notify | `agentd-notify/notify.db` | | +| `memories` | agentd-memory | `agentd-memory/memory.db` | | + +> **Upgrade note (v0.15.0):** The `projects` table was removed from the +> orchestrator database and added to the core database. Run +> `agent admin backfill-projects` after upgrading to assign `organization_id` +> to project rows that were created before multi-tenancy was introduced. + +--- + ## xtask Commands Three `cargo xtask` sub-commands help manage databases during development: diff --git a/ui/src/components/knowledge/ProjectPicker.tsx b/ui/src/components/knowledge/ProjectPicker.tsx index f664ef6a..3f483bb6 100644 --- a/ui/src/components/knowledge/ProjectPicker.tsx +++ b/ui/src/components/knowledge/ProjectPicker.tsx @@ -6,7 +6,7 @@ import { ChevronDown, FolderOpen } from "lucide-react"; import { useEffect, useRef, useState } from "react"; import type { Project } from "@/types/orchestrator"; -import { orchestratorClient } from "@/services/orchestrator"; +import { coreClient } from "@/services/core"; interface ProjectPickerProps { selectedId: string | null; @@ -27,7 +27,7 @@ export function ProjectPicker({ useEffect(() => { let cancelled = false; setLoading(true); - orchestratorClient + coreClient .listProjects({ limit: 200 }) .then((page) => { if (!cancelled) setProjects(page.items); diff --git a/ui/src/services/core.ts b/ui/src/services/core.ts new file mode 100644 index 00000000..8e8f89cf --- /dev/null +++ b/ui/src/services/core.ts @@ -0,0 +1,32 @@ +/** + * Client for the Core service (default port 17000). + * + * The core service is the API gateway and owns the canonical project entity. + * Project CRUD operations are served from `/api/v1/projects`. + */ + +import type { PaginatedResponse } from "@/types/common"; +import type { ListProjectsParams, Project } from "@/types/orchestrator"; +import { ApiClient, withAuth } from "./base"; +import { serviceConfig } from "./config"; + +export class CoreClient extends ApiClient { + // ------------------------------------------------------------------------- + // Projects + // ------------------------------------------------------------------------- + + /** `GET /api/v1/projects` — list all projects. */ + listProjects( + params?: ListProjectsParams, + ): Promise> { + return this.get>( + "/api/v1/projects", + params as Record, + ); + } +} + +/** Singleton client instance using the configured core service URL */ +export const coreClient = new CoreClient( + withAuth({ baseUrl: serviceConfig.coreServiceUrl }), +);