use std::collections::BTreeMap;
use super::naming::MIGRATION_FILE_EXT;
use super::projection::BucketKey;
use super::schema::AppliedSchema;
use super::target::FilesystemBucket;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DriftDiagnostic {
pub bucket: BucketKey,
pub kind: DriftKind,
pub text: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum DriftKind {
Outcome2ComposedNotApplied,
Outcome3Drift,
Outcome4PendingInvalid,
D004FilesystemUnregistered,
D004RegisteredMissingFolder,
}
impl DriftKind {
pub fn is_outcome3_drift(self) -> bool {
matches!(self, DriftKind::Outcome3Drift)
}
}
pub fn classify_bucket(
bucket: &BucketKey,
models: Option<&AppliedSchema>,
pending: Option<&AppliedSchema>,
snapshot: Option<&AppliedSchema>,
) -> Option<DriftDiagnostic> {
classify_bucket_with_pending(bucket, models, pending, snapshot, None)
}
pub fn classify_bucket_with_pending(
bucket: &BucketKey,
models: Option<&AppliedSchema>,
pending: Option<&AppliedSchema>,
snapshot: Option<&AppliedSchema>,
pending_version: Option<&str>,
) -> Option<DriftDiagnostic> {
let models_eq_pending = match (models, pending) {
(Some(m), Some(p)) => schema_equiv(m, p),
(None, None) => true,
_ => false,
};
let models_eq_snapshot = match (models, snapshot) {
(Some(m), Some(s)) => schema_equiv(m, s),
(None, None) => true,
_ => false,
};
let pending_eq_snapshot = match (pending, snapshot) {
(Some(p), Some(s)) => schema_equiv(p, s),
(None, None) => true,
_ => false,
};
if pending.is_none() && models_eq_snapshot {
return None;
}
if pending.is_some() && models_eq_pending && models_eq_snapshot {
return None;
}
if pending.is_some() && models_eq_pending && !models_eq_snapshot {
return Some(DriftDiagnostic {
bucket: bucket.clone(),
kind: DriftKind::Outcome2ComposedNotApplied,
text: format_warning_outcome2(bucket, pending_version),
});
}
if pending.is_some() && !models_eq_pending && !pending_eq_snapshot {
return Some(DriftDiagnostic {
bucket: bucket.clone(),
kind: DriftKind::Outcome4PendingInvalid,
text: format_warning_outcome4(bucket),
});
}
if !models_eq_snapshot && (pending.is_none() || pending_eq_snapshot) {
return Some(DriftDiagnostic {
bucket: bucket.clone(),
kind: DriftKind::Outcome3Drift,
text: format_warning_outcome3(bucket),
});
}
Some(DriftDiagnostic {
bucket: bucket.clone(),
kind: DriftKind::Outcome3Drift,
text: format_warning_outcome3(bucket),
})
}
pub fn classify_filesystem_drift(
filesystem: &std::collections::BTreeSet<FilesystemBucket>,
snapshots: &BTreeMap<BucketKey, AppliedSchema>,
) -> Vec<DriftDiagnostic> {
let mut diagnostics = Vec::new();
let mut registered_per_db: BTreeMap<String, std::collections::BTreeSet<String>> =
BTreeMap::new();
for (key, snap) in snapshots {
let entry = registered_per_db.entry(key.database.clone()).or_default();
for app in &snap.registered_apps {
entry.insert(app.clone());
}
}
for fs_bucket in filesystem {
let known_apps = registered_per_db.get(&fs_bucket.database);
let is_registered = known_apps
.map(|set| set.contains(&fs_bucket.app))
.unwrap_or(false);
if !is_registered {
let bucket = BucketKey {
database: fs_bucket.database.clone(),
app: fs_bucket.app.clone(),
};
diagnostics.push(DriftDiagnostic {
kind: DriftKind::D004FilesystemUnregistered,
text: format_warning_d004_unregistered(&bucket),
bucket,
});
}
}
let fs_lookup: std::collections::BTreeSet<(String, String)> = filesystem
.iter()
.map(|b| (b.database.clone(), b.app.clone()))
.collect();
for (database, apps) in ®istered_per_db {
for app in apps {
if !fs_lookup.contains(&(database.clone(), app.clone())) {
let bucket = BucketKey {
database: database.clone(),
app: app.clone(),
};
diagnostics.push(DriftDiagnostic {
kind: DriftKind::D004RegisteredMissingFolder,
text: format_warning_d004_missing(&bucket),
bucket,
});
}
}
}
diagnostics
}
pub fn format_warning_outcome2(bucket: &BucketKey, pending_version: Option<&str>) -> String {
let (filename, version) = match pending_version {
Some(v) => (format!("{v}{ext}", ext = MIGRATION_FILE_EXT), v.to_string()),
None => (
"<unknown>".to_string() + MIGRATION_FILE_EXT,
"<unknown>".to_string(),
),
};
format!(
"composed migration not yet applied: {filename} (version {version}; bucket {database}/{app})",
database = bucket.database,
app = super::target::app_dirname(&bucket.app),
)
}
pub fn format_warning_outcome3(bucket: &BucketKey) -> String {
format!(
"model drift detected for {database}/{app}; run `djogi migrations compose` to stage the delta",
database = bucket.database,
app = super::target::app_dirname(&bucket.app),
)
}
pub fn format_warning_outcome4(bucket: &BucketKey) -> String {
format!(
"pending compose for {database}/{app} is stale relative to model state; re-run `djogi migrations compose`",
database = bucket.database,
app = super::target::app_dirname(&bucket.app),
)
}
pub fn format_warning_d004_unregistered(bucket: &BucketKey) -> String {
format!(
"D004: filesystem app \"{database}/{app}\" not registered in snapshot",
database = bucket.database,
app = super::target::app_dirname(&bucket.app),
)
}
pub fn format_warning_d004_missing(bucket: &BucketKey) -> String {
format!(
"D004: registered app \"{database}/{app}\" missing from filesystem",
database = bucket.database,
app = super::target::app_dirname(&bucket.app),
)
}
fn schema_equiv(a: &AppliedSchema, b: &AppliedSchema) -> bool {
a.djogi_version == b.djogi_version
&& a.enums == b.enums
&& a.format_version == b.format_version
&& a.indexes == b.indexes
&& a.models == b.models
&& a.registered_apps == b.registered_apps
}
#[cfg(test)]
mod tests {
use super::*;
use crate::migrate::schema::SNAPSHOT_FORMAT_VERSION;
use std::collections::BTreeMap;
fn empty_schema() -> AppliedSchema {
AppliedSchema {
djogi_version: "0.1.0".to_string(),
enums: BTreeMap::new(),
format_version: SNAPSHOT_FORMAT_VERSION.to_string(),
generated_at: "2026-04-25T00:00:00Z".to_string(),
indexes: Vec::new(),
models: BTreeMap::new(),
registered_apps: vec!["".to_string()],
}
}
fn drifted_schema() -> AppliedSchema {
AppliedSchema {
djogi_version: "9.9.9".to_string(),
..empty_schema()
}
}
fn global_bucket() -> BucketKey {
BucketKey {
database: "main".into(),
app: "".into(),
}
}
#[test]
fn outcome1_no_pending_models_match_snapshot_returns_none() {
let m = empty_schema();
let s = empty_schema();
let diag = classify_bucket(&global_bucket(), Some(&m), None, Some(&s));
assert!(diag.is_none());
}
#[test]
fn outcome1_all_three_match_returns_none() {
let m = empty_schema();
let p = empty_schema();
let s = empty_schema();
let diag = classify_bucket(&global_bucket(), Some(&m), Some(&p), Some(&s));
assert!(diag.is_none());
}
#[test]
fn outcome2_composed_not_applied() {
let m = drifted_schema();
let p = drifted_schema();
let s = empty_schema();
let diag = classify_bucket(&global_bucket(), Some(&m), Some(&p), Some(&s)).expect("diag");
assert_eq!(diag.kind, DriftKind::Outcome2ComposedNotApplied);
assert_eq!(
diag.text,
"composed migration not yet applied: <unknown>.sdjql (version <unknown>; bucket main/_global_)"
);
}
#[test]
fn outcome2_composed_not_applied_with_pending_version() {
let m = drifted_schema();
let p = drifted_schema();
let s = empty_schema();
let diag = classify_bucket_with_pending(
&global_bucket(),
Some(&m),
Some(&p),
Some(&s),
Some("V20260425010203__add_widgets"),
)
.expect("diag");
assert_eq!(diag.kind, DriftKind::Outcome2ComposedNotApplied);
assert_eq!(
diag.text,
"composed migration not yet applied: V20260425010203__add_widgets.sdjql \
(version V20260425010203__add_widgets; bucket main/_global_)"
);
}
#[test]
fn outcome3_drift_no_pending() {
let m = drifted_schema();
let s = empty_schema();
let diag = classify_bucket(&global_bucket(), Some(&m), None, Some(&s)).expect("diag");
assert_eq!(diag.kind, DriftKind::Outcome3Drift);
assert_eq!(
diag.text,
"model drift detected for main/_global_; run `djogi migrations compose` to stage the delta"
);
}
#[test]
fn outcome4_pending_invalid() {
let m = drifted_schema();
let p = AppliedSchema {
djogi_version: "5.5.5".to_string(),
..empty_schema()
};
let s = empty_schema();
let diag = classify_bucket(&global_bucket(), Some(&m), Some(&p), Some(&s)).expect("diag");
assert_eq!(diag.kind, DriftKind::Outcome4PendingInvalid);
assert_eq!(
diag.text,
"pending compose for main/_global_ is stale relative to model state; re-run `djogi migrations compose`"
);
}
#[test]
fn d004_filesystem_unregistered() {
let mut snapshots = BTreeMap::new();
let mut snap = empty_schema();
snap.registered_apps = vec!["".to_string(), "billing".to_string()];
snapshots.insert(global_bucket(), snap);
let fs: std::collections::BTreeSet<FilesystemBucket> = [
FilesystemBucket {
database: "main".into(),
app: "".into(),
},
FilesystemBucket {
database: "main".into(),
app: "billing".into(),
},
FilesystemBucket {
database: "main".into(),
app: "ghost".into(),
},
]
.into_iter()
.collect();
let diags = classify_filesystem_drift(&fs, &snapshots);
assert_eq!(diags.len(), 1);
assert_eq!(diags[0].kind, DriftKind::D004FilesystemUnregistered);
assert_eq!(
diags[0].text,
"D004: filesystem app \"main/ghost\" not registered in snapshot"
);
}
#[test]
fn d004_registered_missing_folder() {
let mut snapshots = BTreeMap::new();
let mut snap = empty_schema();
snap.registered_apps = vec!["".to_string(), "billing".to_string()];
snapshots.insert(global_bucket(), snap);
let fs: std::collections::BTreeSet<FilesystemBucket> = [FilesystemBucket {
database: "main".into(),
app: "".into(),
}]
.into_iter()
.collect();
let diags = classify_filesystem_drift(&fs, &snapshots);
assert_eq!(diags.len(), 1);
assert_eq!(diags[0].kind, DriftKind::D004RegisteredMissingFolder);
assert_eq!(
diags[0].text,
"D004: registered app \"main/billing\" missing from filesystem"
);
}
#[test]
fn schema_equiv_ignores_generated_at() {
let mut a = empty_schema();
a.generated_at = "2020-01-01T00:00:00Z".into();
let b = empty_schema();
assert!(schema_equiv(&a, &b));
}
#[test]
fn fresh_snapshot_no_pending_no_models_synced() {
let diag = classify_bucket(&global_bucket(), None, None, None);
assert!(diag.is_none());
}
}