Skip to main content

type_bridge_migration/state/
mod.rs

1//! Migration state backend seam.
2//!
3//! [`MigrationStateStore`] is the trait that abstracts applied-state storage.
4//! Two implementations exist:
5//!
6//! - [`memory::InMemoryStateStore`] — insertion-ordered, mutex-backed store
7//!   used for Rust unit tests with no TypeDB connection.
8//! - [`typedb::TypeDbStateStore`] — persists state over the ORM
9//!   [`Database`][type_bridge_orm::session::Database] seam.
10//!
11//! [`schema`] is the public, canonical description of the TypeDB types owned by
12//! the default store. Bootstrap and language bindings both consume that same
13//! [`SchemaInfo`][type_bridge_orm::_schema::SchemaInfo].
14
15pub mod memory;
16pub mod schema;
17pub mod typedb;
18
19pub use memory::InMemoryStateStore;
20pub use schema::{
21    MigrationStateSchemaKind, applied_migration_entity_label, is_migration_state_type,
22    migration_state_schema,
23};
24pub use typedb::{
25    LEGACY_CUTOVER_SENTINEL_APP_LABEL, LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
26    LEGACY_CUTOVER_SENTINEL_MIGRATION_ID, LEGACY_CUTOVER_SENTINEL_NAME,
27    LEGACY_WRITER_CUTOVER_MESSAGE, LegacyCutoverSentinelError, LegacyCutoverSentinelExpectation,
28    TypeDbStateStore, VerifiedLegacyAppliedPartition, require_legacy_writer_open,
29    require_legacy_writer_open_in_transaction,
30};
31
32use std::net::UdpSocket;
33
34use chrono::Utc;
35use network_interface::{NetworkInterface, NetworkInterfaceConfig};
36use type_bridge_orm::session::backend::BoxFuture;
37use uuid::Uuid;
38
39use crate::plan::{MigrationAction, MigrationExecution};
40use crate::{AppliedMigrationRecord, Result};
41
42const TIMESTAMP_FORMAT: &str = "%Y-%m-%dT%H:%M:%S.%6f";
43
44/// Best-effort executor identity stored with migration run-log rows.
45#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
46pub struct MigrationExecutorInfo {
47    /// Best-effort local IP address of the process that executed the migration.
48    #[serde(default)]
49    pub ip: Option<String>,
50    /// Best-effort MAC address of the process that executed the migration.
51    #[serde(default)]
52    pub mac: Option<String>,
53}
54
55/// Append/update record for one migration execution attempt.
56#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
57pub struct MigrationRunRecord {
58    /// Unique execution-attempt identifier.
59    pub run_id: String,
60    /// Application or migration package label.
61    pub app_label: String,
62    /// Migration file stem, such as `0001_initial`.
63    pub name: String,
64    /// Migration checksum observed when this run started.
65    pub checksum: String,
66    /// Execution direction: `apply` or `rollback`.
67    pub direction: String,
68    /// Execution status: `started`, `succeeded`, or `failed`.
69    pub status: String,
70    /// UTC timestamp when the run started.
71    pub started_at: String,
72    /// UTC timestamp when the run finished, when available.
73    #[serde(default)]
74    pub finished_at: Option<String>,
75    /// Error message captured for failed runs, when available.
76    #[serde(default)]
77    pub error: Option<String>,
78    /// Best-effort executor IP address, when available.
79    #[serde(default)]
80    pub executor_ip: Option<String>,
81    /// Best-effort executor MAC address, when available.
82    #[serde(default)]
83    pub executor_mac: Option<String>,
84}
85
86/// Return the current UTC timestamp in the TypeDB datetime literal format.
87pub fn migration_timestamp_now() -> String {
88    Utc::now().format(TIMESTAMP_FORMAT).to_string()
89}
90
91/// Collect best-effort executor identity for audit logging.
92pub fn collect_executor_info() -> MigrationExecutorInfo {
93    MigrationExecutorInfo {
94        ip: local_ip(),
95        mac: local_mac(),
96    }
97}
98
99/// Create a started run-log record for a planned migration execution.
100pub fn started_run_record(
101    migration: &MigrationExecution,
102    checksum: String,
103    executor: &MigrationExecutorInfo,
104) -> MigrationRunRecord {
105    MigrationRunRecord {
106        run_id: Uuid::new_v4().to_string(),
107        app_label: migration.app_label.clone(),
108        name: migration.name.clone(),
109        checksum,
110        direction: direction_label(migration.action).to_string(),
111        status: "started".to_string(),
112        started_at: migration_timestamp_now(),
113        finished_at: None,
114        error: None,
115        executor_ip: executor.ip.clone(),
116        executor_mac: executor.mac.clone(),
117    }
118}
119
120/// Return a finished copy of a migration run-log record.
121pub fn finished_run_record(
122    mut record: MigrationRunRecord,
123    status: &str,
124    error: Option<String>,
125) -> MigrationRunRecord {
126    record.status = status.to_string();
127    record.finished_at = Some(migration_timestamp_now());
128    record.error = error;
129    record
130}
131
132fn direction_label(action: MigrationAction) -> &'static str {
133    match action {
134        MigrationAction::Apply => "apply",
135        MigrationAction::Rollback => "rollback",
136    }
137}
138
139fn local_ip() -> Option<String> {
140    let socket = UdpSocket::bind("0.0.0.0:0").ok()?;
141    socket.connect("8.8.8.8:80").ok()?;
142    let addr = socket.local_addr().ok()?;
143    Some(addr.ip().to_string()).filter(|ip| ip != "0.0.0.0")
144}
145
146fn local_mac() -> Option<String> {
147    NetworkInterface::show()
148        .ok()?
149        .into_iter()
150        .filter(|interface| !interface.internal)
151        .find_map(|interface| {
152            interface
153                .mac_addr
154                .as_deref()
155                .and_then(normalize_mac_address)
156        })
157}
158
159fn normalize_mac_address(value: &str) -> Option<String> {
160    let mut bytes = [0_u8; 6];
161    let mut parts = value.split([':', '-']);
162    for byte in &mut bytes {
163        let part = parts.next()?;
164        if part.len() != 2 {
165            return None;
166        }
167        *byte = u8::from_str_radix(part, 16).ok()?;
168    }
169    if parts.next().is_some() || bytes.iter().all(|byte| *byte == 0) {
170        return None;
171    }
172    Some(format_mac_address(bytes))
173}
174
175fn format_mac_address(bytes: [u8; 6]) -> String {
176    format!(
177        "{:02x}:{:02x}:{:02x}:{:02x}:{:02x}:{:02x}",
178        bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5]
179    )
180}
181
182/// Seam trait for migration applied-state storage.
183///
184/// Implementations must be [`Send`] + [`Sync`] and expose four operations that
185/// map 1-to-1 to the existing Python `MigrationStateManager` surface:
186///
187/// - [`ensure_schema`][Self::ensure_schema] — idempotent schema bootstrap.
188/// - [`load_applied`][Self::load_applied] — full applied-state read.
189/// - [`record_applied`][Self::record_applied] — insert/replace one record.
190/// - [`record_unapplied`][Self::record_unapplied] — remove one record (absent
191///   → `Ok`, matching Python delete semantics).
192///
193/// Method futures are returned as [`BoxFuture`] so the trait is object-safe
194/// and can be used as `Box<dyn MigrationStateStore>`.
195pub trait MigrationStateStore: Send + Sync {
196    /// Ensure the state-storage schema exists, creating it if absent.
197    ///
198    /// This is idempotent: calling it on an already-initialised store is a
199    /// no-op `Ok(())`.
200    fn ensure_schema(&self) -> BoxFuture<'_, Result<()>>;
201
202    /// Load all applied migration records in stable insertion order.
203    fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>>;
204
205    /// Load all migration execution run-log records.
206    fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>>;
207
208    /// Record a migration as applied, inserting or replacing by `(app_label,
209    /// name)` identity.
210    ///
211    /// Calling this twice with the same identity is idempotent: the second
212    /// call replaces the first; no duplicates are stored.
213    fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>>;
214
215    /// Remove the applied record identified by `(app_label, name)`.
216    ///
217    /// If no such record exists the call succeeds silently, matching the
218    /// Python `record_unapplied` delete-absent no-op.
219    fn record_unapplied<'a>(
220        &'a self,
221        app_label: &'a str,
222        name: &'a str,
223    ) -> BoxFuture<'a, Result<()>>;
224
225    /// Insert or replace one migration execution run-log record by `run_id`.
226    fn record_run(&self, record: MigrationRunRecord) -> BoxFuture<'_, Result<()>>;
227}
228
229#[cfg(test)]
230mod tests {
231    use super::{format_mac_address, normalize_mac_address};
232
233    #[test]
234    fn formats_mac_addresses_in_lowercase_colon_notation() {
235        assert_eq!(
236            format_mac_address([0x00, 0x11, 0xAB, 0xCD, 0xEF, 0x42]),
237            "00:11:ab:cd:ef:42"
238        );
239    }
240
241    #[test]
242    fn normalizes_supported_mac_address_strings() {
243        assert_eq!(
244            normalize_mac_address("00:11:AB:CD:EF:42").as_deref(),
245            Some("00:11:ab:cd:ef:42")
246        );
247        assert_eq!(
248            normalize_mac_address("00-11-ab-cd-ef-42").as_deref(),
249            Some("00:11:ab:cd:ef:42")
250        );
251    }
252
253    #[test]
254    fn rejects_invalid_or_zero_mac_address_strings() {
255        assert_eq!(normalize_mac_address("00:00:00:00:00:00"), None);
256        assert_eq!(normalize_mac_address("00:11:22:33:44"), None);
257        assert_eq!(normalize_mac_address("00:11:22:33:44:zz"), None);
258    }
259}