pub mod memory;
pub mod schema;
pub mod typedb;
pub use memory::InMemoryStateStore;
pub use schema::{
MigrationStateSchemaKind, applied_migration_entity_label, is_migration_state_type,
migration_state_schema,
};
pub use typedb::{
LEGACY_CUTOVER_SENTINEL_APP_LABEL, LEGACY_CUTOVER_SENTINEL_APPLIED_AT,
LEGACY_CUTOVER_SENTINEL_MIGRATION_ID, LEGACY_CUTOVER_SENTINEL_NAME,
LEGACY_WRITER_CUTOVER_MESSAGE, LegacyCutoverSentinelError, LegacyCutoverSentinelExpectation,
TypeDbStateStore, VerifiedLegacyAppliedPartition, require_legacy_writer_open,
require_legacy_writer_open_in_transaction,
};
use std::net::UdpSocket;
use chrono::Utc;
use network_interface::{NetworkInterface, NetworkInterfaceConfig};
use type_bridge_orm::session::backend::BoxFuture;
use uuid::Uuid;
use crate::plan::{MigrationAction, MigrationExecution};
use crate::{AppliedMigrationRecord, Result};
const TIMESTAMP_FORMAT: &str = "%Y-%m-%dT%H:%M:%S.%6f";
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize, Default)]
pub struct MigrationExecutorInfo {
#[serde(default)]
pub ip: Option<String>,
#[serde(default)]
pub mac: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct MigrationRunRecord {
pub run_id: String,
pub app_label: String,
pub name: String,
pub checksum: String,
pub direction: String,
pub status: String,
pub started_at: String,
#[serde(default)]
pub finished_at: Option<String>,
#[serde(default)]
pub error: Option<String>,
#[serde(default)]
pub executor_ip: Option<String>,
#[serde(default)]
pub executor_mac: Option<String>,
}
pub fn migration_timestamp_now() -> String {
Utc::now().format(TIMESTAMP_FORMAT).to_string()
}
pub fn collect_executor_info() -> MigrationExecutorInfo {
MigrationExecutorInfo {
ip: local_ip(),
mac: local_mac(),
}
}
pub fn started_run_record(
migration: &MigrationExecution,
checksum: String,
executor: &MigrationExecutorInfo,
) -> MigrationRunRecord {
MigrationRunRecord {
run_id: Uuid::new_v4().to_string(),
app_label: migration.app_label.clone(),
name: migration.name.clone(),
checksum,
direction: direction_label(migration.action).to_string(),
status: "started".to_string(),
started_at: migration_timestamp_now(),
finished_at: None,
error: None,
executor_ip: executor.ip.clone(),
executor_mac: executor.mac.clone(),
}
}
pub fn finished_run_record(
mut record: MigrationRunRecord,
status: &str,
error: Option<String>,
) -> MigrationRunRecord {
record.status = status.to_string();
record.finished_at = Some(migration_timestamp_now());
record.error = error;
record
}
fn direction_label(action: MigrationAction) -> &'static str {
match action {
MigrationAction::Apply => "apply",
MigrationAction::Rollback => "rollback",
}
}
fn local_ip() -> Option<String> {
let socket = UdpSocket::bind("0.0.0.0:0").ok()?;
socket.connect("8.8.8.8:80").ok()?;
let addr = socket.local_addr().ok()?;
Some(addr.ip().to_string()).filter(|ip| ip != "0.0.0.0")
}
fn local_mac() -> Option<String> {
NetworkInterface::show()
.ok()?
.into_iter()
.filter(|interface| !interface.internal)
.find_map(|interface| {
interface
.mac_addr
.as_deref()
.and_then(normalize_mac_address)
})
}
fn normalize_mac_address(value: &str) -> Option<String> {
let mut bytes = [0_u8; 6];
let mut parts = value.split([':', '-']);
for byte in &mut bytes {
let part = parts.next()?;
if part.len() != 2 {
return None;
}
*byte = u8::from_str_radix(part, 16).ok()?;
}
if parts.next().is_some() || bytes.iter().all(|byte| *byte == 0) {
return None;
}
Some(format_mac_address(bytes))
}
fn format_mac_address(bytes: [u8; 6]) -> String {
format!(
"{:02x}:{:02x}:{:02x}:{:02x}:{:02x}:{:02x}",
bytes[0], bytes[1], bytes[2], bytes[3], bytes[4], bytes[5]
)
}
pub trait MigrationStateStore: Send + Sync {
fn ensure_schema(&self) -> BoxFuture<'_, Result<()>>;
fn load_applied(&self) -> BoxFuture<'_, Result<Vec<AppliedMigrationRecord>>>;
fn load_runs(&self) -> BoxFuture<'_, Result<Vec<MigrationRunRecord>>>;
fn record_applied(&self, record: AppliedMigrationRecord) -> BoxFuture<'_, Result<()>>;
fn record_unapplied<'a>(
&'a self,
app_label: &'a str,
name: &'a str,
) -> BoxFuture<'a, Result<()>>;
fn record_run(&self, record: MigrationRunRecord) -> BoxFuture<'_, Result<()>>;
}
#[cfg(test)]
mod tests {
use super::{format_mac_address, normalize_mac_address};
#[test]
fn formats_mac_addresses_in_lowercase_colon_notation() {
assert_eq!(
format_mac_address([0x00, 0x11, 0xAB, 0xCD, 0xEF, 0x42]),
"00:11:ab:cd:ef:42"
);
}
#[test]
fn normalizes_supported_mac_address_strings() {
assert_eq!(
normalize_mac_address("00:11:AB:CD:EF:42").as_deref(),
Some("00:11:ab:cd:ef:42")
);
assert_eq!(
normalize_mac_address("00-11-ab-cd-ef-42").as_deref(),
Some("00:11:ab:cd:ef:42")
);
}
#[test]
fn rejects_invalid_or_zero_mac_address_strings() {
assert_eq!(normalize_mac_address("00:00:00:00:00:00"), None);
assert_eq!(normalize_mac_address("00:11:22:33:44"), None);
assert_eq!(normalize_mac_address("00:11:22:33:44:zz"), None);
}
}