From fa3a6b71fb7123a51bb13c562596ebae3f49d681 Mon Sep 17 00:00:00 2001 From: Sharang Parnerkar <30073382+mighty840@users.noreply.github.com> Date: Fri, 10 Jul 2026 18:25:46 +0200 Subject: [PATCH 1/2] feat(migrate): id-preserving mapping core for onboarding backfill (#132 part 1) Pure transforms folding legacy records into the unified model, preserving _id so every downstream collection keyed by that hex id keeps resolving: - repo_to_target: TrackedRepository -> OnboardedTarget (GitRepo artifact with branch/watermark/auth; tracker config -> scan_config; webhook/schedule/counts). - dast_to_artifact / dast_to_standalone_target: DastTarget -> LiveUrl artifact (crawl config + auth), or a standalone target when repo-less. - fold_dast_into_target: append a LiveUrl to a repo-derived target (idempotent), promoting BackendService -> WebApp for a browser app. 5 unit tests. The per-tenant DB orchestration (idempotent backfill + revert) and the `migrate onboarding` CLI subcommand are the next increment. Refs #132. Co-Authored-By: Claude Fable 5 --- compliance-agent/src/lib.rs | 1 + compliance-agent/src/migrate/mod.rs | 8 + compliance-agent/src/migrate/onboarding.rs | 234 +++++++++++++++++++++ 3 files changed, 243 insertions(+) create mode 100644 compliance-agent/src/migrate/mod.rs create mode 100644 compliance-agent/src/migrate/onboarding.rs diff --git a/compliance-agent/src/lib.rs b/compliance-agent/src/lib.rs index 2cc5179..a0a07bf 100644 --- a/compliance-agent/src/lib.rs +++ b/compliance-agent/src/lib.rs @@ -8,6 +8,7 @@ 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/migrate/mod.rs b/compliance-agent/src/migrate/mod.rs new file mode 100644 index 0000000..2546602 --- /dev/null +++ b/compliance-agent/src/migrate/mod.rs @@ -0,0 +1,8 @@ +//! 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 new file mode 100644 index 0000000..9c1b421 --- /dev/null +++ b/compliance-agent/src/migrate/onboarding.rs @@ -0,0 +1,234 @@ +//! 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, +}; + +/// 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 +} + +#[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); + } +} -- 2.54.0 From 0d83859bcf75a747a977e0ded617b13eca557d15 Mon Sep 17 00:00:00 2001 From: Sharang Parnerkar <30073382+mighty840@users.noreply.github.com> Date: Fri, 10 Jul 2026 19:49:52 +0200 Subject: [PATCH 2/2] feat(migrate): per-tenant onboarding backfill + revert + CLI (#132 part 2) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire the id-preserving mappers into a runnable, idempotent, reversible migration: - backfill_onboarded_targets(db, dry_run): repositories -> onboarded_targets; dast_targets fold into the linked repo target (append LiveUrl, promote type, relink dast_scan_runs/pentest_sessions history) or migrate standalone. Skips existing (by _id), writes a schema_migrations marker; dry-run computes the report without writing. Legacy collections are never deleted. - revert(db): replay the relink log, drop onboarded_targets + the log, clear the marker — restores the pre-migration state. - CLI: `compliance-agent migrate onboarding [--all | --tenant ] [--dry-run] [--revert]`, with DatabasePool::list_tenant_ids for --all. - database.rs: collection_named accessor + list_tenant_ids helper. Integration test (real Mongo, local-only — CI is --lib) covers fold + relink + idempotency + revert end to end. 5 mapper unit tests run in CI. Closes #132. Co-Authored-By: Claude Fable 5 --- compliance-agent/src/database.rs | 24 +++ compliance-agent/src/main.rs | 55 +++++- compliance-agent/src/migrate/onboarding.rs | 172 ++++++++++++++++++ .../tests/integration/migration.rs | 156 ++++++++++++++++ compliance-agent/tests/integration/mod.rs | 1 + 5 files changed, 407 insertions(+), 1 deletion(-) create mode 100644 compliance-agent/tests/integration/migration.rs diff --git a/compliance-agent/src/database.rs b/compliance-agent/src/database.rs index 1105a58..cdadf3b 100644 --- a/compliance-agent/src/database.rs +++ b/compliance-agent/src/database.rs @@ -179,6 +179,23 @@ impl DatabasePool { .collect()) } + /// Tenant ids for every provisioned tenant database, derived by stripping + /// the `_` from the database names. Skips the admin database + /// (`__admin`). Hash-fallback names (very long tenant_ids) are lost + /// at the cluster level and cannot be recovered here — in practice tenant + /// ids are UUIDs and never hit that path. Used by the migration CLI's + /// `--all` mode. + pub async fn list_tenant_ids(&self) -> Result, AgentError> { + let prefix = format!("{}_", self.db_prefix); + Ok(self + .list_tenant_db_names() + .await? + .into_iter() + .filter_map(|n| n.strip_prefix(&prefix).map(str::to_string)) + .filter(|id| !id.starts_with('_')) + .collect()) + } + /// Drop the database for a specific tenant. Used by GDPR delete /// and tenant offboarding. Idempotent — dropping a non-existent /// database is a no-op at the driver level. @@ -521,6 +538,13 @@ impl Database { self.inner.collection("onboarded_targets") } + /// A typed handle to an arbitrary collection by name. For bookkeeping + /// collections without a dedicated model (e.g. `schema_migrations`, + /// `onboarding_migration_log`). + pub fn collection_named(&self, name: &str) -> Collection { + self.inner.collection(name) + } + pub fn dast_scan_runs(&self) -> Collection { self.inner.collection("dast_scan_runs") } diff --git a/compliance-agent/src/main.rs b/compliance-agent/src/main.rs index 110634f..95323a9 100644 --- a/compliance-agent/src/main.rs +++ b/compliance-agent/src/main.rs @@ -1,4 +1,50 @@ -use compliance_agent::{agent, api, config, database, scheduler, ssh, webhooks}; +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(()) +} #[tokio::main] async fn main() -> Result<(), Box> { @@ -31,6 +77,13 @@ 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/onboarding.rs b/compliance-agent/src/migrate/onboarding.rs index 9c1b421..7652316 100644 --- a/compliance-agent/src/migrate/onboarding.rs +++ b/compliance-agent/src/migrate/onboarding.rs @@ -12,6 +12,14 @@ 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)] @@ -119,6 +127,170 @@ pub fn dast_to_standalone_target(dast: &DastTarget) -> OnboardedTarget { 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 { diff --git a/compliance-agent/tests/integration/migration.rs b/compliance-agent/tests/integration/migration.rs new file mode 100644 index 0000000..9b6a116 --- /dev/null +++ b/compliance-agent/tests/integration/migration.rs @@ -0,0 +1,156 @@ +// 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 1baec19..5f9870e 100644 --- a/compliance-agent/tests/integration/mod.rs +++ b/compliance-agent/tests/integration/mod.rs @@ -7,3 +7,4 @@ // Or nightly: (via CI with MongoDB service container) mod api; +mod migration; -- 2.54.0