From 17f664014482d17d1d77cc8065a4bd25e3e65020 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:09:44 -0700 Subject: [PATCH 01/18] epic: Phase 1 - Stand up projects entity in core service From b705a8a26af07a26e9a0ecacc6bd9ed9289a3295 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:10:14 -0700 Subject: [PATCH 02/18] feat: add project entity, migration, and storage to core crate From ab5f982ba833eebdda9f1a284e376f30cd7b1e10 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:12:56 -0700 Subject: [PATCH 03/18] feat(core): add project entity, migration, and storage Adds the canonical `projects` entity to the core crate following the existing `organizations` pattern: - `entity/project.rs`: SeaORM Model with id, name, description, organization_id, created_at, updated_at fields - `migration/m20260616_000004_create_projects_table.rs`: creates the projects table with unique index on name and index on organization_id - `project_storage.rs`: ProjectStorage with create/get_by_id/get_by_name/ list/list_org/list_paginated/update/delete; list_org includes NULL org rows for tenant-transition compatibility - Registers project in entity/mod.rs, migration/mod.rs, storage.rs, and lib.rs Closes #1308 --- crates/core/src/entity/mod.rs | 4 +- crates/core/src/entity/project.rs | 29 ++ crates/core/src/lib.rs | 1 + .../m20260616_000004_create_projects_table.rs | 72 +++++ crates/core/src/migration/mod.rs | 2 + crates/core/src/project_storage.rs | 296 ++++++++++++++++++ crates/core/src/storage.rs | 6 + 7 files changed, 409 insertions(+), 1 deletion(-) create mode 100644 crates/core/src/entity/project.rs create mode 100644 crates/core/src/migration/m20260616_000004_create_projects_table.rs create mode 100644 crates/core/src/project_storage.rs 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()) From 280c41e9240751962e19abe3a5124b88e2f4fc3f Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:14:25 -0700 Subject: [PATCH 04/18] feat: add project CRUD API routes to core crate From 9d6ba58a4376afa3fdf286c26a791cd3c2fd5236 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:17:01 -0700 Subject: [PATCH 05/18] feat(core): add project CRUD API routes (address review feedback) Add AuthUser extractor to all five handlers to enforce authentication on every endpoint. Without it any unauthenticated caller could read, create, update, and delete projects. Also add empty-name validation to update_project_handler (mirrors the guard already present in create_project_handler) so whitespace-only names are rejected with 400 on PUT as well as POST. Update all 19 tests to include Bearer tokens obtained via /auth/register and add explicit 401-unauthenticated tests for each HTTP method. Note: in-memory pagination in list_projects_handler (list_org + skip/take) is a known scaling cliff flagged by the reviewer; a follow-up list_org_paginated storage method is needed. Closes #1309 --- crates/core/src/api.rs | 7 + crates/core/src/api/projects.rs | 850 ++++++++++++++++++++++++++++++++ 2 files changed, 857 insertions(+) create mode 100644 crates/core/src/api/projects.rs 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..4569ffeb --- /dev/null +++ b/crates/core/src/api/projects.rs @@ -0,0 +1,850 @@ +//! 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 { 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); + } + + // ----------------------------------------------------------------------- + // 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); + } +} From d07161ee66f0524e88db47a16806182a4ec370d5 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Fri, 19 Jun 2026 10:35:41 -0700 Subject: [PATCH 06/18] epic: Add backfill-projects admin CLI command From 8dd7904975243891446bf76bfdd6f4026e97ea9a Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Fri, 19 Jun 2026 10:36:09 -0700 Subject: [PATCH 07/18] feat: add backfill-projects admin CLI command From 421b06441777c228274732b88b133f7db15275cd Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Fri, 19 Jun 2026 10:42:32 -0700 Subject: [PATCH 08/18] feat(admin): add backfill-projects CLI command (address review feedback) --- crates/cli/src/commands/admin.rs | 473 ++++++++++++++++++++++++++++++- 1 file changed, 472 insertions(+), 1 deletion(-) diff --git a/crates/cli/src/commands/admin.rs b/crates/cli/src/commands/admin.rs index aa92da63..812f983c 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 + } } } } @@ -239,6 +281,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 +674,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"); + } } From c55d2499788e6c3213509ba80d385d095de581e9 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:26:00 -0700 Subject: [PATCH 09/18] epic: Remove project CRUD from orchestrator, keep association endpoints From a1c45c755fadd41ba9b5e0c81f2a06467f1ffd5f Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Thu, 18 Jun 2026 17:39:15 -0700 Subject: [PATCH 10/18] feat(orchestrator): remove project CRUD, keep association endpoints - Delete ProjectStorage struct and all methods from storage.rs - Delete model_to_project helper - Remove project_entity import from storage.rs - Remove project_storage() accessor from manager.rs - Remove GET/POST /projects and GET/PUT/DELETE /projects/{id} routes - Remove list_projects, create_project, get_project, update_project, delete_project handlers - Remove ps.get() existence checks from all association handlers - Remove pub mod project from entity/mod.rs (unused after dropping checks) - Add allow(dead_code) to project types in types.rs (kept for client.rs compat until #1311 lands) - Update CLI admin SCOPED_TABLES: move projects from orchestrator to core --- crates/cli/src/commands/admin.rs | 3 +- crates/orchestrator/src/api.rs | 147 +------------ crates/orchestrator/src/entity/mod.rs | 1 - crates/orchestrator/src/manager.rs | 7 +- crates/orchestrator/src/storage.rs | 290 +------------------------- crates/orchestrator/src/types.rs | 10 + 6 files changed, 23 insertions(+), 435 deletions(-) diff --git a/crates/cli/src/commands/admin.rs b/crates/cli/src/commands/admin.rs index 812f983c..ddc9d295 100644 --- a/crates/cli/src/commands/admin.rs +++ b/crates/cli/src/commands/admin.rs @@ -124,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"] }, 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/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")] From 88f8ddcc2ae673c2c8ad8d36abe18c8fd3ee0855 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Sun, 21 Jun 2026 09:46:43 -0700 Subject: [PATCH 11/18] epic: Phase 4 - Repoint project consumers to core From c77961437833ad075f4ae9bb7d0f17d623d5946b Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Sun, 21 Jun 2026 09:46:48 -0700 Subject: [PATCH 12/18] feat: repoint Rust client and CLI project commands to core From b986853a16e71459bc5fcf25935c41532ccaac37 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Sun, 21 Jun 2026 09:55:30 -0700 Subject: [PATCH 13/18] feat(orchestrator,cli): repoint project CRUD to core service --- crates/cli/src/commands/project.rs | 2 - crates/cli/src/main.rs | 12 ++- crates/core/src/api/projects.rs | 2 +- crates/orchestrator/src/client.rs | 131 ++++++++++++++++++++++++++--- 4 files changed, 127 insertions(+), 20 deletions(-) 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/core/src/api/projects.rs b/crates/core/src/api/projects.rs index 4569ffeb..925f8139 100644 --- a/crates/core/src/api/projects.rs +++ b/crates/core/src/api/projects.rs @@ -270,7 +270,7 @@ mod tests { async fn test_app() -> (Router, tempfile::TempDir) { let (conn, tmp) = create_test_connection().await; let storage = Storage::new(conn).await.unwrap(); - let state = AppState { storage }; + let state = AppState::new(storage); let app = crate::api::create_router(state); (app, tmp) } diff --git a/crates/orchestrator/src/client.rs b/crates/orchestrator/src/client.rs index af4cea5e..92b2577b 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,53 @@ 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). + /// + /// 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 + self.delete_core(&format!("/projects/{id}")).await } /// List agents associated with a project. @@ -553,6 +600,64 @@ impl OrchestratorClient { self.get(&format!("/agents/{agent_id}/conversation/{event_id}")).await } + // -- Core-routing HTTP helpers -- + // These helpers target `core_base_url` when set, falling back to `base_url`. + + fn core_url(&self, path: &str) -> String { + let base = self.core_base_url.as_deref().unwrap_or(&self.base_url); + format!("{base}{path}") + } + + async fn get_core(&self, path: &str) -> Result { + let url = self.core_url(path); + 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 GET {url}"))?; + Self::handle_response(response).await + } + + async fn post_core( + &self, + path: &str, + body: &T, + ) -> Result { + let url = self.core_url(path); + 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 POST {url}"))?; + Self::handle_response(response).await + } + + async fn put_core(&self, path: &str, body: &T) -> Result { + let url = self.core_url(path); + 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_core(&self, path: &str) -> Result<()> { + let url = self.core_url(path); + let mut req = self.client.delete(&url); + if let Some(t) = &self.token { + req = req.bearer_auth(t); + } + let response = req.send().await.context(format!("Failed to DELETE {url}"))?; + if response.status().is_success() { + Ok(()) + } else { + let status = response.status(); + let error_text = response.text().await.unwrap_or_default(); + Err(anyhow::anyhow!("Request failed with status {status}: {error_text}")) + } + } + // -- Private HTTP helpers -- async fn get(&self, path: &str) -> Result { From 60ac86185aef846596238061938472d27ac4e5a5 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Sun, 21 Jun 2026 09:56:41 -0700 Subject: [PATCH 14/18] feat: repoint MCP tools and UI project calls to core From 069295c5f1eeb1e73fcdaffa6c6c6556dd603ef1 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Sun, 21 Jun 2026 10:01:26 -0700 Subject: [PATCH 15/18] feat(mcp,ui): repoint MCP tools and UI project calls to core --- crates/common/src/config.rs | 12 +++++++ crates/mcp/src/client.rs | 5 +++ crates/mcp/src/config.rs | 9 ++++++ crates/mcp/src/tools/orchestrator_debug.rs | 16 ++++------ crates/mcp/tests/common/mod.rs | 1 + ui/src/components/knowledge/ProjectPicker.tsx | 4 +-- ui/src/services/core.ts | 32 +++++++++++++++++++ 7 files changed, 67 insertions(+), 12 deletions(-) create mode 100644 ui/src/services/core.ts 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/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/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 }), +); From cc6ec74564184f75d89d33a85e73b096547f2b98 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Mon, 22 Jun 2026 10:11:38 -0700 Subject: [PATCH 16/18] fix(orchestrator): address review feedback on delete_project and helpers - Fix delete_project docstring: core does not enforce association constraints server-side, so add a client-side pre-check that queries list_project_agents and list_project_workflows before deleting - Replace core_url(path) helper with url_for(path, use_core: bool) to reduce duplication; extract URL-based implementations (get_url, post_url, put_url, delete_url, delete_with_response_url) that both *_core and regular wrappers delegate to - Add unit tests for url_for fallback logic using a struct-literal url_only_client helper to avoid reqwest TLS init in test threads Co-Authored-By: Claude Sonnet 4.6 --- crates/orchestrator/src/client.rs | 213 ++++++++++++++++++++---------- 1 file changed, 144 insertions(+), 69 deletions(-) diff --git a/crates/orchestrator/src/client.rs b/crates/orchestrator/src/client.rs index 92b2577b..5c6b1142 100644 --- a/crates/orchestrator/src/client.rs +++ b/crates/orchestrator/src/client.rs @@ -481,10 +481,39 @@ impl OrchestratorClient { 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<()> { + // 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 } @@ -600,22 +629,29 @@ impl OrchestratorClient { self.get(&format!("/agents/{agent_id}/conversation/{event_id}")).await } - // -- Core-routing HTTP helpers -- - // These helpers target `core_base_url` when set, falling back to `base_url`. - - fn core_url(&self, path: &str) -> String { - let base = self.core_base_url.as_deref().unwrap_or(&self.base_url); + // -- 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 { - let url = self.core_url(path); - 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 GET {url}"))?; - Self::handle_response(response).await + self.get_url(self.url_for(path, true)).await } async fn post_core( @@ -623,45 +659,66 @@ impl OrchestratorClient { path: &str, body: &T, ) -> Result { - let url = self.core_url(path); - 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 POST {url}"))?; - Self::handle_response(response).await + self.post_url(self.url_for(path, true), body).await } async fn put_core(&self, path: &str, body: &T) -> Result { - let url = self.core_url(path); - let mut req = self.client.put(&url).json(body); + 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 { + 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 PUT {url}"))?; + let response = req.send().await.context(format!("Failed to PATCH {url}"))?; Self::handle_response(response).await } - async fn delete_core(&self, path: &str) -> Result<()> { - let url = self.core_url(path); - let mut req = self.client.delete(&url); + 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 DELETE {url}"))?; - if response.status().is_success() { - Ok(()) - } else { - let status = response.status(); - let error_text = response.text().await.unwrap_or_default(); - Err(anyhow::anyhow!("Request failed with status {status}: {error_text}")) - } + Self::handle_response(response).await } - // -- Private HTTP helpers -- + async fn delete_with_response(&self, path: &str) -> Result { + self.delete_with_response_url(self.url_for(path, false)).await + } - async fn get(&self, path: &str) -> Result { - let url = format!("{}{}", self.base_url, path); + // -- 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); @@ -670,8 +727,11 @@ impl OrchestratorClient { Self::handle_response(response).await } - async fn post(&self, path: &str, body: &T) -> Result { - let url = format!("{}{}", self.base_url, path); + 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); @@ -680,8 +740,7 @@ impl OrchestratorClient { Self::handle_response(response).await } - async fn put(&self, path: &str, body: &T) -> 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); @@ -690,18 +749,7 @@ impl OrchestratorClient { 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); - if let Some(t) = &self.token { - req = req.bearer_auth(t); - } - let response = req.send().await.context(format!("Failed to PATCH {url}"))?; - Self::handle_response(response).await - } - - async fn delete(&self, path: &str) -> Result<()> { - let url = format!("{}{}", self.base_url, path); + 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); @@ -716,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); @@ -755,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() { @@ -775,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"); + } } From 404e413bba0086b17c5d8620f3892572ac214bb4 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Mon, 22 Jun 2026 15:22:11 -0700 Subject: [PATCH 17/18] epic: Add project migration tests and update documentation From d494d4176d48c85a94681dd47932e8e5541215d0 Mon Sep 17 00:00:00 2001 From: Geoff Johnson Date: Mon, 22 Jun 2026 15:29:47 -0700 Subject: [PATCH 18/18] feat(tests,docs): add project migration tests and update documentation Core handler tests: - Add test_list_projects_with_tenant_header_scopes_results to verify X-Tenant-ID scoping on GET /api/v1/projects (org-filtered + NULL rows vs. unscoped list returning all) Orchestrator association tests (new file): - crates/orchestrator/tests/project_associations_http.rs: 9 tests covering POST/DELETE agent and workflow associations, 404 on missing agent/workflow, and behavior when project ID does not exist in core CLI integration tests: - Add project_tests module to crates/cli/tests/integration_test.rs with mockito tests verifying CRUD routes to core URL, association routes to orchestrator URL, and delete_project pre-check behavior Documentation: - README: add agentd-core to port table (17000/7000), update description to mention project management, add project management CLI entry, add v0.15.0 migration status section with operator upgrade note - docs/storage.md: add service-to-table ownership table with note that projects moved from orchestrator to core in v0.15.0 Co-Authored-By: Claude Sonnet 4.6 --- README.md | 12 +- crates/cli/tests/integration_test.rs | 208 +++++++++ crates/core/src/api/projects.rs | 61 +++ .../tests/project_associations_http.rs | 414 ++++++++++++++++++ docs/storage.md | 22 + 5 files changed, 715 insertions(+), 2 deletions(-) create mode 100644 crates/orchestrator/tests/project_associations_http.rs 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/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/core/src/api/projects.rs b/crates/core/src/api/projects.rs index 925f8139..db2f186a 100644 --- a/crates/core/src/api/projects.rs +++ b/crates/core/src/api/projects.rs @@ -643,6 +643,67 @@ mod tests { 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 // ----------------------------------------------------------------------- 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: