1use std::collections::BTreeSet;
17use std::sync::Arc;
18use std::sync::atomic::{AtomicBool, Ordering};
19
20use chrono::Utc;
21pub use type_bridge_contract::reserved::{
22 LEGACY_CUTOVER_SENTINEL_APP_LABEL, LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
23 LEGACY_CUTOVER_SENTINEL_MIGRATION_ID, LEGACY_CUTOVER_SENTINEL_NAME,
24 LEGACY_WRITER_CUTOVER_MESSAGE,
25};
26use type_bridge_orm::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
43const APPLIED_AT_FORMAT: &str = "%Y-%m-%dT%H:%M:%S.%6f";
49
50const LEGACY_STATE_SCHEMA_PROBE_QUERY_TAG: &str =
51 "# typebridge-internal-legacy-state-schema-probe/v1\n";
52
53#[derive(Debug, Clone, Copy)]
56pub enum LegacyCutoverSentinelExpectation<'a> {
57 Absent,
59 OptionalExact(&'a str),
61 RequiredExact(&'a str),
63}
64
65#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct VerifiedLegacyAppliedPartition {
72 applied: Vec<AppliedMigrationRecord>,
73 sentinel_fingerprint: Option<String>,
74}
75
76impl VerifiedLegacyAppliedPartition {
77 pub fn applied(&self) -> &[AppliedMigrationRecord] {
80 &self.applied
81 }
82
83 pub fn into_applied(self) -> Vec<AppliedMigrationRecord> {
85 self.applied
86 }
87
88 pub fn sentinel_fingerprint(&self) -> Option<&str> {
90 self.sentinel_fingerprint.as_deref()
91 }
92}
93
94#[derive(Debug, thiserror::Error)]
96pub enum LegacyCutoverSentinelError {
97 #[error("legacy cutover sentinel storage inspection failed: {0}")]
99 Storage(#[source] MigrationError),
100 #[error("legacy cutover sentinel contract violation: {message}")]
102 Contract {
103 message: String,
105 },
106}
107
108impl LegacyCutoverSentinelError {
109 pub fn is_contract_violation(&self) -> bool {
112 matches!(self, Self::Contract { .. })
113 }
114}
115
116pub struct TypeDbStateStore {
121 db: Arc<Database>,
122 schema_ensured: AtomicBool,
125}
126
127impl TypeDbStateStore {
128 pub fn new(db: Arc<Database>) -> Self {
130 Self {
131 db,
132 schema_ensured: AtomicBool::new(false),
133 }
134 }
135
136 async fn type_exists(&self, type_name: &str) -> Result<bool> {
137 let kind = state_type_kind(type_name).ok_or_else(|| MigrationError::State {
138 message: format!("unknown migration-state schema label: {type_name}"),
139 })?;
140 let check_query = format!(
141 "\n match {kind} $t;\n fetch {{ \"label\": label($t) }};\n "
142 );
143
144 let ctx = self
145 .db
146 .transaction_context(TxType::Read)
147 .await
148 .map_err(map_orm_error)?;
149 let checked = ctx
150 .query(&check_query)
151 .await
152 .map_err(map_orm_error)
153 .and_then(|result| schema_labels_from_result(result, "migration state type probe"))
154 .map(|labels| labels.contains(type_name));
155 let closed = ctx.close().await.map_err(map_orm_error);
156 match (checked, closed) {
157 (Ok(exists), Ok(())) => Ok(exists),
158 (Err(primary), Ok(())) => Err(primary),
159 (Ok(_), Err(cleanup)) => Err(cleanup),
160 (Err(primary), Err(_)) => Err(primary),
161 }
162 }
163
164 async fn ensure_type(&self, type_name: &str, define_typeql: &str) -> Result<()> {
165 if self.type_exists(type_name).await? {
166 return Ok(());
167 }
168
169 let ctx = self
170 .db
171 .transaction_context(TxType::Schema)
172 .await
173 .map_err(map_orm_error)?;
174 if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
175 let _ = ctx.rollback().await;
176 return Err(error);
177 }
178 match ctx.query(define_typeql).await {
179 Ok(_) => {
180 ctx.commit().await.map_err(map_orm_error)?;
181 Ok(())
182 }
183 Err(error) => {
184 if let Err(cleanup) = ctx.rollback().await {
185 return Err(schema_rollback_cleanup_failure(&error, &cleanup));
186 }
187 if self.type_exists(type_name).await? {
188 Ok(())
189 } else {
190 Err(map_orm_error(error))
191 }
192 }
193 }
194 }
195
196 async fn ensure_schema_for_read(&self) -> Result<()> {
197 if self.schema_ensured.load(Ordering::Acquire) {
198 return Ok(());
199 }
200
201 let mut transaction = self.db.read_transaction().await.map_err(map_orm_error)?;
206 let inspected = legacy_state_schema_presence(&mut transaction)
207 .await
208 .map_err(legacy_sentinel_error_into_migration_error);
209 let closed = transaction.close().await.map_err(map_orm_error);
210 let presence = match (inspected, closed) {
211 (Ok(presence), Ok(())) => presence,
212 (Err(primary), Ok(())) => return Err(primary),
213 (Ok(_), Err(cleanup)) => return Err(cleanup),
214 (Err(primary), Err(_)) => return Err(primary),
215 };
216 if presence == LegacyStateSchemaPresence::Complete {
217 self.schema_ensured.store(true, Ordering::Release);
218 return Ok(());
219 }
220
221 self.ensure_schema().await
226 }
227
228 async fn query_documents(&self, query: &str) -> Result<Vec<serde_json::Value>> {
229 let ctx = self
230 .db
231 .transaction_context(TxType::Read)
232 .await
233 .map_err(map_orm_error)?;
234 let queried = ctx
235 .query(query)
236 .await
237 .map_err(map_orm_error)
238 .map(query_result_values);
239 let closed = ctx.close().await.map_err(map_orm_error);
240 match (queried, closed) {
241 (Ok(values), Ok(())) => Ok(values),
242 (Err(primary), Ok(())) => Err(primary),
243 (Ok(_), Err(cleanup)) => Err(cleanup),
244 (Err(primary), Err(_)) => Err(primary),
245 }
246 }
247
248 pub async fn load_applied_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 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 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
440pub 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
455pub 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
615fn 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
642fn 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
672fn extract_value(doc: &serde_json::Value, key: &str) -> Option<String> {
684 let value = doc.get(key)?;
685 extract_scalar(value)
686}
687
688fn 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
701pub 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 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
733pub 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 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 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 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 let delete_existing = format!(
952 "\nmatch\n$m isa {APPLIED_ENTITY},\n has {MIGRATION_ID} {migration_id};\ndelete $m;\n"
953 );
954
955 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 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 #[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 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 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 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 #[test]
1428 fn format_applied_at_matches_python_strftime() {
1429 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 let dt = Utc
1443 .with_ymd_and_hms(2026, 1, 2, 3, 4, 5)
1444 .unwrap()
1445 .with_nanosecond(7_000)
1446 .unwrap();
1447 assert_eq!(format_applied_at(dt), "2026-01-02T03:04:05.000007");
1449 }
1450}