use std::error::Error;
use std::fmt;
use std::io::{Read as _, Write};
use std::path::PathBuf;
use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
use type_bridge_contract::limits::StructuralLimits;
use type_bridge_contract::migration::{
MigrationId, MigrationStep, MigrationStepId, SchemaDeltaStep,
};
use type_bridge_contract::migration_assertion::AssertionExpectation;
use type_bridge_contract::migration_backfill::AttributeBackfillPlan;
use type_bridge_contract::schema::{DeclaredSchema, SchemaDelta};
use type_bridge_query::{MigrationAssertionValidationContext, lower_condition_to_plan};
use type_bridge_schema::{
ManagedDeltaContext, SafetyClass, SafetyConditionDomainIndex, SafetyDerivationProfile,
apply_delta, derive_safety_conditions_with_domain_index, diff_managed, inverse_delta,
managed_schema_state, resolve,
};
use crate::history::MigrationHistoryGraph;
use crate::lowering::{
SchemaFactCatalog, SchemaLoweringBinding, SchemaLoweringDiagnostic,
lower_schema_delta_with_verified_assertions,
};
use crate::manifest::{
SchemaMigrationDraft, VerifiedSchemaMigrationManifest, build_verified_manifest,
delta_diagnostic, encode_verified_manifest, verify_assertion_coverage,
};
use crate::profile::schema_lowering_profile_binding;
use crate::{MigrationAuthoringLock, MigrationDirectory};
use type_bridge_contract::managed_scope::SemanticProfileBinding;
const REVERSE_REJECTION_CODES: [&str; 4] = [
"migration_manifest_reverse_unresolved_safety",
"migration_manifest_reverse_requires_assertions",
"migration_manifest_inverse_replay_mismatch",
"migration_manifest_inverse_plan_invalid",
];
#[derive(Clone, Copy, Debug)]
pub struct MigrationGenerationRequest<'a> {
pub app_label: &'a str,
pub base_name: &'a str,
pub genesis_source: &'a DeclaredSchema,
pub desired: &'a DeclaredSchema,
pub context: &'a ManagedDeltaContext,
}
#[derive(Clone, Copy, Debug)]
pub struct BackfillMigrationGenerationRequest<'a> {
pub app_label: &'a str,
pub base_name: &'a str,
pub genesis_source: &'a DeclaredSchema,
pub desired: &'a DeclaredSchema,
pub plan: &'a AttributeBackfillPlan,
pub context: &'a ManagedDeltaContext,
}
#[derive(Clone, Debug)]
pub enum MigrationGenerationOutcome {
UpToDate,
Generated(Box<GeneratedMigration>),
}
#[derive(Clone, Debug)]
pub struct GeneratedMigration {
manifest: VerifiedSchemaMigrationManifest,
canonical_bytes: Vec<u8>,
}
impl GeneratedMigration {
pub const fn manifest(&self) -> &VerifiedSchemaMigrationManifest {
&self.manifest
}
pub fn canonical_bytes(&self) -> &[u8] {
&self.canonical_bytes
}
pub fn file_name(&self) -> String {
format!("{}.tbmigration.json", self.manifest.id().name().as_str())
}
pub fn preview_file_name(&self) -> String {
format!("{}.typeql", self.manifest.id().name().as_str())
}
}
#[derive(Clone, Debug)]
pub enum MigrationPreviewError {
Diagnostic(Diagnostic),
Lowering(SchemaLoweringDiagnostic),
}
impl fmt::Display for MigrationPreviewError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Diagnostic(diagnostic) => diagnostic.fmt(formatter),
Self::Lowering(diagnostic) => diagnostic.fmt(formatter),
}
}
}
impl Error for MigrationPreviewError {}
impl From<Diagnostic> for MigrationPreviewError {
fn from(value: Diagnostic) -> Self {
Self::Diagnostic(value)
}
}
impl From<SchemaLoweringDiagnostic> for MigrationPreviewError {
fn from(value: SchemaLoweringDiagnostic) -> Self {
Self::Lowering(value)
}
}
pub fn generate_next_migration(
graph: &MigrationHistoryGraph,
request: &MigrationGenerationRequest<'_>,
) -> Result<MigrationGenerationOutcome, Diagnostic> {
for (id, _) in graph.manifests() {
if id.app_label().as_str() != request.app_label {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_generation_foreign_app_label",
"history contains a migration from a different application lineage",
));
}
}
let (source, parents) = match graph.default_head()? {
None => (request.genesis_source, Vec::new()),
Some(head) => (
graph
.manifest(head)
.expect("a graph head is always a graph member")
.target_schema(),
vec![head.clone()],
),
};
let source_state = managed_schema_state(source, request.context).map_err(delta_diagnostic)?;
let desired_state =
managed_schema_state(request.desired, request.context).map_err(delta_diagnostic)?;
if source_state == desired_state {
return Ok(MigrationGenerationOutcome::UpToDate);
}
let delta = diff_managed(source, request.desired, request.context).map_err(delta_diagnostic)?;
let id = MigrationId::new(request.app_label, allocate_name(graph, request.base_name))?;
if graph.manifest(&id).is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_generation_duplicate_name",
"allocated migration identity already exists in verified history",
));
}
let steps = author_steps(&delta, source, request.desired, request.context, true)?;
let draft = SchemaMigrationDraft::new(id.clone(), parents.clone(), steps)?;
let manifest = match build_verified_manifest(draft, (source, request.context)) {
Ok(manifest) => manifest,
Err(diagnostic) if REVERSE_REJECTION_CODES.contains(&diagnostic.code().as_str()) => {
let steps = author_steps(&delta, source, request.desired, request.context, false)?;
let draft = SchemaMigrationDraft::new(id, parents, steps)?;
build_verified_manifest(draft, (source, request.context))?
}
Err(diagnostic) => return Err(diagnostic),
};
let canonical_bytes = encode_verified_manifest(&manifest)?;
Ok(MigrationGenerationOutcome::Generated(Box::new(
GeneratedMigration {
manifest,
canonical_bytes,
},
)))
}
pub fn generate_backfill_migration(
graph: &MigrationHistoryGraph,
request: &BackfillMigrationGenerationRequest<'_>,
) -> Result<GeneratedMigration, Diagnostic> {
for (id, _) in graph.manifests() {
if id.app_label().as_str() != request.app_label {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_generation_foreign_app_label",
"history contains a migration from a different application lineage",
));
}
}
let (source, parents) = match graph.default_head()? {
None => (request.genesis_source, Vec::new()),
Some(head) => (
graph
.manifest(head)
.expect("a graph head is always a graph member")
.target_schema(),
vec![head.clone()],
),
};
let source_state = managed_schema_state(source, request.context).map_err(delta_diagnostic)?;
let desired_state =
managed_schema_state(request.desired, request.context).map_err(delta_diagnostic)?;
if source_state != desired_state {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_backfill_generation_schema_drift",
"backfill authoring requires current Split-YAML to equal the committed head",
));
}
let id = MigrationId::new(request.app_label, allocate_name(graph, request.base_name))?;
if graph.manifest(&id).is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_generation_duplicate_name",
"allocated migration identity already exists in verified history",
));
}
let step = MigrationStep::backfill(MigrationStepId::new("backfill")?, request.plan.clone())?;
let draft = SchemaMigrationDraft::new(id, parents, vec![step])?;
let manifest = build_verified_manifest(draft, (source, request.context))?;
let canonical_bytes = encode_verified_manifest(&manifest)?;
Ok(GeneratedMigration {
manifest,
canonical_bytes,
})
}
pub fn render_migration_preview(
manifest: &VerifiedSchemaMigrationManifest,
context: &ManagedDeltaContext,
) -> Result<String, MigrationPreviewError> {
let binding = SchemaLoweringBinding::current(context.available_capabilities().clone())?;
let safety_profile = SafetyDerivationProfile::new(
manifest.semantic_profile().clone(),
manifest.lowering_profile().clone(),
)?;
let mut current = manifest.source_schema().clone();
let mut pending = Vec::new();
let mut queries = Vec::new();
for step in manifest.steps() {
let Some(schema_step) = step.as_schema_delta() else {
pending.push(step);
continue;
};
let target =
apply_delta(¤t, schema_step.delta(), context).map_err(delta_diagnostic)?;
let coverage = verify_assertion_coverage(
&pending,
schema_step.delta(),
¤t,
&target,
&safety_profile,
)?;
pending.clear();
let source_catalog = SchemaFactCatalog::new(current.facts().cloned())?;
let target_catalog = SchemaFactCatalog::new(target.facts().cloned())?;
let lowering = lower_schema_delta_with_verified_assertions(
schema_step.delta(),
&source_catalog,
&target_catalog,
&binding,
coverage.discharged_operation_indices(),
true,
)?;
for unit in lowering.units() {
for statement in unit.statements() {
queries.push(statement.query().to_owned());
}
}
current = target;
}
Ok(format!("{}\n", queries.join("\n\n")))
}
pub fn try_acquire_migration_authoring_lock(
directory: &MigrationDirectory,
) -> Result<MigrationAuthoringLock<'_>, Diagnostic> {
directory.try_acquire_authoring_lock().map_err(|error| {
if error.kind() == std::io::ErrorKind::WouldBlock {
write_conflict()
} else {
write_failed("authoring lock acquisition", &error)
}
})
}
pub fn write_generated_migration_under_lock(
lock: &MigrationAuthoringLock<'_>,
generated: &GeneratedMigration,
preview: &str,
) -> Result<PathBuf, Diagnostic> {
let directory = lock.directory;
let manifest_name = generated.file_name();
let preview_name = generated.preview_file_name();
if directory
.entry_exists(manifest_name.as_ref())
.map_err(|error| write_failed("manifest presence probe", &error))?
{
return Err(write_conflict());
}
if directory
.entry_exists(preview_name.as_ref())
.map_err(|error| write_failed("preview presence probe", &error))?
{
directory
.remove_file(preview_name.as_ref())
.map_err(|error| write_failed("interrupted preview recovery", &error))?;
}
let manifest_temp = write_unique_temporary(
directory,
&generated.file_name(),
generated.canonical_bytes(),
)?;
let preview_temp = match write_unique_temporary(
directory,
&generated.preview_file_name(),
preview.as_bytes(),
) {
Ok(path) => path,
Err(error) => {
let _ = directory.remove_file(manifest_temp.as_ref());
return Err(error);
}
};
let cleanup = |published_preview: bool| {
let _ = directory.remove_file(manifest_temp.as_ref());
let _ = directory.remove_file(preview_temp.as_ref());
if published_preview {
let _ = directory.remove_file(preview_name.as_ref());
}
};
let published_preview = match publish_no_replace(directory, &preview_temp, &preview_name) {
Ok(()) => true,
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
match read_existing_preview(directory, &preview_name) {
Ok(existing) if existing == preview.as_bytes() => false,
_ => {
cleanup(false);
return Err(write_conflict());
}
}
}
Err(error) => {
cleanup(false);
return Err(write_failed("preview publication link", &error));
}
};
if let Err(error) = publish_no_replace(directory, &manifest_temp, &manifest_name) {
cleanup(published_preview);
return Err(if error.kind() == std::io::ErrorKind::AlreadyExists {
write_conflict()
} else {
write_failed("manifest publication link", &error)
});
}
if let Err(error) = sync_authoring_directory(directory) {
let _ = directory.remove_file(manifest_temp.as_ref());
let _ = directory.remove_file(preview_temp.as_ref());
return Err(error);
}
let _ = directory.remove_file(manifest_temp.as_ref());
let _ = directory.remove_file(preview_temp.as_ref());
sync_authoring_directory(directory)?;
Ok(PathBuf::from(manifest_name))
}
fn read_existing_preview(directory: &MigrationDirectory, name: &str) -> std::io::Result<Vec<u8>> {
let limit = type_bridge_contract::limits::MAX_CANONICAL_BYTES;
let file = directory.open_regular_readonly(name.as_ref())?;
let mut bytes = Vec::new();
file.take(u64::try_from(limit).unwrap_or(u64::MAX).saturating_add(1))
.read_to_end(&mut bytes)?;
if bytes.len() > limit {
return Err(std::io::Error::other(
"migration preview exceeds byte ceiling",
));
}
Ok(bytes)
}
fn write_unique_temporary(
directory: &MigrationDirectory,
final_name: &str,
bytes: &[u8],
) -> Result<String, Diagnostic> {
use std::sync::atomic::{AtomicU64, Ordering};
static NEXT_TEMPORARY: AtomicU64 = AtomicU64::new(1);
for attempt in 0..128_u64 {
let nonce = NEXT_TEMPORARY.fetch_add(1, Ordering::Relaxed);
let name = format!(
".{final_name}.{}.{}.{}.tmp",
std::process::id(),
nonce,
attempt
);
match directory.create_new(name.as_ref()) {
Ok(mut file) => {
if let Err(error) = file.write_all(bytes).and_then(|()| file.sync_all()) {
let _ = directory.remove_file(name.as_ref());
return Err(write_failed("temporary write", &error));
}
return Ok(name);
}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
Err(error) => return Err(write_failed("temporary creation", &error)),
}
}
Err(write_failed(
"temporary name allocation",
&std::io::Error::other("temporary name allocation exhausted"),
))
}
fn publish_no_replace(
directory: &MigrationDirectory,
temporary: &str,
target: &str,
) -> std::io::Result<()> {
directory.hard_link(temporary.as_ref(), target.as_ref())
}
fn sync_authoring_directory(directory: &MigrationDirectory) -> Result<(), Diagnostic> {
directory
.sync_all()
.map_err(|error| write_failed("directory sync", &error))
}
fn write_conflict() -> Diagnostic {
failure(
DiagnosticCategory::InvalidContract,
"migration_generation_write_conflict",
"generated migration file already exists in the target directory",
)
}
fn write_failed(operation: &'static str, error: &std::io::Error) -> Diagnostic {
failure(
DiagnosticCategory::Integrity,
"migration_generation_write_failed",
format!("generated migration file could not be created: {operation}: {error}"),
)
}
fn allocate_name(graph: &MigrationHistoryGraph, base_name: &str) -> String {
let next = graph
.manifests()
.filter_map(|(id, _)| leading_ordinal(id.name().as_str()))
.max()
.unwrap_or(0)
.saturating_add(1);
format!("{next:04}_{base_name}")
}
fn leading_ordinal(name: &str) -> Option<u64> {
let (digits, _) = name.split_once('_')?;
if digits.is_empty() || !digits.bytes().all(|byte| byte.is_ascii_digit()) {
return None;
}
digits.parse().ok()
}
fn author_steps(
delta: &SchemaDelta,
source: &DeclaredSchema,
target: &DeclaredSchema,
context: &ManagedDeltaContext,
with_reverse: bool,
) -> Result<Vec<MigrationStep>, Diagnostic> {
let safety_profile = SafetyDerivationProfile::new(
SemanticProfileBinding::resolve(context.semantic_profile().clone())?,
schema_lowering_profile_binding()?,
)?;
let resolved = resolve(source, context.semantic_profile()).map_err(|diagnostics| {
diagnostics
.iter()
.next()
.map(|diagnostic| diagnostic.diagnostic().clone())
.unwrap_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_generation_resolution_failed",
"assertion source resolution failed without a diagnostic",
)
})
})?;
let source_state = managed_schema_state(source, context).map_err(delta_diagnostic)?;
let validation_context = MigrationAssertionValidationContext::new(&resolved, &source_state);
let mut steps = Vec::new();
let domain = SafetyConditionDomainIndex::new(source, target);
for (ordinal, operation) in delta.operations().iter().enumerate() {
let derived = derive_safety_conditions_with_domain_index(
ordinal,
operation,
source,
target,
&safety_profile,
&domain,
)?;
let required = derived
.conditions()
.iter()
.filter(|condition| condition.policy() == SafetyClass::Conditional);
for (index, condition) in required.enumerate() {
let validated = lower_condition_to_plan(
condition,
&validation_context,
StructuralLimits::CANONICAL,
)?;
steps.push(MigrationStep::assertion(
MigrationStepId::new(format!("assert-{ordinal}-{index}"))?,
validated.plan().clone(),
AssertionExpectation::NoRows,
)?);
}
}
let reverse = if with_reverse {
Some(inverse_delta(delta).map_err(delta_diagnostic)?)
} else {
None
};
steps.push(MigrationStep::from(SchemaDeltaStep::new(
MigrationStepId::new("schema-delta")?,
delta.clone(),
reverse,
)?));
Ok(steps)
}
fn failure(
category: DiagnosticCategory,
code: &'static str,
message: impl Into<String>,
) -> Diagnostic {
Diagnostic::new(
category,
DiagnosticCode::new(code).expect("static generation diagnostic code is canonical"),
message,
)
}