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::{
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#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
46pub struct MigrationExecutorInfo {
47 #[serde(default)]
49 pub ip: Option<String>,
50 #[serde(default)]
52 pub mac: Option<String>,
53}
54
55#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
57pub struct MigrationRunRecord {
58 pub run_id: String,
60 pub app_label: String,
62 pub name: String,
64 pub checksum: String,
66 pub direction: String,
68 pub status: String,
70 pub started_at: String,
72 #[serde(default)]
74 pub finished_at: Option<String>,
75 #[serde(default)]
77 pub error: Option<String>,
78 #[serde(default)]
80 pub executor_ip: Option<String>,
81 #[serde(default)]
83 pub executor_mac: Option<String>,
84}
85
86pub fn migration_timestamp_now() -> String {
88 Utc::now().format(TIMESTAMP_FORMAT).to_string()
89}
90
91pub fn collect_executor_info() -> MigrationExecutorInfo {
93 MigrationExecutorInfo {
94 ip: local_ip(),
95 mac: local_mac(),
96 }
97}
98
99pub 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
120pub 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
182pub trait MigrationStateStore: Send + Sync {
196 fn ensure_schema(&self) -> BoxFuture<'_, Result<()>>;
201
202 fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>>;
204
205 fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>>;
207
208 fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>>;
214
215 fn record_unapplied<'a>(
220 &'a self,
221 app_label: &'a str,
222 name: &'a str,
223 ) -> BoxFuture<'a, Result<()>>;
224
225 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}