Skip to main content

type_bridge_migration/state/
typedb.rs

1//! TypeDB-backed [`MigrationStateStore`] implementation.
2//!
3//! [`TypeDbStateStore`] persists the established migration projection and run
4//! log over the ORM [`Database`][type_bridge_orm::Database] seam. Its per-type
5//! bootstrap definitions are rendered from the canonical migration-state
6//! [`SchemaInfo`][type_bridge_orm::schema::SchemaInfo], and its row queries use
7//! the same semantic label constants. Existing labels, value types, keys, and
8//! storage behavior remain unchanged.
9//!
10//! All work runs through [`Database::transaction_context`], mirroring the
11//! Phase-05 executor: open a context for the right [`TxType`], run the query,
12//! and commit write/schema contexts. No Python transaction crosses this
13//! boundary (invariant 4); only the public ORM `session` API is used
14//! (invariant 7).
15
16use 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::OrmError;
27use type_bridge_orm::schema::{SchemaError, SchemaInfo};
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
43/// Timestamp format mirroring Python's `strftime("%Y-%m-%dT%H:%M:%S.%f")`.
44///
45/// Python `%f` always emits exactly six digits (microseconds, zero-padded);
46/// chrono's `%6f` is the byte-identical equivalent. This is the format written
47/// into the `insert` TypeQL — `record_applied` (`state.py:250`).
48const 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/// Required state of the V2-only legacy-ledger sentinel while a caller reads a
54/// legacy applied projection.
55#[derive(Debug, Clone, Copy)]
56pub enum LegacyCutoverSentinelExpectation<'a> {
57    /// No cutover sentinel may exist.
58    Absent,
59    /// The sentinel may be absent, or must exactly carry this fingerprint.
60    OptionalExact(&'a str),
61    /// One exact sentinel carrying this fingerprint must exist.
62    RequiredExact(&'a str),
63}
64
65/// A verified legacy applied projection with the V2-only sentinel removed.
66///
67/// Construction is possible only after the sentinel candidate probes and all
68/// stored sentinel fields have passed the requested expectation.  Callers must
69/// never filter the reserved identity directly from an unverified row list.
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct VerifiedLegacyAppliedPartition {
72    applied: Vec<AppliedMigrationRecord>,
73    sentinel_fingerprint: Option<String>,
74}
75
76impl VerifiedLegacyAppliedPartition {
77    /// Borrow the released user migration records, excluding the verified
78    /// V2-only sentinel partition.
79    pub fn applied(&self) -> &[AppliedMigrationRecord] {
80        &self.applied
81    }
82
83    /// Consume the partition and return released user migration records.
84    pub fn into_applied(self) -> Vec<AppliedMigrationRecord> {
85        self.applied
86    }
87
88    /// Return the exact sentinel fingerprint when the verified pair is present.
89    pub fn sentinel_fingerprint(&self) -> Option<&str> {
90        self.sentinel_fingerprint.as_deref()
91    }
92}
93
94/// Failure to observe and validate the reserved V2 cutover sentinel.
95#[derive(Debug, thiserror::Error)]
96pub enum LegacyCutoverSentinelError {
97    /// The ledger could not be queried through the retained transaction.
98    #[error("legacy cutover sentinel storage inspection failed: {0}")]
99    Storage(#[source] MigrationError),
100    /// Stored rows violate the exact singleton contract.
101    #[error("legacy cutover sentinel contract violation: {message}")]
102    Contract {
103        /// Stable human-readable contract failure.
104        message: String,
105    },
106}
107
108impl LegacyCutoverSentinelError {
109    /// Return whether this is durable stored drift rather than provider
110    /// infrastructure failure.
111    pub fn is_contract_violation(&self) -> bool {
112        matches!(self, Self::Contract { .. })
113    }
114}
115
116/// TypeDB-backed migration state store over the ORM session seam.
117///
118/// Holds a shared [`Arc<Database>`] (the same handle the rest of the Rust ORM
119/// path uses) and opens its own transaction contexts per operation.
120pub struct TypeDbStateStore {
121    db: Arc<Database>,
122    /// Idempotency latch for [`ensure_schema`](Self::ensure_schema), mirroring
123    /// the Python `_schema_ensured` flag (`state.py:121`).
124    schema_ensured: AtomicBool,
125}
126
127impl TypeDbStateStore {
128    /// Construct a store bound to a shared database handle.
129    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        // An adopted database already has the complete frozen legacy state
202        // schema.  Archival reads must remain available there and must expose
203        // the sentinel to ordinary legacy planning; only a genuinely missing
204        // type enters the writer-guarded bootstrap path.
205        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        // Released readers repaired interrupted incremental bootstraps. Keep
222        // that behavior for absent and partial unadopted schemas by entering
223        // the ordinary writer-guarded bootstrap. An adopted target is rejected
224        // by the sentinel before any repair mutation.
225        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    /// Read the complete released applied ledger through an already-retained
249    /// transaction.
250    ///
251    /// Legacy-frontier cutover uses a managed schema transaction as an
252    /// exclusive guard against V1 write transactions. The caller must ensure
253    /// the released state schema already exists before opening that guard;
254    /// this method performs no bootstrap and never commits or closes it.
255    pub async fn load_applied_in_transaction(
256        transaction: &mut Transaction,
257    ) -> Result<Vec<AppliedMigrationRecord>> {
258        let result = transaction
259            .query(&applied_query())
260            .await
261            .map_err(map_orm_error)?;
262        parse_applied_documents(&query_result_values(result))
263    }
264
265    /// Read the released applied projection and validate the reserved V2
266    /// sentinel through the same retained transaction snapshot.
267    ///
268    /// The sentinel is removed only after both independent identity probes,
269    /// singleton cardinality, every stored field, and the expected anchor
270    /// fingerprint have been checked.  This is the only supported filtering
271    /// boundary for V2 legacy-frontier continuity and digest calculations.
272    pub async fn load_verified_legacy_partition_in_transaction(
273        transaction: &mut Transaction,
274        expectation: LegacyCutoverSentinelExpectation<'_>,
275    ) -> std::result::Result<VerifiedLegacyAppliedPartition, LegacyCutoverSentinelError> {
276        match legacy_state_schema_presence(transaction).await? {
277            LegacyStateSchemaPresence::Absent => {
278                if matches!(
279                    expectation,
280                    LegacyCutoverSentinelExpectation::RequiredExact(_)
281                ) {
282                    return Err(sentinel_contract_error(
283                        "the V2 bridge is active but the frozen legacy ledger schema is absent",
284                    ));
285                }
286                return Ok(VerifiedLegacyAppliedPartition {
287                    applied: Vec::new(),
288                    sentinel_fingerprint: None,
289                });
290            }
291            LegacyStateSchemaPresence::Partial => {
292                return Err(sentinel_contract_error(
293                    "the frozen legacy ledger schema is partially present",
294                ));
295            }
296            LegacyStateSchemaPresence::Complete => {}
297        }
298        let result = transaction
299            .query(&applied_query())
300            .await
301            .map_err(map_sentinel_storage_error)?;
302        let mut applied = parse_applied_documents(&query_result_values(result))
303            .map_err(LegacyCutoverSentinelError::Storage)?;
304
305        let id_candidates = sentinel_query_values(
306            transaction,
307            &format!(
308                "match $m isa {APPLIED_ENTITY}, has {MIGRATION_ID} {}; fetch {{ \"exists\": true }};",
309                typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
310            ),
311        )
312        .await?;
313        let name_candidates = sentinel_query_values(
314            transaction,
315            &format!(
316                "match $m isa {APPLIED_ENTITY}, has {NAME} {}; fetch {{ \"exists\": true }};",
317                typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
318            ),
319        )
320        .await?;
321
322        if id_candidates.is_empty() && name_candidates.is_empty() {
323            if matches!(
324                expectation,
325                LegacyCutoverSentinelExpectation::RequiredExact(_)
326            ) {
327                return Err(sentinel_contract_error(
328                    "the V2 bridge is active but its legacy-writer sentinel is missing",
329                ));
330            }
331            return Ok(VerifiedLegacyAppliedPartition {
332                applied,
333                sentinel_fingerprint: None,
334            });
335        }
336
337        if matches!(expectation, LegacyCutoverSentinelExpectation::Absent) {
338            return Err(sentinel_contract_error(
339                "a legacy-writer sentinel exists without an active or pending V2 bridge",
340            ));
341        }
342        if id_candidates.len() != 1 || name_candidates.len() != 1 {
343            return Err(sentinel_contract_error(
344                "the legacy-writer sentinel is duplicated or has split identity rows",
345            ));
346        }
347
348        let details = sentinel_query_values(
349            transaction,
350            &format!(
351                "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 }};",
352                typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
353                typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
354            ),
355        )
356        .await?;
357        if details.len() != 1 {
358            return Err(sentinel_contract_error(
359                "the legacy-writer sentinel is missing required exact fields",
360            ));
361        }
362        let detail = &details[0];
363        let app = extract_value(detail, "app").ok_or_else(|| {
364            sentinel_contract_error("the legacy-writer sentinel app label is malformed")
365        })?;
366        let applied_at = extract_value(detail, "applied").ok_or_else(|| {
367            sentinel_contract_error("the legacy-writer sentinel applied timestamp is malformed")
368        })?;
369        let fingerprint = extract_value(detail, "checksum").ok_or_else(|| {
370            sentinel_contract_error("the legacy-writer sentinel checksum is malformed")
371        })?;
372        if app != LEGACY_CUTOVER_SENTINEL_APP_LABEL {
373            return Err(sentinel_contract_error(
374                "the legacy-writer sentinel carries a foreign application label",
375            ));
376        }
377        if applied_at != LEGACY_CUTOVER_SENTINEL_APPLIED_AT {
378            return Err(sentinel_contract_error(
379                "the legacy-writer sentinel carries a foreign applied timestamp",
380            ));
381        }
382        if !is_lower_hex_fingerprint(&fingerprint) {
383            return Err(sentinel_contract_error(
384                "the legacy-writer sentinel checksum is not a lowercase 64-hex fingerprint",
385            ));
386        }
387        let expected = match expectation {
388            LegacyCutoverSentinelExpectation::OptionalExact(expected)
389            | LegacyCutoverSentinelExpectation::RequiredExact(expected) => expected,
390            LegacyCutoverSentinelExpectation::Absent => unreachable!("handled above"),
391        };
392        if fingerprint != expected {
393            return Err(sentinel_contract_error(
394                "the legacy-writer sentinel checksum differs from the managed cutover anchor",
395            ));
396        }
397
398        let original_len = applied.len();
399        applied.retain(|record| {
400            record.app_label != LEGACY_CUTOVER_SENTINEL_APP_LABEL
401                || record.name != LEGACY_CUTOVER_SENTINEL_NAME
402        });
403        if original_len.saturating_sub(applied.len()) != 1 {
404            return Err(sentinel_contract_error(
405                "the exact legacy-writer sentinel is absent from the released applied projection",
406            ));
407        }
408
409        Ok(VerifiedLegacyAppliedPartition {
410            applied,
411            sentinel_fingerprint: Some(fingerprint),
412        })
413    }
414
415    /// Stage the complete V2 cutover sentinel in a caller-retained managed
416    /// transaction.  The caller commits this in the same transaction as the
417    /// managed cutover anchor.
418    pub async fn insert_legacy_cutover_sentinel_in_transaction(
419        transaction: &mut Transaction,
420        anchor_fingerprint: &str,
421    ) -> Result<()> {
422        if !is_lower_hex_fingerprint(anchor_fingerprint) {
423            return Err(MigrationError::State {
424                message: "legacy cutover sentinel requires a lowercase 64-hex anchor fingerprint"
425                    .to_owned(),
426            });
427        }
428        let query = format!(
429            "insert $m isa {APPLIED_ENTITY}, has {MIGRATION_ID} {}, has {APP_LABEL} {}, has {NAME} {}, has {APPLIED_AT} {LEGACY_CUTOVER_SENTINEL_APPLIED_AT}, has {CHECKSUM} {};",
430            typeql_string_literal(LEGACY_CUTOVER_SENTINEL_MIGRATION_ID),
431            typeql_string_literal(LEGACY_CUTOVER_SENTINEL_APP_LABEL),
432            typeql_string_literal(LEGACY_CUTOVER_SENTINEL_NAME),
433            typeql_string_literal(anchor_fingerprint),
434        );
435        transaction.query(&query).await.map_err(map_orm_error)?;
436        Ok(())
437    }
438}
439
440/// Fail before a legacy writer uses an already-open transaction when an exact,
441/// managed-anchor-bound V2 cutover pair is present.
442pub async fn require_legacy_writer_open_in_transaction(
443    transaction: &TransactionContext,
444) -> Result<()> {
445    require_orm_legacy_writer_open_in_transaction(transaction)
446        .await
447        .map_err(map_legacy_guard_error)
448}
449
450#[cfg(test)]
451pub(crate) fn is_legacy_state_schema_probe_query(query: &str) -> bool {
452    query.starts_with(LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG)
453}
454
455/// Read-only entry guard for legacy writer surfaces whose external side
456/// effects cannot share a TypeDB transaction.
457pub async fn require_legacy_writer_open(database: &Database) -> Result<()> {
458    require_orm_legacy_writer_open(database)
459        .await
460        .map_err(map_legacy_guard_error)
461}
462
463#[derive(Debug, Clone, Copy, PartialEq, Eq)]
464enum LegacyStateSchemaPresence {
465    Absent,
466    Partial,
467    Complete,
468}
469
470async fn legacy_state_schema_presence(
471    transaction: &mut Transaction,
472) -> std::result::Result<LegacyStateSchemaPresence, LegacyCutoverSentinelError> {
473    let state_schema = migration_state_schema();
474    let mut present = 0_usize;
475    let expected_by_root = [
476        (
477            "attribute",
478            state_schema
479                .attributes
480                .keys()
481                .map(String::as_str)
482                .collect::<Vec<_>>(),
483        ),
484        (
485            "entity",
486            state_schema
487                .entities
488                .keys()
489                .map(String::as_str)
490                .collect::<Vec<_>>(),
491        ),
492        (
493            "relation",
494            state_schema
495                .relations
496                .keys()
497                .map(String::as_str)
498                .collect::<Vec<_>>(),
499        ),
500    ];
501    let total = expected_by_root
502        .iter()
503        .map(|(_, labels)| labels.len())
504        .sum::<usize>();
505    let all_expected = expected_by_root
506        .iter()
507        .flat_map(|(_, labels)| labels.iter().copied())
508        .collect::<BTreeSet<_>>();
509    let mut expected_labels_seen_in_any_kind = 0_usize;
510    for (kind, expected) in expected_by_root {
511        let result = transaction
512            .query(&format!(
513                "{LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG}match {kind} $t; fetch {{ \"label\": label($t) }};"
514            ))
515            .await
516            .map_err(map_sentinel_storage_error)?;
517        let observed = schema_labels_from_result(result, "legacy state schema probe")
518            .map_err(LegacyCutoverSentinelError::Storage)?;
519        expected_labels_seen_in_any_kind += observed
520            .iter()
521            .filter(|label| all_expected.contains(label.as_str()))
522            .count();
523        present += expected
524            .into_iter()
525            .filter(|label| observed.contains(*label))
526            .count();
527    }
528    if present == 0 && expected_labels_seen_in_any_kind == 0 {
529        return Ok(LegacyStateSchemaPresence::Absent);
530    }
531    if present != total || expected_labels_seen_in_any_kind != total {
532        return Ok(LegacyStateSchemaPresence::Partial);
533    }
534    Ok(LegacyStateSchemaPresence::Complete)
535}
536
537fn state_type_kind(type_name: &str) -> Option<&'static str> {
538    let schema = migration_state_schema();
539    if schema.attributes.contains_key(type_name) {
540        Some("attribute")
541    } else if schema.entities.contains_key(type_name) {
542        Some("entity")
543    } else if schema.relations.contains_key(type_name) {
544        Some("relation")
545    } else {
546        None
547    }
548}
549
550fn schema_labels_from_result(result: QueryResult, operation: &str) -> Result<BTreeSet<String>> {
551    let values = match result {
552        QueryResult::Documents(values) | QueryResult::Rows(values) => values,
553        QueryResult::Ok => {
554            return Err(MigrationError::State {
555                message: format!("{operation} returned no document result"),
556            });
557        }
558    };
559    let mut labels = BTreeSet::new();
560    for value in &values {
561        let label = extract_value(value, "label").ok_or_else(|| MigrationError::State {
562            message: format!("{operation} returned a malformed schema label"),
563        })?;
564        labels.insert(label);
565    }
566    Ok(labels)
567}
568
569fn legacy_sentinel_error_into_migration_error(error: LegacyCutoverSentinelError) -> MigrationError {
570    match error {
571        LegacyCutoverSentinelError::Storage(error) => error,
572        LegacyCutoverSentinelError::Contract { message } => MigrationError::State { message },
573    }
574}
575
576fn applied_query() -> String {
577    format!(
578        "\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"
579    )
580}
581
582async fn sentinel_query_values(
583    transaction: &mut Transaction,
584    query: &str,
585) -> std::result::Result<Vec<serde_json::Value>, LegacyCutoverSentinelError> {
586    let result = transaction
587        .query(query)
588        .await
589        .map_err(map_sentinel_storage_error)?;
590    match result {
591        QueryResult::Documents(values) | QueryResult::Rows(values) => Ok(values),
592        QueryResult::Ok => Err(LegacyCutoverSentinelError::Storage(MigrationError::State {
593            message: "legacy cutover sentinel fetch returned no document result".to_owned(),
594        })),
595    }
596}
597
598fn map_sentinel_storage_error(error: OrmError) -> LegacyCutoverSentinelError {
599    LegacyCutoverSentinelError::Storage(map_orm_error(error))
600}
601
602fn sentinel_contract_error(message: impl Into<String>) -> LegacyCutoverSentinelError {
603    LegacyCutoverSentinelError::Contract {
604        message: message.into(),
605    }
606}
607
608fn is_lower_hex_fingerprint(value: &str) -> bool {
609    value.len() == 64
610        && value
611            .bytes()
612            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
613}
614
615/// Format the current UTC time as the Python-compatible applied-at string.
616///
617/// Isolated as a pure function (parameterised on the instant) so the
618/// `%Y-%m-%dT%H:%M:%S.%6f` parity with Python's `%f` is unit-testable without a
619/// live clock or TypeDB.
620fn format_applied_at(now: chrono::DateTime<Utc>) -> String {
621    now.format(APPLIED_AT_FORMAT).to_string()
622}
623
624fn query_result_values(result: QueryResult) -> Vec<serde_json::Value> {
625    match result {
626        QueryResult::Documents(docs) => docs,
627        QueryResult::Rows(rows) => rows,
628        QueryResult::Ok => Vec::new(),
629    }
630}
631
632fn typeql_string_literal(value: &str) -> String {
633    let escaped = value
634        .replace('\\', "\\\\")
635        .replace('"', "\\\"")
636        .replace('\n', "\\n")
637        .replace('\r', "\\r")
638        .replace('\t', "\\t");
639    format!("\"{escaped}\"")
640}
641
642/// Map an ORM-layer error into the migration error hierarchy.
643fn map_orm_error(error: OrmError) -> MigrationError {
644    MigrationError::State {
645        message: error.to_string(),
646    }
647}
648
649fn map_legacy_guard_error(error: OrmError) -> MigrationError {
650    match error {
651        OrmError::Transaction(message) if message == LEGACY_WRITER_CUTOVER_MESSAGE => {
652            MigrationError::State { message }
653        }
654        error => map_orm_error(error),
655    }
656}
657
658fn map_schema_error(error: SchemaError) -> MigrationError {
659    MigrationError::State {
660        message: error.to_string(),
661    }
662}
663
664fn schema_rollback_cleanup_failure(primary: &OrmError, cleanup: &OrmError) -> MigrationError {
665    MigrationError::State {
666        message: format!(
667            "schema bootstrap query failed and rollback was not acknowledged; primary: {primary}; cleanup: {cleanup}"
668        ),
669    }
670}
671
672/// Unwrap a single fetched field into its scalar string.
673///
674/// The Phase-05 ORM backend renders a `fetch { "k": $attr }` document so that
675/// each attribute-bound variable becomes a bare JSON scalar (the attribute's
676/// value), not a `{"value": ...}` wrapper (typedb-driver `json_value`,
677/// `concept_document.rs`). The original Python `_extract_value`
678/// (`state.py:221-237`) defensively handled BOTH shapes — a bare scalar and a
679/// `{"value": ...}` object — so this parser does the same: prefer the inner
680/// `"value"` when the field is an object carrying one, otherwise take the
681/// scalar directly. Returns `None` for a missing/null field or an object with
682/// no usable `"value"`.
683fn extract_value(doc: &serde_json::Value, key: &str) -> Option<String> {
684    let value = doc.get(key)?;
685    extract_scalar(value)
686}
687
688/// Reduce a fetched JSON value to its string form, unwrapping a `{"value":
689/// ...}` envelope when present (the `_extract_value` dict branch).
690fn extract_scalar(value: &serde_json::Value) -> Option<String> {
691    match value {
692        serde_json::Value::Null => None,
693        serde_json::Value::Object(map) => map.get("value").and_then(extract_scalar),
694        serde_json::Value::String(s) => Some(s.clone()),
695        serde_json::Value::Bool(b) => Some(b.to_string()),
696        serde_json::Value::Number(n) => Some(n.to_string()),
697        serde_json::Value::Array(_) => None,
698    }
699}
700
701/// Parse a list of TypeDB fetch documents into applied migration records.
702///
703/// Pure and TypeDB-free so the document-shape contract is unit-testable. Each
704/// document is the `fetch { "app": ..., "name": ..., "applied": ...,
705/// "checksum": ... }` shape from the `load_applied` query (ported from
706/// `state.py:179-191`). A document missing any of `app` / `name` / `checksum`
707/// is skipped, matching the Python `if all([...])` guard (`state.py:206`);
708/// `applied` is carried through as the optional ISO string.
709pub fn parse_applied_documents(
710    values: &[serde_json::Value],
711) -> Result<Vec<AppliedMigrationRecord>> {
712    let mut records = Vec::with_capacity(values.len());
713    for doc in values {
714        let (Some(app_label), Some(name), Some(checksum)) = (
715            extract_value(doc, "app"),
716            extract_value(doc, "name"),
717            extract_value(doc, "checksum"),
718        ) else {
719            // Skip incomplete rows (mirrors the Python `all([...])` guard).
720            continue;
721        };
722        let applied_at = extract_value(doc, "applied");
723        records.push(AppliedMigrationRecord {
724            app_label,
725            name,
726            checksum,
727            applied_at,
728        });
729    }
730    Ok(records)
731}
732
733/// Parse TypeDB fetch documents into migration run-log records.
734pub fn parse_run_documents(values: &[serde_json::Value]) -> Result<Vec<MigrationRunRecord>> {
735    let mut records = Vec::with_capacity(values.len());
736    for doc in values {
737        let (
738            Some(run_id),
739            Some(app_label),
740            Some(name),
741            Some(checksum),
742            Some(direction),
743            Some(status),
744            Some(started_at),
745        ) = (
746            extract_value(doc, "run_id"),
747            extract_value(doc, "app"),
748            extract_value(doc, "name"),
749            extract_value(doc, "checksum"),
750            extract_value(doc, "direction"),
751            extract_value(doc, "status"),
752            extract_value(doc, "started"),
753        )
754        else {
755            continue;
756        };
757        records.push(MigrationRunRecord {
758            run_id,
759            app_label,
760            name,
761            checksum,
762            direction,
763            status,
764            started_at,
765            finished_at: None,
766            error: None,
767            executor_ip: None,
768            executor_mac: None,
769        });
770    }
771    Ok(records)
772}
773
774fn optional_run_field_query(attribute: &str, alias: &str) -> String {
775    format!(
776        "\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"
777    )
778}
779
780fn merge_optional_run_field(
781    runs: &mut [MigrationRunRecord],
782    docs: &[serde_json::Value],
783    target: &str,
784    source: &str,
785) {
786    for doc in docs {
787        let (Some(run_id), Some(value)) =
788            (extract_value(doc, "run_id"), extract_value(doc, source))
789        else {
790            continue;
791        };
792        let Some(run) = runs.iter_mut().find(|run| run.run_id == run_id) else {
793            continue;
794        };
795        match target {
796            "finished_at" => run.finished_at = Some(value),
797            "error" => run.error = Some(value),
798            "executor_ip" => run.executor_ip = Some(value),
799            "executor_mac" => run.executor_mac = Some(value),
800            _ => {}
801        }
802    }
803}
804
805fn run_insert_query(record: &MigrationRunRecord) -> String {
806    let mut fields = vec![
807        format!("    has {RUN_ID} {}", typeql_string_literal(&record.run_id)),
808        format!(
809            "    has {APP_LABEL} {}",
810            typeql_string_literal(&record.app_label)
811        ),
812        format!("    has {NAME} {}", typeql_string_literal(&record.name)),
813        format!(
814            "    has {CHECKSUM} {}",
815            typeql_string_literal(&record.checksum)
816        ),
817        format!(
818            "    has {DIRECTION} {}",
819            typeql_string_literal(&record.direction)
820        ),
821        format!("    has {STATUS} {}", typeql_string_literal(&record.status)),
822        format!("    has {STARTED_AT} {}", record.started_at),
823    ];
824
825    if let Some(finished_at) = &record.finished_at {
826        fields.push(format!("    has {FINISHED_AT} {finished_at}"));
827    }
828    if let Some(error) = &record.error {
829        fields.push(format!("    has {ERROR} {}", typeql_string_literal(error)));
830    }
831    if let Some(executor_ip) = &record.executor_ip {
832        fields.push(format!(
833            "    has {EXECUTOR_IP} {}",
834            typeql_string_literal(executor_ip)
835        ));
836    }
837    if let Some(executor_mac) = &record.executor_mac {
838        fields.push(format!(
839            "    has {EXECUTOR_MAC} {}",
840            typeql_string_literal(executor_mac)
841        ));
842    }
843
844    format!("\ninsert $r isa {RUN_ENTITY},\n{};\n", fields.join(",\n"))
845}
846
847impl MigrationStateStore for TypeDbStateStore {
848    fn ensure_schema(&self) -> BoxFuture<'_, Result<()>> {
849        Box::pin(async move {
850            // Even the latched fast path is a legacy writer entry point.  The
851            // read guard rejects an already-adopted target; each actual schema
852            // transaction repeats the guard for race-free mutation ordering.
853            require_legacy_writer_open(self.db.as_ref()).await?;
854            if self.schema_ensured.load(Ordering::Acquire) {
855                return Ok(());
856            }
857
858            let state_schema = migration_state_schema();
859
860            for (name, attribute) in &state_schema.attributes {
861                let mut definition = SchemaInfo::default();
862                definition
863                    .attributes
864                    .insert(name.clone(), attribute.clone());
865                let define = definition.to_typeql().map_err(map_schema_error)?;
866                self.ensure_type(name, &define).await?;
867            }
868
869            for (name, entity) in &state_schema.entities {
870                let mut definition = SchemaInfo::default();
871                definition.entities.insert(name.clone(), entity.clone());
872                let define = definition.to_typeql().map_err(map_schema_error)?;
873                self.ensure_type(name, &define).await?;
874            }
875
876            for (name, relation) in &state_schema.relations {
877                let mut definition = SchemaInfo::default();
878                definition.relations.insert(name.clone(), relation.clone());
879                let define = definition.to_typeql().map_err(map_schema_error)?;
880                self.ensure_type(name, &define).await?;
881            }
882
883            self.schema_ensured.store(true, Ordering::Release);
884            Ok(())
885        })
886    }
887
888    fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>> {
889        Box::pin(async move {
890            self.ensure_schema_for_read().await?;
891
892            // Ported verbatim from state.py:179-191.
893            let query = applied_query();
894
895            let values = self.query_documents(&query).await?;
896            parse_applied_documents(&values)
897        })
898    }
899
900    fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>> {
901        Box::pin(async move {
902            self.ensure_schema_for_read().await?;
903
904            let query = format!(
905                "\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"
906            );
907            let mut runs = parse_run_documents(&self.query_documents(&query).await?)?;
908
909            let finished_query = optional_run_field_query(FINISHED_AT, "finished");
910            let finished_docs = self.query_documents(&finished_query).await?;
911            merge_optional_run_field(&mut runs, &finished_docs, "finished_at", "finished");
912
913            let error_query = optional_run_field_query(ERROR, "error");
914            let error_docs = self.query_documents(&error_query).await?;
915            merge_optional_run_field(&mut runs, &error_docs, "error", "error");
916
917            let ip_query = optional_run_field_query(EXECUTOR_IP, "executor_ip");
918            let ip_docs = self.query_documents(&ip_query).await?;
919            merge_optional_run_field(&mut runs, &ip_docs, "executor_ip", "executor_ip");
920
921            let mac_query = optional_run_field_query(EXECUTOR_MAC, "executor_mac");
922            let mac_docs = self.query_documents(&mac_query).await?;
923            merge_optional_run_field(&mut runs, &mac_docs, "executor_mac", "executor_mac");
924
925            Ok(runs)
926        })
927    }
928
929    fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>> {
930        Box::pin(async move {
931            self.ensure_schema().await?;
932
933            // Rust stamps applied_at when the record carries none, mirroring
934            // state.py:249-250 (`datetime.now(UTC).strftime(...)`).
935            let applied_at = record
936                .applied_at
937                .clone()
938                .unwrap_or_else(|| format_applied_at(Utc::now()));
939
940            let migration_id = format!("{}:{}", record.app_label, record.name);
941            let migration_id = typeql_string_literal(&migration_id);
942            let app = typeql_string_literal(&record.app_label);
943            let name = typeql_string_literal(&record.name);
944            let checksum = typeql_string_literal(&record.checksum);
945
946            // Idempotent replace: delete any existing row for this migration_id
947            // (the @key) before inserting, so re-recording an already-applied
948            // migration updates in place instead of failing the @key constraint.
949            // This gives the TypeDB store the same dedup semantics the in-memory
950            // store has behind the shared seam.
951            let delete_existing = format!(
952                "\nmatch\n$m isa {APPLIED_ENTITY},\n    has {MIGRATION_ID} {migration_id};\ndelete $m;\n"
953            );
954
955            // Field set + storage schema match state.py:253-260. `applied_at` is
956            // emitted unquoted (a TypeQL datetime literal); every other field quoted.
957            let insert = format!(
958                "\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",
959            );
960
961            let ctx = self
962                .db
963                .transaction_context(TxType::Write)
964                .await
965                .map_err(map_orm_error)?;
966            if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
967                let _ = ctx.rollback().await;
968                return Err(error);
969            }
970            ctx.query(&delete_existing).await.map_err(map_orm_error)?;
971            ctx.query(&insert).await.map_err(map_orm_error)?;
972            ctx.commit().await.map_err(map_orm_error)?;
973            Ok(())
974        })
975    }
976
977    fn record_unapplied<'a>(
978        &'a self,
979        app_label: &'a str,
980        name: &'a str,
981    ) -> BoxFuture<'a, Result<()>> {
982        Box::pin(async move {
983            self.ensure_schema().await?;
984
985            // Storage schema matches state.py:288-293; the delete clause uses the
986            // TypeDB 3.x form `delete $m;` (the older `delete $m isa <type>;` that
987            // the Python code carried is a parse error on TypeDB 3.x).
988            let app_label = typeql_string_literal(app_label);
989            let name = typeql_string_literal(name);
990            let query = format!(
991                "\nmatch\n$m isa {APPLIED_ENTITY},\n    has {APP_LABEL} {app_label},\n    has {NAME} {name};\ndelete $m;\n"
992            );
993
994            let ctx = self
995                .db
996                .transaction_context(TxType::Write)
997                .await
998                .map_err(map_orm_error)?;
999            if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
1000                let _ = ctx.rollback().await;
1001                return Err(error);
1002            }
1003            ctx.query(&query).await.map_err(map_orm_error)?;
1004            ctx.commit().await.map_err(map_orm_error)?;
1005            Ok(())
1006        })
1007    }
1008
1009    fn record_run(&self, record: MigrationRunRecord) -> BoxFuture<'_, Result<()>> {
1010        Box::pin(async move {
1011            self.ensure_schema().await?;
1012
1013            let run_id = typeql_string_literal(&record.run_id);
1014            let delete_existing =
1015                format!("\nmatch\n$r isa {RUN_ENTITY},\n    has {RUN_ID} {run_id};\ndelete $r;\n");
1016            let insert = run_insert_query(&record);
1017
1018            let ctx = self
1019                .db
1020                .transaction_context(TxType::Write)
1021                .await
1022                .map_err(map_orm_error)?;
1023            if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
1024                let _ = ctx.rollback().await;
1025                return Err(error);
1026            }
1027            ctx.query(&delete_existing).await.map_err(map_orm_error)?;
1028            ctx.query(&insert).await.map_err(map_orm_error)?;
1029            ctx.commit().await.map_err(map_orm_error)?;
1030            Ok(())
1031        })
1032    }
1033}
1034
1035#[cfg(test)]
1036mod tests {
1037    use super::*;
1038    use crate::testing::{MockEvent, MockMigrationBackend};
1039    use chrono::{TimeZone, Timelike};
1040
1041    // ── document parsing (no TypeDB) ────────────────────────────────────────
1042    //
1043    // Feeds a hand-built fetch-document list in the REAL TypeDB shape: each
1044    // attribute-bound variable renders as a bare JSON scalar (typedb-driver
1045    // 3.8.1 `json_value`), NOT a `{"value": ...}` wrapper. This is the P0
1046    // silent-failure guard — a shape mismatch yields empty state with no error.
1047
1048    #[test]
1049    fn parse_applied_documents_extracts_bare_scalar_fields() {
1050        let docs = vec![serde_json::json!({
1051            "app": "myapp",
1052            "name": "0001_initial",
1053            "applied": "2026-06-05T00:00:00.000000000",
1054            "checksum": "abc123"
1055        })];
1056
1057        let records = parse_applied_documents(&docs).unwrap();
1058        assert_eq!(records.len(), 1);
1059        assert_eq!(records[0].app_label, "myapp");
1060        assert_eq!(records[0].name, "0001_initial");
1061        assert_eq!(records[0].checksum, "abc123");
1062        assert_eq!(
1063            records[0].applied_at.as_deref(),
1064            Some("2026-06-05T00:00:00.000000000")
1065        );
1066    }
1067
1068    #[tokio::test]
1069    async fn state_readers_close_every_read_context() {
1070        let responses = vec![
1071            QueryResult::Documents(vec![serde_json::json!({
1072                "app": "myapp",
1073                "name": "0001_initial",
1074                "applied": "2026-06-05T00:00:00.000000000",
1075                "checksum": "abc123"
1076            })]),
1077            QueryResult::Documents(Vec::new()),
1078            QueryResult::Documents(Vec::new()),
1079            QueryResult::Documents(Vec::new()),
1080            QueryResult::Documents(Vec::new()),
1081            QueryResult::Documents(Vec::new()),
1082        ];
1083        let (backend, log) = MockMigrationBackend::with_state_read_responses(responses);
1084        let store =
1085            TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1086
1087        assert_eq!(store.load_applied().await.unwrap().len(), 1);
1088        assert!(store.load_runs().await.unwrap().is_empty());
1089
1090        let events = log.lock().unwrap();
1091        let opens = events
1092            .iter()
1093            .filter(|event| matches!(event, MockEvent::OpenTx(TxType::Read)))
1094            .count();
1095        let closes = events
1096            .iter()
1097            .filter(|event| matches!(event, MockEvent::Close))
1098            .count();
1099        assert_eq!(opens, 7, "one schema inspection plus six ledger reads");
1100        assert_eq!(closes, opens, "every read context must be acknowledged");
1101    }
1102
1103    #[tokio::test]
1104    async fn load_applied_preserves_query_error_when_close_also_fails() {
1105        // Close 0 terminates the successful schema-presence inspection. The
1106        // applied-ledger query and close then fail together at indexes 0 and 1.
1107        let (backend, log) = MockMigrationBackend::with_state_read_and_close_failure(0, 1);
1108        let store =
1109            TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1110
1111        let error = store
1112            .load_applied()
1113            .await
1114            .expect_err("the ledger query must fail");
1115        let message = error.to_string();
1116        assert!(message.contains("injected query failure for testing"));
1117        assert!(!message.contains("injected close failure for testing"));
1118        assert_eq!(
1119            log.lock()
1120                .unwrap()
1121                .iter()
1122                .filter(|event| matches!(event, MockEvent::Close))
1123                .count(),
1124            2,
1125            "the failed ledger read must still acknowledge close"
1126        );
1127    }
1128
1129    #[tokio::test]
1130    async fn load_runs_preserves_query_error_when_close_also_fails() {
1131        let (backend, log) = MockMigrationBackend::with_state_read_and_close_failure(0, 1);
1132        let store =
1133            TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1134
1135        let error = store
1136            .load_runs()
1137            .await
1138            .expect_err("the base run-log query must fail");
1139        let message = error.to_string();
1140        assert!(message.contains("injected query failure for testing"));
1141        assert!(!message.contains("injected close failure for testing"));
1142        assert_eq!(
1143            log.lock()
1144                .unwrap()
1145                .iter()
1146                .filter(|event| matches!(event, MockEvent::Close))
1147                .count(),
1148            2,
1149            "the failed run-log read must still acknowledge close"
1150        );
1151    }
1152
1153    #[tokio::test]
1154    async fn unadopted_partial_state_schema_is_repaired_before_read() {
1155        let (backend, log, labels) =
1156            MockMigrationBackend::with_partial_state_schema(&[RUN_ENTITY], false);
1157        let store =
1158            TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1159
1160        assert!(store.load_applied().await.unwrap().is_empty());
1161        assert!(labels.lock().unwrap().contains(RUN_ENTITY));
1162        assert!(log.lock().unwrap().iter().any(|event| {
1163            matches!(event, MockEvent::Query(TxType::Schema, query) if query.contains(&format!("entity {RUN_ENTITY}")))
1164        }));
1165    }
1166
1167    #[tokio::test]
1168    async fn adopted_partial_state_schema_fails_before_repair() {
1169        let (backend, log, labels) =
1170            MockMigrationBackend::with_partial_state_schema(&[RUN_ENTITY], true);
1171        let store =
1172            TypeDbStateStore::new(Arc::new(Database::with_backend(Box::new(backend), "test")));
1173
1174        let error = store
1175            .load_applied()
1176            .await
1177            .expect_err("the sentinel must block partial-schema repair");
1178        assert!(error.to_string().contains(LEGACY_WRITER_CUTOVER_MESSAGE));
1179        assert!(!labels.lock().unwrap().contains(RUN_ENTITY));
1180        assert!(
1181            !log.lock()
1182                .unwrap()
1183                .iter()
1184                .any(|event| matches!(event, MockEvent::Query(TxType::Schema, _))),
1185            "adopted schema repair must not reach a define query"
1186        );
1187    }
1188
1189    #[tokio::test]
1190    async fn exact_sentinel_partition_is_verified_before_filtering() {
1191        let fingerprint = "a".repeat(64);
1192        let sentinel = serde_json::json!({
1193            "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1194            "name": LEGACY_CUTOVER_SENTINEL_NAME,
1195            "applied": LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
1196            "checksum": fingerprint,
1197        });
1198        let responses = vec![
1199            QueryResult::Documents(vec![sentinel]),
1200            QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1201            QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1202            QueryResult::Documents(vec![serde_json::json!({
1203                "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1204                "applied": LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
1205                "checksum": fingerprint,
1206            })]),
1207        ];
1208        let (backend, _) = MockMigrationBackend::with_state_read_responses(responses);
1209        let database = Database::with_backend(Box::new(backend), "test");
1210        let mut transaction = database.read_transaction().await.unwrap();
1211
1212        let partition = TypeDbStateStore::load_verified_legacy_partition_in_transaction(
1213            &mut transaction,
1214            LegacyCutoverSentinelExpectation::RequiredExact(&fingerprint),
1215        )
1216        .await
1217        .expect("verify exact sentinel");
1218        transaction.close().await.unwrap();
1219
1220        assert!(partition.applied().is_empty());
1221        assert_eq!(partition.sentinel_fingerprint(), Some(fingerprint.as_str()));
1222    }
1223
1224    #[tokio::test]
1225    async fn malformed_sentinel_timestamp_is_not_filtered() {
1226        let fingerprint = "b".repeat(64);
1227        let malformed_timestamp = "1970-01-01T00:00:00";
1228        let responses = vec![
1229            QueryResult::Documents(vec![serde_json::json!({
1230                "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1231                "name": LEGACY_CUTOVER_SENTINEL_NAME,
1232                "applied": malformed_timestamp,
1233                "checksum": fingerprint,
1234            })]),
1235            QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1236            QueryResult::Documents(vec![serde_json::json!({"exists": true})]),
1237            QueryResult::Documents(vec![serde_json::json!({
1238                "app": LEGACY_CUTOVER_SENTINEL_APP_LABEL,
1239                "applied": malformed_timestamp,
1240                "checksum": fingerprint,
1241            })]),
1242        ];
1243        let (backend, _) = MockMigrationBackend::with_state_read_responses(responses);
1244        let database = Database::with_backend(Box::new(backend), "test");
1245        let mut transaction = database.read_transaction().await.unwrap();
1246
1247        let error = TypeDbStateStore::load_verified_legacy_partition_in_transaction(
1248            &mut transaction,
1249            LegacyCutoverSentinelExpectation::RequiredExact(&fingerprint),
1250        )
1251        .await
1252        .expect_err("malformed sentinel must fail closed");
1253        transaction.close().await.unwrap();
1254
1255        assert!(error.is_contract_violation());
1256        assert!(error.to_string().contains("foreign applied timestamp"));
1257    }
1258
1259    #[test]
1260    fn sentinel_name_is_outside_the_released_numbered_loader_namespace() {
1261        assert!(
1262            !LEGACY_CUTOVER_SENTINEL_NAME
1263                .as_bytes()
1264                .first()
1265                .is_some_and(u8::is_ascii_digit)
1266        );
1267    }
1268
1269    #[test]
1270    fn parse_applied_documents_also_unwraps_value_envelope() {
1271        // Defensive: the original `_extract_value` accepted a `{"value": ...}`
1272        // object form too. Mixed shapes in one list must both parse.
1273        let docs = vec![serde_json::json!({
1274            "app": {"value": "myapp"},
1275            "name": {"value": "0002_next"},
1276            "applied": {"value": "2026-06-05T01:02:03.000000000"},
1277            "checksum": {"value": "def456"}
1278        })];
1279
1280        let records = parse_applied_documents(&docs).unwrap();
1281        assert_eq!(records.len(), 1);
1282        assert_eq!(records[0].app_label, "myapp");
1283        assert_eq!(records[0].name, "0002_next");
1284        assert_eq!(records[0].checksum, "def456");
1285        assert_eq!(
1286            records[0].applied_at.as_deref(),
1287            Some("2026-06-05T01:02:03.000000000")
1288        );
1289    }
1290
1291    #[test]
1292    fn parse_applied_documents_empty_list_is_empty() {
1293        let records = parse_applied_documents(&[]).unwrap();
1294        assert!(records.is_empty());
1295    }
1296
1297    #[test]
1298    fn parse_applied_documents_skips_incomplete_rows() {
1299        // Missing `checksum` → skipped, mirroring the Python `all([...])` guard.
1300        let docs = vec![
1301            serde_json::json!({
1302                "app": "myapp",
1303                "name": "0001_initial",
1304                "applied": "2026-06-05T00:00:00.000000000"
1305            }),
1306            serde_json::json!({
1307                "app": "myapp",
1308                "name": "0002_next",
1309                "applied": "2026-06-05T00:00:00.000000000",
1310                "checksum": "ok"
1311            }),
1312        ];
1313
1314        let records = parse_applied_documents(&docs).unwrap();
1315        assert_eq!(records.len(), 1);
1316        assert_eq!(records[0].name, "0002_next");
1317    }
1318
1319    #[test]
1320    fn parse_applied_documents_carries_missing_applied_as_none() {
1321        let docs = vec![serde_json::json!({
1322            "app": "myapp",
1323            "name": "0001_initial",
1324            "checksum": "abc123"
1325        })];
1326
1327        let records = parse_applied_documents(&docs).unwrap();
1328        assert_eq!(records.len(), 1);
1329        assert!(records[0].applied_at.is_none());
1330    }
1331
1332    #[test]
1333    fn parse_run_documents_extracts_required_fields() {
1334        let docs = vec![serde_json::json!({
1335            "run_id": "run-1",
1336            "app": "app",
1337            "name": "0001_initial",
1338            "checksum": "abc123",
1339            "direction": "apply",
1340            "status": "started",
1341            "started": "2026-06-05T00:00:00.000000"
1342        })];
1343
1344        let records = parse_run_documents(&docs).unwrap();
1345
1346        assert_eq!(records.len(), 1);
1347        assert_eq!(records[0].run_id, "run-1");
1348        assert_eq!(records[0].direction, "apply");
1349        assert_eq!(records[0].status, "started");
1350        assert_eq!(records[0].finished_at, None);
1351    }
1352
1353    #[test]
1354    fn merge_optional_run_field_updates_matching_record_only() {
1355        let mut records = vec![MigrationRunRecord {
1356            run_id: "run-1".to_string(),
1357            app_label: "app".to_string(),
1358            name: "0001_initial".to_string(),
1359            checksum: "abc123".to_string(),
1360            direction: "apply".to_string(),
1361            status: "started".to_string(),
1362            started_at: "2026-06-05T00:00:00.000000".to_string(),
1363            finished_at: None,
1364            error: None,
1365            executor_ip: None,
1366            executor_mac: None,
1367        }];
1368        let docs = vec![serde_json::json!({
1369            "run_id": "run-1",
1370            "finished": "2026-06-05T00:00:01.000000"
1371        })];
1372
1373        merge_optional_run_field(&mut records, &docs, "finished_at", "finished");
1374
1375        assert_eq!(
1376            records[0].finished_at.as_deref(),
1377            Some("2026-06-05T00:00:01.000000")
1378        );
1379    }
1380
1381    #[test]
1382    fn typeql_string_literal_escapes_user_controlled_text() {
1383        assert_eq!(typeql_string_literal("a\"b\\c\n"), "\"a\\\"b\\\\c\\n\"");
1384    }
1385
1386    #[test]
1387    fn schema_bootstrap_rollback_failure_preserves_primary_and_cleanup() {
1388        let primary = OrmError::QueryExecution("define failed".to_owned());
1389        let cleanup = OrmError::Transaction("rollback failed".to_owned());
1390        let error = schema_rollback_cleanup_failure(&primary, &cleanup).to_string();
1391        assert!(error.contains("define failed"), "{error}");
1392        assert!(error.contains("rollback failed"), "{error}");
1393        assert!(error.contains("rollback was not acknowledged"), "{error}");
1394    }
1395
1396    #[test]
1397    fn run_insert_query_includes_optional_fields_when_present() {
1398        let record = MigrationRunRecord {
1399            run_id: "run-1".to_string(),
1400            app_label: "app".to_string(),
1401            name: "0001_initial".to_string(),
1402            checksum: "abc123".to_string(),
1403            direction: "apply".to_string(),
1404            status: "failed".to_string(),
1405            started_at: "2026-06-05T00:00:00.000000".to_string(),
1406            finished_at: Some("2026-06-05T00:00:01.000000".to_string()),
1407            error: Some("quote: \"boom\"".to_string()),
1408            executor_ip: Some("127.0.0.1".to_string()),
1409            executor_mac: Some("00:11:22:33:44:55".to_string()),
1410        };
1411
1412        let query = run_insert_query(&record);
1413
1414        assert!(query.contains("has migration_run_id \"run-1\""));
1415        assert!(query.contains("has migration_finished_at 2026-06-05T00:00:01.000000"));
1416        assert!(query.contains("has migration_error \"quote: \\\"boom\\\"\""));
1417        assert!(query.contains("has migration_executor_ip \"127.0.0.1\""));
1418        assert!(query.contains("has migration_executor_mac \"00:11:22:33:44:55\""));
1419    }
1420
1421    // ── timestamp parity (no TypeDB) ────────────────────────────────────────
1422    //
1423    // Asserts the Rust format string is byte-identical to Python's
1424    // `strftime("%Y-%m-%dT%H:%M:%S.%f")`: 6 fractional digits, zero-padded,
1425    // no timezone suffix.
1426
1427    #[test]
1428    fn format_applied_at_matches_python_strftime() {
1429        // 2026-06-05 14:09:08.123456 UTC.
1430        let dt = Utc
1431            .with_ymd_and_hms(2026, 6, 5, 14, 9, 8)
1432            .unwrap()
1433            .with_nanosecond(123_456_000)
1434            .unwrap();
1435        assert_eq!(format_applied_at(dt), "2026-06-05T14:09:08.123456");
1436    }
1437
1438    #[test]
1439    fn format_applied_at_zero_pads_microseconds() {
1440        // Python `%f` zero-pads to exactly 6 digits; a sub-microsecond value
1441        // truncates and pads identically under chrono `%6f`.
1442        let dt = Utc
1443            .with_ymd_and_hms(2026, 1, 2, 3, 4, 5)
1444            .unwrap()
1445            .with_nanosecond(7_000)
1446            .unwrap();
1447        // 7000 ns = 7 microseconds → ".000007".
1448        assert_eq!(format_applied_at(dt), "2026-01-02T03:04:05.000007");
1449    }
1450}