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