1use 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
43const 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#[derive(Clone, Copy, Debug)]
57pub struct MigrationGenerationRequest<'a> {
58 pub app_label: &'a str,
60 pub base_name: &'a str,
62 pub genesis_source: &'a DeclaredSchema,
64 pub desired: &'a DeclaredSchema,
66 pub context: &'a ManagedDeltaContext,
68}
69
70#[derive(Clone, Copy, Debug)]
72pub struct BackfillMigrationGenerationRequest<'a> {
73 pub app_label: &'a str,
75 pub base_name: &'a str,
77 pub genesis_source: &'a DeclaredSchema,
79 pub desired: &'a DeclaredSchema,
81 pub plan: &'a AttributeBackfillPlan,
83 pub context: &'a ManagedDeltaContext,
85}
86
87#[derive(Clone, Debug)]
89pub enum MigrationGenerationOutcome {
90 UpToDate,
92 Generated(Box<GeneratedMigration>),
94}
95
96#[derive(Clone, Debug)]
98pub struct GeneratedMigration {
99 manifest: VerifiedSchemaMigrationManifest,
100 canonical_bytes: Vec<u8>,
101}
102
103impl GeneratedMigration {
104 pub const fn manifest(&self) -> &VerifiedSchemaMigrationManifest {
106 &self.manifest
107 }
108
109 pub fn canonical_bytes(&self) -> &[u8] {
111 &self.canonical_bytes
112 }
113
114 pub fn file_name(&self) -> String {
116 format!("{}.tbmigration.json", self.manifest.id().name().as_str())
117 }
118
119 pub fn preview_file_name(&self) -> String {
121 format!("{}.typeql", self.manifest.id().name().as_str())
122 }
123}
124
125#[derive(Clone, Debug)]
127pub enum MigrationPreviewError {
128 Diagnostic(Diagnostic),
130 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
157pub 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
226pub 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
283pub 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(¤t, schema_step.delta(), context).map_err(delta_diagnostic)?;
309 let coverage = verify_assertion_coverage(
310 &pending,
311 schema_step.delta(),
312 ¤t,
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
337pub 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
355pub 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 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 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
498fn 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
525fn 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
537fn 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
559fn 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 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}