1use std::collections::BTreeSet;
17use std::sync::Arc;
18use std::sync::atomic::{AtomicBool, Ordering};
19
20use chrono::Utc;
21pub use type_bridge_contract::reserved::{
22 LEGACY_CUTOVER_SENTINEL_APP_LABEL, LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
23 LEGACY_CUTOVER_SENTINEL_MIGRATION_ID, LEGACY_CUTOVER_SENTINEL_NAME,
24 LEGACY_WRITER_CUTOVER_MESSAGE,
25};
26use type_bridge_orm::_schema::{SchemaError, SchemaInfo};
27use type_bridge_orm::OrmError;
28use type_bridge_orm::session::backend::{BoxFuture, QueryResult, TxType};
29use type_bridge_orm::{
30 Database, Transaction, TransactionContext,
31 require_legacy_writer_open as require_orm_legacy_writer_open,
32 require_legacy_writer_open_in_transaction as require_orm_legacy_writer_open_in_transaction,
33};
34
35use crate::state::schema::labels::{
36 APP_LABEL, APPLIED_AT, APPLIED_ENTITY, CHECKSUM, DIRECTION, ERROR, EXECUTOR_IP, EXECUTOR_MAC,
37 FINISHED_AT, MIGRATION_ID, NAME, RUN_ENTITY, RUN_ID, STARTED_AT, STATUS,
38};
39use crate::state::schema::migration_state_schema;
40use crate::state::{MigrationRunRecord, MigrationStateStore};
41use crate::{AppliedMigrationRecord, MigrationError, Result};
42
43const APPLIED_AT_FORMAT: &str = "%Y-%m-%dT%H:%M:%S.%6f";
49
50const LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG: &str =
51 "# typebridge-internal-legacy-state-schema-probe/v1\n";
52
53#[derive(Debug, Clone, Copy)]
56pub enum LegacyCutoverSentinelExpectation<'a> {
57 Absent,
59 OptionalExact(&'a str),
61 RequiredExact(&'a str),
63}
64
65#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct VerifiedLegacyAppliedPartition {
72 applied: Vec<AppliedMigrationRecord>,
73 sentinel_fingerprint: Option<String>,
74}
75
76impl VerifiedLegacyAppliedPartition {
77 pub fn applied(&self) -> &[AppliedMigrationRecord] {
80 &self.applied
81 }
82
83 pub fn into_applied(self) -> Vec<AppliedMigrationRecord> {
85 self.applied
86 }
87
88 pub fn sentinel_fingerprint(&self) -> Option<&str> {
90 self.sentinel_fingerprint.as_deref()
91 }
92}
93
94#[derive(Debug, thiserror::Error)]
96pub enum LegacyCutoverSentinelError {
97 #[error("legacy cutover sentinel storage inspection failed: {0}")]
99 Storage(#[source] MigrationError),
100 #[error("legacy cutover sentinel contract violation: {message}")]
102 Contract {
103 message: String,
105 },
106}
107
108impl LegacyCutoverSentinelError {
109 pub fn is_contract_violation(&self) -> bool {
112 matches!(self, Self::Contract { .. })
113 }
114}
115
116pub struct TypeDbStateStore {
121 db: Arc<Database>,
122 schema_ensured: AtomicBool,
125}
126
127impl TypeDbStateStore {
128 pub fn new(db: Arc<Database>) -> Self {
130 Self {
131 db,
132 schema_ensured: AtomicBool::new(false),
133 }
134 }
135
136 async fn type_exists(&self, type_name: &str) -> Result<bool> {
137 let kind = state_type_kind(type_name).ok_or_else(|| MigrationError::State {
138 message: format!("unknown migration-state schema label: {type_name}"),
139 })?;
140 let check_query = format!(
141 "\n match {kind} $t;\n fetch {{ \"label\": label($t) }};\n "
142 );
143
144 let ctx = self
145 .db
146 .transaction_context(TxType::Read)
147 .await
148 .map_err(map_orm_error)?;
149 let checked = ctx
150 .query(&check_query)
151 .await
152 .map_err(map_orm_error)
153 .and_then(|result| schema_labels_from_result(result, "migration state type probe"))
154 .map(|labels| labels.contains(type_name));
155 let closed = ctx.close().await.map_err(map_orm_error);
156 match (checked, closed) {
157 (Ok(exists), Ok(())) => Ok(exists),
158 (Err(primary), Ok(())) => Err(primary),
159 (Ok(_), Err(cleanup)) => Err(cleanup),
160 (Err(primary), Err(_)) => Err(primary),
161 }
162 }
163
164 async fn ensure_type(&self, type_name: &str, define_typeql: &str) -> Result<()> {
165 if self.type_exists(type_name).await? {
166 return Ok(());
167 }
168
169 let ctx = self
170 .db
171 .transaction_context(TxType::Schema)
172 .await
173 .map_err(map_orm_error)?;
174 if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
175 let _ = ctx.rollback().await;
176 return Err(error);
177 }
178 match ctx.query(define_typeql).await {
179 Ok(_) => {
180 ctx.commit().await.map_err(map_orm_error)?;
181 Ok(())
182 }
183 Err(error) => {
184 if let Err(cleanup) = ctx.rollback().await {
185 return Err(schema_rollback_cleanup_failure(&error, &cleanup));
186 }
187 if self.type_exists(type_name).await? {
188 Ok(())
189 } else {
190 Err(map_orm_error(error))
191 }
192 }
193 }
194 }
195
196 async fn ensure_schema_for_read(&self) -> Result<()> {
197 if self.schema_ensured.load(Ordering::Acquire) {
198 return Ok(());
199 }
200
201 let mut transaction = self.db.read_transaction().await.map_err(map_orm_error)?;
206 let inspected = legacy_state_schema_presence(&mut transaction)
207 .await
208 .map_err(legacy_sentinel_error_into_migration_error);
209 let closed = transaction.close().await.map_err(map_orm_error);
210 let presence = match (inspected, closed) {
211 (Ok(presence), Ok(())) => presence,
212 (Err(primary), Ok(())) => return Err(primary),
213 (Ok(_), Err(cleanup)) => return Err(cleanup),
214 (Err(primary), Err(_)) => return Err(primary),
215 };
216 if presence == LegacyStateSchemaPresence::Complete {
217 self.schema_ensured.store(true, Ordering::Release);
218 return Ok(());
219 }
220
221 self.ensure_schema().await
226 }
227
228 async fn query_documents(&self, query: &str) -> Result<Vec<serde_json::Value>> {
229 let ctx = self
230 .db
231 .transaction_context(TxType::Read)
232 .await
233 .map_err(map_orm_error)?;
234 let queried = ctx
235 .query(query)
236 .await
237 .map_err(map_orm_error)
238 .map(query_result_values);
239 let closed = ctx.close().await.map_err(map_orm_error);
240 match (queried, closed) {
241 (Ok(values), Ok(())) => Ok(values),
242 (Err(primary), Ok(())) => Err(primary),
243 (Ok(_), Err(cleanup)) => Err(cleanup),
244 (Err(primary), Err(_)) => Err(primary),
245 }
246 }
247
248 pub async fn load_applied_archival(&self) -> Result<Vec<AppliedMigrationRecord>> {
254 if !self.archival_schema_is_complete().await? {
255 return Ok(Vec::new());
256 }
257 parse_applied_documents(&self.query_documents(&applied_query()).await?)
258 }
259
260 pub async fn load_runs_archival(&self) -> Result<Vec<MigrationRunRecord>> {
262 if !self.archival_schema_is_complete().await? {
263 return Ok(Vec::new());
264 }
265 self.load_runs_documents().await
266 }
267
268 async fn archival_schema_is_complete(&self) -> Result<bool> {
269 let mut transaction = self.db.read_transaction().await.map_err(map_orm_error)?;
270 let inspected = legacy_state_schema_presence(&mut transaction)
271 .await
272 .map_err(legacy_sentinel_error_into_migration_error);
273 let closed = transaction.close().await.map_err(map_orm_error);
274 let presence = match (inspected, closed) {
275 (Ok(presence), Ok(())) => presence,
276 (Err(primary), Ok(())) => return Err(primary),
277 (Ok(_), Err(cleanup)) => return Err(cleanup),
278 (Err(primary), Err(_)) => return Err(primary),
279 };
280 match presence {
281 LegacyStateSchemaPresence::Absent => Ok(false),
282 LegacyStateSchemaPresence::Complete => Ok(true),
283 LegacyStateSchemaPresence::Partial => Err(MigrationError::State {
284 message: "the frozen legacy ledger schema is partially present".to_owned(),
285 }),
286 }
287 }
288
289 async fn load_runs_documents(&self) -> Result<Vec<MigrationRunRecord>> {
290 let query = format!(
291 "\nmatch\n$r isa {RUN_ENTITY},\n has {RUN_ID} $run_id,\n has {APP_LABEL} $app,\n has {NAME} $name,\n has {CHECKSUM} $checksum,\n has {DIRECTION} $direction,\n has {STATUS} $status,\n has {STARTED_AT} $started;\nfetch {{\n \"run_id\": $run_id,\n \"app\": $app,\n \"name\": $name,\n \"checksum\": $checksum,\n \"direction\": $direction,\n \"status\": $status,\n \"started\": $started\n}};\n"
292 );
293 let mut runs = parse_run_documents(&self.query_documents(&query).await?)?;
294
295 let finished_query = optional_run_field_query(FINISHED_AT, "finished");
296 let finished_docs = self.query_documents(&finished_query).await?;
297 merge_optional_run_field(&mut runs, &finished_docs, "finished_at", "finished");
298
299 let error_query = optional_run_field_query(ERROR, "error");
300 let error_docs = self.query_documents(&error_query).await?;
301 merge_optional_run_field(&mut runs, &error_docs, "error", "error");
302
303 let ip_query = optional_run_field_query(EXECUTOR_IP, "executor_ip");
304 let ip_docs = self.query_documents(&ip_query).await?;
305 merge_optional_run_field(&mut runs, &ip_docs, "executor_ip", "executor_ip");
306
307 let mac_query = optional_run_field_query(EXECUTOR_MAC, "executor_mac");
308 let mac_docs = self.query_documents(&mac_query).await?;
309 merge_optional_run_field(&mut runs, &mac_docs, "executor_mac", "executor_mac");
310
311 Ok(runs)
312 }
313
314 pub async fn load_applied_in_transaction(
322 transaction: &mut Transaction,
323 ) -> Result<Vec<AppliedMigrationRecord>> {
324 let result = transaction
325 .query(&applied_query())
326 .await
327 .map_err(map_orm_error)?;
328 parse_applied_documents(&query_result_values(result))
329 }
330
331 pub async fn load_verified_legacy_partition_in_transaction(
339 transaction: &mut Transaction,
340 expectation: LegacyCutoverSentinelExpectation<'_>,
341 ) -> std::result::Result<VerifiedLegacyAppliedPartition, LegacyCutoverSentinelError> {
342 match legacy_state_schema_presence(transaction).await? {
343 LegacyStateSchemaPresence::Absent => {
344 if matches!(
345 expectation,
346 LegacyCutoverSentinelExpectation::RequiredExact(_)
347 ) {
348 return Err(sentinel_contract_error(
349 "the V2 bridge is active but the frozen legacy ledger schema is absent",
350 ));
351 }
352 return Ok(VerifiedLegacyAppliedPartition {
353 applied: Vec::new(),
354 sentinel_fingerprint: None,
355 });
356 }
357 LegacyStateSchemaPresence::Partial => {
358 return Err(sentinel_contract_error(
359 "the frozen legacy ledger schema is partially present",
360 ));
361 }
362 LegacyStateSchemaPresence::Complete => {}
363 }
364 let result = transaction
365 .query(&applied_query())
366 .await
367 .map_err(map_sentinel_storage_error)?;
368 let mut applied = parse_applied_documents(&query_result_values(result))
369 .map_err(LegacyCutoverSentinelError::Storage)?;
370
371 let id_candidates = sentinel_query_values(
372 transaction,
373 &format!(
374 "match $m isa {APPLIED_ENTITY}, has {MIGRATION_ID} {}; fetch {{ \"exists\": true }};",
375 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
376 ),
377 )
378 .await?;
379 let name_candidates = sentinel_query_values(
380 transaction,
381 &format!(
382 "match $m isa {APPLIED_ENTITY}, has {NAME} {}; fetch {{ \"exists\": true }};",
383 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
384 ),
385 )
386 .await?;
387
388 if id_candidates.is_empty() && name_candidates.is_empty() {
389 if matches!(
390 expectation,
391 LegacyCutoverSentinelExpectation::RequiredExact(_)
392 ) {
393 return Err(sentinel_contract_error(
394 "the V2 bridge is active but its legacy-writer sentinel is missing",
395 ));
396 }
397 return Ok(VerifiedLegacyAppliedPartition {
398 applied,
399 sentinel_fingerprint: None,
400 });
401 }
402
403 if matches!(expectation, LegacyCutoverSentinelExpectation::Absent) {
404 return Err(sentinel_contract_error(
405 "a legacy-writer sentinel exists without an active or pending V2 bridge",
406 ));
407 }
408 if id_candidates.len() != 1 || name_candidates.len() != 1 {
409 return Err(sentinel_contract_error(
410 "the legacy-writer sentinel is duplicated or has split identity rows",
411 ));
412 }
413
414 let details = sentinel_query_values(
415 transaction,
416 &format!(
417 "match $m isa {APPLIED_ENTITY}, has {MIGRATION_ID} {}, has {APP_LABEL} $app, has {NAME} {}, has {APPLIED_AT} $applied, has {CHECKSUM} $checksum; fetch {{ \"app\": $app, \"applied\": $applied, \"checksum\": $checksum }};",
418 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
419 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
420 ),
421 )
422 .await?;
423 if details.len() != 1 {
424 return Err(sentinel_contract_error(
425 "the legacy-writer sentinel is missing required exact fields",
426 ));
427 }
428 let detail = &details[0];
429 let app = extract_value(detail, "app").ok_or_else(|| {
430 sentinel_contract_error("the legacy-writer sentinel app label is malformed")
431 })?;
432 let applied_at = extract_value(detail, "applied").ok_or_else(|| {
433 sentinel_contract_error("the legacy-writer sentinel applied timestamp is malformed")
434 })?;
435 let fingerprint = extract_value(detail, "checksum").ok_or_else(|| {
436 sentinel_contract_error("the legacy-writer sentinel checksum is malformed")
437 })?;
438 if app != LEGACY_CUTOVER_SENTINEL_APP_LABEL {
439 return Err(sentinel_contract_error(
440 "the legacy-writer sentinel carries a foreign application label",
441 ));
442 }
443 if applied_at != LEGACY_CUTOVER_SENTINEL_APPLIED_AT {
444 return Err(sentinel_contract_error(
445 "the legacy-writer sentinel carries a foreign applied timestamp",
446 ));
447 }
448 if !is_lower_hex_fingerprint(&fingerprint) {
449 return Err(sentinel_contract_error(
450 "the legacy-writer sentinel checksum is not a lowercase 64-hex fingerprint",
451 ));
452 }
453 let expected = match expectation {
454 LegacyCutoverSentinelExpectation::OptionalExact(expected)
455 | LegacyCutoverSentinelExpectation::RequiredExact(expected) => expected,
456 LegacyCutoverSentinelExpectation::Absent => unreachable!("handled above"),
457 };
458 if fingerprint != expected {
459 return Err(sentinel_contract_error(
460 "the legacy-writer sentinel checksum differs from the managed cutover anchor",
461 ));
462 }
463
464 let original_len = applied.len();
465 applied.retain(|record| {
466 record.app_label != LEGACY_CUTOVER_SENTINEL_APP_LABEL
467 || record.name != LEGACY_CUTOVER_SENTINEL_NAME
468 });
469 if original_len.saturating_sub(applied.len()) != 1 {
470 return Err(sentinel_contract_error(
471 "the exact legacy-writer sentinel is absent from the released applied projection",
472 ));
473 }
474
475 Ok(VerifiedLegacyAppliedPartition {
476 applied,
477 sentinel_fingerprint: Some(fingerprint),
478 })
479 }
480
481 pub async fn insert_legacy_cutover_sentinel_in_transaction(
485 transaction: &mut Transaction,
486 anchor_fingerprint: &str,
487 ) -> Result<()> {
488 if !is_lower_hex_fingerprint(anchor_fingerprint) {
489 return Err(MigrationError::State {
490 message: "legacy cutover sentinel requires a lowercase 64-hex anchor fingerprint"
491 .to_owned(),
492 });
493 }
494 let query = format!(
495 "insert $m isa {APPLIED_ENTITY}, has {MIGRATION_ID} {}, has {APP_LABEL} {}, has {NAME} {}, has {APPLIED_AT} {LEGACY_CUTOVER_SENTINEL_APPLIED_AT}, has {CHECKSUM} {};",
496 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
497 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_APP_LABEL),
498 typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
499 typeql_string_literal(anchor_fingerprint),
500 );
501 transaction.query(&query).await.map_err(map_orm_error)?;
502 Ok(())
503 }
504}
505
506pub async fn require_legacy_writer_open_in_transaction(
509 transaction: &TransactionContext,
510) -> Result<()> {
511 require_orm_legacy_writer_open_in_transaction(transaction)
512 .await
513 .map_err(map_legacy_guard_error)
514}
515
516#[cfg(test)]
517pub(crate) fn is_legacy_state_schema_probe_query(query: &str) -> bool {
518 query.starts_with(LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG)
519}
520
521pub async fn require_legacy_writer_open(database: &Database) -> Result<()> {
524 require_orm_legacy_writer_open(database)
525 .await
526 .map_err(map_legacy_guard_error)
527}
528
529#[derive(Debug, Clone, Copy, PartialEq, Eq)]
530enum LegacyStateSchemaPresence {
531 Absent,
532 Partial,
533 Complete,
534}
535
536async fn legacy_state_schema_presence(
537 transaction: &mut Transaction,
538) -> std::result::Result<LegacyStateSchemaPresence, LegacyCutoverSentinelError> {
539 let state_schema = migration_state_schema();
540 let mut present = 0_usize;
541 let expected_by_root = [
542 (
543 "attribute",
544 state_schema
545 .attributes
546 .keys()
547 .map(String::as_str)
548 .collect::<Vec<_>>(),
549 ),
550 (
551 "entity",
552 state_schema
553 .entities
554 .keys()
555 .map(String::as_str)
556 .collect::<Vec<_>>(),
557 ),
558 (
559 "relation",
560 state_schema
561 .relations
562 .keys()
563 .map(String::as_str)
564 .collect::<Vec<_>>(),
565 ),
566 ];
567 let total = expected_by_root
568 .iter()
569 .map(|(_, labels)| labels.len())
570 .sum::<usize>();
571 let all_expected = expected_by_root
572 .iter()
573 .flat_map(|(_, labels)| labels.iter().copied())
574 .collect::<BTreeSet<_>>();
575 let mut expected_labels_seen_in_any_kind = 0_usize;
576 for (kind, expected) in expected_by_root {
577 let result = transaction
578 .query(&format!(
579 "{LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG}match {kind} $t; fetch {{ \"label\": label($t) }};"
580 ))
581 .await
582 .map_err(map_sentinel_storage_error)?;
583 let observed = schema_labels_from_result(result, "legacy state schema probe")
584 .map_err(LegacyCutoverSentinelError::Storage)?;
585 expected_labels_seen_in_any_kind += observed
586 .iter()
587 .filter(|label| all_expected.contains(label.as_str()))
588 .count();
589 present += expected
590 .into_iter()
591 .filter(|label| observed.contains(*label))
592 .count();
593 }
594 if present == 0 && expected_labels_seen_in_any_kind == 0 {
595 return Ok(LegacyStateSchemaPresence::Absent);
596 }
597 if present != total || expected_labels_seen_in_any_kind != total {
598 return Ok(LegacyStateSchemaPresence::Partial);
599 }
600 Ok(LegacyStateSchemaPresence::Complete)
601}
602
603fn state_type_kind(type_name: &str) -> Option<&'static str> {
604 let schema = migration_state_schema();
605 if schema.attributes.contains_key(type_name) {
606 Some("attribute")
607 } else if schema.entities.contains_key(type_name) {
608 Some("entity")
609 } else if schema.relations.contains_key(type_name) {
610 Some("relation")
611 } else {
612 None
613 }
614}
615
616fn schema_labels_from_result(result: QueryResult, operation: &str) -> Result<BTreeSet<String>> {
617 let values = match result {
618 QueryResult::Documents(values) | QueryResult::Rows(values) => values,
619 QueryResult::Ok => {
620 return Err(MigrationError::State {
621 message: format!("{operation} returned no document result"),
622 });
623 }
624 };
625 let mut labels = BTreeSet::new();
626 for value in &values {
627 let label = extract_value(value, "label").ok_or_else(|| MigrationError::State {
628 message: format!("{operation} returned a malformed schema label"),
629 })?;
630 labels.insert(label);
631 }
632 Ok(labels)
633}
634
635fn legacy_sentinel_error_into_migration_error(error: LegacyCutoverSentinelError) -> MigrationError {
636 match error {
637 LegacyCutoverSentinelError::Storage(error) => error,
638 LegacyCutoverSentinelError::Contract { message } => MigrationError::State { message },
639 }
640}
641
642fn applied_query() -> String {
643 format!(
644 "\nmatch\n$m isa {APPLIED_ENTITY},\n has {APP_LABEL} $app,\n has {NAME} $name,\n has {APPLIED_AT} $applied,\n has {CHECKSUM} $checksum;\nfetch {{\n \"app\": $app,\n \"name\": $name,\n \"applied\": $applied,\n \"checksum\": $checksum\n}};\n"
645 )
646}
647
648async fn sentinel_query_values(
649 transaction: &mut Transaction,
650 query: &str,
651) -> std::result::Result<Vec<serde_json::Value>, LegacyCutoverSentinelError> {
652 let result = transaction
653 .query(query)
654 .await
655 .map_err(map_sentinel_storage_error)?;
656 match result {
657 QueryResult::Documents(values) | QueryResult::Rows(values) => Ok(values),
658 QueryResult::Ok => Err(LegacyCutoverSentinelError::Storage(MigrationError::State {
659 message: "legacy cutover sentinel fetch returned no document result".to_owned(),
660 })),
661 }
662}
663
664fn map_sentinel_storage_error(error: OrmError) -> LegacyCutoverSentinelError {
665 LegacyCutoverSentinelError::Storage(map_orm_error(error))
666}
667
668fn sentinel_contract_error(message: impl Into<String>) -> LegacyCutoverSentinelError {
669 LegacyCutoverSentinelError::Contract {
670 message: message.into(),
671 }
672}
673
674fn is_lower_hex_fingerprint(value: &str) -> bool {
675 value.len() == 64
676 && value
677 .bytes()
678 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
679}
680
681fn format_applied_at(now: chrono::DateTime<Utc>) -> String {
687 now.format(APPLIED_AT_FORMAT).to_string()
688}
689
690fn query_result_values(result: QueryResult) -> Vec<serde_json::Value> {
691 match result {
692 QueryResult::Documents(docs) => docs,
693 QueryResult::Rows(rows) => rows,
694 QueryResult::Ok => Vec::new(),
695 }
696}
697
698fn typeql_string_literal(value: &str) -> String {
699 let escaped = value
700 .replace('\\', "\\\\")
701 .replace('"', "\\\"")
702 .replace('\n', "\\n")
703 .replace('\r', "\\r")
704 .replace('\t', "\\t");
705 format!("\"{escaped}\"")
706}
707
708fn map_orm_error(error: OrmError) -> MigrationError {
710 MigrationError::State {
711 message: error.to_string(),
712 }
713}
714
715fn map_legacy_guard_error(error: OrmError) -> MigrationError {
716 match error {
717 OrmError::Transaction(message) if message == LEGACY_WRITER_CUTOVER_MESSAGE => {
718 MigrationError::State { message }
719 }
720 error => map_orm_error(error),
721 }
722}
723
724fn map_schema_error(error: SchemaError) -> MigrationError {
725 MigrationError::State {
726 message: error.to_string(),
727 }
728}
729
730fn schema_rollback_cleanup_failure(primary: &OrmError, cleanup: &OrmError) -> MigrationError {
731 MigrationError::State {
732 message: format!(
733 "schema bootstrap query failed and rollback was not acknowledged; primary: {primary}; cleanup: {cleanup}"
734 ),
735 }
736}
737
738fn extract_value(doc: &serde_json::Value, key: &str) -> Option<String> {
750 let value = doc.get(key)?;
751 extract_scalar(value)
752}
753
754fn extract_scalar(value: &serde_json::Value) -> Option<String> {
757 match value {
758 serde_json::Value::Null => None,
759 serde_json::Value::Object(map) => map.get("value").and_then(extract_scalar),
760 serde_json::Value::String(s) => Some(s.clone()),
761 serde_json::Value::Bool(b) => Some(b.to_string()),
762 serde_json::Value::Number(n) => Some(n.to_string()),
763 serde_json::Value::Array(_) => None,
764 }
765}
766
767pub fn parse_applied_documents(
776 values: &[serde_json::Value],
777) -> Result<Vec<AppliedMigrationRecord>> {
778 let mut records = Vec::with_capacity(values.len());
779 for doc in values {
780 let (Some(app_label), Some(name), Some(checksum)) = (
781 extract_value(doc, "app"),
782 extract_value(doc, "name"),
783 extract_value(doc, "checksum"),
784 ) else {
785 continue;
787 };
788 let applied_at = extract_value(doc, "applied");
789 records.push(AppliedMigrationRecord {
790 app_label,
791 name,
792 checksum,
793 applied_at,
794 });
795 }
796 Ok(records)
797}
798
799pub fn parse_run_documents(values: &[serde_json::Value]) -> Result<Vec<MigrationRunRecord>> {
801 let mut records = Vec::with_capacity(values.len());
802 for doc in values {
803 let (
804 Some(run_id),
805 Some(app_label),
806 Some(name),
807 Some(checksum),
808 Some(direction),
809 Some(status),
810 Some(started_at),
811 ) = (
812 extract_value(doc, "run_id"),
813 extract_value(doc, "app"),
814 extract_value(doc, "name"),
815 extract_value(doc, "checksum"),
816 extract_value(doc, "direction"),
817 extract_value(doc, "status"),
818 extract_value(doc, "started"),
819 )
820 else {
821 continue;
822 };
823 records.push(MigrationRunRecord {
824 run_id,
825 app_label,
826 name,
827 checksum,
828 direction,
829 status,
830 started_at,
831 finished_at: None,
832 error: None,
833 executor_ip: None,
834 executor_mac: None,
835 });
836 }
837 Ok(records)
838}
839
840fn optional_run_field_query(attribute: &str, alias: &str) -> String {
841 format!(
842 "\nmatch\n$r isa {RUN_ENTITY},\n has {RUN_ID} $run_id,\n has {attribute} ${alias};\nfetch {{\n \"run_id\": $run_id,\n \"{alias}\": ${alias}\n}};\n"
843 )
844}
845
846fn merge_optional_run_field(
847 runs: &mut [MigrationRunRecord],
848 docs: &[serde_json::Value],
849 target: &str,
850 source: &str,
851) {
852 for doc in docs {
853 let (Some(run_id), Some(value)) =
854 (extract_value(doc, "run_id"), extract_value(doc, source))
855 else {
856 continue;
857 };
858 let Some(run) = runs.iter_mut().find(|run| run.run_id == run_id) else {
859 continue;
860 };
861 match target {
862 "finished_at" => run.finished_at = Some(value),
863 "error" => run.error = Some(value),
864 "executor_ip" => run.executor_ip = Some(value),
865 "executor_mac" => run.executor_mac = Some(value),
866 _ => {}
867 }
868 }
869}
870
871fn run_insert_query(record: &MigrationRunRecord) -> String {
872 let mut fields = vec![
873 format!(" has {RUN_ID} {}", typeql_string_literal(&record.run_id)),
874 format!(
875 " has {APP_LABEL} {}",
876 typeql_string_literal(&record.app_label)
877 ),
878 format!(" has {NAME} {}", typeql_string_literal(&record.name)),
879 format!(
880 " has {CHECKSUM} {}",
881 typeql_string_literal(&record.checksum)
882 ),
883 format!(
884 " has {DIRECTION} {}",
885 typeql_string_literal(&record.direction)
886 ),
887 format!(" has {STATUS} {}", typeql_string_literal(&record.status)),
888 format!(" has {STARTED_AT} {}", record.started_at),
889 ];
890
891 if let Some(finished_at) = &record.finished_at {
892 fields.push(format!(" has {FINISHED_AT} {finished_at}"));
893 }
894 if let Some(error) = &record.error {
895 fields.push(format!(" has {ERROR} {}", typeql_string_literal(error)));
896 }
897 if let Some(executor_ip) = &record.executor_ip {
898 fields.push(format!(
899 " has {EXECUTOR_IP} {}",
900 typeql_string_literal(executor_ip)
901 ));
902 }
903 if let Some(executor_mac) = &record.executor_mac {
904 fields.push(format!(
905 " has {EXECUTOR_MAC} {}",
906 typeql_string_literal(executor_mac)
907 ));
908 }
909
910 format!("\ninsert $r isa {RUN_ENTITY},\n{};\n", fields.join(",\n"))
911}
912
913impl MigrationStateStore for TypeDbStateStore {
914 fn ensure_schema(&self) -> BoxFuture<'_, Result<()>> {
915 Box::pin(async move {
916 require_legacy_writer_open(self.db.as_ref()).await?;
920 if self.schema_ensured.load(Ordering::Acquire) {
921 return Ok(());
922 }
923
924 let state_schema = migration_state_schema();
925
926 for (name, attribute) in &state_schema.attributes {
927 let mut definition = SchemaInfo::default();
928 definition
929 .attributes
930 .insert(name.clone(), attribute.clone());
931 let define = definition.to_typeql().map_err(map_schema_error)?;
932 self.ensure_type(name, &define).await?;
933 }
934
935 for (name, entity) in &state_schema.entities {
936 let mut definition = SchemaInfo::default();
937 definition.entities.insert(name.clone(), entity.clone());
938 let define = definition.to_typeql().map_err(map_schema_error)?;
939 self.ensure_type(name, &define).await?;
940 }
941
942 for (name, relation) in &state_schema.relations {
943 let mut definition = SchemaInfo::default();
944 definition.relations.insert(name.clone(), relation.clone());
945 let define = definition.to_typeql().map_err(map_schema_error)?;
946 self.ensure_type(name, &define).await?;
947 }
948
949 self.schema_ensured.store(true, Ordering::Release);
950 Ok(())
951 })
952 }
953
954 fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>> {
955 Box::pin(async move {
956 self.ensure_schema_for_read().await?;
957
958 let query = applied_query();
960
961 let values = self.query_documents(&query).await?;
962 parse_applied_documents(&values)
963 })
964 }
965
966 fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>> {
967 Box::pin(async move {
968 self.ensure_schema_for_read().await?;
969 self.load_runs_documents().await
970 })
971 }
972
973 fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>> {
974 Box::pin(async move {
975 self.ensure_schema().await?;
976
977 let applied_at = record
980 .applied_at
981 .clone()
982 .unwrap_or_else(|| format_applied_at(Utc::now()));
983
984 let migration_id = format!("{}:{}", record.app_label, record.name);
985 let migration_id = typeql_string_literal(&migration_id);
986 let app = typeql_string_literal(&record.app_label);
987 let name = typeql_string_literal(&record.name);
988 let checksum = typeql_string_literal(&record.checksum);
989
990 let delete_existing = format!(
996 "\nmatch\n$m isa {APPLIED_ENTITY},\n has {MIGRATION_ID} {migration_id};\ndelete $m;\n"
997 );
998
999 let insert = format!(
1002 "\ninsert $m isa {APPLIED_ENTITY},\n has {MIGRATION_ID} {migration_id},\n has {APP_LABEL} {app},\n has {NAME} {name},\n has {APPLIED_AT} {applied_at},\n has {CHECKSUM} {checksum};\n",
1003 );
1004
1005 let ctx = self
1006 .db
1007 .transaction_context(TxType::Write)
1008 .await
1009 .map_err(map_orm_error)?;
1010 if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
1011 let _ = ctx.rollback().await;
1012 return Err(error);
1013 }
1014 ctx.query(&delete_existing).await.map_err(map_orm_error)?;
1015 ctx.query(&insert).await.map_err(map_orm_error)?;
1016 ctx.commit().await.map_err(map_orm_error)?;
1017 Ok(())
1018 })
1019 }
1020
1021 fn record_unapplied<'a>(
1022 &'a self,
1023 app_label: &'a str,
1024 name: &'a str,
1025 ) -> BoxFuture<'a, Result<()>> {
1026 Box::pin(async move {
1027 self.ensure_schema().await?;
1028
1029 let app_label = typeql_string_literal(app_label);
1033 let name = typeql_string_literal(name);
1034 let query = format!(
1035 "\nmatch\n$m isa {APPLIED_ENTITY},\n has {APP_LABEL} {app_label},\n has {NAME} {name};\ndelete $m;\n"
1036 );
1037
1038 let ctx = self
1039 .db
1040 .transaction_context(TxType::Write)
1041 .await
1042 .map_err(map_orm_error)?;
1043 if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
1044 let _ = ctx.rollback().await;
1045 return Err(error);
1046 }
1047 ctx.query(&query).await.map_err(map_orm_error)?;
1048 ctx.commit().await.map_err(map_orm_error)?;
1049 Ok(())
1050 })
1051 }
1052
1053 fn record_run(&self, record: MigrationRunRecord) -> BoxFuture<'_, Result<()>> {
1054 Box::pin(async move {
1055 self.ensure_schema().await?;
1056
1057 let run_id = typeql_string_literal(&record.run_id);
1058 let delete_existing =
1059 format!("\nmatch\n$r isa {RUN_ENTITY},\n has {RUN_ID} {run_id};\ndelete $r;\n");
1060 let insert = run_insert_query(&record);
1061
1062 let ctx = self
1063 .db
1064 .transaction_context(TxType::Write)
1065 .await
1066 .map_err(map_orm_error)?;
1067 if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
1068 let _ = ctx.rollback().await;
1069 return Err(error);
1070 }
1071 ctx.query(&delete_existing).await.map_err(map_orm_error)?;
1072 ctx.query(&insert).await.map_err(map_orm_error)?;
1073 ctx.commit().await.map_err(map_orm_error)?;
1074 Ok(())
1075 })
1076 }
1077}
1078
1079#[cfg(test)]
1080mod tests {
1081 use super::*;
1082 use crate::testing::{MockEvent, MockMigrationBackend};
1083 use chrono::{TimeZone, Timelike};
1084
1085 #[test]
1093 fn parse_applied_documents_extracts_bare_scalar_fields() {
1094 let docs = vec![serde_json::json!({
1095 "app": "myapp",
1096 "name": "0001_initial",
1097 "applied": "2026-06-05T00:00:00.000000000",
1098 "checksum": "abc123"
1099 })];
1100
1101 let records = parse_applied_documents(&docs).unwrap();
1102 assert_eq!(records.len(), 1);
1103 assert_eq!(records[0].app_label, "myapp");
1104 assert_eq!(records[0].name, "0001_initial");
1105 assert_eq!(records[0].checksum, "abc123");
1106 assert_eq!(
1107 records[0].applied_at.as_deref(),
1108 Some("2026-06-05T00:00:00.000000000")
1109 );
1110 }
1111
1112 #[tokio::test]
1113 async fn state_readers_close_every_read_context() {
1114 let responses = vec![
1115 QueryResult::Documents(vec![serde_json::json!({
1116 "app": "myapp",
1117 "name": "0001_initial",
1118 "applied": "2026-06-05T00:00:00.000000000",
1119 "checksum": "abc123"
1120 })]),
1121 QueryResult::Documents(Vec::new()),
1122 QueryResult::Documents(Vec::new()),
1123 QueryResult::Documents(Vec::new()),
1124 QueryResult::Documents(Vec::new()),
1125 QueryResult::Documents(Vec::new()),
1126 ];
1127 let (backend, log) = MockMigrationBackend::with_state_read_responses(responses);
1128 let store =
1129 TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1130
1131 assert_eq!(store.load_applied().await.unwrap().len(), 1);
1132 assert!(store.load_runs().await.unwrap().is_empty());
1133
1134 let events = log.lock().unwrap();
1135 let opens = events
1136 .iter()
1137 .filter(|event| matches!(event, MockEvent::OpenTx(TxType::Read)))
1138 .count();
1139 let closes = events
1140 .iter()
1141 .filter(|event| matches!(event, MockEvent::Close))
1142 .count();
1143 assert_eq!(opens, 7, "one schema inspection plus six ledger reads");
1144 assert_eq!(closes, opens, "every read context must be acknowledged");
1145 }
1146
1147 #[tokio::test]
1148 async fn load_applied_preserves_query_error_when_close_also_fails() {
1149 let (backend, log) = MockMigrationBackend::with_state_read_and_close_failure(0, 1);
1152 let store =
1153 TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1154
1155 let error = store
1156 .load_applied()
1157 .await
1158 .expect_err("the ledger query must fail");
1159 let message = error.to_string();
1160 assert!(message.contains("injected query failure for testing"));
1161 assert!(!message.contains("injected close failure for testing"));
1162 assert_eq!(
1163 log.lock()
1164 .unwrap()
1165 .iter()
1166 .filter(|event| matches!(event, MockEvent::Close))
1167 .count(),
1168 2,
1169 "the failed ledger read must still acknowledge close"
1170 );
1171 }
1172
1173 #[tokio::test]
1174 async fn load_runs_preserves_query_error_when_close_also_fails() {
1175 let (backend, log) = MockMigrationBackend::with_state_read_and_close_failure(0, 1);
1176 let store =
1177 TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1178
1179 let error = store
1180 .load_runs()
1181 .await
1182 .expect_err("the base run-log query must fail");
1183 let message = error.to_string();
1184 assert!(message.contains("injected query failure for testing"));
1185 assert!(!message.contains("injected close failure for testing"));
1186 assert_eq!(
1187 log.lock()
1188 .unwrap()
1189 .iter()
1190 .filter(|event| matches!(event, MockEvent::Close))
1191 .count(),
1192 2,
1193 "the failed run-log read must still acknowledge close"
1194 );
1195 }
1196
1197 #[tokio::test]
1198 async fn unadopted_partial_state_schema_is_repaired_before_read() {
1199 let (backend, log, labels) =
1200 MockMigrationBackend::with_partial_state_schema(&[RUN_ENTITY], false);
1201 let store =
1202 TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1203
1204 assert!(store.load_applied().await.unwrap().is_empty());
1205 assert!(labels.lock().unwrap().contains(RUN_ENTITY));
1206 assert!(log.lock().unwrap().iter().any(|event| {
1207 matches!(event, MockEvent::Query(TxType::Schema, query) if query.contains(&format!("entity {RUN_ENTITY}")))
1208 }));
1209 }
1210
1211 #[tokio::test]
1212 async fn adopted_partial_state_schema_fails_before_repair() {
1213 let (backend, log, labels) =
1214 MockMigrationBackend::with_partial_state_schema(&[RUN_ENTITY], true);
1215 let store =
1216 TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1217
1218 let error = store
1219 .load_applied()
1220 .await
1221 .expect_err("the sentinel must block partial-schema repair");
1222 assert!(error.to_string().contains(LEGACY_WRITER_CUTOVER_MESSAGE));
1223 assert!(!labels.lock().unwrap().contains(RUN_ENTITY));
1224 assert!(
1225 !log.lock()
1226 .unwrap()
1227 .iter()
1228 .any(|event| matches!(event, MockEvent::Query(TxType::Schema, _))),
1229 "adopted schema repair must not reach a define query"
1230 );
1231 }
1232
1233 #[tokio::test]
1234 async fn exact_sentinel_partition_is_verified_before_filtering() {
1235 let fingerprint = "a".repeat(64);
1236 let sentinel = serde_json::json!({
1237 "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1238 "name": LEGACY_CUTOVER_SENTINEL_NAME,
1239 "applied": LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
1240 "checksum": fingerprint,
1241 });
1242 let responses = vec![
1243 QueryResult::Documents(vec![sentinel]),
1244 QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1245 QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1246 QueryResult::Documents(vec![serde_json::json!({
1247 "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1248 "applied": LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
1249 "checksum": fingerprint,
1250 })]),
1251 ];
1252 let (backend, _) = MockMigrationBackend::with_state_read_responses(responses);
1253 let database = Database::with_backend(Box::new(backend), "test");
1254 let mut transaction = database.read_transaction().await.unwrap();
1255
1256 let partition = TypeDbStateStore::load_verified_legacy_partition_in_transaction(
1257 &mut transaction,
1258 LegacyCutoverSentinelExpectation::RequiredExact(&fingerprint),
1259 )
1260 .await
1261 .expect("verify exact sentinel");
1262 transaction.close().await.unwrap();
1263
1264 assert!(partition.applied().is_empty());
1265 assert_eq!(partition.sentinel_fingerprint(), Some(fingerprint.as_str()));
1266 }
1267
1268 #[tokio::test]
1269 async fn malformed_sentinel_timestamp_is_not_filtered() {
1270 let fingerprint = "b".repeat(64);
1271 let malformed_timestamp = "1970-01-01T00:00:00";
1272 let responses = vec![
1273 QueryResult::Documents(vec![serde_json::json!({
1274 "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1275 "name": LEGACY_CUTOVER_SENTINEL_NAME,
1276 "applied": malformed_timestamp,
1277 "checksum": fingerprint,
1278 })]),
1279 QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1280 QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1281 QueryResult::Documents(vec![serde_json::json!({
1282 "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1283 "applied": malformed_timestamp,
1284 "checksum": fingerprint,
1285 })]),
1286 ];
1287 let (backend, _) = MockMigrationBackend::with_state_read_responses(responses);
1288 let database = Database::with_backend(Box::new(backend), "test");
1289 let mut transaction = database.read_transaction().await.unwrap();
1290
1291 let error = TypeDbStateStore::load_verified_legacy_partition_in_transaction(
1292 &mut transaction,
1293 LegacyCutoverSentinelExpectation::RequiredExact(&fingerprint),
1294 )
1295 .await
1296 .expect_err("malformed sentinel must fail closed");
1297 transaction.close().await.unwrap();
1298
1299 assert!(error.is_contract_violation());
1300 assert!(error.to_string().contains("foreign applied timestamp"));
1301 }
1302
1303 #[test]
1304 fn sentinel_name_is_outside_the_released_numbered_loader_namespace() {
1305 assert!(
1306 !LEGACY_CUTOVER_SENTINEL_NAME
1307 .as_bytes()
1308 .first()
1309 .is_some_and(u8::is_ascii_digit)
1310 );
1311 }
1312
1313 #[test]
1314 fn parse_applied_documents_also_unwraps_value_envelope() {
1315 let docs = vec![serde_json::json!({
1318 "app": {"value": "myapp"},
1319 "name": {"value": "0002_next"},
1320 "applied": {"value": "2026-06-05T01:02:03.000000000"},
1321 "checksum": {"value": "def456"}
1322 })];
1323
1324 let records = parse_applied_documents(&docs).unwrap();
1325 assert_eq!(records.len(), 1);
1326 assert_eq!(records[0].app_label, "myapp");
1327 assert_eq!(records[0].name, "0002_next");
1328 assert_eq!(records[0].checksum, "def456");
1329 assert_eq!(
1330 records[0].applied_at.as_deref(),
1331 Some("2026-06-05T01:02:03.000000000")
1332 );
1333 }
1334
1335 #[test]
1336 fn parse_applied_documents_empty_list_is_empty() {
1337 let records = parse_applied_documents(&[]).unwrap();
1338 assert!(records.is_empty());
1339 }
1340
1341 #[test]
1342 fn parse_applied_documents_skips_incomplete_rows() {
1343 let docs = vec![
1345 serde_json::json!({
1346 "app": "myapp",
1347 "name": "0001_initial",
1348 "applied": "2026-06-05T00:00:00.000000000"
1349 }),
1350 serde_json::json!({
1351 "app": "myapp",
1352 "name": "0002_next",
1353 "applied": "2026-06-05T00:00:00.000000000",
1354 "checksum": "ok"
1355 }),
1356 ];
1357
1358 let records = parse_applied_documents(&docs).unwrap();
1359 assert_eq!(records.len(), 1);
1360 assert_eq!(records[0].name, "0002_next");
1361 }
1362
1363 #[test]
1364 fn parse_applied_documents_carries_missing_applied_as_none() {
1365 let docs = vec![serde_json::json!({
1366 "app": "myapp",
1367 "name": "0001_initial",
1368 "checksum": "abc123"
1369 })];
1370
1371 let records = parse_applied_documents(&docs).unwrap();
1372 assert_eq!(records.len(), 1);
1373 assert!(records[0].applied_at.is_none());
1374 }
1375
1376 #[test]
1377 fn parse_run_documents_extracts_required_fields() {
1378 let docs = vec![serde_json::json!({
1379 "run_id": "run-1",
1380 "app": "app",
1381 "name": "0001_initial",
1382 "checksum": "abc123",
1383 "direction": "apply",
1384 "status": "started",
1385 "started": "2026-06-05T00:00:00.000000"
1386 })];
1387
1388 let records = parse_run_documents(&docs).unwrap();
1389
1390 assert_eq!(records.len(), 1);
1391 assert_eq!(records[0].run_id, "run-1");
1392 assert_eq!(records[0].direction, "apply");
1393 assert_eq!(records[0].status, "started");
1394 assert_eq!(records[0].finished_at, None);
1395 }
1396
1397 #[test]
1398 fn merge_optional_run_field_updates_matching_record_only() {
1399 let mut records = vec![MigrationRunRecord {
1400 run_id: "run-1".to_string(),
1401 app_label: "app".to_string(),
1402 name: "0001_initial".to_string(),
1403 checksum: "abc123".to_string(),
1404 direction: "apply".to_string(),
1405 status: "started".to_string(),
1406 started_at: "2026-06-05T00:00:00.000000".to_string(),
1407 finished_at: None,
1408 error: None,
1409 executor_ip: None,
1410 executor_mac: None,
1411 }];
1412 let docs = vec![serde_json::json!({
1413 "run_id": "run-1",
1414 "finished": "2026-06-05T00:00:01.000000"
1415 })];
1416
1417 merge_optional_run_field(&mut records, &docs, "finished_at", "finished");
1418
1419 assert_eq!(
1420 records[0].finished_at.as_deref(),
1421 Some("2026-06-05T00:00:01.000000")
1422 );
1423 }
1424
1425 #[test]
1426 fn typeql_string_literal_escapes_user_controlled_text() {
1427 assert_eq!(typeql_string_literal("a\"b\\c\n"), "\"a\\\"b\\\\c\\n\"");
1428 }
1429
1430 #[test]
1431 fn schema_bootstrap_rollback_failure_preserves_primary_and_cleanup() {
1432 let primary = OrmError::QueryExecution("define failed".to_owned());
1433 let cleanup = OrmError::Transaction("rollback failed".to_owned());
1434 let error = schema_rollback_cleanup_failure(&primary, &cleanup).to_string();
1435 assert!(error.contains("define failed"), "{error}");
1436 assert!(error.contains("rollback failed"), "{error}");
1437 assert!(error.contains("rollback was not acknowledged"), "{error}");
1438 }
1439
1440 #[test]
1441 fn run_insert_query_includes_optional_fields_when_present() {
1442 let record = MigrationRunRecord {
1443 run_id: "run-1".to_string(),
1444 app_label: "app".to_string(),
1445 name: "0001_initial".to_string(),
1446 checksum: "abc123".to_string(),
1447 direction: "apply".to_string(),
1448 status: "failed".to_string(),
1449 started_at: "2026-06-05T00:00:00.000000".to_string(),
1450 finished_at: Some("2026-06-05T00:00:01.000000".to_string()),
1451 error: Some("quote: \"boom\"".to_string()),
1452 executor_ip: Some("127.0.0.1".to_string()),
1453 executor_mac: Some("00:11:22:33:44:55".to_string()),
1454 };
1455
1456 let query = run_insert_query(&record);
1457
1458 assert!(query.contains("has migration_run_id \"run-1\""));
1459 assert!(query.contains("has migration_finished_at 2026-06-05T00:00:01.000000"));
1460 assert!(query.contains("has migration_error \"quote: \\\"boom\\\"\""));
1461 assert!(query.contains("has migration_executor_ip \"127.0.0.1\""));
1462 assert!(query.contains("has migration_executor_mac \"00:11:22:33:44:55\""));
1463 }
1464
1465 #[test]
1472 fn format_applied_at_matches_python_strftime() {
1473 let dt = Utc
1475 .with_ymd_and_hms(2026, 6, 5, 14, 9, 8)
1476 .unwrap()
1477 .with_nanosecond(123_456_000)
1478 .unwrap();
1479 assert_eq!(format_applied_at(dt), "2026-06-05T14:09:08.123456");
1480 }
1481
1482 #[test]
1483 fn format_applied_at_zero_pads_microseconds() {
1484 let dt = Utc
1487 .with_ymd_and_hms(2026, 1, 2, 3, 4, 5)
1488 .unwrap()
1489 .with_nanosecond(7_000)
1490 .unwrap();
1491 assert_eq!(format_applied_at(dt), "2026-01-02T03:04:05.000007");
1493 }
1494}