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