407 lines
15 KiB
Rust
407 lines
15 KiB
Rust
//! 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<bool, AgentError> {
|
|
let found = db
|
|
.collection_named::<Document>("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<MigrationReport, AgentError> {
|
|
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::<Document>("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::<Document>("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::<Document>("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::<Document>("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<String>, 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);
|
|
}
|
|
}
|