diff --git a/compliance-agent/src/agent.rs b/compliance-agent/src/agent.rs index 92d12f8..fac9035 100644 --- a/compliance-agent/src/agent.rs +++ b/compliance-agent/src/agent.rs @@ -63,21 +63,12 @@ impl ComplianceAgent { let db = self.db_pool.for_tenant_id(tenant_id).await?; let orchestrator = PipelineOrchestrator::new(self.config.clone(), db, self.llm.clone(), self.http.clone()); - if self.config.unified_pipeline { - orchestrator.run_target(repo_id, trigger).await - } else { - orchestrator.run(repo_id, trigger).await - } + orchestrator.run_target(repo_id, trigger).await } - /// Run a scan for an onboarded target through the unified pipeline, - /// unconditionally. - /// - /// Unlike [`Self::run_scan`], this does *not* consult the - /// `unified_pipeline` transition flag: the caller (the `/targets/{id}/scan` - /// endpoint) operates on `onboarded_targets` by construction, so it must - /// always dispatch to `run_target` regardless of how the legacy paths - /// (scheduler, webhooks, `/repositories/{id}/scan`) are configured. + /// Alias for [`Self::run_scan`] — every scan runs the unified onboarded-target + /// pipeline. Kept as a distinct name for the `/targets/{id}/scan` endpoint's + /// intent. pub async fn run_target_scan( &self, tenant_id: &str, @@ -100,16 +91,19 @@ impl ComplianceAgent { head_sha: &str, ) -> Result<(), crate::error::AgentError> { let db = self.db_pool.for_tenant_id(tenant_id).await?; - let repo = db - .repositories() - .find_one(mongodb::bson::doc! { - "_id": mongodb::bson::oid::ObjectId::parse_str(repo_id) - .map_err(|e| crate::error::AgentError::Other(e.to_string()))? - }) + let oid = mongodb::bson::oid::ObjectId::parse_str(repo_id) + .map_err(|e| crate::error::AgentError::Other(e.to_string()))?; + let target = db + .onboarded_targets() + .find_one(mongodb::bson::doc! { "_id": oid }) .await? .ok_or_else(|| { - crate::error::AgentError::Other(format!("Repository {repo_id} not found")) + crate::error::AgentError::Other(format!("Target {repo_id} not found")) })?; + let code = target.code_artifact().ok_or_else(|| { + crate::error::AgentError::Other(format!("Target {repo_id} has no code artifact")) + })?; + let repo = crate::pipeline::repo_view::RepoView::from_target(&target, code); let orchestrator = PipelineOrchestrator::new(self.config.clone(), db, self.llm.clone(), self.http.clone()); diff --git a/compliance-agent/src/api/handlers/chat.rs b/compliance-agent/src/api/handlers/chat.rs index 9583ed8..a402cd2 100644 --- a/compliance-agent/src/api/handlers/chat.rs +++ b/compliance-agent/src/api/handlers/chat.rs @@ -146,7 +146,7 @@ pub async fn build_embeddings( let agent_clone = (*agent).clone(); tokio::spawn(async move { let repo = match db - .repositories() + .onboarded_targets() .find_one(doc! { "_id": mongodb::bson::oid::ObjectId::parse_str(&repo_id).ok() }) .await { @@ -194,14 +194,22 @@ pub async fn build_embeddings( } }; + let code = match repo.code_artifact() { + Some(c) => c, + None => { + tracing::error!("Target {repo_id} has no code artifact for embedding build"); + return; + } + }; + let view = crate::pipeline::repo_view::RepoView::from_target(&repo, code); let creds = crate::pipeline::git::RepoCredentials { ssh_key_path: Some(agent_clone.config.ssh_key_path.clone()), - auth_token: repo.auth_token.clone(), - auth_username: repo.auth_username.clone(), + auth_token: view.auth_token.clone(), + auth_username: view.auth_username.clone(), }; let git_ops = crate::pipeline::git::GitOps::new(&agent_clone.config.git_clone_base_path, creds); - let repo_path = match git_ops.clone_or_fetch(&repo.git_url, &repo.name) { + let repo_path = match git_ops.clone_or_fetch(&view.git_url, &view.name) { Ok(p) => p, Err(e) => { tracing::error!("Failed to clone repo for embedding build: {e}"); diff --git a/compliance-agent/src/api/handlers/graph.rs b/compliance-agent/src/api/handlers/graph.rs index 8c10c9c..937efd3 100644 --- a/compliance-agent/src/api/handlers/graph.rs +++ b/compliance-agent/src/api/handlers/graph.rs @@ -255,7 +255,7 @@ pub async fn get_file_content( // Look up the repository to get repo name let repo = db - .repositories() + .onboarded_targets() .find_one(doc! { "_id": mongodb::bson::oid::ObjectId::parse_str(&repo_id).ok() }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)? @@ -317,7 +317,7 @@ pub async fn trigger_build( let agent_clone = (*agent).clone(); tokio::spawn(async move { let repo = match db - .repositories() + .onboarded_targets() .find_one(doc! { "_id": mongodb::bson::oid::ObjectId::parse_str(&repo_id).ok() }) .await { @@ -328,14 +328,22 @@ pub async fn trigger_build( } }; + let code = match repo.code_artifact() { + Some(c) => c, + None => { + tracing::error!("Target {repo_id} has no code artifact for graph build"); + return; + } + }; + let view = crate::pipeline::repo_view::RepoView::from_target(&repo, code); let creds = crate::pipeline::git::RepoCredentials { ssh_key_path: Some(agent_clone.config.ssh_key_path.clone()), - auth_token: repo.auth_token.clone(), - auth_username: repo.auth_username.clone(), + auth_token: view.auth_token.clone(), + auth_username: view.auth_username.clone(), }; let git_ops = crate::pipeline::git::GitOps::new(&agent_clone.config.git_clone_base_path, creds); - let repo_path = match git_ops.clone_or_fetch(&repo.git_url, &repo.name) { + let repo_path = match git_ops.clone_or_fetch(&view.git_url, &view.name) { Ok(p) => p, Err(e) => { tracing::error!("Failed to clone repo for graph build: {e}"); diff --git a/compliance-agent/src/api/handlers/health.rs b/compliance-agent/src/api/handlers/health.rs index ef3b030..fbdfd26 100644 --- a/compliance-agent/src/api/handlers/health.rs +++ b/compliance-agent/src/api/handlers/health.rs @@ -10,6 +10,18 @@ pub async fn health() -> Json { Json(serde_json::json!({ "status": "ok" })) } +/// GET /api/v1/settings/ssh-public-key — the agent's SSH deploy public key, +/// for adding as a read-only deploy key on private git targets. +#[tracing::instrument(skip_all)] +pub async fn get_ssh_public_key( + axum::extract::Extension(agent): AgentExt, +) -> Result, axum::http::StatusCode> { + let public_path = format!("{}.pub", agent.config.ssh_key_path); + let public_key = + std::fs::read_to_string(&public_path).map_err(|_| axum::http::StatusCode::NOT_FOUND)?; + Ok(Json(serde_json::json!({ "public_key": public_key.trim() }))) +} + #[tracing::instrument(skip_all)] pub async fn stats_overview( axum::extract::Extension(agent): AgentExt, @@ -19,7 +31,7 @@ pub async fn stats_overview( let db = &db; let total_repositories = db - .repositories() + .onboarded_targets() .count_documents(doc! {}) .await .unwrap_or(0); diff --git a/compliance-agent/src/api/handlers/mod.rs b/compliance-agent/src/api/handlers/mod.rs index 8b4631c..9478e63 100644 --- a/compliance-agent/src/api/handlers/mod.rs +++ b/compliance-agent/src/api/handlers/mod.rs @@ -12,7 +12,6 @@ pub mod notifications; pub mod onboarding; pub mod pentest_handlers; pub use pentest_handlers as pentest; -pub mod repos; pub mod sbom; pub mod scans; @@ -21,6 +20,5 @@ pub use dto::*; pub use findings::*; pub use health::*; pub use issues::*; -pub use repos::*; pub use sbom::*; pub use scans::*; diff --git a/compliance-agent/src/api/handlers/onboarding.rs b/compliance-agent/src/api/handlers/onboarding.rs index 55ab297..03eb635 100644 --- a/compliance-agent/src/api/handlers/onboarding.rs +++ b/compliance-agent/src/api/handlers/onboarding.rs @@ -246,15 +246,116 @@ pub async fn delete_target( .delete_one(doc! { "_id": oid }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; - // Cascade the collections keyed by repo_id == target id (best-effort). - let by_repo = doc! { "repo_id": &id }; - let _ = db.findings().delete_many(by_repo.clone()).await; - let _ = db.scan_runs().delete_many(by_repo.clone()).await; - let _ = db.sbom_entries().delete_many(by_repo.clone()).await; - let _ = db.cve_alerts().delete_many(by_repo).await; + // Cascade all data keyed by repo_id == target id (best-effort). + let db = &db; + let _ = db.findings().delete_many(doc! { "repo_id": &id }).await; + let _ = db.sbom_entries().delete_many(doc! { "repo_id": &id }).await; + let _ = db.scan_runs().delete_many(doc! { "repo_id": &id }).await; + let _ = db.cve_alerts().delete_many(doc! { "repo_id": &id }).await; + let _ = db + .tracker_issues() + .delete_many(doc! { "repo_id": &id }) + .await; + let _ = db.graph_nodes().delete_many(doc! { "repo_id": &id }).await; + let _ = db.graph_edges().delete_many(doc! { "repo_id": &id }).await; + let _ = db.graph_builds().delete_many(doc! { "repo_id": &id }).await; + let _ = db + .impact_analyses() + .delete_many(doc! { "repo_id": &id }) + .await; + let _ = db + .code_embeddings() + .delete_many(doc! { "repo_id": &id }) + .await; + let _ = db + .embedding_builds() + .delete_many(doc! { "repo_id": &id }) + .await; + + // DAST targets linked to this target, and all their downstream data. + if let Ok(mut cursor) = db.dast_targets().find(doc! { "repo_id": &id }).await { + use futures_util::StreamExt; + while let Some(Ok(dt)) = cursor.next().await { + let dast_target_id = dt.id.map(|oid| oid.to_hex()).unwrap_or_default(); + if !dast_target_id.is_empty() { + cascade_delete_dast_target(db, &dast_target_id).await; + } + } + } + + // Pentest sessions linked directly to this target (not via a DAST target). + if let Ok(mut cursor) = db.pentest_sessions().find(doc! { "repo_id": &id }).await { + use futures_util::StreamExt; + while let Some(Ok(session)) = cursor.next().await { + let session_id = session.id.map(|oid| oid.to_hex()).unwrap_or_default(); + if !session_id.is_empty() { + let _ = db + .attack_chain_nodes() + .delete_many(doc! { "session_id": &session_id }) + .await; + let _ = db + .pentest_messages() + .delete_many(doc! { "session_id": &session_id }) + .await; + let _ = db + .dast_findings() + .delete_many(doc! { "session_id": &session_id }) + .await; + } + } + } + let _ = db + .pentest_sessions() + .delete_many(doc! { "repo_id": &id }) + .await; + Ok(Json(serde_json::json!({ "status": "deleted" }))) } +/// Delete a DAST target and everything downstream of it (pentest sessions + +/// their attack chains / messages / findings, DAST scan runs + findings). +async fn cascade_delete_dast_target(db: &crate::database::Database, target_id: &str) { + use futures_util::StreamExt; + if let Ok(mut cursor) = db + .pentest_sessions() + .find(doc! { "target_id": target_id }) + .await + { + while let Some(Ok(session)) = cursor.next().await { + let session_id = session.id.map(|oid| oid.to_hex()).unwrap_or_default(); + if !session_id.is_empty() { + let _ = db + .attack_chain_nodes() + .delete_many(doc! { "session_id": &session_id }) + .await; + let _ = db + .pentest_messages() + .delete_many(doc! { "session_id": &session_id }) + .await; + let _ = db + .dast_findings() + .delete_many(doc! { "session_id": &session_id }) + .await; + } + } + } + let _ = db + .pentest_sessions() + .delete_many(doc! { "target_id": target_id }) + .await; + let _ = db + .dast_findings() + .delete_many(doc! { "target_id": target_id }) + .await; + let _ = db + .dast_scan_runs() + .delete_many(doc! { "target_id": target_id }) + .await; + if let Ok(oid) = mongodb::bson::oid::ObjectId::parse_str(target_id) { + let _ = db.dast_targets().delete_one(doc! { "_id": oid }).await; + } +} + /// POST /api/v1/targets/{id}/artifacts — attach an artifact (by reference). #[tracing::instrument(skip_all, fields(target_id = %id))] pub async fn add_artifact( diff --git a/compliance-agent/src/api/handlers/pentest_handlers/session.rs b/compliance-agent/src/api/handlers/pentest_handlers/session.rs index 91c8c58..579d781 100644 --- a/compliance-agent/src/api/handlers/pentest_handlers/session.rs +++ b/compliance-agent/src/api/handlers/pentest_handlers/session.rs @@ -113,14 +113,14 @@ pub async fn create_session( session.config = Some(config.clone()); session.repo_id = target.repo_id.clone(); - // Resolve repo_id from git_repo_url if provided + // Resolve repo_id (target id) from git_repo_url if provided if let Some(ref git_url) = config.git_repo_url { - if let Ok(Some(repo)) = db - .repositories() - .find_one(doc! { "git_url": git_url }) + if let Ok(Some(target)) = db + .onboarded_targets() + .find_one(doc! { "artifacts.source_ref": git_url }) .await { - session.repo_id = repo.id.map(|oid| oid.to_hex()); + session.repo_id = target.id.map(|oid| oid.to_hex()); } } @@ -380,17 +380,20 @@ pub async fn lookup_repo( ) -> Result>, StatusCode> { let db = tenant_db(&agent, &tenant).await?; let repo = db - .repositories() - .find_one(doc! { "git_url": ¶ms.url }) + .onboarded_targets() + .find_one(doc! { "artifacts.source_ref": ¶ms.url }) .await .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; let data = match repo { - Some(r) => serde_json::json!({ - "name": r.name, - "default_branch": r.default_branch, - "last_scanned_commit": r.last_scanned_commit, - }), + Some(r) => { + let git = r.code_artifact().and_then(|c| c.git.as_ref()); + serde_json::json!({ + "name": r.name, + "default_branch": git.map(|g| g.default_branch.clone()), + "last_scanned_commit": git.and_then(|g| g.last_scanned_commit.clone()), + }) + } None => serde_json::Value::Null, }; diff --git a/compliance-agent/src/api/handlers/repos.rs b/compliance-agent/src/api/handlers/repos.rs deleted file mode 100644 index 5c58167..0000000 --- a/compliance-agent/src/api/handlers/repos.rs +++ /dev/null @@ -1,339 +0,0 @@ -use axum::extract::{Extension, Path, Query}; -use axum::http::StatusCode; -use axum::Json; -use mongodb::bson::doc; - -use super::dto::*; -use compliance_core::models::*; -use compliance_core::tenant_ctx::TenantCtx; - -#[tracing::instrument(skip_all)] -pub async fn list_repositories( - Extension(agent): AgentExt, - tenant: TenantCtx, - Query(params): Query, -) -> ApiResult> { - let db = tenant_db(&agent, &tenant).await?; - let db = &db; - let skip = (params.page.saturating_sub(1)) * params.limit as u64; - let total = db - .repositories() - .count_documents(doc! {}) - .await - .unwrap_or(0); - - let repos = match db - .repositories() - .find(doc! {}) - .skip(skip) - .limit(params.limit) - .await - { - Ok(cursor) => collect_cursor_async(cursor).await, - Err(e) => { - tracing::warn!("Failed to fetch repositories: {e}"); - Vec::new() - } - }; - - Ok(Json(ApiResponse { - data: repos, - total: Some(total), - page: Some(params.page), - })) -} - -#[tracing::instrument(skip_all)] -pub async fn add_repository( - Extension(agent): AgentExt, - tenant: TenantCtx, - Json(req): Json, -) -> Result>, (StatusCode, String)> { - // Validate repository access before saving - let creds = crate::pipeline::git::RepoCredentials { - ssh_key_path: Some(agent.config.ssh_key_path.clone()), - auth_token: req.auth_token.clone(), - auth_username: req.auth_username.clone(), - }; - - if let Err(e) = crate::pipeline::git::GitOps::test_access(&req.git_url, &creds) { - return Err(( - StatusCode::BAD_REQUEST, - format!("Cannot access repository: {e}"), - )); - } - - let mut repo = TrackedRepository::new(req.name, req.git_url); - repo.default_branch = req.default_branch; - repo.auth_token = req.auth_token; - repo.auth_username = req.auth_username; - repo.tracker_type = req.tracker_type; - repo.tracker_owner = req.tracker_owner; - repo.tracker_repo = req.tracker_repo; - repo.tracker_token = req.tracker_token; - repo.scan_schedule = req.scan_schedule; - - let db = tenant_db(&agent, &tenant) - .await - .map_err(|s| (s, "failed to acquire tenant database".to_string()))?; - db.repositories().insert_one(&repo).await.map_err(|_| { - ( - StatusCode::CONFLICT, - "Repository already exists".to_string(), - ) - })?; - - Ok(Json(ApiResponse { - data: repo, - total: None, - page: None, - })) -} - -#[tracing::instrument(skip_all, fields(repo_id = %id))] -pub async fn update_repository( - Extension(agent): AgentExt, - tenant: TenantCtx, - Path(id): Path, - Json(req): Json, -) -> Result, StatusCode> { - let oid = mongodb::bson::oid::ObjectId::parse_str(&id).map_err(|_| StatusCode::BAD_REQUEST)?; - let db = tenant_db(&agent, &tenant).await?; - - let mut set_doc = doc! { "updated_at": mongodb::bson::DateTime::now() }; - - if let Some(name) = &req.name { - set_doc.insert("name", name); - } - if let Some(branch) = &req.default_branch { - set_doc.insert("default_branch", branch); - } - if let Some(token) = &req.auth_token { - set_doc.insert("auth_token", token); - } - if let Some(username) = &req.auth_username { - set_doc.insert("auth_username", username); - } - if let Some(tracker_type) = &req.tracker_type { - set_doc.insert("tracker_type", tracker_type.to_string()); - } - if let Some(owner) = &req.tracker_owner { - set_doc.insert("tracker_owner", owner); - } - if let Some(repo) = &req.tracker_repo { - set_doc.insert("tracker_repo", repo); - } - if let Some(token) = &req.tracker_token { - set_doc.insert("tracker_token", token); - } - if let Some(schedule) = &req.scan_schedule { - set_doc.insert("scan_schedule", schedule); - } - - let result = db - .repositories() - .update_one(doc! { "_id": oid }, doc! { "$set": set_doc }) - .await - .map_err(|e| { - tracing::warn!("Failed to update repository: {e}"); - StatusCode::INTERNAL_SERVER_ERROR - })?; - - if result.matched_count == 0 { - return Err(StatusCode::NOT_FOUND); - } - - Ok(Json(serde_json::json!({ "status": "updated" }))) -} - -#[tracing::instrument(skip_all)] -pub async fn get_ssh_public_key( - Extension(agent): AgentExt, -) -> Result, StatusCode> { - let public_path = format!("{}.pub", agent.config.ssh_key_path); - let public_key = std::fs::read_to_string(&public_path).map_err(|_| StatusCode::NOT_FOUND)?; - Ok(Json(serde_json::json!({ "public_key": public_key.trim() }))) -} - -#[tracing::instrument(skip_all, fields(repo_id = %id))] -pub async fn trigger_scan( - Extension(agent): AgentExt, - tenant: TenantCtx, - Path(id): Path, -) -> Result, StatusCode> { - let agent_clone = (*agent).clone(); - let tenant_id = tenant.0.tenant_id.clone(); - tokio::spawn(async move { - if let Err(e) = agent_clone - .run_scan(&tenant_id, &id, ScanTrigger::Manual) - .await - { - tracing::error!("Manual scan failed for {id}: {e}"); - } - }); - - Ok(Json(serde_json::json!({ "status": "scan_triggered" }))) -} - -/// Return the webhook secret for a repository (used by dashboard to display it) -pub async fn get_webhook_config( - Extension(agent): AgentExt, - tenant: TenantCtx, - Path(id): Path, -) -> Result, StatusCode> { - let oid = mongodb::bson::oid::ObjectId::parse_str(&id).map_err(|_| StatusCode::BAD_REQUEST)?; - let db = tenant_db(&agent, &tenant).await?; - let repo = db - .repositories() - .find_one(doc! { "_id": oid }) - .await - .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)? - .ok_or(StatusCode::NOT_FOUND)?; - - let tracker_type = repo - .tracker_type - .as_ref() - .map(|t| t.to_string()) - .unwrap_or_else(|| "gitea".to_string()); - - Ok(Json(serde_json::json!({ - "webhook_secret": repo.webhook_secret, - "tracker_type": tracker_type, - }))) -} - -#[tracing::instrument(skip_all, fields(repo_id = %id))] -pub async fn delete_repository( - Extension(agent): AgentExt, - tenant: TenantCtx, - Path(id): Path, -) -> Result, StatusCode> { - let oid = mongodb::bson::oid::ObjectId::parse_str(&id).map_err(|_| StatusCode::BAD_REQUEST)?; - let db = tenant_db(&agent, &tenant).await?; - let db = &db; - - // Delete the repository - let result = db - .repositories() - .delete_one(doc! { "_id": oid }) - .await - .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?; - - if result.deleted_count == 0 { - return Err(StatusCode::NOT_FOUND); - } - - // Cascade delete all related data - let _ = db.findings().delete_many(doc! { "repo_id": &id }).await; - let _ = db.sbom_entries().delete_many(doc! { "repo_id": &id }).await; - let _ = db.scan_runs().delete_many(doc! { "repo_id": &id }).await; - let _ = db.cve_alerts().delete_many(doc! { "repo_id": &id }).await; - let _ = db - .tracker_issues() - .delete_many(doc! { "repo_id": &id }) - .await; - let _ = db.graph_nodes().delete_many(doc! { "repo_id": &id }).await; - let _ = db.graph_edges().delete_many(doc! { "repo_id": &id }).await; - let _ = db.graph_builds().delete_many(doc! { "repo_id": &id }).await; - let _ = db - .impact_analyses() - .delete_many(doc! { "repo_id": &id }) - .await; - let _ = db - .code_embeddings() - .delete_many(doc! { "repo_id": &id }) - .await; - let _ = db - .embedding_builds() - .delete_many(doc! { "repo_id": &id }) - .await; - - // Cascade delete DAST targets linked to this repo, and all their downstream data - // (scan runs, findings, pentest sessions, attack chains, messages) - if let Ok(mut cursor) = db.dast_targets().find(doc! { "repo_id": &id }).await { - use futures_util::StreamExt; - while let Some(Ok(target)) = cursor.next().await { - let target_id = target.id.map(|oid| oid.to_hex()).unwrap_or_default(); - if !target_id.is_empty() { - cascade_delete_dast_target(db, &target_id).await; - } - } - } - - // Also delete pentest sessions linked directly to this repo (not via target) - if let Ok(mut cursor) = db.pentest_sessions().find(doc! { "repo_id": &id }).await { - use futures_util::StreamExt; - while let Some(Ok(session)) = cursor.next().await { - let session_id = session.id.map(|oid| oid.to_hex()).unwrap_or_default(); - if !session_id.is_empty() { - let _ = db - .attack_chain_nodes() - .delete_many(doc! { "session_id": &session_id }) - .await; - let _ = db - .pentest_messages() - .delete_many(doc! { "session_id": &session_id }) - .await; - // Delete DAST findings produced by this session - let _ = db - .dast_findings() - .delete_many(doc! { "session_id": &session_id }) - .await; - } - } - } - let _ = db - .pentest_sessions() - .delete_many(doc! { "repo_id": &id }) - .await; - - Ok(Json(serde_json::json!({ "status": "deleted" }))) -} - -/// Cascade-delete a DAST target and all its downstream data. -async fn cascade_delete_dast_target(db: &crate::database::Database, target_id: &str) { - // Delete pentest sessions for this target (and their attack chains + messages) - if let Ok(mut cursor) = db - .pentest_sessions() - .find(doc! { "target_id": target_id }) - .await - { - use futures_util::StreamExt; - while let Some(Ok(session)) = cursor.next().await { - let session_id = session.id.map(|oid| oid.to_hex()).unwrap_or_default(); - if !session_id.is_empty() { - let _ = db - .attack_chain_nodes() - .delete_many(doc! { "session_id": &session_id }) - .await; - let _ = db - .pentest_messages() - .delete_many(doc! { "session_id": &session_id }) - .await; - let _ = db - .dast_findings() - .delete_many(doc! { "session_id": &session_id }) - .await; - } - } - } - let _ = db - .pentest_sessions() - .delete_many(doc! { "target_id": target_id }) - .await; - - // Delete DAST scan runs and their findings - let _ = db - .dast_findings() - .delete_many(doc! { "target_id": target_id }) - .await; - let _ = db - .dast_scan_runs() - .delete_many(doc! { "target_id": target_id }) - .await; - - // Delete the target itself - if let Ok(oid) = mongodb::bson::oid::ObjectId::parse_str(target_id) { - let _ = db.dast_targets().delete_one(doc! { "_id": oid }).await; - } -} diff --git a/compliance-agent/src/api/routes.rs b/compliance-agent/src/api/routes.rs index 3f1cf07..bf73519 100644 --- a/compliance-agent/src/api/routes.rs +++ b/compliance-agent/src/api/routes.rs @@ -11,20 +11,6 @@ pub fn build_router() -> Router { "/api/v1/settings/ssh-public-key", get(handlers::get_ssh_public_key), ) - .route("/api/v1/repositories", get(handlers::list_repositories)) - .route("/api/v1/repositories", post(handlers::add_repository)) - .route( - "/api/v1/repositories/{id}/scan", - post(handlers::trigger_scan), - ) - .route( - "/api/v1/repositories/{id}", - delete(handlers::delete_repository).patch(handlers::update_repository), - ) - .route( - "/api/v1/repositories/{id}/webhook-config", - get(handlers::get_webhook_config), - ) // Unified onboarding targets (#131). .route( "/api/v1/targets", diff --git a/compliance-agent/src/config.rs b/compliance-agent/src/config.rs index ad86563..a754975 100644 --- a/compliance-agent/src/config.rs +++ b/compliance-agent/src/config.rs @@ -47,12 +47,6 @@ pub fn load_config() -> Result { .unwrap_or_else(|| "/tmp/compliance-scanner/repos".to_string()), artifact_store_base_path: env_var_opt("ARTIFACT_STORE_BASE_PATH") .unwrap_or_else(|| "/data/compliance-scanner/artifacts".to_string()), - // Defaults ON: the unified onboarded-target pipeline is now the primary - // path (no legacy `repositories` data in production). Set - // `UNIFIED_PIPELINE=0` to fall back to the legacy repository pipeline. - unified_pipeline: env_var_opt("UNIFIED_PIPELINE") - .map(|v| v == "1" || v.eq_ignore_ascii_case("true")) - .unwrap_or(true), ssh_key_path: env_var_opt("SSH_KEY_PATH") .unwrap_or_else(|| "/data/compliance-scanner/ssh/id_ed25519".to_string()), keycloak_url: env_var_opt("KEYCLOAK_URL"), diff --git a/compliance-agent/src/database.rs b/compliance-agent/src/database.rs index cdadf3b..e10f3bf 100644 --- a/compliance-agent/src/database.rs +++ b/compliance-agent/src/database.rs @@ -249,16 +249,6 @@ impl Database { } pub async fn ensure_indexes(&self) -> Result<(), AgentError> { - // repositories: unique git_url - self.repositories() - .create_index( - IndexModel::builder() - .keys(doc! { "git_url": 1 }) - .options(IndexOptions::builder().unique(true).build()) - .build(), - ) - .await?; - // findings: unique fingerprint self.findings() .create_index( @@ -479,10 +469,6 @@ impl Database { Ok(()) } - pub fn repositories(&self) -> Collection { - self.inner.collection("repositories") - } - pub fn findings(&self) -> Collection { self.inner.collection("findings") } diff --git a/compliance-agent/src/lib.rs b/compliance-agent/src/lib.rs index a0a07bf..2cc5179 100644 --- a/compliance-agent/src/lib.rs +++ b/compliance-agent/src/lib.rs @@ -8,7 +8,6 @@ pub mod database; pub mod error; pub mod ingest; pub mod llm; -pub mod migrate; pub mod pentest; pub mod pipeline; pub mod rag; diff --git a/compliance-agent/src/main.rs b/compliance-agent/src/main.rs index 95323a9..110634f 100644 --- a/compliance-agent/src/main.rs +++ b/compliance-agent/src/main.rs @@ -1,50 +1,4 @@ -use compliance_agent::{agent, api, config, database, migrate, scheduler, ssh, webhooks}; - -/// Run the `migrate onboarding` subcommand and exit. Backfills (or reverts) the -/// unified `onboarded_targets` collection per tenant. -/// -/// Usage: `compliance-agent migrate onboarding [--all | --tenant ] [--dry-run] [--revert]` -async fn run_migration( - args: &[String], - pool: &database::DatabasePool, -) -> Result<(), compliance_agent::error::AgentError> { - if args.get(2).map(String::as_str) != Some("onboarding") { - eprintln!( - "usage: compliance-agent migrate onboarding [--all | --tenant ] [--dry-run] [--revert]" - ); - std::process::exit(2); - } - let has = |flag: &str| args.iter().any(|a| a == flag); - let dry_run = has("--dry-run"); - let revert = has("--revert"); - let tenant = args - .iter() - .position(|a| a == "--tenant") - .and_then(|i| args.get(i + 1)) - .cloned(); - - let tenants: Vec = if has("--all") { - pool.list_tenant_ids().await? - } else if let Some(t) = tenant { - vec![t] - } else { - eprintln!("specify --all or --tenant "); - std::process::exit(2); - }; - - for tenant_id in tenants { - let db = pool.for_tenant_id(&tenant_id).await?; - if revert { - migrate::onboarding::revert(&db).await?; - println!("[{tenant_id}] reverted onboarding backfill"); - } else { - let report = migrate::onboarding::backfill_onboarded_targets(&db, dry_run).await?; - let prefix = if dry_run { "(dry-run) " } else { "" }; - println!("[{tenant_id}] {prefix}{report:?}"); - } - } - Ok(()) -} +use compliance_agent::{agent, api, config, database, scheduler, ssh, webhooks}; #[tokio::main] async fn main() -> Result<(), Box> { @@ -77,13 +31,6 @@ async fn main() -> Result<(), Box> { let db_pool = database::DatabasePool::connect(&config.mongodb_uri, &config.mongodb_database).await?; - // One-shot subcommands run and exit without starting the servers. - let args: Vec = std::env::args().collect(); - if args.get(1).map(String::as_str) == Some("migrate") { - run_migration(&args, &db_pool).await?; - return Ok(()); - } - let agent = agent::ComplianceAgent::new(config.clone(), db_pool); tracing::info!("Starting scheduler..."); diff --git a/compliance-agent/src/migrate/mod.rs b/compliance-agent/src/migrate/mod.rs deleted file mode 100644 index 2546602..0000000 --- a/compliance-agent/src/migrate/mod.rs +++ /dev/null @@ -1,8 +0,0 @@ -//! One-time data migrations. -//! -//! Currently just the onboarding backfill ([`onboarding`]), which folds the -//! legacy `repositories` and `dast_targets` collections into the unified -//! `onboarded_targets` collection, preserving `_id` so every downstream record -//! keyed by `repo_id` / `target_id` keeps resolving. - -pub mod onboarding; diff --git a/compliance-agent/src/migrate/onboarding.rs b/compliance-agent/src/migrate/onboarding.rs deleted file mode 100644 index 7652316..0000000 --- a/compliance-agent/src/migrate/onboarding.rs +++ /dev/null @@ -1,406 +0,0 @@ -//! Backfill: legacy `repositories` + `dast_targets` → `onboarded_targets`. -//! -//! The transforms here are **id-preserving**: an [`OnboardedTarget`] keeps the -//! same `_id` as the `TrackedRepository` / `DastTarget` it came from, so every -//! downstream collection keyed by that hex id (findings, sbom, scan_runs, -//! graph, dast_*, pentest_*) keeps resolving with zero row rewrites, and -//! existing webhook URLs keep working. The mapping functions are pure and unit -//! tested; the DB orchestration (idempotent per-tenant backfill + revert) is a -//! thin driver over them. - -use compliance_core::models::{ - Artifact, ArtifactKind, DastTarget, DastTargetType, GitArtifactConfig, IssueTrackerConfig, - OnboardedTarget, TargetType, TrackedRepository, WebArtifactConfig, -}; -use futures_util::TryStreamExt; -use mongodb::bson::{doc, Document}; - -use crate::database::Database; -use crate::error::AgentError; - -/// Marker id in `schema_migrations` recording that the backfill has run. -const MIGRATION_MARKER: &str = "onboarding_backfill_v1"; - -/// Summary of a backfill run. -#[derive(Debug, Clone, Default, PartialEq, Eq)] -pub struct MigrationReport { - /// Repositories turned into onboarded targets. - pub repos_migrated: u64, - /// DAST targets folded into an existing (repo-linked) target as a LiveUrl. - pub dast_targets_folded: u64, - /// DAST targets with no repo link, migrated as standalone targets. - pub dast_targets_standalone: u64, - /// Records skipped because a target with that `_id` already existed. - pub skipped_existing: u64, -} - -/// Map a legacy `DastTargetType` to a unified [`TargetType`]. REST/GraphQL APIs -/// are backend services; a browser app is a web app. -fn target_type_for_dast(kind: &DastTargetType) -> TargetType { - match kind { - DastTargetType::WebApp => TargetType::WebApp, - DastTargetType::RestApi | DastTargetType::GraphQl => TargetType::BackendService, - } -} - -/// Build the LiveUrl artifact for a DAST target (its base URL + crawl config + -/// auth). Shared by fold-in and standalone migration. -pub fn dast_to_artifact(dast: &DastTarget) -> Artifact { - let mut artifact = Artifact::live_url(dast.base_url.clone()); - artifact.web = Some(WebArtifactConfig { - target_kind: dast.target_type.clone(), - excluded_paths: dast.excluded_paths.clone(), - max_crawl_depth: dast.max_crawl_depth, - rate_limit: dast.rate_limit, - allow_destructive: dast.allow_destructive, - }); - artifact.auth = dast.auth_config.clone().map(Into::into); - artifact -} - -/// Map a `TrackedRepository` to an onboarded target, preserving `_id`. The git -/// remote becomes a `GitRepo` artifact carrying the repo's branch, watermark, -/// and auth; tracker config folds into `scan_config`. -/// -/// `target_type` is a safe default (`BackendService`) — the classifier can -/// refine it later; `classification` is left `None` (unconfirmed). -pub fn repo_to_target(repo: &TrackedRepository) -> OnboardedTarget { - let mut target = OnboardedTarget::new(repo.name.clone(), TargetType::BackendService); - target.id = repo.id; - - let mut artifact = Artifact::git_repo(repo.git_url.clone(), repo.default_branch.clone()); - artifact.git = Some(GitArtifactConfig { - default_branch: repo.default_branch.clone(), - last_scanned_commit: repo.last_scanned_commit.clone(), - local_path: repo.local_path.clone(), - }); - if repo.auth_token.is_some() || repo.auth_username.is_some() { - artifact.auth = Some(compliance_core::models::ArtifactAuth { - method: "token".to_string(), - username: repo.auth_username.clone(), - secret: repo.auth_token.clone(), - ..Default::default() - }); - } - target.artifacts.push(artifact); - - if repo.tracker_type.is_some() { - target.scan_config.issue_tracker = Some(IssueTrackerConfig { - tracker_type: repo.tracker_type.clone(), - owner: repo.tracker_owner.clone(), - repo: repo.tracker_repo.clone(), - token: repo.tracker_token.clone(), - }); - } - - target.scan_schedule = repo.scan_schedule.clone(); - target.webhook_enabled = repo.webhook_enabled; - target.webhook_secret = repo.webhook_secret.clone(); - target.findings_count = repo.findings_count; - target.created_at = repo.created_at; - target.updated_at = repo.updated_at; - target -} - -/// Append a DAST target's LiveUrl artifact onto an existing (repo-derived) -/// target. If the repo default was `BackendService` but the DAST target is a -/// browser web app, promote the type to `WebApp`. -pub fn fold_dast_into_target(target: &mut OnboardedTarget, dast: &DastTarget) { - if matches!(dast.target_type, DastTargetType::WebApp) - && target.target_type == TargetType::BackendService - { - target.target_type = TargetType::WebApp; - } - if !target.has(ArtifactKind::LiveUrl) { - target.artifacts.push(dast_to_artifact(dast)); - } -} - -/// Map a repo-less DAST target to a standalone onboarded target, preserving `_id`. -pub fn dast_to_standalone_target(dast: &DastTarget) -> OnboardedTarget { - let mut target = - OnboardedTarget::new(dast.name.clone(), target_type_for_dast(&dast.target_type)); - target.id = dast.id; - target.artifacts.push(dast_to_artifact(dast)); - target.created_at = dast.created_at; - target.updated_at = dast.updated_at; - target -} - -/// Whether the onboarding backfill has already been applied to this database. -pub async fn already_applied(db: &Database) -> Result { - let found = db - .collection_named::("schema_migrations") - .find_one(doc! { "_id": MIGRATION_MARKER }) - .await?; - Ok(found.is_some()) -} - -/// Backfill `onboarded_targets` from `repositories` + `dast_targets` for one -/// tenant database. -/// -/// Id-preserving and **idempotent**: targets that already exist (by `_id`) are -/// skipped, so re-running is safe. With `dry_run`, computes the report without -/// writing. The legacy collections are never deleted; the only mutation outside -/// `onboarded_targets` is the history relink of folded DAST targets, which is -/// logged so [`revert`] can undo it. -pub async fn backfill_onboarded_targets( - db: &Database, - dry_run: bool, -) -> Result { - let mut report = MigrationReport::default(); - - // 1. repositories -> onboarded_targets (preserve _id, skip existing). - let mut repos = db.repositories().find(doc! {}).await?; - while let Some(repo) = repos.try_next().await? { - let Some(id) = repo.id else { continue }; - if db - .onboarded_targets() - .find_one(doc! { "_id": id }) - .await? - .is_some() - { - report.skipped_existing += 1; - continue; - } - if !dry_run { - db.onboarded_targets() - .insert_one(repo_to_target(&repo)) - .await?; - } - report.repos_migrated += 1; - } - - // 2. dast_targets -> fold into the linked repo target, or migrate standalone. - let mut dasts = db.dast_targets().find(doc! {}).await?; - while let Some(dast) = dasts.try_next().await? { - let Some(dast_id) = dast.id else { continue }; - let repo_oid = dast - .repo_id - .as_deref() - .and_then(|r| mongodb::bson::oid::ObjectId::parse_str(r).ok()); - let linked = match repo_oid { - Some(oid) => db.onboarded_targets().find_one(doc! { "_id": oid }).await?, - None => None, - }; - - match (linked, repo_oid) { - // Fold into an existing repo-derived target. - (Some(mut target), Some(oid)) => { - if target.has(ArtifactKind::LiveUrl) { - report.skipped_existing += 1; // already folded on a prior run - continue; - } - fold_dast_into_target(&mut target, &dast); - if !dry_run { - db.onboarded_targets() - .replace_one(doc! { "_id": oid }, &target) - .await?; - relink_history(db, &dast_id.to_hex(), &oid.to_hex()).await?; - } - report.dast_targets_folded += 1; - } - // No linked repo target: migrate as a standalone target (keeps _id). - _ => { - if db - .onboarded_targets() - .find_one(doc! { "_id": dast_id }) - .await? - .is_some() - { - report.skipped_existing += 1; - continue; - } - if !dry_run { - db.onboarded_targets() - .insert_one(dast_to_standalone_target(&dast)) - .await?; - } - report.dast_targets_standalone += 1; - } - } - } - - if !dry_run { - db.collection_named::("schema_migrations") - .update_one( - doc! { "_id": MIGRATION_MARKER }, - doc! { "$set": { "applied_at": mongodb::bson::DateTime::now() } }, - ) - .upsert(true) - .await?; - } - Ok(report) -} - -/// Relink DAST scan runs and pentest sessions from the old DAST target id to the -/// unified target id, logging each move so [`revert`] can undo it. -/// -/// Note: if multiple DAST targets fold into the same repo target, revert -/// restores only the last-logged mapping — a rare edge. The source collections -/// (`repositories`, `dast_targets`) are never deleted, so no data is lost. -async fn relink_history(db: &Database, old_id: &str, new_id: &str) -> Result<(), AgentError> { - db.dast_scan_runs() - .update_many( - doc! { "target_id": old_id }, - doc! { "$set": { "target_id": new_id } }, - ) - .await?; - db.pentest_sessions() - .update_many( - doc! { "target_id": old_id }, - doc! { "$set": { "target_id": new_id } }, - ) - .await?; - db.collection_named::("onboarding_migration_log") - .insert_one(doc! { "old_target_id": old_id, "new_target_id": new_id }) - .await?; - Ok(()) -} - -/// Undo the backfill: replay the relink log in reverse, drop `onboarded_targets` -/// and the log, and clear the marker. The legacy collections are untouched, so -/// this restores the pre-migration state. -pub async fn revert(db: &Database) -> Result<(), AgentError> { - let log = db.collection_named::("onboarding_migration_log"); - let mut cursor = log.find(doc! {}).await?; - while let Some(entry) = cursor.try_next().await? { - if let (Ok(old), Ok(new)) = ( - entry.get_str("old_target_id"), - entry.get_str("new_target_id"), - ) { - db.dast_scan_runs() - .update_many( - doc! { "target_id": new }, - doc! { "$set": { "target_id": old } }, - ) - .await?; - db.pentest_sessions() - .update_many( - doc! { "target_id": new }, - doc! { "$set": { "target_id": old } }, - ) - .await?; - } - } - db.onboarded_targets().drop().await?; - log.drop().await?; - db.collection_named::("schema_migrations") - .delete_one(doc! { "_id": MIGRATION_MARKER }) - .await?; - Ok(()) -} - -#[cfg(test)] -#[allow(clippy::expect_used, clippy::unwrap_used)] -mod tests { - use super::*; - use compliance_core::models::{DastAuthConfig, TrackerType}; - - fn repo() -> TrackedRepository { - let mut r = TrackedRepository::new("acme".to_string(), "https://git/acme.git".to_string()); - r.id = Some(mongodb::bson::oid::ObjectId::new()); - r.default_branch = "develop".to_string(); - r.last_scanned_commit = Some("abc123".to_string()); - r.auth_token = Some("pat".to_string()); - r.auth_username = Some("bob".to_string()); - r.tracker_type = Some(TrackerType::Gitea); - r.tracker_owner = Some("acme".to_string()); - r.findings_count = 7; - r - } - - fn dast(repo_id: Option, kind: DastTargetType) -> DastTarget { - let mut d = DastTarget::new( - "acme-web".to_string(), - "https://acme.example.com".to_string(), - kind, - ); - d.id = Some(mongodb::bson::oid::ObjectId::new()); - d.repo_id = repo_id; - d.max_crawl_depth = 5; - d.auth_config = Some(DastAuthConfig { - method: "bearer".to_string(), - login_url: None, - username: None, - password: None, - token: Some("tok".to_string()), - headers: None, - }); - d - } - - #[test] - fn repo_maps_preserving_id_and_git_artifact() { - let r = repo(); - let t = repo_to_target(&r); - assert_eq!(t.id, r.id); // id preserved - assert_eq!(t.findings_count, 7); - assert_eq!(t.scan_schedule, r.scan_schedule); - let git = t.code_artifact().expect("git artifact"); - assert_eq!(git.kind, ArtifactKind::GitRepo); - assert_eq!(git.source_ref, "https://git/acme.git"); - let gc = git.git.as_ref().expect("git config"); - assert_eq!(gc.default_branch, "develop"); - assert_eq!(gc.last_scanned_commit.as_deref(), Some("abc123")); - let auth = git.auth.as_ref().expect("auth"); - assert_eq!(auth.secret.as_deref(), Some("pat")); - assert_eq!(auth.username.as_deref(), Some("bob")); - assert_eq!( - t.scan_config - .issue_tracker - .as_ref() - .and_then(|it| it.tracker_type.clone()), - Some(TrackerType::Gitea) - ); - } - - #[test] - fn standalone_dast_maps_preserving_id_and_live_url() { - let d = dast(None, DastTargetType::WebApp); - let t = dast_to_standalone_target(&d); - assert_eq!(t.id, d.id); - assert_eq!(t.target_type, TargetType::WebApp); - let url = t.live_url().expect("live url"); - assert_eq!(url.source_ref, "https://acme.example.com"); - let web = url.web.as_ref().expect("web config"); - assert_eq!(web.max_crawl_depth, 5); - assert_eq!( - url.auth.as_ref().and_then(|a| a.secret.clone()), - Some("tok".to_string()) - ); - } - - #[test] - fn rest_api_dast_maps_to_backend_service() { - let d = dast(None, DastTargetType::RestApi); - assert_eq!( - dast_to_standalone_target(&d).target_type, - TargetType::BackendService - ); - } - - #[test] - fn fold_adds_live_url_and_promotes_webapp() { - let mut t = repo_to_target(&repo()); - assert_eq!(t.target_type, TargetType::BackendService); - fold_dast_into_target(&mut t, &dast(Some("x".to_string()), DastTargetType::WebApp)); - assert_eq!(t.target_type, TargetType::WebApp); // promoted - assert!(t.has(ArtifactKind::LiveUrl)); - assert!(t.has(ArtifactKind::GitRepo)); - } - - #[test] - fn fold_is_idempotent_on_live_url() { - let mut t = repo_to_target(&repo()); - let d = dast(Some("x".to_string()), DastTargetType::WebApp); - fold_dast_into_target(&mut t, &d); - fold_dast_into_target(&mut t, &d); - let live_urls = t - .artifacts - .iter() - .filter(|a| a.kind == ArtifactKind::LiveUrl) - .count(); - assert_eq!(live_urls, 1); - } -} diff --git a/compliance-agent/src/pentest/cleanup.rs b/compliance-agent/src/pentest/cleanup.rs index 739b2f9..dce1c61 100644 --- a/compliance-agent/src/pentest/cleanup.rs +++ b/compliance-agent/src/pentest/cleanup.rs @@ -342,7 +342,6 @@ mod tests { pentest_imap_password: None, admin_api_token: None, tenant_registry_url: None, - unified_pipeline: false, } } diff --git a/compliance-agent/src/pipeline/git.rs b/compliance-agent/src/pipeline/git.rs index 8d9c8a5..f29d5e0 100644 --- a/compliance-agent/src/pipeline/git.rs +++ b/compliance-agent/src/pipeline/git.rs @@ -138,7 +138,7 @@ impl GitOps { /// Build credentials from agent config + per-repo overrides pub fn make_repo_credentials( config: &compliance_core::AgentConfig, - repo: &compliance_core::models::TrackedRepository, + repo: &crate::pipeline::repo_view::RepoView, ) -> RepoCredentials { RepoCredentials { ssh_key_path: Some(config.ssh_key_path.clone()), diff --git a/compliance-agent/src/pipeline/issue_creation.rs b/compliance-agent/src/pipeline/issue_creation.rs index d9bdd9b..00c18cf 100644 --- a/compliance-agent/src/pipeline/issue_creation.rs +++ b/compliance-agent/src/pipeline/issue_creation.rs @@ -1,5 +1,6 @@ use mongodb::bson::doc; +use crate::pipeline::repo_view::RepoView; use compliance_core::models::*; use super::orchestrator::{extract_base_url, PipelineOrchestrator}; @@ -10,7 +11,7 @@ use crate::trackers; impl PipelineOrchestrator { /// Build an issue tracker client from a repository's tracker configuration. /// Returns `None` if the repo has no tracker configured. - pub(super) fn build_tracker(&self, repo: &TrackedRepository) -> Option { + pub(super) fn build_tracker(&self, repo: &RepoView) -> Option { let tracker_type = repo.tracker_type.as_ref()?; // Per-repo token takes precedence, fall back to global config match tracker_type { @@ -81,7 +82,7 @@ impl PipelineOrchestrator { #[tracing::instrument(skip_all, fields(repo_id = %repo_id))] pub(super) async fn create_tracker_issues( &self, - repo: &TrackedRepository, + repo: &RepoView, repo_id: &str, new_findings: &[Finding], ) -> Result<(), AgentError> { diff --git a/compliance-agent/src/pipeline/mod.rs b/compliance-agent/src/pipeline/mod.rs index c157fde..2a0f2be 100644 --- a/compliance-agent/src/pipeline/mod.rs +++ b/compliance-agent/src/pipeline/mod.rs @@ -11,6 +11,7 @@ pub mod orchestrator; pub mod patterns; pub mod plan; mod pr_review; +pub mod repo_view; pub mod sbom; pub mod semgrep; mod tracker_dispatch; diff --git a/compliance-agent/src/pipeline/orchestrator.rs b/compliance-agent/src/pipeline/orchestrator.rs index 4135568..d5b8ad4 100644 --- a/compliance-agent/src/pipeline/orchestrator.rs +++ b/compliance-agent/src/pipeline/orchestrator.rs @@ -16,6 +16,7 @@ use crate::pipeline::gitleaks::GitleaksScanner; use crate::pipeline::lint::LintScanner; use crate::pipeline::patterns::{GdprPatternScanner, OAuthPatternScanner}; use crate::pipeline::plan::build_scan_plan; +use crate::pipeline::repo_view::RepoView; use crate::pipeline::sbom::SbomScanner; use crate::pipeline::semgrep::SemgrepScanner; @@ -51,72 +52,8 @@ impl PipelineOrchestrator { } } - #[tracing::instrument(skip_all, fields(repo_id = %repo_id, trigger = ?trigger))] - pub async fn run(&self, repo_id: &str, trigger: ScanTrigger) -> Result<(), AgentError> { - // Look up the repository - let repo = self - .db - .repositories() - .find_one(doc! { "_id": mongodb::bson::oid::ObjectId::parse_str(repo_id).map_err(|e| AgentError::Other(e.to_string()))? }) - .await? - .ok_or_else(|| AgentError::Other(format!("Repository {repo_id} not found")))?; - - // Create scan run - let scan_run = ScanRun::new(repo_id.to_string(), trigger); - let insert = self.db.scan_runs().insert_one(&scan_run).await?; - let scan_run_id = insert - .inserted_id - .as_object_id() - .map(|id| id.to_hex()) - .unwrap_or_default(); - - let result = self.run_pipeline(&repo, &scan_run_id).await; - - // Update scan run status - match &result { - Ok(count) => { - self.db - .scan_runs() - .update_one( - doc! { "_id": &insert.inserted_id }, - doc! { - "$set": { - "status": "completed", - "current_phase": "completed", - "new_findings_count": *count as i64, - "completed_at": mongodb::bson::DateTime::now(), - } - }, - ) - .await?; - } - Err(e) => { - tracing::error!(repo_id, error = %e, "Scan pipeline failed"); - self.db - .scan_runs() - .update_one( - doc! { "_id": &insert.inserted_id }, - doc! { - "$set": { - "status": "failed", - "error_message": e.to_string(), - "completed_at": mongodb::bson::DateTime::now(), - } - }, - ) - .await?; - } - } - - result.map(|_| ()) - } - #[tracing::instrument(skip_all, fields(repo_id = repo.name.as_str()))] - async fn run_pipeline( - &self, - repo: &TrackedRepository, - scan_run_id: &str, - ) -> Result { + async fn run_pipeline(&self, repo: &RepoView, scan_run_id: &str) -> Result { let repo_id = repo.id.as_ref().map(|id| id.to_hex()).unwrap_or_default(); // Stage 0: Change detection @@ -130,7 +67,6 @@ impl PipelineOrchestrator { return Ok(0); } - let current_sha = GitOps::get_head_sha(&repo_path)?; let mut all_findings: Vec = Vec::new(); // Stage 1: Semgrep SAST @@ -396,20 +332,9 @@ impl PipelineOrchestrator { tracing::warn!("[{repo_id}] Issue creation failed: {e}"); } - // Stage 7: Update repository - self.db - .repositories() - .update_one( - doc! { "_id": repo.id }, - doc! { - "$set": { - "last_scanned_commit": ¤t_sha, - "updated_at": mongodb::bson::DateTime::now(), - }, - "$inc": { "findings_count": new_count as i64 }, - }, - ) - .await?; + // The onboarded target's findings_count and the git artifact's + // last_scanned_commit watermark are persisted by `finalize_target` after + // `run_pipeline` returns. // Stage 8: DAST (async, optional — only if a DastTarget is configured) tracing::info!("[{repo_id}] Stage 8: Checking for DAST targets"); @@ -526,7 +451,7 @@ impl PipelineOrchestrator { match target.code_artifact() { Some(code) if code.kind == ArtifactKind::GitRepo => { - let repo = repo_view_from_target(target, code); + let repo = RepoView::from_target(target, code); let new_count = self.run_pipeline(&repo, scan_run_id).await?; self.finalize_target(target, &repo, new_count).await?; Ok(new_count) @@ -703,7 +628,7 @@ impl PipelineOrchestrator { async fn finalize_target( &self, target: &OnboardedTarget, - repo: &TrackedRepository, + repo: &RepoView, new_count: u32, ) -> Result<(), AgentError> { let oid = match target.id { @@ -751,37 +676,6 @@ impl PipelineOrchestrator { } } -/// Build a legacy `TrackedRepository` view from an onboarded target's code -/// artifact, so the unified pipeline can reuse the existing repo pipeline. The -/// inverse of the migration's `repo_to_target`. `_id` is preserved so findings -/// and DAST lookups resolve against the same key. -fn repo_view_from_target(target: &OnboardedTarget, code: &Artifact) -> TrackedRepository { - let mut repo = TrackedRepository::new(target.name.clone(), code.source_ref.clone()); - repo.id = target.id; - if let Some(git) = &code.git { - repo.default_branch = git.default_branch.clone(); - repo.last_scanned_commit = git.last_scanned_commit.clone(); - repo.local_path = git.local_path.clone(); - } - if let Some(auth) = &code.auth { - repo.auth_token = auth.secret.clone(); - repo.auth_username = auth.username.clone(); - } - if let Some(it) = &target.scan_config.issue_tracker { - repo.tracker_type = it.tracker_type.clone(); - repo.tracker_owner = it.owner.clone(); - repo.tracker_repo = it.repo.clone(); - repo.tracker_token = it.token.clone(); - } - repo.scan_schedule = target.scan_schedule.clone(); - repo.webhook_enabled = target.webhook_enabled; - repo.webhook_secret = target.webhook_secret.clone(); - repo.findings_count = target.findings_count; - repo.created_at = target.created_at; - repo.updated_at = target.updated_at; - repo -} - /// Extract the scheme + host from a git URL. /// e.g. "https://gitea.example.com/owner/repo.git" -> "https://gitea.example.com" /// e.g. "ssh://git@gitea.example.com:22/owner/repo.git" -> "https://gitea.example.com" @@ -840,7 +734,7 @@ mod tests { target.artifacts.push(artifact); let code = target.code_artifact().expect("code artifact"); - let repo = repo_view_from_target(&target, code); + let repo = RepoView::from_target(&target, code); assert_eq!(repo.id, target.id); // preserved assert_eq!(repo.git_url, "https://git/acme.git"); diff --git a/compliance-agent/src/pipeline/pr_review.rs b/compliance-agent/src/pipeline/pr_review.rs index 0ea7b2d..c3ea5d7 100644 --- a/compliance-agent/src/pipeline/pr_review.rs +++ b/compliance-agent/src/pipeline/pr_review.rs @@ -1,3 +1,4 @@ +use crate::pipeline::repo_view::RepoView; use compliance_core::models::*; use super::dedup::compute_fingerprint; @@ -14,7 +15,7 @@ impl PipelineOrchestrator { #[tracing::instrument(skip_all, fields(repo_id = %repo_id, pr_number))] pub async fn run_pr_review( &self, - repo: &TrackedRepository, + repo: &RepoView, repo_id: &str, pr_number: u64, base_sha: &str, diff --git a/compliance-agent/src/pipeline/repo_view.rs b/compliance-agent/src/pipeline/repo_view.rs new file mode 100644 index 0000000..5a4aa92 --- /dev/null +++ b/compliance-agent/src/pipeline/repo_view.rs @@ -0,0 +1,74 @@ +//! `RepoView` — an internal, non-persisted view of a code target for the scan +//! pipeline. +//! +//! It replaces the old persisted `TrackedRepository` model. The pipeline +//! (SAST → SBOM → CVE → triage → issues → DAST, and PR review) only ever needs a +//! flat bundle of git + issue-tracker + auth fields; those are projected from an +//! [`OnboardedTarget`] and its code [`Artifact`] by [`RepoView::from_target`]. +//! Nothing here is written to Mongo — onboarded targets are the sole persisted +//! entity. + +use compliance_core::models::{Artifact, OnboardedTarget, TrackerType}; + +/// A flat, pipeline-facing view of a code target. Built from an onboarded +/// target; never persisted. +#[derive(Debug, Clone)] +pub struct RepoView { + /// The onboarded target's id (used as `repo_id` across findings/sbom/etc.). + pub id: Option, + pub name: String, + pub git_url: String, + pub default_branch: String, + pub local_path: Option, + pub scan_schedule: Option, + pub webhook_enabled: bool, + pub webhook_secret: Option, + pub tracker_type: Option, + pub tracker_owner: Option, + pub tracker_repo: Option, + pub tracker_token: Option, + pub auth_token: Option, + pub auth_username: Option, + pub last_scanned_commit: Option, + pub findings_count: u32, +} + +impl RepoView { + /// Project an onboarded target + its code artifact into a pipeline view. + pub fn from_target(target: &OnboardedTarget, code: &Artifact) -> Self { + let mut view = Self { + id: target.id, + name: target.name.clone(), + git_url: code.source_ref.clone(), + default_branch: "main".to_string(), + local_path: None, + scan_schedule: target.scan_schedule.clone(), + webhook_enabled: target.webhook_enabled, + webhook_secret: target.webhook_secret.clone(), + tracker_type: None, + tracker_owner: None, + tracker_repo: None, + tracker_token: None, + auth_token: None, + auth_username: None, + last_scanned_commit: None, + findings_count: target.findings_count, + }; + if let Some(git) = &code.git { + view.default_branch = git.default_branch.clone(); + view.last_scanned_commit = git.last_scanned_commit.clone(); + view.local_path = git.local_path.clone(); + } + if let Some(auth) = &code.auth { + view.auth_token = auth.secret.clone(); + view.auth_username = auth.username.clone(); + } + if let Some(it) = &target.scan_config.issue_tracker { + view.tracker_type = it.tracker_type.clone(); + view.tracker_owner = it.owner.clone(); + view.tracker_repo = it.repo.clone(); + view.tracker_token = it.token.clone(); + } + view + } +} diff --git a/compliance-agent/src/scheduler.rs b/compliance-agent/src/scheduler.rs index 98f8a87..17fdd3c 100644 --- a/compliance-agent/src/scheduler.rs +++ b/compliance-agent/src/scheduler.rs @@ -348,7 +348,7 @@ async fn monitor_cves(agent: &ComplianceAgent, tenant_id: &str) { std::collections::HashMap::new(); for rid in &repo_ids { if let Ok(oid) = mongodb::bson::oid::ObjectId::parse_str(rid) { - if let Ok(Some(repo)) = db.repositories().find_one(doc! { "_id": oid }).await { + if let Ok(Some(repo)) = db.onboarded_targets().find_one(doc! { "_id": oid }).await { repo_names.insert(rid.clone(), repo.name.clone()); } } diff --git a/compliance-agent/src/webhooks/gitea.rs b/compliance-agent/src/webhooks/gitea.rs index a77c161..375b74c 100644 --- a/compliance-agent/src/webhooks/gitea.rs +++ b/compliance-agent/src/webhooks/gitea.rs @@ -31,7 +31,7 @@ pub async fn handle_gitea_webhook( } }; let repo = match db - .repositories() + .onboarded_targets() .find_one(mongodb::bson::doc! { "_id": oid }) .await { diff --git a/compliance-agent/src/webhooks/github.rs b/compliance-agent/src/webhooks/github.rs index a9b7048..f10ee12 100644 --- a/compliance-agent/src/webhooks/github.rs +++ b/compliance-agent/src/webhooks/github.rs @@ -31,7 +31,7 @@ pub async fn handle_github_webhook( } }; let repo = match db - .repositories() + .onboarded_targets() .find_one(mongodb::bson::doc! { "_id": oid }) .await { diff --git a/compliance-agent/src/webhooks/gitlab.rs b/compliance-agent/src/webhooks/gitlab.rs index c830a28..fd12ff9 100644 --- a/compliance-agent/src/webhooks/gitlab.rs +++ b/compliance-agent/src/webhooks/gitlab.rs @@ -27,7 +27,7 @@ pub async fn handle_gitlab_webhook( } }; let repo = match db - .repositories() + .onboarded_targets() .find_one(mongodb::bson::doc! { "_id": oid }) .await { diff --git a/compliance-agent/tests/common/mod.rs b/compliance-agent/tests/common/mod.rs index e85a475..9d0934e 100644 --- a/compliance-agent/tests/common/mod.rs +++ b/compliance-agent/tests/common/mod.rs @@ -70,7 +70,6 @@ impl TestServer { pentest_imap_password: None, admin_api_token: None, tenant_registry_url: None, - unified_pipeline: false, }; let agent = ComplianceAgent::new(config, db_pool); diff --git a/compliance-agent/tests/integration/api/cascade_delete.rs b/compliance-agent/tests/integration/api/cascade_delete.rs index 0ba9cc2..07cb647 100644 --- a/compliance-agent/tests/integration/api/cascade_delete.rs +++ b/compliance-agent/tests/integration/api/cascade_delete.rs @@ -113,15 +113,16 @@ async fn delete_repo_cascades_to_dast_and_pentest_data() { // Create a repo let resp = server .post( - "/api/v1/repositories", + "/api/v1/targets", &json!({ "name": "cascade-test", - "git_url": "https://github.com/example/cascade-test.git", + "target_type": "web_app", + "artifacts": [{ "kind": "git_repo", "source_ref": "https://github.com/example/cascade-test.git", "branch": "main" }], }), ) .await; let body: serde_json::Value = resp.json().await.unwrap(); - let repo_id = body["data"]["id"].as_str().unwrap().to_string(); + let repo_id = body["data"]["_id"]["$oid"].as_str().unwrap().to_string(); // Insert DAST target linked to repo let target_id = insert_dast_target(&server, &repo_id, "cascade-target").await; @@ -140,9 +141,7 @@ async fn delete_repo_cascades_to_dast_and_pentest_data() { assert_eq!(count_docs(&server, "dast_findings").await, 1); // Delete the repo - let resp = server - .delete(&format!("/api/v1/repositories/{repo_id}")) - .await; + let resp = server.delete(&format!("/api/v1/targets/{repo_id}")).await; assert_eq!(resp.status(), 200); // All downstream data should be gone @@ -161,15 +160,16 @@ async fn delete_repo_cascades_sast_findings_and_sbom() { // Create a repo let resp = server .post( - "/api/v1/repositories", + "/api/v1/targets", &json!({ "name": "sast-cascade", - "git_url": "https://github.com/example/sast-cascade.git", + "target_type": "web_app", + "artifacts": [{ "kind": "git_repo", "source_ref": "https://github.com/example/sast-cascade.git", "branch": "main" }], }), ) .await; let body: serde_json::Value = resp.json().await.unwrap(); - let repo_id = body["data"]["id"].as_str().unwrap().to_string(); + let repo_id = body["data"]["_id"]["$oid"].as_str().unwrap().to_string(); // Insert SAST finding and SBOM entry let mongodb_uri = std::env::var("TEST_MONGODB_URI") @@ -209,9 +209,7 @@ async fn delete_repo_cascades_sast_findings_and_sbom() { assert_eq!(count_docs(&server, "sbom_entries").await, 1); // Delete repo - server - .delete(&format!("/api/v1/repositories/{repo_id}")) - .await; + server.delete(&format!("/api/v1/targets/{repo_id}")).await; // Both should be gone assert_eq!(count_docs(&server, "findings").await, 0); diff --git a/compliance-agent/tests/integration/api/mod.rs b/compliance-agent/tests/integration/api/mod.rs index 960394b..23f1594 100644 --- a/compliance-agent/tests/integration/api/mod.rs +++ b/compliance-agent/tests/integration/api/mod.rs @@ -3,5 +3,4 @@ mod dast; mod findings; mod health; mod onboarding; -mod repositories; mod stats; diff --git a/compliance-agent/tests/integration/api/repositories.rs b/compliance-agent/tests/integration/api/repositories.rs deleted file mode 100644 index 7cf476f..0000000 --- a/compliance-agent/tests/integration/api/repositories.rs +++ /dev/null @@ -1,110 +0,0 @@ -use crate::common::TestServer; -use serde_json::json; - -#[tokio::test] -async fn add_and_list_repository() { - let server = TestServer::start().await; - - // Initially empty - let resp = server.get("/api/v1/repositories").await; - assert_eq!(resp.status(), 200); - let body: serde_json::Value = resp.json().await.unwrap(); - assert_eq!(body["data"].as_array().unwrap().len(), 0); - - // Add a repository - let resp = server - .post( - "/api/v1/repositories", - &json!({ - "name": "test-repo", - "git_url": "https://github.com/example/test-repo.git", - }), - ) - .await; - assert_eq!(resp.status(), 200); - let body: serde_json::Value = resp.json().await.unwrap(); - let repo_id = body["data"]["id"].as_str().unwrap().to_string(); - assert!(!repo_id.is_empty()); - - // List should now return 1 - let resp = server.get("/api/v1/repositories").await; - let body: serde_json::Value = resp.json().await.unwrap(); - let repos = body["data"].as_array().unwrap(); - assert_eq!(repos.len(), 1); - assert_eq!(repos[0]["name"], "test-repo"); - - server.cleanup().await; -} - -#[tokio::test] -async fn add_duplicate_repository_fails() { - let server = TestServer::start().await; - - let payload = json!({ - "name": "dup-repo", - "git_url": "https://github.com/example/dup-repo.git", - }); - - // First add succeeds - let resp = server.post("/api/v1/repositories", &payload).await; - assert_eq!(resp.status(), 200); - - // Second add with same git_url should fail (unique index) - let resp = server.post("/api/v1/repositories", &payload).await; - assert_ne!(resp.status(), 200); - - server.cleanup().await; -} - -#[tokio::test] -async fn delete_repository() { - let server = TestServer::start().await; - - // Add a repo - let resp = server - .post( - "/api/v1/repositories", - &json!({ - "name": "to-delete", - "git_url": "https://github.com/example/to-delete.git", - }), - ) - .await; - let body: serde_json::Value = resp.json().await.unwrap(); - let repo_id = body["data"]["id"].as_str().unwrap(); - - // Delete it - let resp = server - .delete(&format!("/api/v1/repositories/{repo_id}")) - .await; - assert_eq!(resp.status(), 200); - - // List should be empty again - let resp = server.get("/api/v1/repositories").await; - let body: serde_json::Value = resp.json().await.unwrap(); - assert_eq!(body["data"].as_array().unwrap().len(), 0); - - server.cleanup().await; -} - -#[tokio::test] -async fn delete_nonexistent_repository_returns_404() { - let server = TestServer::start().await; - - let resp = server - .delete("/api/v1/repositories/000000000000000000000000") - .await; - assert_eq!(resp.status(), 404); - - server.cleanup().await; -} - -#[tokio::test] -async fn delete_invalid_id_returns_400() { - let server = TestServer::start().await; - - let resp = server.delete("/api/v1/repositories/not-a-valid-id").await; - assert_eq!(resp.status(), 400); - - server.cleanup().await; -} diff --git a/compliance-agent/tests/integration/api/stats.rs b/compliance-agent/tests/integration/api/stats.rs index 969dc3f..3997f8c 100644 --- a/compliance-agent/tests/integration/api/stats.rs +++ b/compliance-agent/tests/integration/api/stats.rs @@ -5,13 +5,14 @@ use serde_json::json; async fn stats_overview_reflects_inserted_data() { let server = TestServer::start().await; - // Add a repo + // Add a target server .post( - "/api/v1/repositories", + "/api/v1/targets", &json!({ "name": "stats-repo", - "git_url": "https://github.com/example/stats-repo.git", + "target_type": "web_app", + "artifacts": [{ "kind": "git_repo", "source_ref": "https://github.com/example/stats-repo.git", "branch": "main" }], }), ) .await; diff --git a/compliance-agent/tests/integration/migration.rs b/compliance-agent/tests/integration/migration.rs deleted file mode 100644 index 9b6a116..0000000 --- a/compliance-agent/tests/integration/migration.rs +++ /dev/null @@ -1,156 +0,0 @@ -// Integration tests for the onboarding backfill migration. -// -// Requires MongoDB (set TEST_MONGODB_URI if not at the default). -// Not run in CI (which is `--lib` only) — run locally: -// cargo test -p compliance-agent --test e2e migration - -use compliance_agent::database::{Database, DatabasePool}; -use compliance_agent::migrate::onboarding; -use compliance_core::models::{ - ArtifactKind, DastTarget, DastTargetType, TargetType, TrackedRepository, -}; -use mongodb::bson::{doc, Document}; - -async fn fresh_db() -> (DatabasePool, String, Database) { - let uri = std::env::var("TEST_MONGODB_URI") - .unwrap_or_else(|_| "mongodb://root:example@localhost:27017/?authSource=admin".into()); - // Prefix must fit the pool's 30-char cap (`_<32 hex>` <= 63). - let prefix = format!("t_{}", &uuid::Uuid::new_v4().simple().to_string()[..16]); - let pool = DatabasePool::connect(&uri, &prefix) - .await - .expect("connect mongo"); - let db = pool.for_tenant_id("t1").await.expect("tenant db"); - (pool, prefix, db) -} - -async fn cleanup(pool: &DatabasePool, prefix: &str) { - if let Ok(names) = pool.client().list_database_names().await { - for n in names { - if n.starts_with(prefix) { - pool.client().database(&n).drop().await.ok(); - } - } - } -} - -#[tokio::test] -async fn backfill_folds_relinks_is_idempotent_and_reversible() { - let (pool, prefix, db) = fresh_db().await; - - // Seed a repo. - let repo = TrackedRepository::new("acme".into(), "https://git/acme.git".into()); - let repo_id = db - .repositories() - .insert_one(repo) - .await - .expect("insert repo") - .inserted_id - .as_object_id() - .expect("repo oid"); - - // A DAST target linked to the repo (folds + promotes to WebApp + relinks). - let mut linked = DastTarget::new( - "acme-web".into(), - "https://acme.example.com".into(), - DastTargetType::WebApp, - ); - linked.repo_id = Some(repo_id.to_hex()); - let linked_id = db - .dast_targets() - .insert_one(linked) - .await - .expect("insert linked dast") - .inserted_id - .as_object_id() - .expect("linked oid"); - - // A repo-less DAST target (standalone). - let standalone = DastTarget::new( - "acme-api".into(), - "https://api.acme.com".into(), - DastTargetType::RestApi, - ); - let standalone_id = db - .dast_targets() - .insert_one(standalone) - .await - .expect("insert standalone dast") - .inserted_id - .as_object_id() - .expect("standalone oid"); - - // A DAST scan run pointing at the linked target — should be relinked to the repo. - db.collection_named::("dast_scan_runs") - .insert_one(doc! { "target_id": linked_id.to_hex(), "status": "completed" }) - .await - .expect("insert dast run"); - - // --- Backfill --- - assert!(!onboarding::already_applied(&db).await.unwrap()); - let report = onboarding::backfill_onboarded_targets(&db, false) - .await - .expect("backfill"); - assert_eq!(report.repos_migrated, 1); - assert_eq!(report.dast_targets_folded, 1); - assert_eq!(report.dast_targets_standalone, 1); - assert!(onboarding::already_applied(&db).await.unwrap()); - - // Repo target: preserved _id, has git + folded live-url, promoted to WebApp. - let repo_target = db - .onboarded_targets() - .find_one(doc! { "_id": repo_id }) - .await - .unwrap() - .expect("repo target"); - assert!(repo_target.has(ArtifactKind::GitRepo)); - assert!(repo_target.has(ArtifactKind::LiveUrl)); - assert_eq!(repo_target.target_type, TargetType::WebApp); - - // Standalone target: preserved _id, live-url, backend service. - let standalone_target = db - .onboarded_targets() - .find_one(doc! { "_id": standalone_id }) - .await - .unwrap() - .expect("standalone target"); - assert!(standalone_target.has(ArtifactKind::LiveUrl)); - assert_eq!(standalone_target.target_type, TargetType::BackendService); - - // The DAST run was relinked from the old dast id to the repo (unified) id. - let run = db - .collection_named::("dast_scan_runs") - .find_one(doc! {}) - .await - .unwrap() - .expect("run"); - assert_eq!(run.get_str("target_id").unwrap(), repo_id.to_hex()); - - // --- Idempotent: re-run migrates nothing new --- - let again = onboarding::backfill_onboarded_targets(&db, false) - .await - .expect("backfill again"); - assert_eq!(again.repos_migrated, 0); - assert_eq!(again.dast_targets_folded, 0); - assert_eq!(again.dast_targets_standalone, 0); - assert!(again.skipped_existing >= 2); - - // --- Revert: onboarded targets gone, relink undone, marker cleared --- - onboarding::revert(&db).await.expect("revert"); - assert_eq!( - db.onboarded_targets() - .count_documents(doc! {}) - .await - .unwrap(), - 0 - ); - let run_after = db - .collection_named::("dast_scan_runs") - .find_one(doc! {}) - .await - .unwrap() - .expect("run"); - assert_eq!(run_after.get_str("target_id").unwrap(), linked_id.to_hex()); - assert!(!onboarding::already_applied(&db).await.unwrap()); - - cleanup(&pool, &prefix).await; -} diff --git a/compliance-agent/tests/integration/mod.rs b/compliance-agent/tests/integration/mod.rs index 5f9870e..1baec19 100644 --- a/compliance-agent/tests/integration/mod.rs +++ b/compliance-agent/tests/integration/mod.rs @@ -7,4 +7,3 @@ // Or nightly: (via CI with MongoDB service container) mod api; -mod migration; diff --git a/compliance-agent/tests/tenant_isolation.rs b/compliance-agent/tests/tenant_isolation.rs index c204f7b..7a32d5d 100644 --- a/compliance-agent/tests/tenant_isolation.rs +++ b/compliance-agent/tests/tenant_isolation.rs @@ -11,7 +11,7 @@ #![allow(clippy::expect_used, clippy::unwrap_used)] use compliance_agent::database::DatabasePool; -use compliance_core::models::TrackedRepository; +use compliance_core::models::{Artifact, OnboardedTarget, TargetType}; use compliance_core::{OrgRole, TenantContext, TenantStatus}; use mongodb::bson::doc; @@ -28,27 +28,12 @@ fn ctx(tenant_id: &str, slug: &str) -> TenantContext { } } -fn fixture_repo(name: &str, git_url: &str) -> TrackedRepository { - TrackedRepository { - id: None, - name: name.to_string(), - git_url: git_url.to_string(), - default_branch: "main".to_string(), - local_path: None, - scan_schedule: None, - webhook_enabled: false, - webhook_secret: None, - tracker_type: None, - tracker_owner: None, - tracker_repo: None, - tracker_token: None, - auth_token: None, - auth_username: None, - last_scanned_commit: None, - findings_count: 0, - created_at: chrono::Utc::now(), - updated_at: chrono::Utc::now(), - } +fn fixture_repo(name: &str, git_url: &str) -> OnboardedTarget { + let mut target = OnboardedTarget::new(name.to_string(), TargetType::WebApp); + target + .artifacts + .push(Artifact::git_repo(git_url.to_string(), "main".to_string())); + target } #[tokio::test] @@ -71,12 +56,12 @@ async fn pool_isolates_tenants_at_driver_level() { // Write distinct repos into each tenant's database. acme_db - .repositories() + .onboarded_targets() .insert_one(fixture_repo("acme-app", "git@example.com:acme/app.git")) .await .expect("insert acme"); globex_db - .repositories() + .onboarded_targets() .insert_one(fixture_repo( "globex-platform", "git@example.com:globex/platform.git", @@ -173,12 +158,12 @@ async fn admin_helpers_list_and_drop_tenant_dbs() { let acme_db = pool.for_tenant(&acme).await.expect("acme db"); let globex_db = pool.for_tenant(&globex).await.expect("globex db"); acme_db - .repositories() + .onboarded_targets() .insert_one(fixture_repo("acme-app", "git@example.com:acme/app.git")) .await .expect("insert acme"); globex_db - .repositories() + .onboarded_targets() .insert_one(fixture_repo("globex-app", "git@example.com:globex/app.git")) .await .expect("insert globex"); @@ -284,9 +269,9 @@ fn short_id() -> String { } /// Drain a `repositories` find cursor on the given tenant database. -async fn collect(db: &compliance_agent::database::Database) -> Vec { +async fn collect(db: &compliance_agent::database::Database) -> Vec { let mut cursor = db - .repositories() + .onboarded_targets() .find(doc! {}) .await .expect("find repositories"); diff --git a/compliance-core/src/config.rs b/compliance-core/src/config.rs index a787cc4..db88fef 100644 --- a/compliance-core/src/config.rs +++ b/compliance-core/src/config.rs @@ -49,11 +49,6 @@ pub struct AgentConfig { /// of tenants to iterate. When `None` or unreachable, scheduler /// falls back to `SCHEDULER_TENANT_IDS` env (M7.2-C). pub tenant_registry_url: Option, - /// When true, `run_scan` dispatches to the unified `run_target` pipeline - /// (reads `onboarded_targets`) instead of the legacy repository pipeline. - /// Env `UNIFIED_PIPELINE`. Defaults on; set `UNIFIED_PIPELINE=0` to use the - /// legacy repository pipeline. - pub unified_pipeline: bool, } #[derive(Clone, Debug, Serialize, Deserialize)] diff --git a/compliance-core/src/models/mod.rs b/compliance-core/src/models/mod.rs index 4e84714..a56241e 100644 --- a/compliance-core/src/models/mod.rs +++ b/compliance-core/src/models/mod.rs @@ -44,6 +44,6 @@ pub use pentest::{ PentestStatus, PentestStrategy, SeverityDistribution, TestUserRecord, TesterInfo, ToolCallRecord, }; -pub use repository::{ScanTrigger, TrackedRepository}; +pub use repository::ScanTrigger; pub use sbom::{SbomEntry, VulnRef}; pub use scan::{ScanPhase, ScanRun, ScanRunStatus, ScanType}; diff --git a/compliance-core/src/models/repository.rs b/compliance-core/src/models/repository.rs index eae5cae..fe5212e 100644 --- a/compliance-core/src/models/repository.rs +++ b/compliance-core/src/models/repository.rs @@ -1,8 +1,6 @@ -use chrono::{DateTime, Utc}; -use serde::{Deserialize, Deserializer, Serialize}; - -use super::issue::TrackerType; +use serde::{Deserialize, Serialize}; +/// What initiated a scan. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum ScanTrigger { @@ -10,92 +8,3 @@ pub enum ScanTrigger { Webhook, Manual, } - -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct TrackedRepository { - #[serde(rename = "_id", skip_serializing_if = "Option::is_none")] - pub id: Option, - #[serde(default)] - pub name: String, - #[serde(default)] - pub git_url: String, - #[serde(default = "default_branch")] - pub default_branch: String, - pub local_path: Option, - pub scan_schedule: Option, - #[serde(default)] - pub webhook_enabled: bool, - /// Auto-generated HMAC secret for verifying incoming webhooks - #[serde(default, skip_serializing_if = "Option::is_none")] - pub webhook_secret: Option, - pub tracker_type: Option, - pub tracker_owner: Option, - pub tracker_repo: Option, - /// Optional per-repo PAT for the issue tracker (GitHub/GitLab/Jira) - #[serde(default, skip_serializing_if = "Option::is_none")] - pub tracker_token: Option, - /// Optional auth token for HTTPS private repos (PAT or password) - #[serde(default, skip_serializing_if = "Option::is_none")] - pub auth_token: Option, - /// Optional username for HTTPS auth (defaults to "x-access-token" for PATs) - #[serde(default, skip_serializing_if = "Option::is_none")] - pub auth_username: Option, - pub last_scanned_commit: Option, - #[serde(default, deserialize_with = "deserialize_findings_count")] - pub findings_count: u32, - #[serde( - default = "chrono::Utc::now", - with = "super::serde_helpers::bson_datetime" - )] - pub created_at: DateTime, - #[serde( - default = "chrono::Utc::now", - with = "super::serde_helpers::bson_datetime" - )] - pub updated_at: DateTime, -} - -fn default_branch() -> String { - "main".to_string() -} - -fn deserialize_findings_count<'de, D>(deserializer: D) -> Result -where - D: Deserializer<'de>, -{ - let bson = bson::Bson::deserialize(deserializer)?; - match &bson { - bson::Bson::Int32(n) => Ok(*n as u32), - bson::Bson::Int64(n) => Ok(*n as u32), - bson::Bson::Double(n) => Ok(*n as u32), - _ => Ok(0), - } -} - -impl TrackedRepository { - pub fn new(name: String, git_url: String) -> Self { - let now = Utc::now(); - // Generate a random webhook secret (hex-encoded UUID v4, no dashes) - let webhook_secret = uuid::Uuid::new_v4().to_string().replace('-', ""); - Self { - id: None, - name, - git_url, - default_branch: "main".to_string(), - local_path: None, - scan_schedule: None, - auth_token: None, - auth_username: None, - webhook_enabled: false, - webhook_secret: Some(webhook_secret), - tracker_type: None, - tracker_owner: None, - tracker_repo: None, - tracker_token: None, - last_scanned_commit: None, - findings_count: 0, - created_at: now, - updated_at: now, - } - } -} diff --git a/compliance-dashboard/src/app.rs b/compliance-dashboard/src/app.rs index 9c47ac4..3771c36 100644 --- a/compliance-dashboard/src/app.rs +++ b/compliance-dashboard/src/app.rs @@ -10,8 +10,6 @@ pub enum Route { #[layout(AppShell)] #[route("/")] OverviewPage {}, - #[route("/repositories")] - RepositoriesPage {}, #[route("/targets")] TargetsPage {}, #[route("/onboard")] diff --git a/compliance-dashboard/src/components/pentest_wizard.rs b/compliance-dashboard/src/components/pentest_wizard.rs index 2f0cc0d..961ccc6 100644 --- a/compliance-dashboard/src/components/pentest_wizard.rs +++ b/compliance-dashboard/src/components/pentest_wizard.rs @@ -4,8 +4,9 @@ use dioxus_free_icons::Icon; use crate::app::Route; use crate::infrastructure::dast::fetch_dast_targets; +use crate::infrastructure::onboarding::fetch_targets; use crate::infrastructure::pentest::{create_pentest_session_wizard, lookup_repo_by_url}; -use crate::infrastructure::repositories::{fetch_repositories, fetch_ssh_public_key}; +use crate::infrastructure::repositories::fetch_ssh_public_key; const DISCLAIMER_TEXT: &str = "I confirm that I have authorization to perform security testing \ against the specified target. I understand that penetration testing may cause disruption to the \ @@ -39,7 +40,7 @@ pub fn PentestWizard(show: Signal) -> Element { let mut show_target_dropdown = use_signal(|| false); let mut show_repo_dropdown = use_signal(|| false); let existing_targets = use_resource(|| async { fetch_dast_targets().await.ok() }); - let existing_repos = use_resource(|| async { fetch_repositories(1).await.ok() }); + let existing_repos = use_resource(|| async { fetch_targets().await.ok() }); // SSH key state for private repos let mut ssh_public_key = use_signal(String::new); @@ -211,7 +212,25 @@ pub fn PentestWizard(show: Signal) -> Element { Some(Some(data)) => data .data .iter() - .map(|r| (r.git_url.clone(), r.name.clone())) + .filter_map(|t| { + let name = t + .get("name") + .and_then(|v| v.as_str()) + .unwrap_or_default() + .to_string(); + let git_url = t + .get("artifacts") + .and_then(|a| a.as_array()) + .and_then(|arr| { + arr.iter().find(|a| { + a.get("kind").and_then(|k| k.as_str()) == Some("git_repo") + }) + }) + .and_then(|a| a.get("source_ref")) + .and_then(|s| s.as_str())? + .to_string(); + Some((git_url, name)) + }) .collect(), _ => Vec::new(), } diff --git a/compliance-dashboard/src/infrastructure/database.rs b/compliance-dashboard/src/infrastructure/database.rs index cc45c24..8bf641d 100644 --- a/compliance-dashboard/src/infrastructure/database.rs +++ b/compliance-dashboard/src/infrastructure/database.rs @@ -19,10 +19,6 @@ impl Database { Ok(Self { inner: db }) } - pub fn repositories(&self) -> Collection { - self.inner.collection("repositories") - } - pub fn findings(&self) -> Collection { self.inner.collection("findings") } diff --git a/compliance-dashboard/src/infrastructure/repositories.rs b/compliance-dashboard/src/infrastructure/repositories.rs index 666c7e4..309a906 100644 --- a/compliance-dashboard/src/infrastructure/repositories.rs +++ b/compliance-dashboard/src/infrastructure/repositories.rs @@ -1,145 +1,10 @@ +//! The agent's SSH deploy public key — shown so a read-only deploy key can be +//! added to private git targets. (The legacy repositories CRUD moved to the +//! unified onboarding/targets API.) + use dioxus::prelude::*; -use serde::{Deserialize, Serialize}; - -use compliance_core::models::TrackedRepository; - -#[derive(Debug, Clone, Serialize, Deserialize, Default)] -pub struct RepositoryListResponse { - pub data: Vec, - pub total: Option, - pub page: Option, -} - -#[server] -pub async fn fetch_repositories(page: u64) -> Result { - let path = format!("/api/v1/repositories?page={page}&limit=20"); - let resp = super::agent_client::agent_get(&path) - .await? - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - let body: RepositoryListResponse = resp - .json() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - Ok(body) -} - -#[server] -pub async fn add_repository( - name: String, - git_url: String, - default_branch: String, - auth_token: Option, - auth_username: Option, - tracker_type: Option, - tracker_owner: Option, - tracker_repo: Option, - tracker_token: Option, -) -> Result<(), ServerFnError> { - let mut body = serde_json::json!({ - "name": name, - "git_url": git_url, - "default_branch": default_branch, - }); - if let Some(token) = auth_token.filter(|t| !t.is_empty()) { - body["auth_token"] = serde_json::Value::String(token); - } - if let Some(username) = auth_username.filter(|u| !u.is_empty()) { - body["auth_username"] = serde_json::Value::String(username); - } - if let Some(tt) = tracker_type.filter(|t| !t.is_empty()) { - body["tracker_type"] = serde_json::Value::String(tt); - } - if let Some(to) = tracker_owner.filter(|t| !t.is_empty()) { - body["tracker_owner"] = serde_json::Value::String(to); - } - if let Some(tr) = tracker_repo.filter(|t| !t.is_empty()) { - body["tracker_repo"] = serde_json::Value::String(tr); - } - if let Some(tk) = tracker_token.filter(|t| !t.is_empty()) { - body["tracker_token"] = serde_json::Value::String(tk); - } - - let resp = super::agent_client::agent_request(reqwest::Method::POST, "/api/v1/repositories") - .await? - .json(&body) - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - - if !resp.status().is_success() { - let body = resp.text().await.unwrap_or_default(); - return Err(ServerFnError::new(format!( - "Failed to add repository: {body}" - ))); - } - - Ok(()) -} - -#[server] -pub async fn update_repository( - repo_id: String, - name: Option, - default_branch: Option, - auth_token: Option, - auth_username: Option, - tracker_type: Option, - tracker_owner: Option, - tracker_repo: Option, - tracker_token: Option, - scan_schedule: Option, -) -> Result<(), ServerFnError> { - let mut body = serde_json::Map::new(); - if let Some(v) = name.filter(|s| !s.is_empty()) { - body.insert("name".into(), serde_json::Value::String(v)); - } - if let Some(v) = default_branch.filter(|s| !s.is_empty()) { - body.insert("default_branch".into(), serde_json::Value::String(v)); - } - if let Some(v) = auth_token { - body.insert("auth_token".into(), serde_json::Value::String(v)); - } - if let Some(v) = auth_username { - body.insert("auth_username".into(), serde_json::Value::String(v)); - } - if let Some(v) = tracker_type { - body.insert("tracker_type".into(), serde_json::Value::String(v)); - } - if let Some(v) = tracker_owner { - body.insert("tracker_owner".into(), serde_json::Value::String(v)); - } - if let Some(v) = tracker_repo { - body.insert("tracker_repo".into(), serde_json::Value::String(v)); - } - if let Some(v) = tracker_token { - body.insert("tracker_token".into(), serde_json::Value::String(v)); - } - if let Some(v) = scan_schedule { - body.insert("scan_schedule".into(), serde_json::Value::String(v)); - } - - let resp = super::agent_client::agent_request( - reqwest::Method::PATCH, - &format!("/api/v1/repositories/{repo_id}"), - ) - .await? - .json(&body) - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - - if !resp.status().is_success() { - let text = resp.text().await.unwrap_or_default(); - return Err(ServerFnError::new(format!( - "Failed to update repository: {text}" - ))); - } - - Ok(()) -} +/// Fetch the agent's SSH deploy public key. #[server] pub async fn fetch_ssh_public_key() -> Result { let resp = super::agent_client::agent_get("/api/v1/settings/ssh-public-key") @@ -163,86 +28,3 @@ pub async fn fetch_ssh_public_key() -> Result { .unwrap_or("") .to_string()) } - -#[server] -pub async fn delete_repository(repo_id: String) -> Result<(), ServerFnError> { - let resp = super::agent_client::agent_request( - reqwest::Method::DELETE, - &format!("/api/v1/repositories/{repo_id}"), - ) - .await? - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - - if !resp.status().is_success() { - let body = resp.text().await.unwrap_or_default(); - return Err(ServerFnError::new(format!( - "Failed to delete repository: {body}" - ))); - } - - Ok(()) -} - -#[server] -pub async fn trigger_repo_scan(repo_id: String) -> Result<(), ServerFnError> { - super::agent_client::agent_request( - reqwest::Method::POST, - &format!("/api/v1/repositories/{repo_id}/scan"), - ) - .await? - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - - Ok(()) -} - -#[derive(Debug, Clone, Serialize, Deserialize, Default)] -pub struct WebhookConfigResponse { - pub webhook_secret: Option, - pub tracker_type: String, -} - -#[server] -pub async fn fetch_webhook_config(repo_id: String) -> Result { - let resp = - super::agent_client::agent_get(&format!("/api/v1/repositories/{repo_id}/webhook-config")) - .await? - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - let body: WebhookConfigResponse = resp - .json() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - Ok(body) -} - -/// Check if a repository has any running scans -#[server] -pub async fn check_repo_scanning(repo_id: String) -> Result { - let resp = super::agent_client::agent_get("/api/v1/scan-runs?page=1&limit=1") - .await? - .send() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - let body: serde_json::Value = resp - .json() - .await - .map_err(|e| ServerFnError::new(e.to_string()))?; - - // Check if the most recent scan for this repo is still running - if let Some(scans) = body.get("data").and_then(|d| d.as_array()) { - for scan in scans { - let scan_repo = scan.get("repo_id").and_then(|v| v.as_str()).unwrap_or(""); - let status = scan.get("status").and_then(|v| v.as_str()).unwrap_or(""); - if scan_repo == repo_id && status == "running" { - return Ok(true); - } - } - } - - Ok(false) -} diff --git a/compliance-dashboard/src/pages/chat_index.rs b/compliance-dashboard/src/pages/chat_index.rs index 5d56141..bdff7cf 100644 --- a/compliance-dashboard/src/pages/chat_index.rs +++ b/compliance-dashboard/src/pages/chat_index.rs @@ -2,11 +2,11 @@ use dioxus::prelude::*; use crate::app::Route; use crate::components::page_header::PageHeader; -use crate::infrastructure::repositories::fetch_repositories; +use crate::infrastructure::onboarding::fetch_targets; #[component] pub fn ChatIndexPage() -> Element { - let repos = use_resource(|| async { fetch_repositories(1).await.ok() }); + let repos = use_resource(|| async { fetch_targets().await.ok() }); rsx! { PageHeader { @@ -28,10 +28,32 @@ pub fn ChatIndexPage() -> Element { div { class: "graph-index-grid", for repo in repo_list { { - let repo_id = repo.id.map(|id| id.to_hex()).unwrap_or_default(); - let name = repo.name.clone(); - let url = repo.git_url.clone(); - let branch = repo.default_branch.clone(); + let repo_id = repo.get("_id").and_then(|o| o.get("$oid")).and_then(|s| s.as_str()).unwrap_or_default().to_string(); + let name = repo.get("name").and_then(|n| n.as_str()).unwrap_or_default().to_string(); + let url = repo + .get("artifacts") + .and_then(|a| a.as_array()) + .and_then(|arr| { + arr.iter().find(|a| { + a.get("kind").and_then(|k| k.as_str()) == Some("git_repo") + }) + }) + .and_then(|a| a.get("source_ref")) + .and_then(|s| s.as_str()) + .unwrap_or_default() + .to_string(); + let branch = repo + .get("artifacts") + .and_then(|a| a.as_array()) + .and_then(|arr| { + arr.iter().find_map(|a| { + a.get("git") + .and_then(|g| g.get("default_branch")) + .and_then(|b| b.as_str()) + }) + }) + .unwrap_or("main") + .to_string(); rsx! { Link { to: Route::ChatPage { repo_id }, diff --git a/compliance-dashboard/src/pages/graph_index.rs b/compliance-dashboard/src/pages/graph_index.rs index fedf10a..e4be9ae 100644 --- a/compliance-dashboard/src/pages/graph_index.rs +++ b/compliance-dashboard/src/pages/graph_index.rs @@ -2,11 +2,11 @@ use dioxus::prelude::*; use crate::app::Route; use crate::components::page_header::PageHeader; -use crate::infrastructure::repositories::fetch_repositories; +use crate::infrastructure::onboarding::fetch_targets; #[component] pub fn GraphIndexPage() -> Element { - let repos = use_resource(|| async { fetch_repositories(1).await.ok() }); + let repos = use_resource(|| async { fetch_targets().await.ok() }); rsx! { PageHeader { @@ -28,27 +28,34 @@ pub fn GraphIndexPage() -> Element { div { class: "graph-index-grid", for repo in repo_list { { - let repo_id = repo.id.map(|id| id.to_hex()).unwrap_or_default(); - let name = repo.name.clone(); - let url = repo.git_url.clone(); - let branch = repo.default_branch.clone(); - let findings = repo.findings_count; + let repo_id = repo.get("_id").and_then(|o| o.get("$oid")).and_then(|s| s.as_str()).unwrap_or_default().to_string(); + let name = repo.get("name").and_then(|n| n.as_str()).unwrap_or_default().to_string(); + let url = repo + .get("artifacts") + .and_then(|a| a.as_array()) + .and_then(|arr| { + arr.iter().find(|a| { + a.get("kind").and_then(|k| k.as_str()) == Some("git_repo") + }) + }) + .and_then(|a| a.get("source_ref")) + .and_then(|s| s.as_str()) + .unwrap_or_default() + .to_string(); + let branch = repo + .get("artifacts") + .and_then(|a| a.as_array()) + .and_then(|arr| { + arr.iter().find_map(|a| { + a.get("git") + .and_then(|g| g.get("default_branch")) + .and_then(|b| b.as_str()) + }) + }) + .unwrap_or("main") + .to_string(); + let findings = repo.get("findings_count").and_then(|n| n.as_u64()).unwrap_or(0); let findings_label = if findings != 1 { format!("{findings} findings") } else { "1 finding".to_string() }; - let updated = { - let now = chrono::Utc::now(); - let diff = now.signed_duration_since(repo.updated_at); - if diff.num_minutes() < 1 { - "just now".to_string() - } else if diff.num_hours() < 1 { - format!("{}m ago", diff.num_minutes()) - } else if diff.num_days() < 1 { - format!("{}h ago", diff.num_hours()) - } else if diff.num_days() < 30 { - format!("{}d ago", diff.num_days()) - } else { - repo.updated_at.format("%Y-%m-%d").to_string() - } - }; rsx! { Link { to: Route::GraphExplorerPage { repo_id }, @@ -67,9 +74,6 @@ pub fn GraphIndexPage() -> Element { span { class: "graph-repo-card-tag graph-repo-card-tag-findings", "{findings_label}" } - span { class: "graph-repo-card-tag", - "Updated {updated}" - } } } } diff --git a/compliance-dashboard/src/pages/mod.rs b/compliance-dashboard/src/pages/mod.rs index bd9722c..c352195 100644 --- a/compliance-dashboard/src/pages/mod.rs +++ b/compliance-dashboard/src/pages/mod.rs @@ -16,7 +16,6 @@ pub mod onboarding; pub mod overview; pub mod pentest_dashboard; pub mod pentest_session; -pub mod repositories; pub mod sbom; pub mod targets; @@ -38,6 +37,5 @@ pub use onboarding::OnboardingPage; pub use overview::OverviewPage; pub use pentest_dashboard::PentestDashboardPage; pub use pentest_session::PentestSessionPage; -pub use repositories::RepositoriesPage; pub use sbom::SbomPage; pub use targets::TargetsPage; diff --git a/compliance-dashboard/src/pages/overview.rs b/compliance-dashboard/src/pages/overview.rs index fc25dc4..cfb2cc6 100644 --- a/compliance-dashboard/src/pages/overview.rs +++ b/compliance-dashboard/src/pages/overview.rs @@ -6,7 +6,7 @@ use crate::app::Route; use crate::components::page_header::PageHeader; use crate::components::stat_card::StatCard; use crate::infrastructure::mcp::fetch_mcp_servers; -use crate::infrastructure::repositories::fetch_repositories; +use crate::infrastructure::onboarding::fetch_targets; #[cfg(feature = "server")] use crate::infrastructure::stats::fetch_overview_stats; @@ -26,7 +26,7 @@ pub fn OverviewPage() -> Element { } }); - let repos = use_resource(|| async { fetch_repositories(1).await.ok() }); + let repos = use_resource(|| async { fetch_targets().await.ok() }); let mcp_servers = use_resource(|| async { fetch_mcp_servers().await.ok() }); rsx! { @@ -94,8 +94,8 @@ pub fn OverviewPage() -> Element { style: "display: grid; grid-template-columns: repeat(3, 1fr); gap: 1rem; padding: 1rem;", for repo in repo_list { { - let repo_id = repo.id.map(|id| id.to_hex()).unwrap_or_default(); - let name = repo.name.clone(); + let repo_id = repo.get("_id").and_then(|o| o.get("$oid")).and_then(|s| s.as_str()).unwrap_or_default().to_string(); + let name = repo.get("name").and_then(|n| n.as_str()).unwrap_or_default().to_string(); rsx! { Link { to: Route::ChatPage { repo_id }, diff --git a/compliance-dashboard/src/pages/repositories.rs b/compliance-dashboard/src/pages/repositories.rs deleted file mode 100644 index 89bf163..0000000 --- a/compliance-dashboard/src/pages/repositories.rs +++ /dev/null @@ -1,428 +0,0 @@ -use dioxus::prelude::*; -use dioxus_free_icons::icons::bs_icons::*; -#[allow(unused_imports)] -use dioxus_free_icons::icons::bs_icons::{BsGear, BsPencil}; -use dioxus_free_icons::Icon; - -use crate::components::page_header::PageHeader; -use crate::components::pagination::Pagination; -use crate::components::toast::{ToastType, Toasts}; -use crate::pages::graph_explorer::GraphExplorerInline; - -async fn async_sleep_5s() { - #[cfg(feature = "web")] - { - gloo_timers::future::TimeoutFuture::new(5_000).await; - } - #[cfg(not(feature = "web"))] - { - tokio::time::sleep(std::time::Duration::from_secs(5)).await; - } -} - -#[component] -pub fn RepositoriesPage() -> Element { - let mut page = use_signal(|| 1u64); - let mut toasts = use_context::(); - let mut confirm_delete = use_signal(|| Option::<(String, String)>::None); // (id, name) - let mut edit_repo_id = use_signal(|| Option::::None); - let mut edit_name = use_signal(String::new); - let mut edit_branch = use_signal(String::new); - let mut edit_tracker_type = use_signal(String::new); - let mut edit_tracker_owner = use_signal(String::new); - let mut edit_tracker_repo = use_signal(String::new); - let mut edit_tracker_token = use_signal(String::new); - let mut edit_saving = use_signal(|| false); - let mut edit_webhook_secret = use_signal(|| Option::::None); - let mut edit_webhook_tracker = use_signal(String::new); - let mut scanning_ids = use_signal(Vec::::new); - let mut graph_repo_id = use_signal(|| Option::::None); - - let mut repos = use_resource(move || { - let p = page(); - async move { - crate::infrastructure::repositories::fetch_repositories(p) - .await - .ok() - } - }); - - rsx! { - PageHeader { - title: "Repositories", - description: "Legacy git repositories. Onboard new targets from Targets / Onboard.", - } - - // ── Delete confirmation dialog ── - if let Some((del_id, del_name)) = confirm_delete() { - div { class: "modal-overlay", - div { class: "modal-dialog", - h3 { "Delete Repository" } - p { - "Are you sure you want to delete " - strong { "{del_name}" } - "?" - } - p { class: "modal-warning", - "This will permanently remove all associated findings, SBOM entries, scan runs, graph data, embeddings, and CVE alerts." - } - div { class: "modal-actions", - button { - class: "btn btn-secondary", - onclick: move |_| confirm_delete.set(None), - "Cancel" - } - button { - class: "btn btn-danger", - onclick: move |_| { - let id = del_id.clone(); - let name = del_name.clone(); - confirm_delete.set(None); - spawn(async move { - match crate::infrastructure::repositories::delete_repository(id).await { - Ok(_) => { - toasts.push(ToastType::Success, format!("{name} deleted")); - repos.restart(); - } - Err(e) => toasts.push(ToastType::Error, e.to_string()), - } - }); - }, - "Delete" - } - } - } - } - } - - // ── Edit repository dialog ── - if let Some(eid) = edit_repo_id() { - div { class: "modal-overlay", - div { class: "modal-dialog", - h3 { "Edit Repository" } - div { class: "form-group", - label { "Name" } - input { - r#type: "text", - value: "{edit_name}", - oninput: move |e| edit_name.set(e.value()), - } - } - div { class: "form-group", - label { "Default Branch" } - input { - r#type: "text", - value: "{edit_branch}", - oninput: move |e| edit_branch.set(e.value()), - } - } - h4 { style: "margin-top: 16px; margin-bottom: 8px; font-size: 14px; color: var(--text-secondary);", "Issue Tracker" } - div { class: "form-group", - label { "Tracker Type" } - select { - value: "{edit_tracker_type}", - onchange: move |e| edit_tracker_type.set(e.value()), - option { value: "", "None" } - option { value: "github", "GitHub" } - option { value: "gitlab", "GitLab" } - option { value: "gitea", "Gitea" } - option { value: "jira", "Jira" } - } - } - div { class: "form-group", - label { "Owner / Namespace" } - input { - r#type: "text", - placeholder: "org-name", - value: "{edit_tracker_owner}", - oninput: move |e| edit_tracker_owner.set(e.value()), - } - } - div { class: "form-group", - label { "Repository / Project" } - input { - r#type: "text", - placeholder: "repo-name", - value: "{edit_tracker_repo}", - oninput: move |e| edit_tracker_repo.set(e.value()), - } - } - div { class: "form-group", - label { "Tracker Token (leave empty to keep existing)" } - input { - r#type: "password", - placeholder: "Enter new token to change", - value: "{edit_tracker_token}", - oninput: move |e| edit_tracker_token.set(e.value()), - } - } - // Webhook configuration section - if let Some(secret) = edit_webhook_secret() { - h4 { - style: "margin-top: 16px; margin-bottom: 8px; font-size: 14px; color: var(--text-secondary);", - "Webhook Configuration" - } - p { - style: "font-size: 12px; color: var(--text-secondary); margin-bottom: 8px;", - "Add this webhook in your repository settings to enable push-triggered scans and PR reviews." - } - div { class: "form-group", - label { "Webhook URL" } - { - #[cfg(feature = "web")] - let origin = web_sys::window() - .and_then(|w: web_sys::Window| w.location().origin().ok()) - .unwrap_or_default(); - #[cfg(not(feature = "web"))] - let origin = String::new(); - let webhook_url = format!("{origin}/webhook/{}/{eid}", edit_webhook_tracker()); - rsx! { - div { class: "copyable", - input { - r#type: "text", - readonly: true, - style: "font-family: monospace; font-size: 12px; flex: 1;", - value: "{webhook_url}", - } - crate::components::copy_button::CopyButton { value: webhook_url.clone() } - } - } - } - } - div { class: "form-group", - label { "Webhook Secret" } - div { class: "copyable", - input { - r#type: "text", - readonly: true, - style: "font-family: monospace; font-size: 12px; flex: 1;", - value: "{secret}", - } - crate::components::copy_button::CopyButton { value: secret.clone() } - } - } - } - div { class: "modal-actions", - button { - class: "btn btn-secondary", - onclick: move |_| edit_repo_id.set(None), - "Cancel" - } - button { - class: "btn btn-primary", - disabled: edit_saving(), - onclick: move |_| { - let id = eid.clone(); - let nm = { let v = edit_name(); if v.is_empty() { None } else { Some(v) } }; - let br = { let v = edit_branch(); if v.is_empty() { None } else { Some(v) } }; - let tt = { let v = edit_tracker_type(); if v.is_empty() { None } else { Some(v) } }; - let t_owner = { let v = edit_tracker_owner(); if v.is_empty() { None } else { Some(v) } }; - let t_repo = { let v = edit_tracker_repo(); if v.is_empty() { None } else { Some(v) } }; - let t_tok = { let v = edit_tracker_token(); if v.is_empty() { None } else { Some(v) } }; - edit_saving.set(true); - spawn(async move { - match crate::infrastructure::repositories::update_repository( - id, nm, br, None, None, tt, t_owner, t_repo, t_tok, None, - ).await { - Ok(_) => { - toasts.push(ToastType::Success, "Repository updated"); - repos.restart(); - } - Err(e) => toasts.push(ToastType::Error, e.to_string()), - } - edit_saving.set(false); - edit_repo_id.set(None); - }); - }, - if edit_saving() { "Saving..." } else { "Save" } - } - } - } - } - } - - match &*repos.read() { - Some(Some(resp)) => { - let total_pages = resp.total.unwrap_or(0).div_ceil(20).max(1); - rsx! { - div { class: "card", - div { class: "table-wrapper", - table { - thead { - tr { - th { "Name" } - th { "Git URL" } - th { "Branch" } - th { "Findings" } - th { "Last Scanned" } - th { "Actions" } - } - } - tbody { - for repo in &resp.data { - { - let repo_id = repo.id.as_ref().map(|id| id.to_hex()).unwrap_or_default(); - let repo_id_scan = repo_id.clone(); - let repo_id_del = repo_id.clone(); - let repo_id_edit = repo_id.clone(); - let repo_name_del = repo.name.clone(); - let edit_repo_data = repo.clone(); - let is_scanning = scanning_ids().contains(&repo_id); - rsx! { - tr { - td { "{repo.name}" } - td { - style: "font-size: 12px; font-family: monospace;", - "{repo.git_url}" - } - td { "{repo.default_branch}" } - td { "{repo.findings_count}" } - td { - { - let now = chrono::Utc::now(); - let diff = now.signed_duration_since(repo.updated_at); - let label = if diff.num_minutes() < 1 { - "just now".to_string() - } else if diff.num_hours() < 1 { - format!("{}m ago", diff.num_minutes()) - } else if diff.num_days() < 1 { - format!("{}h ago", diff.num_hours()) - } else if diff.num_days() < 30 { - format!("{}d ago", diff.num_days()) - } else { - repo.updated_at.format("%Y-%m-%d").to_string() - }; - rsx! { span { style: "font-size: 12px;", "{label}" } } - } - } - td { style: "display: flex; gap: 4px;", - button { - class: if graph_repo_id().as_deref() == Some(repo_id.as_str()) { "btn btn-ghost btn-active" } else { "btn btn-ghost" }, - title: "View graph", - onclick: { - let rid = repo_id.clone(); - move |_| { - if graph_repo_id().as_deref() == Some(rid.as_str()) { - graph_repo_id.set(None); - } else { - graph_repo_id.set(Some(rid.clone())); - } - } - }, - Icon { icon: BsDiagram3, width: 16, height: 16 } - } - button { - class: "btn btn-ghost", - title: "Edit repository", - onclick: move |_| { - edit_name.set(edit_repo_data.name.clone()); - edit_branch.set(edit_repo_data.default_branch.clone()); - edit_tracker_type.set( - edit_repo_data.tracker_type.as_ref().map(|t| t.to_string()).unwrap_or_default() - ); - edit_tracker_owner.set(edit_repo_data.tracker_owner.clone().unwrap_or_default()); - edit_tracker_repo.set(edit_repo_data.tracker_repo.clone().unwrap_or_default()); - edit_tracker_token.set(String::new()); - edit_webhook_secret.set(None); - edit_webhook_tracker.set(String::new()); - edit_repo_id.set(Some(repo_id_edit.clone())); - // Fetch webhook config in background - let rid = repo_id_edit.clone(); - spawn(async move { - if let Ok(cfg) = crate::infrastructure::repositories::fetch_webhook_config(rid).await { - edit_webhook_secret.set(cfg.webhook_secret); - edit_webhook_tracker.set(cfg.tracker_type); - } - }); - }, - Icon { icon: BsPencil, width: 16, height: 16 } - } - button { - class: if is_scanning { "btn btn-ghost btn-scanning" } else { "btn btn-ghost" }, - title: "Trigger scan", - disabled: is_scanning, - onclick: move |_| { - let id = repo_id_scan.clone(); - // Add to scanning set - let mut ids = scanning_ids(); - ids.push(id.clone()); - scanning_ids.set(ids); - spawn(async move { - match crate::infrastructure::repositories::trigger_repo_scan(id.clone()).await { - Ok(_) => { - toasts.push(ToastType::Success, "Scan triggered"); - // Poll until scan completes - loop { - async_sleep_5s().await; - match crate::infrastructure::repositories::check_repo_scanning(id.clone()).await { - Ok(false) => break, - Ok(true) => continue, - Err(_) => break, - } - } - toasts.push(ToastType::Success, "Scan complete"); - repos.restart(); - } - Err(e) => toasts.push(ToastType::Error, e.to_string()), - } - // Remove from scanning set - let mut ids = scanning_ids(); - ids.retain(|i| i != &id); - scanning_ids.set(ids); - }); - }, - if is_scanning { - span { class: "spinner" } - } else { - Icon { icon: BsPlayCircle, width: 16, height: 16 } - } - } - button { - class: "btn btn-ghost btn-ghost-danger", - title: "Delete repository", - onclick: move |_| { - confirm_delete.set(Some((repo_id_del.clone(), repo_name_del.clone()))); - }, - Icon { icon: BsTrash, width: 16, height: 16 } - } - } - } - } - } - } - } - } - } - Pagination { - current_page: page(), - total_pages: total_pages, - on_page_change: move |p| page.set(p), - } - } - - // Inline graph explorer - if let Some(rid) = graph_repo_id() { - div { class: "card", style: "margin-top: 16px;", - div { class: "card-header", style: "display: flex; justify-content: space-between; align-items: center;", - span { "Code Graph" } - button { - class: "btn btn-sm btn-ghost", - title: "Close graph", - onclick: move |_| { graph_repo_id.set(None); }, - Icon { icon: BsX, width: 18, height: 18 } - } - } - GraphExplorerInline { repo_id: rid } - } - } - } - }, - Some(None) => rsx! { - div { class: "card", p { "Failed to load repositories." } } - }, - None => rsx! { - div { class: "loading", "Loading repositories..." } - }, - } - } -}