use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::{BTreeMap, BTreeSet};
use std::env;
use std::fs;
use std::io;
use std::path::{Path, PathBuf};
use crate::ast::ProtoSchema;
use crate::generation::{
CatalogManifest, SqlGenerationConfig, generate_bootstrap_sql, generate_delta_sql,
};
use crate::migration::diff::diff_manifests;
const ARTIFACT_EXTENSIONS: &[&str] = &["sql", "json", "yaml", "cypher"];
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BackendSyncTarget {
Postgres,
Qdrant,
Minio,
Redis,
Mongodb,
Neo4j,
Clickhouse,
Other(crate::backend::BackendKind),
All,
}
impl BackendSyncTarget {
pub fn from_token(token: &str) -> Option<Self> {
let token = token.trim();
if token.eq_ignore_ascii_case("all") {
return Some(Self::All);
}
let kind = crate::backend::BackendKind::from_token(token).or_else(|| match token {
"mongo" => Some(crate::backend::BackendKind::Mongodb),
"mssql" => Some(crate::backend::BackendKind::Mssql),
_ => None,
})?;
Some(match kind {
crate::backend::BackendKind::Postgres => Self::Postgres,
crate::backend::BackendKind::Qdrant => Self::Qdrant,
crate::backend::BackendKind::Minio => Self::Minio,
crate::backend::BackendKind::Redis => Self::Redis,
crate::backend::BackendKind::Mongodb => Self::Mongodb,
crate::backend::BackendKind::Neo4j => Self::Neo4j,
crate::backend::BackendKind::Clickhouse => Self::Clickhouse,
other => Self::Other(other),
})
}
pub fn includes(&self, kind: &crate::backend::BackendKind) -> bool {
use crate::backend::BackendKind as K;
match (self, kind) {
(Self::All, _) => true,
(Self::Postgres, K::Postgres) => true,
(Self::Qdrant, K::Qdrant) => true,
(Self::Minio, K::Minio) => true,
(Self::Redis, K::Redis) => true,
(Self::Mongodb, K::Mongodb) => true,
(Self::Neo4j, K::Neo4j) => true,
(Self::Clickhouse, K::Clickhouse) => true,
(Self::Other(target), kind) if target == kind => true,
_ => false,
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MultiBackendSyncReport {
pub proto_checksum: String,
pub backends: BTreeMap<String, MigrationSyncReport>,
pub clean: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum FileSyncStatus {
Written,
Verified,
Stale,
NoHeader,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MigrationFileRecord {
pub rel_path: String,
pub status: FileSyncStatus,
pub file_checksum: String,
pub proto_checksum: String,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct MigrationSyncReport {
pub db_ops_root: String,
pub proto_checksum: String,
pub total_files: usize,
pub written: usize,
pub verified: usize,
pub stale: usize,
pub no_header: usize,
pub bootstrap_written: usize,
pub files: Vec<MigrationFileRecord>,
pub bootstrap_files: Vec<MigrationFileRecord>,
pub warnings: Vec<String>,
pub clean: bool,
}
#[derive(Debug, Clone)]
pub struct DbOpsSyncConfig {
pub db_ops_root: Option<PathBuf>,
pub sql_config: SqlGenerationConfig,
pub force_bootstrap: bool,
pub backend: BackendSyncTarget,
}
impl Default for DbOpsSyncConfig {
fn default() -> Self {
Self {
db_ops_root: None,
sql_config: SqlGenerationConfig::default(),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
}
}
}
impl DbOpsSyncConfig {
pub fn from_env() -> Self {
let backend = env::var("UDB_DB_OPS_BACKEND")
.ok()
.and_then(|raw| BackendSyncTarget::from_token(&raw))
.unwrap_or(BackendSyncTarget::Postgres);
Self {
db_ops_root: env::var("UDB_DB_OPS_ROOT").ok().map(PathBuf::from),
sql_config: SqlGenerationConfig::default(),
force_bootstrap: env::var("UDB_DB_OPS_FORCE_BOOTSTRAP")
.map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
.unwrap_or(false),
backend,
}
}
}
#[derive(Debug)]
pub enum SyncError {
Io(PathBuf, io::Error),
Checksum(serde_json::Error),
SqlGeneration(String),
}
impl std::fmt::Display for SyncError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
SyncError::Io(path, err) => write!(f, "I/O error at {}: {err}", path.display()),
SyncError::Checksum(err) => write!(f, "checksum serialisation error: {err}"),
SyncError::SqlGeneration(err) => write!(f, "SQL generation error: {err}"),
}
}
}
impl std::error::Error for SyncError {}
pub fn sync_db_ops(
schemas: &[ProtoSchema],
config: &DbOpsSyncConfig,
) -> Result<MigrationSyncReport, SyncError> {
let manifest = CatalogManifest::from_schemas(schemas)
.map_err(|err| SyncError::SqlGeneration(err.to_string()))?;
let proto_checksum = manifest.checksum_sha256.clone();
let db_ops_root = match &config.db_ops_root {
Some(p) => p.clone(),
None => discover_db_ops_root()?,
};
let seeds_dir = resolve_seeders_dir(&db_ops_root);
ensure_seeds_dir(&seeds_dir)?;
let pg_root = db_ops_root.join("postgres");
let artifacts = generate_bootstrap_sql(schemas, &config.sql_config)
.map_err(|err| SyncError::SqlGeneration(err.to_string()))?;
sync_backend_artifacts(
&artifacts,
Some(&manifest),
&pg_root,
&proto_checksum,
&config.sql_config,
config.force_bootstrap,
)
}
pub fn sync_all_backends(
schemas: &[ProtoSchema],
config: &DbOpsSyncConfig,
) -> Result<MultiBackendSyncReport, SyncError> {
let manifest = CatalogManifest::from_schemas(schemas)
.map_err(|err| SyncError::SqlGeneration(err.to_string()))?;
let proto_checksum = manifest.checksum_sha256.clone();
let db_ops_root = match &config.db_ops_root {
Some(p) => p.clone(),
None => discover_db_ops_root()?,
};
let seeds_dir = resolve_seeders_dir(&db_ops_root);
ensure_seeds_dir(&seeds_dir)?;
let run_postgres = config
.backend
.includes(&crate::backend::BackendKind::Postgres);
let mut reports: BTreeMap<String, MigrationSyncReport> = BTreeMap::new();
if run_postgres {
let artifacts = generate_bootstrap_sql(schemas, &config.sql_config)
.map_err(|err| SyncError::SqlGeneration(err.to_string()))?;
if !artifacts.is_empty() {
let report = sync_backend_artifacts(
&artifacts,
Some(&manifest),
&db_ops_root.join("postgres"),
&proto_checksum,
&config.sql_config,
config.force_bootstrap,
)?;
reports.insert("postgres".to_string(), report);
}
}
for plugin in crate::backend::all_plugins() {
let kind = plugin.kind();
if kind == crate::backend::BackendKind::Postgres {
continue;
}
if !config.backend.includes(&kind) {
continue;
}
let artifacts = plugin
.generate_artifacts(&manifest, &config.sql_config)
.map_err(SyncError::SqlGeneration)?;
if artifacts.is_empty() {
if kind.capabilities().supports_schema_migration {
tracing::warn!(
backend = %kind.as_str(),
"backend advertises schema-migration support but generated zero artifacts; \
no DDL was synced (generate_artifacts is not implemented for this backend)"
);
}
continue;
}
let subdir = plugin.sync_subdir();
let report = sync_backend_artifacts(
&artifacts,
None,
&db_ops_root.join(subdir),
&proto_checksum,
&config.sql_config,
config.force_bootstrap,
)?;
reports.insert(subdir.to_string(), report);
}
let clean = reports.values().all(|r| r.clean);
Ok(MultiBackendSyncReport {
proto_checksum,
backends: reports,
clean,
})
}
fn sync_backend_artifacts(
artifacts: &[crate::generation::GeneratedArtifact],
manifest: Option<&CatalogManifest>,
backend_root: &Path,
proto_checksum: &str,
sql_config: &SqlGenerationConfig,
force_bootstrap: bool,
) -> Result<MigrationSyncReport, SyncError> {
let migrations_dir = backend_root.join("migrations");
let bootstrap_dir = backend_root.join("bootstrap");
let sidecar_path = migrations_dir.join(".last_manifest.json");
for dir in [&migrations_dir, &bootstrap_dir] {
fs::create_dir_all(dir).map_err(|e| SyncError::Io(dir.clone(), e))?;
}
let existing = scan_artifact_files_recursive(&migrations_dir)?;
let expected_rel_paths = artifacts
.iter()
.map(|artifact| artifact.rel_path.clone())
.collect::<BTreeSet<_>>();
let mut warnings: Vec<String> = Vec::new();
let mut file_records: Vec<MigrationFileRecord> = Vec::new();
let mut bootstrap_records: Vec<MigrationFileRecord> = Vec::new();
if existing.is_empty() {
for artifact in artifacts {
let dest = migrations_dir.join(&artifact.rel_path);
ensure_parent(&dest)?;
fs::write(&dest, artifact.content.as_bytes())
.map_err(|e| SyncError::Io(dest.clone(), e))?;
file_records.push(MigrationFileRecord {
rel_path: format!("migrations/{}", artifact.rel_path),
status: FileSyncStatus::Written,
file_checksum: String::new(),
proto_checksum: proto_checksum.to_string(),
});
}
if let Some(mf) = manifest
&& let Ok(json) = serde_json::to_string_pretty(mf)
{
let _ = fs::write(&sidecar_path, json.as_bytes());
}
} else {
let mut stale_schemas: BTreeSet<String> = BTreeSet::new();
let mut pending_new_artifacts = Vec::new();
for artifact in artifacts {
let dest = migrations_dir.join(&artifact.rel_path);
if dest.exists() {
let content =
fs::read_to_string(&dest).map_err(|e| SyncError::Io(dest.clone(), e))?;
let embedded = extract_proto_checksum_header(&content);
let body_checksum = artifact_body_checksum(&content);
let expected_body_checksum = artifact_body_checksum(&artifact.content);
let status = if embedded.is_empty() {
warnings.push(format!(
"migrations/{}: missing proto_manifest_checksum header — cannot verify",
artifact.rel_path
));
FileSyncStatus::NoHeader
} else if embedded == proto_checksum {
if body_checksum == expected_body_checksum {
FileSyncStatus::Verified
} else {
warnings.push(format!(
"migrations/{}: artifact body checksum mismatch — expected {}, found {}",
artifact.rel_path, expected_body_checksum, body_checksum
));
stale_schemas.insert(artifact.schema.clone());
FileSyncStatus::Stale
}
} else {
stale_schemas.insert(artifact.schema.clone());
FileSyncStatus::Stale
};
file_records.push(MigrationFileRecord {
rel_path: format!("migrations/{}", artifact.rel_path),
status,
file_checksum: embedded,
proto_checksum: proto_checksum.to_string(),
});
} else {
pending_new_artifacts.push(artifact.clone());
file_records.push(MigrationFileRecord {
rel_path: format!("migrations/{}", artifact.rel_path),
status: FileSyncStatus::Written,
file_checksum: String::new(),
proto_checksum: proto_checksum.to_string(),
});
}
}
let extra_artifacts = existing
.iter()
.filter_map(|path| relative_artifact_path(&migrations_dir, path))
.filter(|rel_path| !expected_rel_paths.contains(rel_path))
.collect::<Vec<_>>();
let has_stale = file_records
.iter()
.any(|r| r.status == FileSyncStatus::Stale);
let has_no_header = file_records
.iter()
.any(|r| r.status == FileSyncStatus::NoHeader);
let has_noncanonical = has_stale || has_no_header || !extra_artifacts.is_empty();
if has_noncanonical && !force_bootstrap {
warnings.push(
"non-canonical migrations detected; refreshing canonical migrations in place"
.to_string(),
);
reset_artifact_dir(&migrations_dir)?;
reset_artifact_dir(&bootstrap_dir)?;
file_records.clear();
for artifact in artifacts {
let dest = migrations_dir.join(&artifact.rel_path);
ensure_parent(&dest)?;
fs::write(&dest, artifact.content.as_bytes())
.map_err(|e| SyncError::Io(dest.clone(), e))?;
file_records.push(MigrationFileRecord {
rel_path: format!("migrations/{}", artifact.rel_path),
status: FileSyncStatus::Written,
file_checksum: proto_checksum.to_string(),
proto_checksum: proto_checksum.to_string(),
});
}
if let Some(mf) = manifest
&& let Ok(json) = serde_json::to_string_pretty(mf)
{
let _ = fs::write(&sidecar_path, json.as_bytes());
}
} else {
for artifact in &pending_new_artifacts {
let dest = migrations_dir.join(&artifact.rel_path);
ensure_parent(&dest)?;
fs::write(&dest, artifact.content.as_bytes())
.map_err(|e| SyncError::Io(dest.clone(), e))?;
}
if has_noncanonical || force_bootstrap {
let delta_artifacts: Vec<crate::generation::GeneratedArtifact> = if let Some(
new_manifest,
) = manifest
{
let prior_manifest = fs::read_to_string(&sidecar_path)
.ok()
.and_then(|json| serde_json::from_str::<CatalogManifest>(&json).ok());
if let Some(ref prior) = prior_manifest {
let changes = diff_manifests(Some(prior), new_manifest);
let delta = generate_delta_sql(new_manifest, &changes, sql_config);
if delta.is_empty() {
warnings.push(
"proto changed but no auto-safe ALTER statements were generated; \
writing baseline bootstrap for review (may contain blocked ops)"
.to_string(),
);
artifacts
.iter()
.filter(|a| {
force_bootstrap
|| stale_schemas.contains(&a.schema)
|| a.schema.is_empty()
})
.cloned()
.collect()
} else {
delta
}
} else {
warnings.push(
"no prior manifest sidecar found; writing baseline bootstrap \
for first-time stale detection"
.to_string(),
);
artifacts
.iter()
.filter(|a| {
force_bootstrap
|| stale_schemas.contains(&a.schema)
|| a.schema.is_empty()
})
.cloned()
.collect()
}
} else {
artifacts
.iter()
.filter(|a| {
force_bootstrap
|| stale_schemas.contains(&a.schema)
|| a.schema.is_empty()
})
.cloned()
.collect()
};
reset_artifact_dir(&bootstrap_dir)?;
for artifact in &delta_artifacts {
let dest = bootstrap_dir.join(&artifact.rel_path);
ensure_parent(&dest)?;
fs::write(&dest, artifact.content.as_bytes())
.map_err(|e| SyncError::Io(dest.clone(), e))?;
bootstrap_records.push(MigrationFileRecord {
rel_path: format!("bootstrap/{}", artifact.rel_path),
status: FileSyncStatus::Written,
file_checksum: String::new(),
proto_checksum: proto_checksum.to_string(),
});
}
if let Some(mf) = manifest
&& let Ok(json) = serde_json::to_string_pretty(mf)
{
let _ = fs::write(&sidecar_path, json.as_bytes());
}
}
}
}
let written = file_records
.iter()
.filter(|r| r.status == FileSyncStatus::Written)
.count();
let verified = file_records
.iter()
.filter(|r| r.status == FileSyncStatus::Verified)
.count();
let stale = file_records
.iter()
.filter(|r| r.status == FileSyncStatus::Stale)
.count();
let no_header = file_records
.iter()
.filter(|r| r.status == FileSyncStatus::NoHeader)
.count();
let clean = stale == 0 && no_header == 0;
Ok(MigrationSyncReport {
db_ops_root: backend_root.display().to_string(),
proto_checksum: proto_checksum.to_string(),
total_files: file_records.len(),
written,
verified,
stale,
no_header,
bootstrap_written: bootstrap_records.len(),
files: file_records,
bootstrap_files: bootstrap_records,
warnings,
clean,
})
}
pub(crate) fn resolve_seeders_dir(db_ops_root: &Path) -> PathBuf {
let repo_root = db_ops_root
.parent()
.map(Path::to_path_buf)
.unwrap_or_else(|| db_ops_root.to_path_buf());
env::var("UDB_SEEDERS_PATH")
.ok()
.or_else(|| env::var("UDB_SEEDER_PATH").ok())
.map(|raw| raw.trim().to_string())
.filter(|raw| !raw.is_empty())
.map(PathBuf::from)
.map(|path| {
if path.is_absolute() {
path
} else {
repo_root.join(path)
}
})
.unwrap_or_else(|| db_ops_root.join("seeds"))
}
fn ensure_seeds_dir(seeds_dir: &Path) -> Result<(), SyncError> {
fs::create_dir_all(seeds_dir).map_err(|e| SyncError::Io(seeds_dir.to_path_buf(), e))?;
let gitkeep = seeds_dir.join(".gitkeep");
if !gitkeep.exists() {
fs::write(
&gitkeep,
b"# Place seed files here (e.g. 001_roles.sql).\n\
# These run after migrations during test/staging bootstraps.\n",
)
.map_err(|e| SyncError::Io(gitkeep.clone(), e))?;
}
Ok(())
}
fn scan_artifact_files_recursive(dir: &Path) -> Result<Vec<PathBuf>, SyncError> {
let mut out = Vec::new();
collect_artifacts_filtered(dir, &mut out)?;
Ok(out)
}
fn collect_artifacts_filtered(dir: &Path, out: &mut Vec<PathBuf>) -> Result<(), SyncError> {
let entries = match fs::read_dir(dir) {
Ok(e) => e,
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(()),
Err(e) => return Err(SyncError::Io(dir.to_path_buf(), e)),
};
for entry in entries {
let entry = entry.map_err(|e| SyncError::Io(dir.to_path_buf(), e))?;
let path = entry.path();
if path.is_dir() {
collect_artifacts_filtered(&path, out)?;
} else if path
.file_name()
.and_then(|n| n.to_str())
.is_some_and(|n| n == ".last_manifest.json")
{
continue;
} else if path
.extension()
.and_then(|s| s.to_str())
.is_some_and(|ext| ARTIFACT_EXTENSIONS.contains(&ext))
{
out.push(path);
}
}
Ok(())
}
fn ensure_parent(dest: &Path) -> Result<(), SyncError> {
if let Some(parent) = dest.parent()
&& !parent.exists()
{
fs::create_dir_all(parent).map_err(|e| SyncError::Io(parent.to_path_buf(), e))?;
}
Ok(())
}
fn reset_artifact_dir(dir: &Path) -> Result<(), SyncError> {
match fs::remove_dir_all(dir) {
Ok(()) => {}
Err(err) if err.kind() == io::ErrorKind::NotFound => {}
Err(err) => return Err(SyncError::Io(dir.to_path_buf(), err)),
}
fs::create_dir_all(dir).map_err(|e| SyncError::Io(dir.to_path_buf(), e))?;
Ok(())
}
fn relative_artifact_path(base: &Path, path: &Path) -> Option<String> {
path.strip_prefix(base)
.ok()
.map(|rel| rel.to_string_lossy().replace('\\', "/"))
}
fn extract_proto_checksum_header(content: &str) -> String {
for line in content.lines() {
let line = line.trim();
if let Some(val) = line.strip_prefix("-- UDB:proto_manifest_checksum=") {
return val.trim().to_string();
}
if let Some(val) = line.strip_prefix("// UDB:proto_manifest_checksum=") {
return val.trim().to_string();
}
if let Some(rest) = line.strip_prefix("# UDB:proto_manifest_checksum:") {
return rest.trim().to_string();
}
if let Some(key_pos) = line.find("\"proto_manifest_checksum\"") {
let after_key = &line[key_pos + "\"proto_manifest_checksum\"".len()..];
if let Some(colon_pos) = after_key.find(':') {
let after_colon = after_key[colon_pos + 1..].trim_start();
if let Some(rest) = after_colon.strip_prefix('"')
&& let Some(end_quote) = rest.find('"')
{
let value = &rest[..end_quote];
if !value.is_empty() {
return value.to_string();
}
}
}
}
if !line.is_empty()
&& !line.starts_with("--")
&& !line.starts_with("//")
&& !line.starts_with('#')
&& !line.starts_with('{')
&& !line.starts_with('"')
&& !line.starts_with('_')
{
break;
}
}
String::new()
}
fn artifact_body_checksum(content: &str) -> String {
let digest = Sha256::digest(content.as_bytes());
format!("sha256:{digest:x}")
}
pub fn discover_db_ops_root() -> Result<PathBuf, SyncError> {
let cwd = env::current_dir().map_err(|e| SyncError::Io(PathBuf::from("<cwd>"), e))?;
let mut dir = cwd.clone();
loop {
if dir.join("Cargo.toml").exists() {
let parent = dir
.parent()
.map(|p| p.to_path_buf())
.unwrap_or_else(|| dir.clone());
return Ok(parent.join("db_ops"));
}
if dir.join(".git").exists() || dir.join("udb").join("Cargo.toml").exists() {
return Ok(dir.join("db_ops"));
}
match dir.parent() {
Some(p) if p != dir => dir = p.to_path_buf(),
_ => break,
}
}
let fallback = cwd
.parent()
.map(|p| p.to_path_buf())
.unwrap_or_else(|| cwd.clone())
.join("db_ops");
Ok(fallback)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::parser::{ParserConfig, parse_file};
use std::fs;
use std::sync::{Mutex, OnceLock};
fn env_lock() -> std::sync::MutexGuard<'static, ()> {
static ENV_LOCK: OnceLock<Mutex<()>> = OnceLock::new();
ENV_LOCK
.get_or_init(|| Mutex::new(()))
.lock()
.expect("lock env mutex")
}
fn parse_source(source: &str, _filename: &str) -> Vec<ProtoSchema> {
use std::sync::atomic::{AtomicU64, Ordering};
static COUNTER: AtomicU64 = AtomicU64::new(0);
let id = COUNTER.fetch_add(1, Ordering::Relaxed);
let tmp_file = std::env::temp_dir().join(format!("udb_test_parse_{}.proto", id));
fs::write(&tmp_file, source.as_bytes()).unwrap();
let schemas = parse_file(&tmp_file, &ParserConfig::default()).expect("parse proto source");
let _ = fs::remove_file(&tmp_file);
schemas
}
#[test]
fn extract_checksum_from_header() {
let sql = "-- UDB:migration_kind=bootstrap\n\
-- UDB:schema=app\n\
-- UDB:table=users\n\
-- UDB:proto_manifest_checksum=abc123\n\
-- UDB:generator=udb\n\n\
CREATE TABLE \"app\".\"users\" ();\n";
assert_eq!(extract_proto_checksum_header(sql), "abc123");
}
#[test]
fn extract_checksum_missing() {
let sql = "CREATE TABLE \"app\".\"users\" ();\n";
assert_eq!(extract_proto_checksum_header(sql), "");
}
#[test]
fn extract_checksum_stops_at_non_comment() {
let sql = "SET lock_timeout = '5s';\n-- UDB:proto_manifest_checksum=never\n";
assert_eq!(extract_proto_checksum_header(sql), "");
}
#[test]
fn extract_checksum_from_yaml_header() {
let yaml = "# UDB:migration_kind: bootstrap\n\
# UDB:backend: redis\n\
# UDB:proto_manifest_checksum: sha256-deadbeef\n\
namespace: ocr\n";
assert_eq!(extract_proto_checksum_header(yaml), "sha256-deadbeef");
}
#[test]
fn extract_checksum_from_cypher_header() {
let cypher = "// UDB:migration_kind=bootstrap\n\
// UDB:proto_manifest_checksum=cafebabe\n\
CREATE CONSTRAINT;\n";
assert_eq!(extract_proto_checksum_header(cypher), "cafebabe");
}
#[test]
fn extract_checksum_from_json_udb_meta() {
let json =
r#"{"_udb_meta":{"proto_manifest_checksum":"sha256-1234"},"collection_name":"docs"}"#;
assert_eq!(extract_proto_checksum_header(json), "sha256-1234");
}
const SIMPLE_PROTO: &str = r#"
syntax = "proto3";
package example.test.entity.v1;
import "udb/core/common/v1/db.proto";
message Widget {
option (udb.core.common.v1.pg_table) = {
table_name: "widgets"
schema_name: "app"
migration_order: 10
is_table: true
};
string widget_id = 1 [(udb.core.common.v1.pg_column) = {
column_name: "widget_id"
sql_type: "UUID"
primary_key: true
not_null: true
default_value: "gen_random_uuid()"
}];
}
"#;
fn parse_simple() -> Vec<ProtoSchema> {
parse_source(SIMPLE_PROTO, "udb_test_simple.proto")
}
#[test]
fn sync_db_ops_writes_baseline_on_empty_dir() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_sync_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
let report = sync_db_ops(&schemas, &cfg).expect("sync_db_ops");
assert!(
report.written > 0,
"expected baseline artifacts to be written"
);
assert!(report.clean, "fresh run should be clean");
let migrations_dir = tmp.join("postgres").join("migrations");
assert!(
migrations_dir.exists(),
"postgres/migrations/ must be created"
);
let entries: Vec<_> = fs::read_dir(&migrations_dir)
.unwrap()
.filter_map(|e| e.ok())
.collect();
assert!(
!entries.is_empty(),
"at least one migration file must be written"
);
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_db_ops_verifies_matching_checksum() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_verify_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
sync_db_ops(&schemas, &cfg).expect("first sync");
let report2 = sync_db_ops(&schemas, &cfg).expect("second sync");
assert_eq!(report2.stale, 0, "no stale files on identical re-run");
assert!(report2.clean, "second run must be clean");
assert!(
report2.verified > 0,
"at least one file verified on second run"
);
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_db_ops_marks_preserved_header_body_edit_stale() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_body_stale_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
let report1 = sync_db_ops(&schemas, &cfg).expect("first sync");
let migrations_dir = tmp.join("postgres").join("migrations");
let target = scan_artifact_files_recursive(&migrations_dir)
.expect("scan migrations")
.into_iter()
.find(|path| {
fs::read_to_string(path)
.map(|content| !extract_proto_checksum_header(&content).is_empty())
.unwrap_or(false)
})
.expect("migration with checksum header");
let rel_path = format!(
"migrations/{}",
relative_artifact_path(&migrations_dir, &target).expect("relative path")
);
let original = fs::read_to_string(&target).unwrap();
assert_eq!(
extract_proto_checksum_header(&original),
report1.proto_checksum,
"test setup must preserve the current header checksum"
);
fs::write(
&target,
format!("{original}\n-- hand-edited body with preserved header\n"),
)
.unwrap();
let review_cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: true,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let report2 = sync_db_ops(&schemas, &review_cfg).expect("second sync");
let record = report2
.files
.iter()
.find(|record| record.rel_path == rel_path)
.expect("edited migration record");
assert_eq!(record.status, FileSyncStatus::Stale);
assert_eq!(
record.file_checksum, report2.proto_checksum,
"the preserved header still matches the current proto checksum"
);
assert!(!report2.clean, "body edit must make review-mode sync dirty");
assert!(
report2
.warnings
.iter()
.any(|warning| warning.contains("artifact body checksum mismatch")),
"body checksum mismatch should be reported"
);
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_db_ops_refreshes_stale_in_place() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_stale_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
sync_db_ops(&schemas, &cfg).expect("first sync");
let migrations_dir = tmp.join("postgres").join("migrations");
fn patch_sql_files(dir: &std::path::Path) {
if let Ok(entries) = fs::read_dir(dir) {
for entry in entries.filter_map(|e| e.ok()) {
let path = entry.path();
if path.is_dir() {
patch_sql_files(&path);
} else if path.extension().and_then(|e| e.to_str()) == Some("sql") {
let content = fs::read_to_string(&path).unwrap();
let patched = content.replace(
"-- UDB:proto_manifest_checksum=",
"-- UDB:proto_manifest_checksum=STALE_",
);
fs::write(&path, patched).unwrap();
}
}
}
}
patch_sql_files(&migrations_dir);
let report2 = sync_db_ops(&schemas, &cfg).expect("second sync");
assert_eq!(report2.stale, 0, "stale files must be refreshed in place");
assert_eq!(report2.no_header, 0, "refreshed files must be verifiable");
assert!(report2.clean, "run must converge on a clean baseline");
assert!(
report2.written > 0,
"refresh must rewrite canonical artifacts"
);
let bootstrap_dir = tmp.join("postgres").join("bootstrap");
assert!(
scan_artifact_files_recursive(&bootstrap_dir)
.expect("scan bootstrap dir")
.is_empty(),
"bootstrap/ should remain empty during in-place refresh"
);
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_db_ops_force_bootstrap_preserves_review_mode() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_bootstrap_review_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
sync_db_ops(&schemas, &cfg).expect("first sync");
let migrations_dir = tmp.join("postgres").join("migrations");
fn patch_sql_files(dir: &std::path::Path) {
if let Ok(entries) = fs::read_dir(dir) {
for entry in entries.filter_map(|e| e.ok()) {
let path = entry.path();
if path.is_dir() {
patch_sql_files(&path);
} else if path.extension().and_then(|e| e.to_str()) == Some("sql") {
let content = fs::read_to_string(&path).unwrap();
let patched = content.replace(
"-- UDB:proto_manifest_checksum=",
"-- UDB:proto_manifest_checksum=STALE_",
);
fs::write(&path, patched).unwrap();
}
}
}
}
patch_sql_files(&migrations_dir);
let review_cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: true,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let report2 = sync_db_ops(&schemas, &review_cfg).expect("second sync");
assert!(
report2.stale > 0,
"review mode should preserve stale detection"
);
assert!(!report2.clean, "review mode should require operator action");
assert!(
report2.bootstrap_written > 0,
"review mode should emit bootstrap artifacts"
);
let _ = fs::remove_dir_all(&tmp);
}
const MULTI_BACKEND_PROTO: &str = r#"
syntax = "proto3";
package example.test.multi.v1;
import "udb/core/common/v1/db.proto";
message Embedding {
option (udb.core.common.v1.pg_table) = {
table_name: "embeddings"
schema_name: "app"
migration_order: 1
is_table: true
};
option (udb.core.common.v1.vector_store) = {
backend: VECTOR_BACKEND_QDRANT
collection_name: "embeddings"
dimension: 768
distance: VECTOR_DISTANCE_COSINE
};
option (udb.core.common.v1.cache) = {
backend: CACHE_BACKEND_REDIS
key_pattern: "app:embed:{id}"
ttl_seconds: 3600
};
string id = 1 [(udb.core.common.v1.pg_column) = {
column_name: "id"
sql_type: "UUID"
primary_key: true
not_null: true
}];
string s3_uri = 2 [
(udb.core.common.v1.pg_column) = { column_name: "s3_uri" sql_type: "TEXT" },
(udb.core.common.v1.storage) = {
backend: STORAGE_BACKEND_MINIO
bucket_env_key: "EMBED_BUCKET"
key_prefix: "embeddings/"
}
];
}
"#;
#[test]
fn sync_all_backends_creates_multi_dir_layout() {
let _env_guard = env_lock();
let prior_seeders = std::env::var("UDB_SEEDERS_PATH").ok();
let prior_seeder = std::env::var("UDB_SEEDER_PATH").ok();
#[allow(unused_unsafe)]
unsafe {
std::env::remove_var("UDB_SEEDERS_PATH");
std::env::remove_var("UDB_SEEDER_PATH");
}
let tmp = std::env::temp_dir().join(format!(
"udb_test_multi_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::All,
..Default::default()
};
let schemas = parse_source(MULTI_BACKEND_PROTO, "udb_test_multi.proto");
let report = sync_all_backends(&schemas, &cfg).expect("sync_all_backends");
assert!(report.clean, "fresh multi-backend run must be clean");
assert!(
report.backends.contains_key("postgres"),
"postgres backend must be synced"
);
assert!(
report.backends.contains_key("qdrant"),
"qdrant backend must be synced"
);
assert!(tmp.join("seeds").exists(), "seeds/ must be created");
for (backend, backend_report) in &report.backends {
let migrations = tmp.join(backend).join("migrations");
assert!(
migrations.exists(),
"{backend}/migrations/ must exist; report={backend_report:?}"
);
}
#[allow(unused_unsafe)]
unsafe {
match prior_seeders {
Some(value) => std::env::set_var("UDB_SEEDERS_PATH", value),
None => std::env::remove_var("UDB_SEEDERS_PATH"),
}
match prior_seeder {
Some(value) => std::env::set_var("UDB_SEEDER_PATH", value),
None => std::env::remove_var("UDB_SEEDER_PATH"),
}
}
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_all_backends_uses_env_seeders_path() {
let _env_guard = env_lock();
let tmp = std::env::temp_dir().join(format!(
"udb_test_seeders_env_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let db_ops_root = tmp.join("db_ops");
fs::create_dir_all(&db_ops_root).unwrap();
let prior_seeders = std::env::var("UDB_SEEDERS_PATH").ok();
#[allow(unused_unsafe)]
unsafe {
std::env::set_var("UDB_SEEDERS_PATH", "database/seeders/postgress");
}
let cfg = DbOpsSyncConfig {
db_ops_root: Some(db_ops_root.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::All,
..Default::default()
};
let schemas = parse_source(MULTI_BACKEND_PROTO, "udb_test_seeders_env.proto");
let report = sync_all_backends(&schemas, &cfg).expect("sync_all_backends");
assert!(report.clean, "fresh multi-backend run must be clean");
let custom_seeders_dir = tmp.join("database").join("seeders").join("postgress");
assert!(
custom_seeders_dir.exists(),
"custom seeders path must be created"
);
assert!(
custom_seeders_dir.join(".gitkeep").exists(),
"custom seeders path must receive .gitkeep"
);
assert!(
!db_ops_root.join("seeds").exists(),
"default db_ops/seeds path should not be used when env override is set"
);
#[allow(unused_unsafe)]
unsafe {
match prior_seeders {
Some(value) => std::env::set_var("UDB_SEEDERS_PATH", value),
None => std::env::remove_var("UDB_SEEDERS_PATH"),
}
}
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn sync_all_backends_postgres_only_skips_other_dirs() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_pg_only_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
force_bootstrap: false,
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_source(MULTI_BACKEND_PROTO, "udb_test_multi_pg.proto");
let report = sync_all_backends(&schemas, &cfg).expect("sync_all_backends postgres");
assert!(report.backends.contains_key("postgres"));
assert!(
!report.backends.contains_key("qdrant"),
"qdrant should not be synced when backend=postgres"
);
assert!(
!tmp.join("qdrant").exists(),
"qdrant/ dir must not be created"
);
let _ = fs::remove_dir_all(&tmp);
}
#[test]
fn multi_backend_report_proto_checksum_is_stable() {
let tmp = std::env::temp_dir().join(format!(
"udb_test_chksum_{}",
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
fs::create_dir_all(&tmp).unwrap();
let cfg = DbOpsSyncConfig {
db_ops_root: Some(tmp.clone()),
backend: BackendSyncTarget::Postgres,
..Default::default()
};
let schemas = parse_simple();
let r1 = sync_all_backends(&schemas, &cfg).unwrap();
let r2 = sync_all_backends(&schemas, &cfg).unwrap();
assert_eq!(
r1.proto_checksum, r2.proto_checksum,
"proto_checksum must be deterministic"
);
assert!(
!r1.proto_checksum.is_empty(),
"proto_checksum must not be empty"
);
let _ = fs::remove_dir_all(&tmp);
}
}