Skip to main content

type_bridge_schema_migration/
generate.rs

1//! Offline desired-state migration generation.
2//!
3//! Generation is the authoring half of the migration engine: the committed
4//! head's verified target schema is the source, the desired schema compiled
5//! from schema sources is the target, and the minimal fact delta between them
6//! becomes the next canonical manifest. Everything derives from the inputs —
7//! no timestamps, no environment — so a clean checkout regenerates identical
8//! bytes.
9
10use std::error::Error;
11use std::fmt;
12use std::io::{Read as _, Write};
13use std::path::PathBuf;
14
15use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
16use type_bridge_contract::limits::StructuralLimits;
17use type_bridge_contract::migration::{
18    MigrationId, MigrationStep, MigrationStepId, SchemaDeltaStep,
19};
20use type_bridge_contract::migration_assertion::AssertionExpectation;
21use type_bridge_contract::schema::{DeclaredSchema, SchemaDelta};
22use type_bridge_query::{MigrationAssertionValidationContext, lower_condition_to_plan};
23use type_bridge_schema::{
24    ManagedDeltaContext, SafetyClass, SafetyDerivationProfile, apply_delta,
25    derive_safety_conditions, diff_managed, inverse_delta, managed_schema_state, resolve,
26};
27
28use crate::history::MigrationHistoryGraph;
29use crate::lowering::{
30    SchemaFactCatalog, SchemaLoweringBinding, SchemaLoweringDiagnostic,
31    lower_schema_delta_with_verified_assertions,
32};
33use crate::manifest::{
34    SchemaMigrationDraft, VerifiedSchemaMigrationManifest, build_verified_manifest,
35    delta_diagnostic, encode_verified_manifest, verify_assertion_coverage,
36};
37use crate::profile::schema_lowering_profile_binding;
38use crate::{MigrationAuthoringLock, MigrationDirectory};
39use type_bridge_contract::managed_scope::SemanticProfileBinding;
40
41/// Manifest build rejections that only indict the claimed reverse program.
42///
43/// A structural inverse always exists, but it is recorded only when the
44/// verifier accepts it as a real semantic rollback; these codes downgrade the
45/// draft to an irreversible manifest instead of failing generation.
46const REVERSE_REJECTION_CODES: [&str; 4] = [
47    "migration_manifest_reverse_unresolved_safety",
48    "migration_manifest_reverse_requires_assertions",
49    "migration_manifest_inverse_replay_mismatch",
50    "migration_manifest_inverse_plan_invalid",
51];
52
53/// Inputs authoring the next migration for one application lineage.
54#[derive(Clone, Copy, Debug)]
55pub struct MigrationGenerationRequest<'a> {
56    /// Durable application label binding the migration lineage.
57    pub app_label: &'a str,
58    /// Author-supplied descriptive name; the ordinal prefix is allocated.
59    pub base_name: &'a str,
60    /// Exact declared source when history is empty.
61    pub genesis_source: &'a DeclaredSchema,
62    /// Desired schema compiled from the current schema sources.
63    pub desired: &'a DeclaredSchema,
64    /// Managed scope, semantic profile, and available capabilities.
65    pub context: &'a ManagedDeltaContext,
66}
67
68/// Result of one offline generation attempt.
69#[derive(Clone, Debug)]
70pub enum MigrationGenerationOutcome {
71    /// The head target already equals the desired schema.
72    UpToDate,
73    /// A new verified manifest was authored.
74    Generated(Box<GeneratedMigration>),
75}
76
77/// A freshly authored manifest with its exact canonical persistence bytes.
78#[derive(Clone, Debug)]
79pub struct GeneratedMigration {
80    manifest: VerifiedSchemaMigrationManifest,
81    canonical_bytes: Vec<u8>,
82}
83
84impl GeneratedMigration {
85    /// Return the verified manifest.
86    pub const fn manifest(&self) -> &VerifiedSchemaMigrationManifest {
87        &self.manifest
88    }
89
90    /// Return the exact canonical manifest bytes to persist.
91    pub fn canonical_bytes(&self) -> &[u8] {
92        &self.canonical_bytes
93    }
94
95    /// Return the stem-bound canonical file name for this manifest.
96    pub fn file_name(&self) -> String {
97        format!("{}.tbmigration.json", self.manifest.id().name().as_str())
98    }
99
100    /// Return the review-only TypeQL preview file name for this manifest.
101    pub fn preview_file_name(&self) -> String {
102        format!("{}.typeql", self.manifest.id().name().as_str())
103    }
104}
105
106/// Failure rendering a review-only TypeQL preview for a generated manifest.
107#[derive(Clone, Debug)]
108pub enum MigrationPreviewError {
109    /// Verification or replay of the manifest steps failed.
110    Diagnostic(Diagnostic),
111    /// Provider lowering rejected a statement unit.
112    Lowering(SchemaLoweringDiagnostic),
113}
114
115impl fmt::Display for MigrationPreviewError {
116    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
117        match self {
118            Self::Diagnostic(diagnostic) => diagnostic.fmt(formatter),
119            Self::Lowering(diagnostic) => diagnostic.fmt(formatter),
120        }
121    }
122}
123
124impl Error for MigrationPreviewError {}
125
126impl From<Diagnostic> for MigrationPreviewError {
127    fn from(value: Diagnostic) -> Self {
128        Self::Diagnostic(value)
129    }
130}
131
132impl From<SchemaLoweringDiagnostic> for MigrationPreviewError {
133    fn from(value: SchemaLoweringDiagnostic) -> Self {
134        Self::Lowering(value)
135    }
136}
137
138/// Author the next verified migration from the discovered history graph.
139///
140/// The source schema is the sole head's verified target (genesis when the
141/// graph is empty); a multi-head history is refused until an explicit merge
142/// migration joins it. Conditional operations receive verifier-derived
143/// assertion steps; the structural inverse is recorded only when the
144/// verifier accepts it as a real rollback program.
145pub fn generate_next_migration(
146    graph: &MigrationHistoryGraph,
147    request: &MigrationGenerationRequest<'_>,
148) -> Result<MigrationGenerationOutcome, Diagnostic> {
149    for (id, _) in graph.manifests() {
150        if id.app_label().as_str() != request.app_label {
151            return Err(failure(
152                DiagnosticCategory::InvalidContract,
153                "migration_generation_foreign_app_label",
154                "history contains a migration from a different application lineage",
155            ));
156        }
157    }
158
159    let (source, parents) = match graph.default_head()? {
160        None => (request.genesis_source, Vec::new()),
161        Some(head) => (
162            graph
163                .manifest(head)
164                .expect("a graph head is always a graph member")
165                .target_schema(),
166            vec![head.clone()],
167        ),
168    };
169
170    let source_state = managed_schema_state(source, request.context).map_err(delta_diagnostic)?;
171    let desired_state =
172        managed_schema_state(request.desired, request.context).map_err(delta_diagnostic)?;
173    if source_state == desired_state {
174        return Ok(MigrationGenerationOutcome::UpToDate);
175    }
176    let delta = diff_managed(source, request.desired, request.context).map_err(delta_diagnostic)?;
177
178    let id = MigrationId::new(request.app_label, allocate_name(graph, request.base_name))?;
179    if graph.manifest(&id).is_some() {
180        return Err(failure(
181            DiagnosticCategory::InvalidContract,
182            "migration_generation_duplicate_name",
183            "allocated migration identity already exists in verified history",
184        ));
185    }
186
187    let steps = author_steps(&delta, source, request.desired, request.context, true)?;
188    let draft = SchemaMigrationDraft::new(id.clone(), parents.clone(), steps)?;
189    let manifest = match build_verified_manifest(draft, (source, request.context)) {
190        Ok(manifest) => manifest,
191        Err(diagnostic) if REVERSE_REJECTION_CODES.contains(&diagnostic.code().as_str()) => {
192            let steps = author_steps(&delta, source, request.desired, request.context, false)?;
193            let draft = SchemaMigrationDraft::new(id, parents, steps)?;
194            build_verified_manifest(draft, (source, request.context))?
195        }
196        Err(diagnostic) => return Err(diagnostic),
197    };
198    let canonical_bytes = encode_verified_manifest(&manifest)?;
199    Ok(MigrationGenerationOutcome::Generated(Box::new(
200        GeneratedMigration {
201            manifest,
202            canonical_bytes,
203        },
204    )))
205}
206
207/// Render the review-only TypeQL preview of a verified manifest.
208///
209/// The preview replays the manifest steps and lowers every schema delta under
210/// the context's capabilities. Destructive units render without an approval —
211/// the preview is exactly what an operator inspects before granting one —
212/// while classes no approval can execute (backfill, opaque) surface their
213/// gate diagnostic instead of producing a misleading partial preview.
214pub fn render_migration_preview(
215    manifest: &VerifiedSchemaMigrationManifest,
216    context: &ManagedDeltaContext,
217) -> Result<String, MigrationPreviewError> {
218    let binding = SchemaLoweringBinding::current(context.available_capabilities().clone())?;
219    let safety_profile = SafetyDerivationProfile::new(
220        manifest.semantic_profile().clone(),
221        manifest.lowering_profile().clone(),
222    )?;
223    let mut current = manifest.source_schema().clone();
224    let mut pending = Vec::new();
225    let mut queries = Vec::new();
226    for step in manifest.steps() {
227        let Some(schema_step) = step.as_schema_delta() else {
228            pending.push(step);
229            continue;
230        };
231        let target =
232            apply_delta(&current, schema_step.delta(), context).map_err(delta_diagnostic)?;
233        let coverage = verify_assertion_coverage(
234            &pending,
235            schema_step.delta(),
236            &current,
237            &target,
238            &safety_profile,
239        )?;
240        pending.clear();
241        let source_catalog = SchemaFactCatalog::new(current.facts().cloned())?;
242        let target_catalog = SchemaFactCatalog::new(target.facts().cloned())?;
243        let lowering = lower_schema_delta_with_verified_assertions(
244            schema_step.delta(),
245            &source_catalog,
246            &target_catalog,
247            &binding,
248            coverage.conditional_operation_indices(),
249            true,
250        )?;
251        for unit in lowering.units() {
252            for statement in unit.statements() {
253                queries.push(statement.query().to_owned());
254            }
255        }
256        current = target;
257    }
258    Ok(format!("{}\n", queries.join("\n\n")))
259}
260
261/// Acquire the shared canonical-history authoring lock without waiting.
262///
263/// Callers that derive a candidate from the directory must acquire this lock
264/// before re-discovery and retain it through
265/// [`write_generated_migration_under_lock`]. This prevents two different
266/// candidates derived from one stale head from becoming sibling authorities.
267pub fn try_acquire_migration_authoring_lock(
268    directory: &MigrationDirectory,
269) -> Result<MigrationAuthoringLock<'_>, Diagnostic> {
270    directory.try_acquire_authoring_lock().map_err(|error| {
271        if error.kind() == std::io::ErrorKind::WouldBlock {
272            write_conflict()
273        } else {
274            write_failed("authoring lock acquisition", &error)
275        }
276    })
277}
278
279/// Publish a generated manifest while its directory's authoring lock is held.
280///
281/// The lock carries the exact retained directory capability, so a caller
282/// cannot accidentally validate one history and publish into another. An
283/// existing manifest is a conflict, never an overwrite. Both files are
284/// written completely to confined, flushed temporaries and published with
285/// no-replace links, preview first and authority manifest last.
286pub fn write_generated_migration_under_lock(
287    lock: &MigrationAuthoringLock<'_>,
288    generated: &GeneratedMigration,
289    preview: &str,
290) -> Result<PathBuf, Diagnostic> {
291    let directory = lock.directory;
292    let manifest_name = generated.file_name();
293    let preview_name = generated.preview_file_name();
294    if directory
295        .entry_exists(manifest_name.as_ref())
296        .map_err(|error| write_failed("manifest presence probe", &error))?
297    {
298        return Err(write_conflict());
299    }
300    // The advisory directory lock proves no live writer owns this preview;
301    // without a final manifest it is an interrupted, non-authoritative
302    // publication and can be recovered safely.
303    if directory
304        .entry_exists(preview_name.as_ref())
305        .map_err(|error| write_failed("preview presence probe", &error))?
306    {
307        directory
308            .remove_file(preview_name.as_ref())
309            .map_err(|error| write_failed("interrupted preview recovery", &error))?;
310    }
311    let manifest_temp = write_unique_temporary(
312        directory,
313        &generated.file_name(),
314        generated.canonical_bytes(),
315    )?;
316    let preview_temp = match write_unique_temporary(
317        directory,
318        &generated.preview_file_name(),
319        preview.as_bytes(),
320    ) {
321        Ok(path) => path,
322        Err(error) => {
323            let _ = directory.remove_file(manifest_temp.as_ref());
324            return Err(error);
325        }
326    };
327    let cleanup = |published_preview: bool| {
328        let _ = directory.remove_file(manifest_temp.as_ref());
329        let _ = directory.remove_file(preview_temp.as_ref());
330        if published_preview {
331            let _ = directory.remove_file(preview_name.as_ref());
332        }
333    };
334
335    let published_preview = match publish_no_replace(directory, &preview_temp, &preview_name) {
336        Ok(()) => true,
337        Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
338            // A prior interrupted writer or concurrent identical writer may
339            // already have published the non-authoritative preview. Reuse it
340            // only when its bounded bytes are exact; never unlink another
341            // writer's candidate.
342            match read_existing_preview(directory, &preview_name) {
343                Ok(existing) if existing == preview.as_bytes() => false,
344                _ => {
345                    cleanup(false);
346                    return Err(write_conflict());
347                }
348            }
349        }
350        Err(error) => {
351            cleanup(false);
352            return Err(write_failed("preview publication link", &error));
353        }
354    };
355    if let Err(error) = publish_no_replace(directory, &manifest_temp, &manifest_name) {
356        cleanup(published_preview);
357        return Err(if error.kind() == std::io::ErrorKind::AlreadyExists {
358            write_conflict()
359        } else {
360            write_failed("manifest publication link", &error)
361        });
362    }
363    if let Err(error) = sync_authoring_directory(directory) {
364        let _ = directory.remove_file(manifest_temp.as_ref());
365        let _ = directory.remove_file(preview_temp.as_ref());
366        return Err(error);
367    }
368    let _ = directory.remove_file(manifest_temp.as_ref());
369    let _ = directory.remove_file(preview_temp.as_ref());
370    sync_authoring_directory(directory)?;
371    Ok(PathBuf::from(manifest_name))
372}
373
374fn read_existing_preview(directory: &MigrationDirectory, name: &str) -> std::io::Result<Vec<u8>> {
375    let limit = type_bridge_contract::limits::MAX_CANONICAL_BYTES;
376    let file = directory.open_regular_readonly(name.as_ref())?;
377    let mut bytes = Vec::new();
378    file.take(u64::try_from(limit).unwrap_or(u64::MAX).saturating_add(1))
379        .read_to_end(&mut bytes)?;
380    if bytes.len() > limit {
381        return Err(std::io::Error::other(
382            "migration preview exceeds byte ceiling",
383        ));
384    }
385    Ok(bytes)
386}
387
388fn write_unique_temporary(
389    directory: &MigrationDirectory,
390    final_name: &str,
391    bytes: &[u8],
392) -> Result<String, Diagnostic> {
393    use std::sync::atomic::{AtomicU64, Ordering};
394    static NEXT_TEMPORARY: AtomicU64 = AtomicU64::new(1);
395
396    for attempt in 0..128_u64 {
397        let nonce = NEXT_TEMPORARY.fetch_add(1, Ordering::Relaxed);
398        let name = format!(
399            ".{final_name}.{}.{}.{}.tmp",
400            std::process::id(),
401            nonce,
402            attempt
403        );
404        match directory.create_new(name.as_ref()) {
405            Ok(mut file) => {
406                if let Err(error) = file.write_all(bytes).and_then(|()| file.sync_all()) {
407                    let _ = directory.remove_file(name.as_ref());
408                    return Err(write_failed("temporary write", &error));
409                }
410                return Ok(name);
411            }
412            Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
413            Err(error) => return Err(write_failed("temporary creation", &error)),
414        }
415    }
416    Err(write_failed(
417        "temporary name allocation",
418        &std::io::Error::other("temporary name allocation exhausted"),
419    ))
420}
421
422/// Publish a completed temporary under its final name without replacing.
423///
424/// A hard link fails when the destination exists, which keeps publication
425/// atomic and conflict-detecting at once; the temporary is removed by the
426/// caller after both publications succeed.
427fn publish_no_replace(
428    directory: &MigrationDirectory,
429    temporary: &str,
430    target: &str,
431) -> std::io::Result<()> {
432    directory.hard_link(temporary.as_ref(), target.as_ref())
433}
434
435fn sync_authoring_directory(directory: &MigrationDirectory) -> Result<(), Diagnostic> {
436    directory
437        .sync_all()
438        .map_err(|error| write_failed("directory sync", &error))
439}
440
441fn write_conflict() -> Diagnostic {
442    failure(
443        DiagnosticCategory::InvalidContract,
444        "migration_generation_write_conflict",
445        "generated migration file already exists in the target directory",
446    )
447}
448
449/// The integrity diagnostic names the failing operation and its
450/// underlying I/O failure so a production operator can distinguish
451/// permissions, exhaustion, and platform quirks without tracing the
452/// filesystem calls.
453fn write_failed(operation: &'static str, error: &std::io::Error) -> Diagnostic {
454    failure(
455        DiagnosticCategory::Integrity,
456        "migration_generation_write_failed",
457        format!("generated migration file could not be created: {operation}: {error}"),
458    )
459}
460
461/// Allocate the next ordinal-prefixed name in the lineage's local convention.
462///
463/// Ordinal prefixes are display/allocation conventions, never execution
464/// order; names without a numeric prefix simply do not participate.
465fn allocate_name(graph: &MigrationHistoryGraph, base_name: &str) -> String {
466    let next = graph
467        .manifests()
468        .filter_map(|(id, _)| leading_ordinal(id.name().as_str()))
469        .max()
470        .unwrap_or(0)
471        .saturating_add(1);
472    format!("{next:04}_{base_name}")
473}
474
475fn leading_ordinal(name: &str) -> Option<u64> {
476    let (digits, _) = name.split_once('_')?;
477    if digits.is_empty() || !digits.bytes().all(|byte| byte.is_ascii_digit()) {
478        return None;
479    }
480    digits.parse().ok()
481}
482
483/// Author verifier-derived assertion steps followed by the single delta step.
484fn author_steps(
485    delta: &SchemaDelta,
486    source: &DeclaredSchema,
487    target: &DeclaredSchema,
488    context: &ManagedDeltaContext,
489    with_reverse: bool,
490) -> Result<Vec<MigrationStep>, Diagnostic> {
491    let safety_profile = SafetyDerivationProfile::new(
492        SemanticProfileBinding::resolve(context.semantic_profile().clone())?,
493        schema_lowering_profile_binding()?,
494    )?;
495    let resolved = resolve(source, context.semantic_profile()).map_err(|diagnostics| {
496        diagnostics
497            .iter()
498            .next()
499            .map(|diagnostic| diagnostic.diagnostic().clone())
500            .unwrap_or_else(|| {
501                failure(
502                    DiagnosticCategory::Integrity,
503                    "migration_generation_resolution_failed",
504                    "assertion source resolution failed without a diagnostic",
505                )
506            })
507    })?;
508    let source_state = managed_schema_state(source, context).map_err(delta_diagnostic)?;
509    let validation_context = MigrationAssertionValidationContext::new(&resolved, &source_state);
510
511    let mut steps = Vec::new();
512    for (ordinal, operation) in delta.operations().iter().enumerate() {
513        let derived =
514            derive_safety_conditions(ordinal, operation, source, target, &safety_profile)?;
515        // Only conditional requirements become assertions. Destructive guard
516        // conditions are deliberately omitted: an approved destructive
517        // migration means the data loss is intended, and a NoRows guard would
518        // convert it into refuse-if-populated.
519        let required = derived
520            .conditions()
521            .iter()
522            .filter(|condition| condition.policy() == SafetyClass::Conditional);
523        for (index, condition) in required.enumerate() {
524            let validated = lower_condition_to_plan(
525                condition,
526                &validation_context,
527                StructuralLimits::CANONICAL,
528            )?;
529            steps.push(MigrationStep::assertion(
530                MigrationStepId::new(format!("assert-{ordinal}-{index}"))?,
531                validated.plan().clone(),
532                AssertionExpectation::NoRows,
533            )?);
534        }
535    }
536    let reverse = if with_reverse {
537        Some(inverse_delta(delta).map_err(delta_diagnostic)?)
538    } else {
539        None
540    };
541    steps.push(MigrationStep::from(SchemaDeltaStep::new(
542        MigrationStepId::new("schema-delta")?,
543        delta.clone(),
544        reverse,
545    )?));
546    Ok(steps)
547}
548
549fn failure(
550    category: DiagnosticCategory,
551    code: &'static str,
552    message: impl Into<String>,
553) -> Diagnostic {
554    Diagnostic::new(
555        category,
556        DiagnosticCode::new(code).expect("static generation diagnostic code is canonical"),
557        message,
558    )
559}