use std::collections::BTreeSet;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use type_bridge_contract::capability::CapabilitySet;
use type_bridge_contract::codec::{from_canonical_json, to_canonical_json};
use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
use type_bridge_contract::limits::{MAX_CANONICAL_COLLECTION_LEN, StructuralLimits};
use type_bridge_contract::managed_scope::{ManagedScopeBinding, SemanticProfileBinding};
use type_bridge_contract::migration::{
MigrationAppLabel, MigrationFormat, MigrationId, MigrationManifestDigest, MigrationName,
MigrationPlanFingerprint, MigrationStep, MigrationStepId, SchemaDeltaStep,
};
use type_bridge_contract::migration_assertion::{
AssertionExpectation, decode_migration_assertion_plan,
};
use type_bridge_contract::migration_backfill::decode_attribute_backfill_plan;
use type_bridge_contract::schema::{
DeclaredIdentityFingerprint, DeclaredSchema, decode_schema_delta,
};
use type_bridge_contract::schema_delta::ManagedSchemaState;
use type_bridge_contract::schema_fingerprint::{
ManagedDeclaredIdentityFingerprint, ManagedSemanticSchemaFingerprint,
};
use type_bridge_contract::schema_lowering::SchemaLoweringProfileBinding;
use type_bridge_query::{
MigrationAssertionValidationContext, ValidatedMigrationAssertionPlan, lower_condition_to_plan,
};
use type_bridge_schema::{
DeltaError, ManagedDeltaContext, RequiredSafetyCondition, SafetyClass, SafetyCondition,
SafetyConditionDomainIndex, SafetyDerivationProfile, apply_delta,
classify_schema_operation_safety, derive_safety_conditions_with_domain_index,
managed_schema_state, plan_schema_operations, resolve,
};
use crate::legacy::{
LEGACY_APPLIED_SET_ALGORITHM, LEGACY_APPLIED_SET_CANONICALIZATION, LEGACY_CHECKSUM_ALGORITHM,
LegacyAppliedSetDigest, LegacyMigrationChecksum, LegacyMigrationId, LegacyMigrationReference,
};
use crate::profile::schema_lowering_profile_binding;
const MANIFEST_SCHEMA_CANONICALIZATION: &str = "typebridge.schema-c14n/v2";
const MANIFEST_CODEC: &str = "typebridge.canonical-json/v1";
const MANIFEST_DELTA_IR: &str = "typebridge.schema-delta/v1";
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct SchemaMigrationDraft {
id: MigrationId,
parents: Vec<MigrationId>,
steps: Vec<MigrationStep>,
legacy_parents: Vec<LegacyMigrationReference>,
legacy_applied_set: Option<LegacyAppliedSetDigest>,
}
impl SchemaMigrationDraft {
pub fn new<S>(
id: MigrationId,
mut parents: Vec<MigrationId>,
steps: Vec<S>,
) -> Result<Self, Diagnostic>
where
S: Into<MigrationStep>,
{
let steps = steps.into_iter().map(Into::into).collect::<Vec<_>>();
if steps.len() > MAX_CANONICAL_COLLECTION_LEN {
return Err(failure(
DiagnosticCategory::ResourceLimit,
"migration_manifest_step_limit",
"migration draft exceeds the canonical step-count ceiling",
));
}
parents.sort();
if parents.windows(2).any(|pair| pair[0] == pair[1]) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_duplicate_parent",
"migration draft contains a duplicate parent identity",
));
}
if parents.iter().any(|parent| parent == &id) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_self_parent",
"migration draft cannot name itself as a parent",
));
}
let mut step_ids = BTreeSet::new();
for step in &steps {
step.validate()?;
if !step_ids.insert(step.id().clone()) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_duplicate_step_id",
"migration draft contains a duplicate step identity",
));
}
}
Ok(Self {
id,
parents,
steps,
legacy_parents: Vec::new(),
legacy_applied_set: None,
})
}
pub fn legacy_bridge(
id: MigrationId,
mut legacy_parents: Vec<LegacyMigrationReference>,
legacy_applied_set: LegacyAppliedSetDigest,
) -> Result<Self, Diagnostic> {
if legacy_parents.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_empty_legacy_frontier",
"a legacy-frontier bridge must name at least one legacy parent",
));
}
if legacy_parents.len() > MAX_CANONICAL_COLLECTION_LEN {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_frontier_too_large",
"legacy-frontier bridge exceeds the canonical collection ceiling",
));
}
legacy_parents.sort();
if legacy_parents
.windows(2)
.any(|pair| pair[0].id() == pair[1].id())
{
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_duplicate_legacy_parent",
"legacy-frontier bridge names a legacy identity twice",
));
}
Ok(Self {
id,
parents: Vec::new(),
steps: Vec::new(),
legacy_parents,
legacy_applied_set: Some(legacy_applied_set),
})
}
pub const fn id(&self) -> &MigrationId {
&self.id
}
pub fn parents(&self) -> &[MigrationId] {
&self.parents
}
pub fn steps(&self) -> &[MigrationStep] {
&self.steps
}
pub fn legacy_parents(&self) -> &[LegacyMigrationReference] {
&self.legacy_parents
}
pub const fn legacy_applied_set(&self) -> Option<&LegacyAppliedSetDigest> {
self.legacy_applied_set.as_ref()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct VerifiedSchemaMigrationManifest {
format: MigrationFormat,
id: MigrationId,
legacy_parents: Vec<LegacyMigrationReference>,
legacy_applied_set: Option<LegacyAppliedSetDigest>,
lowering_profile: SchemaLoweringProfileBinding,
managed_scope: ManagedScopeBinding,
parents: Vec<MigrationId>,
plan_fingerprint: MigrationPlanFingerprint,
required_capabilities: CapabilitySet,
reversible: bool,
safety: SafetyClass,
semantic_profile: SemanticProfileBinding,
source_schema: DeclaredSchema,
source_state: ManagedSchemaState,
steps: Vec<MigrationStep>,
target_schema: DeclaredSchema,
target_state: ManagedSchemaState,
}
impl VerifiedSchemaMigrationManifest {
pub const fn format(&self) -> &MigrationFormat {
&self.format
}
pub const fn id(&self) -> &MigrationId {
&self.id
}
pub fn parents(&self) -> &[MigrationId] {
&self.parents
}
pub fn legacy_parents(&self) -> &[LegacyMigrationReference] {
&self.legacy_parents
}
pub const fn legacy_applied_set(&self) -> Option<&LegacyAppliedSetDigest> {
self.legacy_applied_set.as_ref()
}
pub fn is_legacy_bridge(&self) -> bool {
!self.legacy_parents.is_empty()
}
pub fn steps(&self) -> &[MigrationStep] {
&self.steps
}
pub const fn managed_scope(&self) -> &ManagedScopeBinding {
&self.managed_scope
}
pub const fn semantic_profile(&self) -> &SemanticProfileBinding {
&self.semantic_profile
}
pub const fn lowering_profile(&self) -> &SchemaLoweringProfileBinding {
&self.lowering_profile
}
pub const fn required_capabilities(&self) -> &CapabilitySet {
&self.required_capabilities
}
pub const fn safety(&self) -> SafetyClass {
self.safety
}
pub const fn reversible(&self) -> bool {
self.reversible
}
pub const fn plan_fingerprint(&self) -> &MigrationPlanFingerprint {
&self.plan_fingerprint
}
pub const fn source_state(&self) -> &ManagedSchemaState {
&self.source_state
}
pub const fn source_schema(&self) -> &DeclaredSchema {
&self.source_schema
}
pub const fn target_state(&self) -> &ManagedSchemaState {
&self.target_state
}
pub const fn target_schema(&self) -> &DeclaredSchema {
&self.target_schema
}
}
pub fn build_verified_manifest(
draft: SchemaMigrationDraft,
context: (&DeclaredSchema, &ManagedDeltaContext),
) -> Result<VerifiedSchemaMigrationManifest, Diagnostic> {
let (source_schema, delta_context) = context;
let source_state =
managed_schema_state(source_schema, delta_context).map_err(delta_diagnostic)?;
let managed_scope = ManagedScopeBinding::exclusive(delta_context.scope_id().clone())?;
if source_state.scope() != &managed_scope {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_scope_mismatch",
"managed source state does not match the verification scope binding",
));
}
let semantic_profile =
SemanticProfileBinding::resolve(delta_context.semantic_profile().clone())?;
let lowering_profile = schema_lowering_profile_binding()?;
let safety_profile =
SafetyDerivationProfile::new(semantic_profile.clone(), lowering_profile.clone())?;
let SchemaMigrationDraft {
id,
parents,
steps,
legacy_parents,
legacy_applied_set,
} = draft;
if legacy_parents.is_empty() {
if legacy_applied_set.is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_applied_set_without_bridge",
"an ordinary migration cannot carry a legacy applied-set digest",
));
}
if steps.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_empty_program",
"a migration without schema steps is valid only as a legacy-frontier bridge",
));
}
} else {
if legacy_applied_set.is_none() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_applied_set_missing",
"a legacy-frontier bridge requires its complete applied-set digest",
));
}
if !steps.is_empty() || !parents.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_bridge_not_zero_operation",
"a legacy-frontier bridge carries no steps and no canonical parents",
));
}
}
let mut current_schema = source_schema.clone();
let mut required_capabilities = source_schema.required_capabilities().clone();
let mut safety = SafetyClass::FormalOnly;
let mut reversible = true;
let mut pending_assertions = Vec::new();
for step in &steps {
if let Some((contract, plan)) = step.as_backfill() {
if !pending_assertions.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_assertion_before_backfill",
"assertion steps must be immediately followed by a schema delta",
));
}
step.validate()?;
step.required_capabilities()
.ensure_supported_by(delta_context.available_capabilities())?;
let intermediate_state =
managed_schema_state(¤t_schema, delta_context).map_err(delta_diagnostic)?;
if plan.managed_semantics() != intermediate_state.managed_semantic_schema()
|| contract.source_semantics() != intermediate_state.managed_semantic_schema()
|| contract.target_semantics() != intermediate_state.managed_semantic_schema()
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_backfill_schema_mismatch",
"backfill plan does not bind the exact historical intermediate schema",
));
}
validate_backfill_historical_schema(plan, ¤t_schema, delta_context)?;
for capability in step.required_capabilities().iter().cloned() {
required_capabilities.insert(capability);
}
safety = safety.max(SafetyClass::BackfillRequired);
reversible &= contract.reverse().is_some();
continue;
}
let Some(schema_step) = step.as_schema_delta() else {
step.validate()?;
step.required_capabilities()
.ensure_supported_by(delta_context.available_capabilities())?;
pending_assertions.push(step);
continue;
};
let delta = schema_step.delta();
if delta.source().scope() != &managed_scope || delta.target().scope() != &managed_scope {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_scope_mismatch",
"schema step crosses the verified managed scope lineage",
));
}
delta
.required_capabilities()
.ensure_supported_by(delta_context.available_capabilities())?;
let target_schema = apply_delta(¤t_schema, delta, delta_context).map_err(|_| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_step_chain_mismatch",
"schema step source does not chain from the preceding verified target",
)
})?;
let planned = plan_schema_operations(¤t_schema, &target_schema).map_err(|_| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_dependency_plan_invalid",
"schema step cannot be reproduced by the dependency planner",
)
})?;
if planned != delta.operations() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_dependency_plan_mismatch",
"schema step operations are not in the canonical dependency plan",
));
}
let coverage = verify_assertion_coverage(
&pending_assertions,
delta,
¤t_schema,
&target_schema,
&safety_profile,
)?;
for assertion in &pending_assertions {
for capability in assertion.required_capabilities().iter().cloned() {
required_capabilities.insert(capability);
}
}
pending_assertions.clear();
safety = safety.max(coverage.effective_safety());
for capability in delta.required_capabilities().iter().cloned() {
required_capabilities.insert(capability);
}
if let Some(reverse) = schema_step.contract().reverse() {
let restored = apply_delta(&target_schema, reverse, delta_context).map_err(|_| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_inverse_replay_mismatch",
"schema step inverse does not replay from its verified target",
)
})?;
let planned_reverse =
plan_schema_operations(&target_schema, &restored).map_err(|_| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_inverse_plan_invalid",
"schema step inverse has no dependency-safe plan",
)
})?;
if planned_reverse != reverse.operations()
|| restored.canonical_identity_bytes()?
!= current_schema.canonical_identity_bytes()?
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_inverse_replay_mismatch",
"schema step inverse does not restore the exact declared source",
));
}
reject_reverse_assertion_requirement(
reverse,
&target_schema,
&restored,
&safety_profile,
)?;
} else {
reversible = false;
}
current_schema = target_schema;
}
if !pending_assertions.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_orphan_assertion",
"assertion steps must be immediately followed by a schema delta",
));
}
let target_state =
managed_schema_state(¤t_schema, delta_context).map_err(delta_diagnostic)?;
let plan_fingerprint = MigrationPlanFingerprint::compute(&steps)?;
Ok(VerifiedSchemaMigrationManifest {
format: MigrationFormat::V1,
id,
legacy_parents,
legacy_applied_set,
lowering_profile,
managed_scope,
parents,
plan_fingerprint,
required_capabilities,
reversible,
safety,
semantic_profile,
source_schema: source_schema.clone(),
source_state,
steps,
target_schema: current_schema,
target_state,
})
}
fn validate_backfill_historical_schema(
plan: &type_bridge_contract::migration_backfill::AttributeBackfillPlan,
schema: &DeclaredSchema,
context: &ManagedDeltaContext,
) -> Result<(), Diagnostic> {
let resolved = resolve(schema, context.semantic_profile()).map_err(|diagnostics| {
diagnostics
.iter()
.next()
.map(|diagnostic| diagnostic.diagnostic().clone())
.unwrap_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_backfill_resolution_failed",
"backfill historical schema resolution failed without a diagnostic",
)
})
})?;
let owner = resolved.types().get(plan.owner()).ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_backfill_owner_missing",
"backfill owner does not exist in the historical intermediate schema",
)
})?;
for (role, attribute) in [
("source", plan.source()),
("destination", plan.destination()),
("partition", plan.partition().stable_attribute()),
] {
if !owner.owns().contains_key(attribute) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_backfill_attribute_not_owned",
"backfill attribute is not effectively owned by the historical owner",
)
.with_detail("attribute_role", role)
.with_detail("attribute", attribute.label().as_str().to_owned()));
}
}
if !owner
.key_attributes()
.contains(plan.partition().stable_attribute())
{
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_backfill_partition_not_key",
"backfill partition attribute must be an effective key of the historical owner",
));
}
Ok(())
}
pub fn decode_verified_manifest(
bytes: &[u8],
context: (&DeclaredSchema, &ManagedDeltaContext),
) -> Result<VerifiedSchemaMigrationManifest, Diagnostic> {
let candidate = from_canonical_json::<ManifestCandidate>(bytes)?;
candidate.validate_header()?;
let draft = candidate.to_draft()?;
let verified = build_verified_manifest(draft, context)?;
if encode_verified_manifest(&verified)? != bytes {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_verification_mismatch",
"manifest claims do not equal the replay-derived verified encoding",
));
}
Ok(verified)
}
pub(crate) fn peek_manifest_identity(
bytes: &[u8],
) -> Result<(MigrationId, Vec<MigrationId>), Diagnostic> {
let candidate = from_canonical_json::<ManifestCandidate>(bytes)?;
candidate.validate_header()?;
let id = candidate.id.rebuild()?;
let parents = candidate
.parents
.iter()
.map(MigrationIdCandidate::rebuild)
.collect::<Result<Vec<_>, _>>()?;
Ok((id, parents))
}
pub(crate) fn peek_manifest_declares_legacy_bridge(bytes: &[u8]) -> Result<bool, Diagnostic> {
let candidate = from_canonical_json::<ManifestCandidate>(bytes)?;
candidate.validate_header()?;
Ok(!candidate.legacy_parents.is_empty())
}
pub fn encode_verified_manifest(
manifest: &VerifiedSchemaMigrationManifest,
) -> Result<Vec<u8>, Diagnostic> {
to_canonical_json(&ManifestWire::from_verified(manifest))
}
pub fn verified_manifest_digest(
manifest: &VerifiedSchemaMigrationManifest,
) -> Result<MigrationManifestDigest, Diagnostic> {
Ok(MigrationManifestDigest::compute(&encode_verified_manifest(
manifest,
)?))
}
fn verified_forward_safety(
safety: SafetyClass,
verifier_resolved: bool,
) -> Result<SafetyClass, Diagnostic> {
if verifier_resolved {
return match safety {
SafetyClass::Conditional => Ok(SafetyClass::Conditional),
SafetyClass::BackfillRequired => Ok(SafetyClass::Conditional),
_ => Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_invalid_resolved_safety",
"verifier resolution targets an operation outside the resolvable safety classes",
)),
};
}
match safety {
SafetyClass::FormalOnly
| SafetyClass::SchemaMetadata
| SafetyClass::Additive
| SafetyClass::Conditional
| SafetyClass::Destructive => Ok(safety),
SafetyClass::BackfillRequired | SafetyClass::Opaque | SafetyClass::Unsupported => {
Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_unresolved_safety",
"migration manifest cannot carry unresolved backfill, opaque, or unsupported work",
))
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(crate) struct VerifiedAssertionCoverage {
discharged_operation_indices: Vec<usize>,
effective_safety: SafetyClass,
validated: Vec<ValidatedMigrationAssertionPlan>,
}
impl VerifiedAssertionCoverage {
pub(crate) fn discharged_operation_indices(&self) -> &[usize] {
&self.discharged_operation_indices
}
pub(crate) const fn effective_safety(&self) -> SafetyClass {
self.effective_safety
}
pub(crate) fn validated(&self) -> &[ValidatedMigrationAssertionPlan] {
&self.validated
}
}
pub(crate) fn verify_assertion_coverage(
assertions: &[&MigrationStep],
delta: &type_bridge_contract::schema::SchemaDelta,
source: &DeclaredSchema,
target: &DeclaredSchema,
profile: &SafetyDerivationProfile,
) -> Result<VerifiedAssertionCoverage, Diagnostic> {
let mut required = Vec::new();
let mut with_destructive_guards = Vec::new();
let mut discharged_operation_indices = Vec::new();
let mut effective_safety = SafetyClass::FormalOnly;
let domain = SafetyConditionDomainIndex::new(source, target);
for (operation_index, operation) in delta.operations().iter().enumerate() {
let mut operation_requires_assertion = false;
let derived = derive_safety_conditions_with_domain_index(
operation_index,
operation,
source,
target,
profile,
&domain,
)?;
for condition in derived.conditions() {
match condition.policy() {
SafetyClass::Conditional => {
operation_requires_assertion = true;
if matches!(condition.condition(), SafetyCondition::Unresolvable { .. }) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_unresolvable_conditional_assertion",
"conditional schema work has no canonical assertion representation",
));
}
required.push(condition.clone());
with_destructive_guards.push(condition.clone());
}
SafetyClass::Destructive if condition.condition().is_resolvable() => {
with_destructive_guards.push(condition.clone());
}
_ => {}
}
}
let discharged = operation_requires_assertion
|| derived.policy() == SafetyClass::Conditional
|| (derived.policy() == SafetyClass::BackfillRequired && derived.is_condition_free());
if discharged {
discharged_operation_indices.push(operation_index);
}
effective_safety = effective_safety.max(verified_forward_safety(
classify_schema_operation_safety(operation),
discharged,
)?);
}
let expected: &[RequiredSafetyCondition] = if assertions.is_empty() {
if !required.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_missing_assertion",
"conditional schema work is missing verifier-derived assertions",
));
}
&[]
} else if assertions.len() == required.len() {
&required
} else if assertions.len() == with_destructive_guards.len() {
&with_destructive_guards
} else {
return Err(failure(
DiagnosticCategory::InvalidContract,
if assertions.len() < required.len() {
"migration_manifest_missing_assertion"
} else {
"migration_manifest_extra_assertion"
},
"assertion count does not equal canonical verifier-derived coverage",
));
};
if expected.is_empty() && !assertions.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_extra_assertion",
"schema delta has no verifier-derived assertion requirement",
));
}
let resolved = resolve(source, profile.semantic().id()).map_err(|diagnostics| {
diagnostics
.iter()
.next()
.map(|diagnostic| diagnostic.diagnostic().clone())
.unwrap_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_assertion_resolution_failed",
"assertion source resolution failed without a diagnostic",
)
})
})?;
let context = MigrationAssertionValidationContext::new(&resolved, delta.source());
let mut validated_plans = Vec::with_capacity(expected.len());
for (actual, condition) in assertions.iter().zip(expected) {
let validated = lower_condition_to_plan(condition, &context, StructuralLimits::CANONICAL)?;
let (contract, plan, expected) = actual.as_assertion().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_assertion_order_mismatch",
"assertion coverage contains a non-assertion step",
)
})?;
if expected != AssertionExpectation::NoRows
|| plan.canonical_bytes()? != validated.plan().canonical_bytes()?
|| plan.fingerprint()? != validated.plan().fingerprint()?
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_assertion_plan_mismatch",
"persisted assertion does not equal verifier-derived canonical plan",
));
}
let rebuilt = MigrationStep::assertion(
contract.id().clone(),
validated.plan().clone(),
AssertionExpectation::NoRows,
)?;
if &rebuilt != *actual {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_assertion_contract_mismatch",
"persisted assertion contract differs from verifier-derived claims",
));
}
validated_plans.push(validated);
}
Ok(VerifiedAssertionCoverage {
discharged_operation_indices,
effective_safety,
validated: validated_plans,
})
}
fn reject_reverse_assertion_requirement(
reverse: &type_bridge_contract::schema::SchemaDelta,
source: &DeclaredSchema,
target: &DeclaredSchema,
profile: &SafetyDerivationProfile,
) -> Result<(), Diagnostic> {
match verify_assertion_coverage(&[], reverse, source, target, profile) {
Ok(_) => Ok(()),
Err(error)
if matches!(
error.code().as_str(),
"migration_manifest_missing_assertion"
| "migration_manifest_unresolvable_conditional_assertion"
) =>
{
Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_reverse_requires_assertions",
"claimed reverse requires assertions that are not represented",
))
}
Err(error) if error.code().as_str() == "migration_manifest_unresolved_safety" => {
Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_reverse_unresolved_safety",
"claimed reverse has unresolved non-assertion migration work",
))
}
Err(error) => Err(error),
}
}
pub(crate) fn delta_diagnostic(error: DeltaError) -> Diagnostic {
match error {
DeltaError::Contract(diagnostic) => diagnostic,
DeltaError::Schema(diagnostics) => diagnostics
.iter()
.next()
.map(|diagnostic| diagnostic.diagnostic().clone())
.unwrap_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_manifest_schema_verification_failed",
"schema verification failed without a diagnostic",
)
}),
}
}
fn failure(category: DiagnosticCategory, code: &'static str, message: &'static str) -> Diagnostic {
Diagnostic::new(
category,
DiagnosticCode::new(code).expect("static manifest diagnostic code is canonical"),
message,
)
}
#[derive(Serialize)]
struct ManifestWire<'a> {
contract: ManifestContractWire<'a>,
fingerprints: ManifestFingerprintsWire<'a>,
format: &'a MigrationFormat,
id: &'a MigrationId,
#[serde(skip_serializing_if = "Option::is_none")]
legacy_applied_set: Option<LegacyAppliedSetWire<'a>>,
#[serde(skip_serializing_if = "Vec::is_empty")]
legacy_parents: Vec<LegacyParentWire<'a>>,
managed_scope: &'a ManagedScopeBinding,
parents: &'a [MigrationId],
required_capabilities: &'a CapabilitySet,
resources: &'static [()],
safety: ManifestSafetyWire,
steps: &'a [MigrationStep],
}
impl<'a> ManifestWire<'a> {
fn from_verified(manifest: &'a VerifiedSchemaMigrationManifest) -> Self {
Self {
contract: ManifestContractWire {
canonicalization: MANIFEST_SCHEMA_CANONICALIZATION,
codec: MANIFEST_CODEC,
delta_ir: MANIFEST_DELTA_IR,
lowering_profile: &manifest.lowering_profile,
semantic_profile: &manifest.semantic_profile,
},
fingerprints: ManifestFingerprintsWire {
plan: &manifest.plan_fingerprint,
source: ManifestEndpointFingerprintsWire {
declared_identity: manifest.source_state.managed_declared_identity(),
resolution_identity: manifest.source_state.declared_identity(),
semantics: manifest.source_state.managed_semantic_schema(),
},
target: ManifestEndpointFingerprintsWire {
declared_identity: manifest.target_state.managed_declared_identity(),
resolution_identity: manifest.target_state.declared_identity(),
semantics: manifest.target_state.managed_semantic_schema(),
},
},
format: &manifest.format,
id: &manifest.id,
legacy_applied_set: manifest.legacy_applied_set.as_ref().map(|digest| {
LegacyAppliedSetWire {
algorithm: digest.algorithm(),
canonicalization: digest.canonicalization(),
digest: digest.as_str(),
}
}),
legacy_parents: manifest
.legacy_parents
.iter()
.map(|reference| LegacyParentWire {
app_label: reference.id().app_label().as_str(),
checksum: LegacyChecksumWire {
algorithm: reference.checksum().algorithm(),
value: reference.checksum().as_str(),
},
name: reference.id().name().as_str(),
})
.collect(),
managed_scope: &manifest.managed_scope,
parents: &manifest.parents,
required_capabilities: &manifest.required_capabilities,
resources: &[],
safety: ManifestSafetyWire {
classification: manifest.safety,
reversible: manifest.reversible,
},
steps: &manifest.steps,
}
}
}
#[derive(Serialize)]
struct ManifestContractWire<'a> {
canonicalization: &'static str,
codec: &'static str,
delta_ir: &'static str,
lowering_profile: &'a SchemaLoweringProfileBinding,
semantic_profile: &'a SemanticProfileBinding,
}
#[derive(Serialize)]
struct ManifestFingerprintsWire<'a> {
plan: &'a MigrationPlanFingerprint,
source: ManifestEndpointFingerprintsWire<'a>,
target: ManifestEndpointFingerprintsWire<'a>,
}
#[derive(Serialize)]
struct ManifestEndpointFingerprintsWire<'a> {
declared_identity: &'a ManagedDeclaredIdentityFingerprint,
resolution_identity: &'a DeclaredIdentityFingerprint,
semantics: &'a ManagedSemanticSchemaFingerprint,
}
#[derive(Serialize)]
struct ManifestSafetyWire {
classification: SafetyClass,
reversible: bool,
}
#[derive(Serialize)]
struct LegacyParentWire<'a> {
app_label: &'a str,
checksum: LegacyChecksumWire<'a>,
name: &'a str,
}
#[derive(Serialize)]
struct LegacyChecksumWire<'a> {
algorithm: &'static str,
value: &'a str,
}
#[derive(Serialize)]
struct LegacyAppliedSetWire<'a> {
algorithm: &'static str,
canonicalization: &'static str,
digest: &'a str,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManifestCandidate {
contract: ManifestContractCandidate,
fingerprints: ManifestFingerprintsCandidate,
format: String,
id: MigrationIdCandidate,
#[serde(default, skip_serializing_if = "Option::is_none")]
legacy_applied_set: Option<LegacyAppliedSetCandidate>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
legacy_parents: Vec<LegacyParentCandidate>,
managed_scope: ManagedScopeCandidate,
parents: Vec<MigrationIdCandidate>,
required_capabilities: CapabilitySet,
resources: Vec<Value>,
safety: ManifestSafetyCandidate,
steps: Vec<Value>,
}
impl ManifestCandidate {
fn validate_header(&self) -> Result<(), Diagnostic> {
MigrationFormat::new(&self.format)?;
if self.contract.canonicalization != MANIFEST_SCHEMA_CANONICALIZATION
|| self.contract.codec != MANIFEST_CODEC
|| self.contract.delta_ir != MANIFEST_DELTA_IR
{
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_contract_mismatch",
"manifest contract metadata is not supported",
));
}
if !self.resources.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_resources_not_empty",
"schema-only manifest resources must be exactly empty",
));
}
parse_safety(&self.safety.classification)?;
Ok(())
}
fn to_draft(&self) -> Result<SchemaMigrationDraft, Diagnostic> {
if !self.legacy_parents.is_empty() {
if !self.parents.is_empty() || !self.steps.is_empty() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_bridge_not_zero_operation",
"a legacy-frontier bridge carries no steps and no canonical parents",
));
}
return SchemaMigrationDraft::legacy_bridge(
self.id.rebuild()?,
self.legacy_parents
.iter()
.map(LegacyParentCandidate::rebuild)
.collect::<Result<Vec<_>, _>>()?,
self.legacy_applied_set
.as_ref()
.ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_applied_set_missing",
"a legacy-frontier bridge requires its complete applied-set digest",
)
})?
.rebuild()?,
);
}
if self.legacy_applied_set.is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_applied_set_without_bridge",
"an ordinary migration cannot carry a legacy applied-set digest",
));
}
SchemaMigrationDraft::new(
self.id.rebuild()?,
self.parents
.iter()
.map(MigrationIdCandidate::rebuild)
.collect::<Result<Vec<_>, _>>()?,
self.steps
.iter()
.map(rebuild_step_candidate)
.collect::<Result<Vec<_>, _>>()?,
)
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManifestContractCandidate {
canonicalization: String,
codec: String,
delta_ir: String,
lowering_profile: ProfileBindingCandidate,
semantic_profile: ProfileBindingCandidate,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ProfileBindingCandidate {
fingerprint: Value,
id: String,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct MigrationIdCandidate {
app_label: String,
name: String,
}
impl MigrationIdCandidate {
fn rebuild(&self) -> Result<MigrationId, Diagnostic> {
Ok(MigrationId::from_components(
MigrationAppLabel::new(self.app_label.clone())?,
MigrationName::new(self.name.clone())?,
))
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct LegacyParentCandidate {
app_label: String,
checksum: LegacyChecksumCandidate,
name: String,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct LegacyChecksumCandidate {
algorithm: String,
value: String,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct LegacyAppliedSetCandidate {
algorithm: String,
canonicalization: String,
digest: String,
}
impl LegacyAppliedSetCandidate {
fn rebuild(&self) -> Result<LegacyAppliedSetDigest, Diagnostic> {
if self.algorithm != LEGACY_APPLIED_SET_ALGORITHM
|| self.canonicalization != LEGACY_APPLIED_SET_CANONICALIZATION
{
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_applied_set_contract",
"legacy applied-set binding carries unsupported digest vocabulary",
));
}
LegacyAppliedSetDigest::new(self.digest.clone())
}
}
impl LegacyParentCandidate {
fn rebuild(&self) -> Result<LegacyMigrationReference, Diagnostic> {
if self.checksum.algorithm != LEGACY_CHECKSUM_ALGORITHM {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_legacy_checksum_algorithm",
"legacy parent checksum carries an unsupported algorithm tag",
));
}
Ok(LegacyMigrationReference::new(
LegacyMigrationId::new(self.app_label.clone(), self.name.clone())?,
LegacyMigrationChecksum::new(self.checksum.value.clone())?,
))
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManagedScopeCandidate {
id: String,
profile: ProfileBindingCandidate,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManifestFingerprintsCandidate {
plan: Value,
source: ManifestEndpointFingerprintsCandidate,
target: ManifestEndpointFingerprintsCandidate,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManifestEndpointFingerprintsCandidate {
declared_identity: Value,
resolution_identity: Value,
semantics: Value,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct ManifestSafetyCandidate {
classification: String,
reversible: bool,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct SchemaStepCandidate {
contract: SchemaStepContractCandidate,
delta: Value,
kind: String,
}
impl SchemaStepCandidate {
fn rebuild(&self) -> Result<MigrationStep, Diagnostic> {
let delta = decode_schema_delta(&to_canonical_json(&self.delta)?)?;
let reverse = self
.contract
.reverse
.as_ref()
.map(|reverse| decode_schema_delta(&to_canonical_json(reverse)?))
.transpose()?;
let trusted = SchemaDeltaStep::new(
MigrationStepId::new(self.contract.id.clone())?,
delta,
reverse,
)?;
if to_canonical_json(self)? != trusted.canonical_bytes()? {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_step_contract_mismatch",
"schema step claims do not match the trusted delta-derived contract",
));
}
Ok(MigrationStep::from(trusted))
}
}
fn rebuild_step_candidate(value: &Value) -> Result<MigrationStep, Diagnostic> {
let kind = value
.as_object()
.and_then(|object| object.get("kind"))
.and_then(Value::as_str)
.ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_missing_step_kind",
"migration step requires a closed kind discriminator",
)
})?;
let bytes = to_canonical_json(value)?;
match kind {
"schema_delta" => from_canonical_json::<SchemaStepCandidate>(&bytes)?.rebuild(),
"assertion" => from_canonical_json::<AssertionStepCandidate>(&bytes)?.rebuild(),
"backfill" => from_canonical_json::<BackfillStepCandidate>(&bytes)?.rebuild(),
_ => Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_unknown_step_kind",
"migration step kind is not in the closed step vocabulary",
)),
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct BackfillStepCandidate {
contract: BackfillStepContractCandidate,
kind: String,
plan: Value,
}
impl BackfillStepCandidate {
fn rebuild(&self) -> Result<MigrationStep, Diagnostic> {
if self.kind != "backfill" {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_backfill_kind_mismatch",
"persisted backfill kind is not supported",
));
}
let plan = decode_attribute_backfill_plan(&to_canonical_json(&self.plan)?)?;
let trusted =
MigrationStep::backfill(MigrationStepId::new(self.contract.id.clone())?, plan)?;
if to_canonical_json(self)? != trusted.canonical_bytes()? {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_backfill_contract_mismatch",
"backfill step claims do not match the trusted plan-derived contract",
));
}
Ok(trusted)
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct AssertionStepCandidate {
contract: AssertionStepContractCandidate,
expected: String,
kind: String,
plan: Value,
}
impl AssertionStepCandidate {
fn rebuild(&self) -> Result<MigrationStep, Diagnostic> {
if self.kind != "assertion" || self.expected != "no_rows" {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_assertion_kind_mismatch",
"persisted assertion kind or expectation is not supported",
));
}
let plan = decode_migration_assertion_plan(&to_canonical_json(&self.plan)?)?;
let trusted = MigrationStep::assertion(
MigrationStepId::new(self.contract.id.clone())?,
plan,
AssertionExpectation::NoRows,
)?;
if to_canonical_json(self)? != trusted.canonical_bytes()? {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_manifest_assertion_contract_mismatch",
"assertion step claims do not match the trusted plan-derived contract",
));
}
Ok(trusted)
}
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct AssertionStepContractCandidate {
id: String,
plan_fingerprint: Value,
recovery: String,
required_capabilities: CapabilitySet,
retry: String,
source_semantics: Value,
target_semantics: Value,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct BackfillStepContractCandidate {
id: String,
plan_fingerprint: Value,
recovery: String,
required_capabilities: CapabilitySet,
retry: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
reverse: Option<String>,
source_semantics: Value,
target_semantics: Value,
}
#[derive(Clone, Deserialize, Serialize)]
#[serde(deny_unknown_fields)]
struct SchemaStepContractCandidate {
delta_fingerprint: Value,
id: String,
recovery: String,
required_capabilities: CapabilitySet,
retry: String,
#[serde(skip_serializing_if = "Option::is_none")]
reverse: Option<Value>,
source_semantics: Value,
target_semantics: Value,
}
fn parse_safety(value: &str) -> Result<SafetyClass, Diagnostic> {
match value {
"formal_only" => Ok(SafetyClass::FormalOnly),
"schema_metadata" => Ok(SafetyClass::SchemaMetadata),
"additive" => Ok(SafetyClass::Additive),
"conditional" => Ok(SafetyClass::Conditional),
"backfill_required" => Ok(SafetyClass::BackfillRequired),
"destructive" => Ok(SafetyClass::Destructive),
"opaque" => Ok(SafetyClass::Opaque),
"unsupported" => Ok(SafetyClass::Unsupported),
_ => Err(failure(
DiagnosticCategory::InvalidContract,
"migration_manifest_unknown_safety",
"manifest safety classification is not in the closed eight-class vocabulary",
)),
}
}