use std::collections::BTreeSet;
use std::io::Read as _;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use type_bridge_contract::capability::CapabilitySet;
use type_bridge_contract::diagnostic::Diagnostic;
use type_bridge_contract::migration::MigrationId;
use type_bridge_contract::schema::{DeclaredSchema, DocumentId};
use type_bridge_schema::SafetyClass;
use type_bridge_schema_compat::{ADOPTED_GENESIS_FILE_NAME, parse_adopted_genesis};
use type_bridge_schema_migration::{
BackfillMigrationGenerationRequest, GeneratedMigration, MigrationDirectory,
MigrationGenerationOutcome, MigrationGenerationRequest, MigrationHistoryGraph,
MigrationPreviewError, VerifiedMigrationHistoryBundle,
canonical_history_declared_legacy_bridge_count_in, discover_verified_migration_chain_in,
encode_verified_migration_history_bundle, generate_backfill_migration, generate_next_migration,
render_migration_preview, require_adoption_authority_pair,
require_adoption_authority_pair_state, try_acquire_migration_authoring_lock,
validate_portable_direct_child, write_generated_migration_under_lock,
};
use crate::{
MAX_BACKFILL_INTENT_BYTES, TypeBridgeWorkspace, TypeBridgeWorkspaceError, WorkspaceConfigError,
WorkspaceConfigErrorCode, parse_backfill_intent,
};
fn migration_directory_escape() -> TypeBridgeWorkspaceError {
TypeBridgeWorkspaceError::Config(
WorkspaceConfigError::new(
WorkspaceConfigErrorCode::PathNotConfined,
"the migration directory must remain beneath the workspace without symbolic links",
)
.with_detail("migration_v2_directory"),
)
}
#[derive(Debug)]
pub struct MigrationDirectoryAuthority {
directory: MigrationDirectory,
configured_path: PathBuf,
owner: Arc<()>,
}
impl MigrationDirectoryAuthority {
#[must_use]
pub const fn directory(&self) -> &MigrationDirectory {
&self.directory
}
#[must_use]
pub fn display_path(&self) -> &Path {
&self.configured_path
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct MigrationPlanEntry {
id: MigrationId,
safety: SafetyClass,
reversible: bool,
}
impl MigrationPlanEntry {
pub const fn id(&self) -> &MigrationId {
&self.id
}
pub const fn safety(&self) -> SafetyClass {
self.safety
}
pub const fn reversible(&self) -> bool {
self.reversible
}
}
impl TypeBridgeWorkspace {
pub fn open_migration_directory(
&self,
) -> Result<MigrationDirectoryAuthority, TypeBridgeWorkspaceError> {
self.open_migration_directory_impl(false)
}
pub fn ensure_migration_directory(
&self,
) -> Result<MigrationDirectoryAuthority, TypeBridgeWorkspaceError> {
self.open_migration_directory_impl(true)
}
fn open_migration_directory_impl(
&self,
create: bool,
) -> Result<MigrationDirectoryAuthority, TypeBridgeWorkspaceError> {
let root = self.config().workspace_root().as_path();
let directory = MigrationDirectory::open_beneath_directory(
self.root_directory(),
self.config().migration_v2_directory().as_path(),
create,
)
.map_err(|_| migration_directory_escape())?;
Ok(MigrationDirectoryAuthority {
directory,
configured_path: root.join(self.config().migration_v2_directory().as_path()),
owner: Arc::clone(self.authority_owner()),
})
}
fn require_migration_directory(
&self,
directory: &MigrationDirectoryAuthority,
) -> Result<(), TypeBridgeWorkspaceError> {
let expected = self
.config()
.workspace_root()
.as_path()
.join(self.config().migration_v2_directory().as_path());
if !Arc::ptr_eq(&directory.owner, self.authority_owner())
|| directory.configured_path != expected
{
return Err(migration_directory_escape());
}
Ok(())
}
pub fn discover_migrations(&self) -> Result<MigrationHistoryGraph, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.discover_migrations_in(&directory)
}
pub fn discover_migrations_in(
&self,
directory: &MigrationDirectoryAuthority,
) -> Result<MigrationHistoryGraph, TypeBridgeWorkspaceError> {
self.discover_migrations_with_genesis_in(directory)
.map(|(graph, _)| graph)
}
pub fn migration_history_bundle(
&self,
) -> Result<VerifiedMigrationHistoryBundle, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.migration_history_bundle_in(&directory)
}
pub fn migration_history_bundle_in(
&self,
directory: &MigrationDirectoryAuthority,
) -> Result<VerifiedMigrationHistoryBundle, TypeBridgeWorkspaceError> {
let graph = self.discover_migrations_in(directory)?;
Ok(VerifiedMigrationHistoryBundle::from_graph(&graph)?)
}
pub fn migration_history_bundle_bytes(&self) -> Result<Vec<u8>, TypeBridgeWorkspaceError> {
let bundle = self.migration_history_bundle()?;
Ok(encode_verified_migration_history_bundle(&bundle)?)
}
fn discover_migrations_with_genesis_in(
&self,
directory: &MigrationDirectoryAuthority,
) -> Result<(MigrationHistoryGraph, DeclaredSchema), TypeBridgeWorkspaceError> {
self.require_migration_directory(directory)?;
let adopted_source = read_adopted_genesis_bounded(directory.directory())?;
let adopted_genesis_present = adopted_source.is_some();
let declared_bridge_count =
canonical_history_declared_legacy_bridge_count_in(directory.directory())?;
require_adoption_authority_pair_state(adopted_genesis_present, declared_bridge_count)?;
let genesis = self.parse_migration_genesis(adopted_source)?;
let graph = discover_verified_migration_chain_in(
directory.directory(),
&genesis,
self.delta_context(),
)?;
require_adoption_authority_pair(&graph, adopted_genesis_present)?;
Ok((graph, genesis))
}
pub fn migration_make(
&self,
base_name: &str,
) -> Result<MigrationGenerationOutcome, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.migration_make_in(&directory, base_name)
}
pub fn migration_make_in(
&self,
directory: &MigrationDirectoryAuthority,
base_name: &str,
) -> Result<MigrationGenerationOutcome, TypeBridgeWorkspaceError> {
let (graph, genesis) = self.discover_migrations_with_genesis_in(directory)?;
let request = MigrationGenerationRequest {
app_label: self.config().app_label().as_str(),
base_name,
genesis_source: &genesis,
desired: self.declared_schema(),
context: self.delta_context(),
};
Ok(generate_next_migration(&graph, &request)?)
}
pub fn author_backfill_migration_in(
&self,
directory: &MigrationDirectoryAuthority,
base_name: &str,
intent_name: &Path,
) -> Result<GeneratedMigration, TypeBridgeWorkspaceError> {
self.require_migration_directory(directory)?;
validate_portable_direct_child(intent_name.as_os_str()).map_err(|_| {
backfill_intent_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::InvalidContract,
"workspace_backfill_intent_path_not_confined",
"backfill intent must be a portable direct-child filename",
)
})?;
let authoring_lock = try_acquire_migration_authoring_lock(directory.directory())?;
let (graph, genesis) = self.discover_migrations_with_genesis_in(directory)?;
let historical_schema = match graph.default_head()? {
Some(head) => graph
.manifest(head)
.expect("the verified default head is present")
.target_schema(),
None => &genesis,
};
let bytes = read_backfill_intent_bounded(directory.directory(), intent_name)?;
let document = DocumentId::new(intent_name.to_str().ok_or_else(|| {
backfill_intent_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::InvalidContract,
"workspace_backfill_intent_non_utf8_name",
"backfill intent filename must be UTF-8",
)
})?)?;
let plan =
parse_backfill_intent(document, &bytes, historical_schema, self.delta_context())?;
let generated = generate_backfill_migration(
&graph,
&BackfillMigrationGenerationRequest {
app_label: self.config().app_label().as_str(),
base_name,
genesis_source: &genesis,
desired: self.declared_schema(),
plan: &plan,
context: self.delta_context(),
},
)?;
let preview = format!(
"-- binding-neutral backfill migration: {}/{}\n-- reviewed source intent: {}\n",
generated.manifest().id().app_label().as_str(),
generated.manifest().id().name().as_str(),
intent_name.display(),
);
write_generated_migration_under_lock(&authoring_lock, &generated, &preview)?;
Ok(generated)
}
pub fn write_generated_migration(
&self,
generated: &GeneratedMigration,
) -> Result<PathBuf, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.write_generated_migration_in(&directory, generated)?;
Ok(directory.display_path().join(generated.file_name()))
}
pub fn write_generated_migration_in(
&self,
directory: &MigrationDirectoryAuthority,
generated: &GeneratedMigration,
) -> Result<PathBuf, TypeBridgeWorkspaceError> {
self.require_migration_directory(directory)?;
let authoring_lock = try_acquire_migration_authoring_lock(directory.directory())?;
let base_name = generated
.manifest()
.id()
.name()
.as_str()
.split_once('_')
.map(|(_, base_name)| base_name)
.filter(|base_name| !base_name.is_empty())
.ok_or_else(generated_migration_stale)?;
let regenerated = self.migration_make_in(directory, base_name)?;
let MigrationGenerationOutcome::Generated(regenerated) = regenerated else {
return Err(generated_migration_stale());
};
if regenerated.canonical_bytes() != generated.canonical_bytes() {
return Err(generated_migration_stale());
}
let preview = render_migration_preview(generated.manifest(), self.delta_context())
.map_err(preview_diagnostic)?;
Ok(write_generated_migration_under_lock(
&authoring_lock,
generated,
&preview,
)?)
}
pub fn migration_plan(
&self,
applied: &BTreeSet<MigrationId>,
) -> Result<Vec<MigrationPlanEntry>, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.migration_plan_in(&directory, applied)
}
pub fn migration_plan_in(
&self,
directory: &MigrationDirectoryAuthority,
applied: &BTreeSet<MigrationId>,
) -> Result<Vec<MigrationPlanEntry>, TypeBridgeWorkspaceError> {
let graph = self.discover_migrations_in(directory)?;
let order = graph.plan_apply_to_default_head(applied)?;
order
.into_iter()
.map(|id| {
let manifest = graph.manifest(&id).ok_or_else(|| {
TypeBridgeWorkspaceError::Contract(Diagnostic::new(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
type_bridge_contract::diagnostic::DiagnosticCode::new(
"workspace_migration_plan_missing_manifest",
)
.expect("static workspace diagnostic code"),
"planned identity has no verified manifest",
))
})?;
Ok(MigrationPlanEntry {
id,
safety: manifest.safety(),
reversible: manifest.reversible(),
})
})
.collect()
}
pub fn migration_genesis(&self) -> Result<DeclaredSchema, TypeBridgeWorkspaceError> {
let directory = self.open_migration_directory()?;
self.migration_genesis_in(&directory)
}
pub fn migration_genesis_in(
&self,
directory: &MigrationDirectoryAuthority,
) -> Result<DeclaredSchema, TypeBridgeWorkspaceError> {
self.require_migration_directory(directory)?;
self.parse_migration_genesis(read_adopted_genesis_bounded(directory.directory())?)
}
fn parse_migration_genesis(
&self,
source: Option<String>,
) -> Result<DeclaredSchema, TypeBridgeWorkspaceError> {
let Some(source) = source else {
return Ok(DeclaredSchema::from_facts(
self.declared_schema().format(),
CapabilitySet::new(),
std::iter::empty(),
)?);
};
let document = DocumentId::new(ADOPTED_GENESIS_FILE_NAME)?;
parse_adopted_genesis(document, &source).map_err(TypeBridgeWorkspaceError::Contract)
}
}
fn read_adopted_genesis_bounded(
directory: &MigrationDirectory,
) -> Result<Option<String>, TypeBridgeWorkspaceError> {
let limit = type_bridge_schema_compat::MAX_TYPEQL_SCHEMA_BYTES;
let file = match directory.open_regular_readonly(ADOPTED_GENESIS_FILE_NAME.as_ref()) {
Ok(file) => file,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(_) => {
return Err(adopted_genesis_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
"workspace_adopted_genesis_unreadable",
"adopted-genesis artifact exists but cannot be read as a regular file",
));
}
};
if !file.metadata().is_ok_and(|metadata| metadata.is_file()) {
return Err(adopted_genesis_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
"workspace_adopted_genesis_not_regular",
"adopted-genesis artifact must be a regular file",
));
}
let mut bytes = Vec::new();
file.take(u64::try_from(limit).unwrap_or(u64::MAX).saturating_add(1))
.read_to_end(&mut bytes)
.map_err(|_| {
adopted_genesis_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
"workspace_adopted_genesis_unreadable",
"adopted-genesis artifact exists but cannot be read",
)
})?;
if bytes.len() > limit {
return Err(adopted_genesis_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::ResourceLimit,
"workspace_adopted_genesis_oversized",
"adopted-genesis artifact exceeds the schema byte ceiling",
));
}
String::from_utf8(bytes).map(Some).map_err(|_| {
adopted_genesis_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::InvalidContract,
"workspace_adopted_genesis_not_utf8",
"adopted-genesis artifact is not valid UTF-8",
)
})
}
fn read_backfill_intent_bounded(
directory: &MigrationDirectory,
name: &Path,
) -> Result<Vec<u8>, TypeBridgeWorkspaceError> {
let file = directory
.open_regular_readonly(name.as_os_str())
.map_err(|_| {
backfill_intent_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
"workspace_backfill_intent_unreadable",
"backfill intent cannot be read as a regular direct-child file",
)
})?;
let mut bytes = Vec::new();
file.take(
u64::try_from(MAX_BACKFILL_INTENT_BYTES)
.unwrap_or(u64::MAX)
.saturating_add(1),
)
.read_to_end(&mut bytes)
.map_err(|_| {
backfill_intent_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
"workspace_backfill_intent_unreadable",
"backfill intent cannot be read",
)
})?;
if bytes.len() > MAX_BACKFILL_INTENT_BYTES {
return Err(backfill_intent_read_error(
type_bridge_contract::diagnostic::DiagnosticCategory::ResourceLimit,
"workspace_backfill_intent_byte_limit",
"backfill intent exceeds the common source-byte ceiling",
));
}
Ok(bytes)
}
fn backfill_intent_read_error(
category: type_bridge_contract::diagnostic::DiagnosticCategory,
code: &'static str,
message: &'static str,
) -> TypeBridgeWorkspaceError {
TypeBridgeWorkspaceError::Contract(Diagnostic::new(
category,
type_bridge_contract::diagnostic::DiagnosticCode::new(code)
.expect("static backfill-intent diagnostic code"),
message,
))
}
fn generated_migration_stale() -> TypeBridgeWorkspaceError {
TypeBridgeWorkspaceError::Contract(Diagnostic::new(
type_bridge_contract::diagnostic::DiagnosticCategory::Integrity,
type_bridge_contract::diagnostic::DiagnosticCode::new(
"workspace_generated_migration_stale",
)
.expect("static workspace diagnostic code"),
"generated migration no longer matches this workspace directory authority",
))
}
fn adopted_genesis_read_error(
category: type_bridge_contract::diagnostic::DiagnosticCategory,
code: &'static str,
message: &'static str,
) -> TypeBridgeWorkspaceError {
TypeBridgeWorkspaceError::Contract(Diagnostic::new(
category,
type_bridge_contract::diagnostic::DiagnosticCode::new(code)
.expect("static workspace diagnostic code"),
message,
))
}
fn preview_diagnostic(error: MigrationPreviewError) -> TypeBridgeWorkspaceError {
match error {
MigrationPreviewError::Diagnostic(diagnostic) => {
TypeBridgeWorkspaceError::Contract(diagnostic)
}
MigrationPreviewError::Lowering(lowering) => {
TypeBridgeWorkspaceError::Contract(Diagnostic::new(
type_bridge_contract::diagnostic::DiagnosticCategory::InvalidContract,
type_bridge_contract::diagnostic::DiagnosticCode::new(lowering.code())
.expect("lowering diagnostic codes are canonical"),
"migration preview lowering failed",
))
}
}
}