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] 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;