type_bridge_migration/state/
mod.rs1pub 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#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
40pub struct MigrationExecutorInfo {
41 #[serde(default)]
43 pub ip: Option<String>,
44 #[serde(default)]
46 pub mac: Option<String>,
47}
48
49#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
51pub struct MigrationRunRecord {
52 pub run_id: String,
54 pub app_label: String,
56 pub name: String,
58 pub checksum: String,
60 pub direction: String,
62 pub status: String,
64 pub started_at: String,
66 #[serde(default)]
68 pub finished_at: Option<String>,
69 #[serde(default)]
71 pub error: Option<String>,
72 #[serde(default)]
74 pub executor_ip: Option<String>,
75 #[serde(default)]
77 pub executor_mac: Option<String>,
78}
79
80pub fn migration_timestamp_now() -> String {
82 Utc::now().format(TIMESTAMP_FORMAT).to_string()
83}
84
85pub fn collect_executor_info() -> MigrationExecutorInfo {
87 MigrationExecutorInfo {
88 ip: local_ip(),
89 mac: local_mac(),
90 }
91}
92
93pub 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
114pub 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
176pub trait MigrationStateStore: Send + Sync {
190 fn ensure_schema(&self) -> BoxFuture<'_, Result<()>>;
195
196 fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>>;
198
199 fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>>;
201
202 fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>>;
208
209 fn record_unapplied<'a>(
214 &'a self,
215 app_label: &'a str,
216 name: &'a str,
217 ) -> BoxFuture<'a, Result<()>>;
218
219 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}