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::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
41const 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#[derive(Clone, Copy, Debug)]
55pub struct MigrationGenerationRequest<'a> {
56 pub app_label: &'a str,
58 pub base_name: &'a str,
60 pub genesis_source: &'a DeclaredSchema,
62 pub desired: &'a DeclaredSchema,
64 pub context: &'a ManagedDeltaContext,
66}
67
68#[derive(Clone, Debug)]
70pub enum MigrationGenerationOutcome {
71 UpToDate,
73 Generated(Box<GeneratedMigration>),
75}
76
77#[derive(Clone, Debug)]
79pub struct GeneratedMigration {
80 manifest: VerifiedSchemaMigrationManifest,
81 canonical_bytes: Vec<u8>,
82}
83
84impl GeneratedMigration {
85 pub const fn manifest(&self) -> &VerifiedSchemaMigrationManifest {
87 &self.manifest
88 }
89
90 pub fn canonical_bytes(&self) -> &[u8] {
92 &self.canonical_bytes
93 }
94
95 pub fn file_name(&self) -> String {
97 format!("{}.tbmigration.json", self.manifest.id().name().as_str())
98 }
99
100 pub fn preview_file_name(&self) -> String {
102 format!("{}.typeql", self.manifest.id().name().as_str())
103 }
104}
105
106#[derive(Clone, Debug)]
108pub enum MigrationPreviewError {
109 Diagnostic(Diagnostic),
111 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
138pub 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
207pub 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(¤t, schema_step.delta(), context).map_err(delta_diagnostic)?;
233 let coverage = verify_assertion_coverage(
234 &pending,
235 schema_step.delta(),
236 ¤t,
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
261pub 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
279pub 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 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 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
422fn 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
449fn 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
461fn 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
483fn 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 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}