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`] seam. Its per-type
5//! bootstrap definitions are rendered from the canonical migration-state
6//! [`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::_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
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 frozen applied ledger without creating or repairing its schema.
249    ///
250    /// A completely absent legacy schema represents an empty history. A
251    /// partial schema is durable drift and fails closed. This is the Python
252    /// archival-reader boundary after the generated-only cutover.
253    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    /// Read the frozen legacy run log without any schema bootstrap or repair.
261    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    /// Read the complete released applied ledger through an already-retained
315    /// transaction.
316    ///
317    /// Legacy-frontier cutover uses a managed schema transaction as an
318    /// exclusive guard against V1 write transactions. The caller must ensure
319    /// the released state schema already exists before opening that guard;
320    /// this method performs no bootstrap and never commits or closes it.
321    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    /// Read the released applied projection and validate the reserved V2
332    /// sentinel through the same retained transaction snapshot.
333    ///
334    /// The sentinel is removed only after both independent identity probes,
335    /// singleton cardinality, every stored field, and the expected anchor
336    /// fingerprint have been checked.  This is the only supported filtering
337    /// boundary for V2 legacy-frontier continuity and digest calculations.
338    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    /// Stage the complete V2 cutover sentinel in a caller-retained managed
482    /// transaction.  The caller commits this in the same transaction as the
483    /// managed cutover anchor.
484    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
506/// Fail before a legacy writer uses an already-open transaction when an exact,
507/// managed-anchor-bound V2 cutover pair is present.
508pub 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
521/// Read-only entry guard for legacy writer surfaces whose external side
522/// effects cannot share a TypeDB transaction.
523pub 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
681/// Format the current UTC time as the Python-compatible applied-at string.
682///
683/// Isolated as a pure function (parameterised on the instant) so the
684/// `%Y-%m-%dT%H:%M:%S.%6f` parity with Python's `%f` is unit-testable without a
685/// live clock or TypeDB.
686fn 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
708/// Map an ORM-layer error into the migration error hierarchy.
709fn 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
738/// Unwrap a single fetched field into its scalar string.
739///
740/// The Phase-05 ORM backend renders a `fetch { "k": $attr }` document so that
741/// each attribute-bound variable becomes a bare JSON scalar (the attribute's
742/// value), not a `{"value": ...}` wrapper (typedb-driver `json_value`,
743/// `concept_document.rs`). The original Python `_extract_value`
744/// (`state.py:221-237`) defensively handled BOTH shapes — a bare scalar and a
745/// `{"value": ...}` object — so this parser does the same: prefer the inner
746/// `"value"` when the field is an object carrying one, otherwise take the
747/// scalar directly. Returns `None` for a missing/null field or an object with
748/// no usable `"value"`.
749fn extract_value(doc: &serde_json::Value, key: &str) -> Option<String> {
750    let value = doc.get(key)?;
751    extract_scalar(value)
752}
753
754/// Reduce a fetched JSON value to its string form, unwrapping a `{"value":
755/// ...}` envelope when present (the `_extract_value` dict branch).
756fn 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
767/// Parse a list of TypeDB fetch documents into applied migration records.
768///
769/// Pure and TypeDB-free so the document-shape contract is unit-testable. Each
770/// document is the `fetch { "app": ..., "name": ..., "applied": ...,
771/// "checksum": ... }` shape from the `load_applied` query (ported from
772/// `state.py:179-191`). A document missing any of `app` / `name` / `checksum`
773/// is skipped, matching the Python `if all([...])` guard (`state.py:206`);
774/// `applied` is carried through as the optional ISO string.
775pub 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            // Skip incomplete rows (mirrors the Python `all([...])` guard).
786            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
799/// Parse TypeDB fetch documents into migration run-log records.
800pub 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            // Even the latched fast path is a legacy writer entry point.  The
917            // read guard rejects an already-adopted target; each actual schema
918            // transaction repeats the guard for race-free mutation ordering.
919            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            // Ported verbatim from state.py:179-191.
959            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            // Rust stamps applied_at when the record carries none, mirroring
978            // state.py:249-250 (`datetime.now(UTC).strftime(...)`).
979            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            // Idempotent replace: delete any existing row for this migration_id
991            // (the @key) before inserting, so re-recording an already-applied
992            // migration updates in place instead of failing the @key constraint.
993            // This gives the TypeDB store the same dedup semantics the in-memory
994            // store has behind the shared seam.
995            let delete_existing = format!(
996                "\nmatch\n$m isa {APPLIED_ENTITY},\n    has {MIGRATION_ID} {migration_id};\ndelete $m;\n"
997            );
998
999            // Field set + storage schema match state.py:253-260. `applied_at` is
1000            // emitted unquoted (a TypeQL datetime literal); every other field quoted.
1001            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            // Storage schema matches state.py:288-293; the delete clause uses the
1030            // TypeDB 3.x form `delete $m;` (the older `delete $m isa <type>;` that
1031            // the Python code carried is a parse error on TypeDB 3.x).
1032            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    // ── document parsing (no TypeDB) ────────────────────────────────────────
1086    //
1087    // Feeds a hand-built fetch-document list in the REAL TypeDB shape: each
1088    // attribute-bound variable renders as a bare JSON scalar (typedb-driver
1089    // 3.8.1 `json_value`), NOT a `{"value": ...}` wrapper. This is the P0
1090    // silent-failure guard — a shape mismatch yields empty state with no error.
1091
1092    #[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        // Close 0 terminates the successful schema-presence inspection. The
1150        // applied-ledger query and close then fail together at indexes 0 and 1.
1151        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        // Defensive: the original `_extract_value` accepted a `{"value": ...}`
1316        // object form too. Mixed shapes in one list must both parse.
1317        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        // Missing `checksum` → skipped, mirroring the Python `all([...])` guard.
1344        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    // ── timestamp parity (no TypeDB) ────────────────────────────────────────
1466    //
1467    // Asserts the Rust format string is byte-identical to Python's
1468    // `strftime("%Y-%m-%dT%H:%M:%S.%f")`: 6 fractional digits, zero-padded,
1469    // no timezone suffix.
1470
1471    #[test]
1472    fn format_applied_at_matches_python_strftime() {
1473        // 2026-06-05 14:09:08.123456 UTC.
1474        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        // Python `%f` zero-pads to exactly 6 digits; a sub-microsecond value
1485        // truncates and pads identically under chrono `%6f`.
1486        let dt = Utc
1487            .with_ymd_and_hms(2026, 1, 2, 3, 4, 5)
1488            .unwrap()
1489            .with_nanosecond(7_000)
1490            .unwrap();
1491        // 7000 ns = 7 microseconds → ".000007".
1492        assert_eq!(format_applied_at(dt), "2026-01-02T03:04:05.000007");
1493    }
1494}