use std::fmt;
use serde::{Serialize, Serializer};
use sha2::{Digest as _, Sha256};
use crate::capability::{CapabilityId, CapabilitySet};
use crate::codec::to_canonical_json;
use crate::diagnostic::{Diagnostic, DiagnosticCategory};
use crate::fingerprint::{
CanonicalizationVersion, Fingerprint, FingerprintDigest, FingerprintDomain,
};
use crate::migration_assertion::{
AssertionExpectation, MigrationAssertionPlan, MigrationAssertionPlanFingerprint,
};
use crate::migration_backfill::AttributeBackfillPlan;
use crate::schema_delta::{SchemaDelta, SchemaOperation};
use crate::schema_fingerprint::ManagedSemanticSchemaFingerprint;
pub const MIGRATION_FORMAT_V1: &str = "typebridge.migration/v1";
pub const MIGRATION_ID_FINGERPRINT_DOMAIN: &str = "typebridge.migration.id";
pub const MIGRATION_ID_CANONICALIZATION: &str = "typebridge.migration-id/v1";
pub const SCHEMA_DELTA_FINGERPRINT_DOMAIN: &str = "typebridge.migration.schema-delta";
pub const SCHEMA_DELTA_FINGERPRINT_CANONICALIZATION: &str = "typebridge.schema-delta/v1";
pub const MIGRATION_PLAN_FINGERPRINT_DOMAIN: &str = "typebridge.migration.plan";
pub const MIGRATION_PLAN_FINGERPRINT_CANONICALIZATION: &str = "typebridge.migration-plan/v1";
pub const CONDITIONAL_RESOLUTION_CAPABILITY: &str = "migration.conditional-resolution";
pub const MAX_MIGRATION_COMPONENT_BYTES: usize = 255;
fn validate_component(
value: String,
kind: &'static str,
allow_leading_digit: bool,
) -> Result<String, Diagnostic> {
let mut bytes = value.bytes();
let valid_first = bytes.next().is_some_and(|byte| {
byte.is_ascii_lowercase() || (allow_leading_digit && byte.is_ascii_digit())
});
let valid_rest = bytes.all(|byte| {
byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-')
});
if value.len() <= MAX_MIGRATION_COMPONENT_BYTES && valid_first && valid_rest {
Ok(value)
} else {
Err(Diagnostic::stable(
DiagnosticCategory::InvalidContract,
"invalid_migration_identity_component",
"migration identity component must be bounded portable lowercase ASCII",
)
.with_detail("component_kind", kind))
}
}
macro_rules! migration_component {
($name:ident, $doc:literal, $kind:literal, $leading_digit:expr) => {
#[doc = $doc]
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
#[serde(transparent)]
pub struct $name(String);
impl $name {
pub fn new(value: impl Into<String>) -> Result<Self, Diagnostic> {
Ok(Self(validate_component(
value.into(),
$kind,
$leading_digit,
)?))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for $name {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(self.as_str())
}
}
};
}
migration_component!(
MigrationAppLabel,
"A portable application label forming the first component of migration identity.",
"app_label",
false
);
migration_component!(
MigrationName,
"A portable filename-stem migration name forming the second identity component.",
"name",
true
);
migration_component!(
MigrationStepId,
"A stable identity unique within one ordered migration plan.",
"step_id",
false
);
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
pub struct MigrationId {
app_label: MigrationAppLabel,
name: MigrationName,
}
impl MigrationId {
pub fn new(app_label: impl Into<String>, name: impl Into<String>) -> Result<Self, Diagnostic> {
Ok(Self {
app_label: MigrationAppLabel::new(app_label)?,
name: MigrationName::new(name)?,
})
}
#[must_use]
pub const fn from_components(app_label: MigrationAppLabel, name: MigrationName) -> Self {
Self { app_label, name }
}
pub const fn app_label(&self) -> &MigrationAppLabel {
&self.app_label
}
pub const fn name(&self) -> &MigrationName {
&self.name
}
pub fn canonical_bytes(&self) -> Result<Vec<u8>, Diagnostic> {
to_canonical_json(self)
}
pub fn ledger_key(&self) -> Result<MigrationLedgerKey, Diagnostic> {
MigrationLedgerKey::compute(self)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct MigrationFormat;
impl MigrationFormat {
pub const V1: Self = Self;
pub fn new(value: &str) -> Result<Self, Diagnostic> {
if value == MIGRATION_FORMAT_V1 {
Ok(Self::V1)
} else {
Err(Diagnostic::stable(
DiagnosticCategory::InvalidContract,
"unsupported_migration_format",
"canonical migration format is not supported",
)
.with_detail("actual", value.to_owned())
.with_detail("supported", MIGRATION_FORMAT_V1))
}
}
pub const fn as_str(self) -> &'static str {
MIGRATION_FORMAT_V1
}
}
impl Serialize for MigrationFormat {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RetryPolicy {
Never,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RecoveryPolicy {
OperatorRequired,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize)]
#[serde(transparent)]
pub struct MigrationManifestDigest(FingerprintDigest);
impl MigrationManifestDigest {
#[must_use]
pub fn compute(canonical_manifest_bytes: &[u8]) -> Self {
let digest = Sha256::digest(canonical_manifest_bytes);
let hex = digest
.iter()
.map(|byte| format!("{byte:02x}"))
.collect::<String>();
Self(
FingerprintDigest::from_hex(&hex)
.expect("SHA-256 always produces a valid lowercase 32-byte digest"),
)
}
pub fn from_hex(value: &str) -> Result<Self, Diagnostic> {
FingerprintDigest::from_hex(value).map(Self)
}
#[must_use]
pub fn to_hex(self) -> String {
self.0.to_hex()
}
#[must_use]
pub const fn bytes(self) -> [u8; 32] {
self.0.bytes()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(transparent)]
pub struct MigrationLedgerKey(Fingerprint);
impl MigrationLedgerKey {
pub fn compute(id: &MigrationId) -> Result<Self, Diagnostic> {
Ok(Self(Fingerprint::compute(
FingerprintDomain::new(MIGRATION_ID_FINGERPRINT_DOMAIN)?,
CanonicalizationVersion::new(MIGRATION_ID_CANONICALIZATION)?,
None,
&id.canonical_bytes()?,
)))
}
pub const fn as_fingerprint(&self) -> &Fingerprint {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(transparent)]
pub struct SchemaDeltaFingerprint(Fingerprint);
impl SchemaDeltaFingerprint {
pub fn compute(delta: &SchemaDelta) -> Result<Self, Diagnostic> {
Ok(Self(Fingerprint::compute(
FingerprintDomain::new(SCHEMA_DELTA_FINGERPRINT_DOMAIN)?,
CanonicalizationVersion::new(SCHEMA_DELTA_FINGERPRINT_CANONICALIZATION)?,
None,
&delta.canonical_bytes()?,
)))
}
pub const fn as_fingerprint(&self) -> &Fingerprint {
&self.0
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum MigrationStepKind {
SchemaDelta,
Assertion,
Backfill,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct BackfillStepContract {
id: MigrationStepId,
plan_fingerprint: Fingerprint,
recovery: RecoveryPolicy,
required_capabilities: CapabilitySet,
retry: RetryPolicy,
#[serde(skip_serializing_if = "Option::is_none")]
reverse: Option<crate::migration_backfill::BackfillReverseProgram>,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
}
impl BackfillStepContract {
fn derive(id: MigrationStepId, plan: &AttributeBackfillPlan) -> Result<Self, Diagnostic> {
Ok(Self {
id,
plan_fingerprint: plan.fingerprint()?,
recovery: RecoveryPolicy::OperatorRequired,
required_capabilities: plan.required_capabilities().clone(),
retry: RetryPolicy::Never,
reverse: plan.reverse(),
source_semantics: plan.managed_semantics().clone(),
target_semantics: plan.managed_semantics().clone(),
})
}
pub const fn id(&self) -> &MigrationStepId {
&self.id
}
pub const fn plan_fingerprint(&self) -> &Fingerprint {
&self.plan_fingerprint
}
pub const fn required_capabilities(&self) -> &CapabilitySet {
&self.required_capabilities
}
pub const fn reverse(&self) -> Option<crate::migration_backfill::BackfillReverseProgram> {
self.reverse
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct SchemaDeltaStepContract {
delta_fingerprint: SchemaDeltaFingerprint,
id: MigrationStepId,
recovery: RecoveryPolicy,
required_capabilities: CapabilitySet,
retry: RetryPolicy,
#[serde(skip_serializing_if = "Option::is_none")]
reverse: Option<SchemaDelta>,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
}
impl SchemaDeltaStepContract {
pub fn new(
id: MigrationStepId,
delta: &SchemaDelta,
reverse: Option<SchemaDelta>,
) -> Result<Self, Diagnostic> {
if let Some(candidate) = &reverse {
let inverse_operations = delta
.operations()
.iter()
.rev()
.flat_map(SchemaOperation::inverse)
.collect();
let expected = SchemaDelta::new(
delta.format(),
delta.target().clone(),
delta.source().clone(),
inverse_operations,
)?;
if candidate != &expected {
return Err(Diagnostic::stable(
DiagnosticCategory::InvalidContract,
"schema_delta_step_inverse_mismatch",
"schema step reverse delta is not the exact operation-wise inverse",
));
}
}
Ok(Self {
delta_fingerprint: SchemaDeltaFingerprint::compute(delta)?,
id,
recovery: RecoveryPolicy::OperatorRequired,
required_capabilities: delta.required_capabilities().clone(),
retry: RetryPolicy::Never,
reverse,
source_semantics: delta.source().managed_semantic_schema().clone(),
target_semantics: delta.target().managed_semantic_schema().clone(),
})
}
pub const fn id(&self) -> &MigrationStepId {
&self.id
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
pub const fn required_capabilities(&self) -> &CapabilitySet {
&self.required_capabilities
}
pub const fn delta_fingerprint(&self) -> &SchemaDeltaFingerprint {
&self.delta_fingerprint
}
pub const fn retry(&self) -> RetryPolicy {
self.retry
}
pub const fn recovery(&self) -> RecoveryPolicy {
self.recovery
}
pub const fn reverse(&self) -> Option<&SchemaDelta> {
self.reverse.as_ref()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct SchemaDeltaStep {
contract: SchemaDeltaStepContract,
delta: SchemaDelta,
kind: MigrationStepKind,
}
impl SchemaDeltaStep {
pub fn new(
id: MigrationStepId,
delta: SchemaDelta,
reverse: Option<SchemaDelta>,
) -> Result<Self, Diagnostic> {
let contract = SchemaDeltaStepContract::new(id, &delta, reverse)?;
Ok(Self {
contract,
delta,
kind: MigrationStepKind::SchemaDelta,
})
}
pub const fn kind(&self) -> MigrationStepKind {
self.kind
}
pub const fn contract(&self) -> &SchemaDeltaStepContract {
&self.contract
}
pub const fn delta(&self) -> &SchemaDelta {
&self.delta
}
pub fn canonical_bytes(&self) -> Result<Vec<u8>, Diagnostic> {
to_canonical_json(self)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct AssertionStepContract {
id: MigrationStepId,
plan_fingerprint: MigrationAssertionPlanFingerprint,
recovery: RecoveryPolicy,
required_capabilities: CapabilitySet,
retry: RetryPolicy,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
}
impl AssertionStepContract {
fn derive(
id: MigrationStepId,
plan: &MigrationAssertionPlan,
expected: AssertionExpectation,
) -> Result<Self, Diagnostic> {
if expected != AssertionExpectation::NoRows
|| plan.expectation() != AssertionExpectation::NoRows
{
return Err(Diagnostic::stable(
DiagnosticCategory::InvalidContract,
"migration_assertion_step_expectation_mismatch",
"persisted migration assertions support only the no-rows expectation",
));
}
let mut required_capabilities = plan.required_capabilities().clone();
required_capabilities.insert(
CapabilityId::new(CONDITIONAL_RESOLUTION_CAPABILITY)
.expect("the fixed conditional-resolution capability is canonical"),
);
Ok(Self {
id,
plan_fingerprint: plan.fingerprint()?,
recovery: RecoveryPolicy::OperatorRequired,
required_capabilities,
retry: RetryPolicy::Never,
source_semantics: plan.managed_semantics().clone(),
target_semantics: plan.managed_semantics().clone(),
})
}
pub const fn id(&self) -> &MigrationStepId {
&self.id
}
pub const fn plan_fingerprint(&self) -> &MigrationAssertionPlanFingerprint {
&self.plan_fingerprint
}
pub const fn required_capabilities(&self) -> &CapabilitySet {
&self.required_capabilities
}
pub const fn retry(&self) -> RetryPolicy {
self.retry
}
pub const fn recovery(&self) -> RecoveryPolicy {
self.recovery
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
pub const fn reverse(&self) -> Option<()> {
None
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MigrationStep {
SchemaDelta(Box<SchemaDeltaStep>),
Assertion {
contract: Box<AssertionStepContract>,
plan: Box<MigrationAssertionPlan>,
expected: AssertionExpectation,
},
Backfill {
contract: Box<BackfillStepContract>,
plan: Box<AttributeBackfillPlan>,
},
}
impl MigrationStep {
pub fn backfill(id: MigrationStepId, plan: AttributeBackfillPlan) -> Result<Self, Diagnostic> {
let contract = BackfillStepContract::derive(id, &plan)?;
Ok(Self::Backfill {
contract: Box::new(contract),
plan: Box::new(plan),
})
}
pub fn assertion(
id: MigrationStepId,
plan: MigrationAssertionPlan,
expected: AssertionExpectation,
) -> Result<Self, Diagnostic> {
let contract = AssertionStepContract::derive(id, &plan, expected)?;
Ok(Self::Assertion {
contract: Box::new(contract),
plan: Box::new(plan),
expected,
})
}
pub const fn kind(&self) -> MigrationStepKind {
match self {
Self::SchemaDelta(_) => MigrationStepKind::SchemaDelta,
Self::Assertion { .. } => MigrationStepKind::Assertion,
Self::Backfill { .. } => MigrationStepKind::Backfill,
}
}
pub const fn id(&self) -> &MigrationStepId {
match self {
Self::SchemaDelta(step) => step.contract().id(),
Self::Assertion { contract, .. } => contract.id(),
Self::Backfill { contract, .. } => contract.id(),
}
}
pub const fn required_capabilities(&self) -> &CapabilitySet {
match self {
Self::SchemaDelta(step) => step.contract().required_capabilities(),
Self::Assertion { contract, .. } => contract.required_capabilities(),
Self::Backfill { contract, .. } => contract.required_capabilities(),
}
}
pub fn as_schema_delta(&self) -> Option<&SchemaDeltaStep> {
match self {
Self::SchemaDelta(step) => Some(step),
Self::Assertion { .. } | Self::Backfill { .. } => None,
}
}
pub fn as_assertion(
&self,
) -> Option<(
&AssertionStepContract,
&MigrationAssertionPlan,
AssertionExpectation,
)> {
match self {
Self::SchemaDelta(_) | Self::Backfill { .. } => None,
Self::Assertion {
contract,
plan,
expected,
} => Some((contract, plan, *expected)),
}
}
pub fn as_backfill(&self) -> Option<(&BackfillStepContract, &AttributeBackfillPlan)> {
match self {
Self::Backfill { contract, plan } => Some((contract, plan)),
Self::SchemaDelta(_) | Self::Assertion { .. } => None,
}
}
pub fn validate(&self) -> Result<(), Diagnostic> {
let rebuilt = match self {
Self::SchemaDelta(step) => Self::SchemaDelta(Box::new(SchemaDeltaStep::new(
step.contract().id().clone(),
step.delta().clone(),
step.contract().reverse().cloned(),
)?)),
Self::Assertion {
contract,
plan,
expected,
} => Self::assertion(contract.id().clone(), plan.as_ref().clone(), *expected)?,
Self::Backfill { contract, plan } => {
Self::backfill(contract.id().clone(), plan.as_ref().clone())?
}
};
if &rebuilt != self {
return Err(Diagnostic::stable(
DiagnosticCategory::Integrity,
"migration_step_contract_mismatch",
"migration step claims differ from constructor-derived claims",
));
}
Ok(())
}
pub fn canonical_bytes(&self) -> Result<Vec<u8>, Diagnostic> {
to_canonical_json(self)
}
}
impl From<SchemaDeltaStep> for MigrationStep {
fn from(step: SchemaDeltaStep) -> Self {
Self::SchemaDelta(Box::new(step))
}
}
#[derive(Serialize)]
struct AssertionStepView<'a> {
contract: &'a AssertionStepContract,
expected: AssertionExpectation,
kind: MigrationStepKind,
plan: &'a MigrationAssertionPlan,
}
#[derive(Serialize)]
struct BackfillStepView<'a> {
contract: &'a BackfillStepContract,
kind: MigrationStepKind,
plan: &'a AttributeBackfillPlan,
}
impl Serialize for MigrationStep {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
match self {
Self::SchemaDelta(step) => step.serialize(serializer),
Self::Assertion {
contract,
plan,
expected,
} => AssertionStepView {
contract,
expected: *expected,
kind: MigrationStepKind::Assertion,
plan,
}
.serialize(serializer),
Self::Backfill { contract, plan } => BackfillStepView {
contract,
kind: MigrationStepKind::Backfill,
plan,
}
.serialize(serializer),
}
}
}
#[derive(Serialize)]
struct MigrationPlanView<'a> {
steps: &'a [MigrationStep],
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
#[serde(transparent)]
pub struct MigrationPlanFingerprint(Fingerprint);
impl MigrationPlanFingerprint {
pub fn canonical_plan_bytes(steps: &[MigrationStep]) -> Result<Vec<u8>, Diagnostic> {
to_canonical_json(&MigrationPlanView { steps })
}
pub fn compute(steps: &[MigrationStep]) -> Result<Self, Diagnostic> {
Ok(Self(Fingerprint::compute(
FingerprintDomain::new(MIGRATION_PLAN_FINGERPRINT_DOMAIN)?,
CanonicalizationVersion::new(MIGRATION_PLAN_FINGERPRINT_CANONICALIZATION)?,
None,
&Self::canonical_plan_bytes(steps)?,
)))
}
pub const fn as_fingerprint(&self) -> &Fingerprint {
&self.0
}
}