#[cfg(test)]
use crate::storage_adapter::SharedStorageAdapterRead;
use std::collections::{BTreeMap, BTreeSet};
use std::ops::Bound;
use bytes::Bytes;
use crate::LixError;
use crate::branch::BranchHeadControlContext;
use crate::changelog::{ChangelogContext, ChangelogReader, CommitLoadRequest};
use crate::init::{
CURRENT_FORMAT_VERSION, REPOSITORY_PROTOCOL_KEY, REPOSITORY_PROTOCOL_SPACE,
RepositoryProtocolStatus, parse_repository_protocol,
};
use crate::storage_adapter::{
Storage, StorageAdapterRead as _, StorageBeginScanOptions,
StorageCoreProjection as CoreProjection, StorageError, StorageGetOptions as GetOptions,
StorageKey as Key, StorageKeyRange, StoragePrecondition as Precondition,
StorageProjectedValue as ProjectedValue, StorageReadOptions as ReadOptions, StorageWrite,
StorageWriteOptions as WriteOptions,
};
use crate::tracked_state::{
CommitStateReplayDebt, TrackedStateContext, TrackedStateFilter, TrackedStateKeyRef,
TrackedStateReadColumns, TrackedStateScanRequest, backfill_row_pk_index_for_commit,
encode_commit_state_manifest_replacement_for_migration,
};
const REPOSITORY_PROTOCOL_V72: &[u8] = b"tracked-default-branch.v72";
const REPOSITORY_PROTOCOL_V73: &[u8] = b"tracked-default-branch.v73";
const REPOSITORY_PROTOCOL_V74: &[u8] = b"tracked-default-branch.v74";
const ACCOUNT_SCHEMA_KEY: &str = "lix_account";
#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct MigrationOptions {
pub max_changes: usize,
pub max_preflight_bytes: usize,
}
impl Default for MigrationOptions {
fn default() -> Self {
Self {
max_changes: 250_000,
max_preflight_bytes: 512 * 1024 * 1024,
}
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) struct MigrationReport {
pub from_version: u32,
pub to_version: u32,
pub changes_rewritten: u64,
pub commit_members_rewritten: u64,
}
use super::inspection::{MigrationStatus, inspect_lix_with_adapter};
pub(crate) async fn migrate_lix_with_adapter<S>(
storage: S,
adapter: crate::storage_adapter::StorageAdapter<S>,
options: MigrationOptions,
) -> Result<MigrationReport, LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let read = super::MigrationPlanningRead::new(&adapter)
.await
.map_err(storage_error)?;
let protocol_status = crate::init::repository_protocol_status(&read).await?;
if let RepositoryProtocolStatus::MigrationRequired { found_version } = protocol_status
&& (found_version < 72
|| !crate::migration::registry::has_complete_migration_path(
found_version,
CURRENT_FORMAT_VERSION,
))
{
return Err(migration_error(format!(
"repository v{found_version} predates the v{CURRENT_FORMAT_VERSION} complete-snapshot commit format; no in-place migration is available"
)));
}
let from_version = match protocol_status {
RepositoryProtocolStatus::MigrationRequired {
found_version: found_version @ (72 | 73 | 74 | 75 | 76 | 77 | 78 | 79 | 80),
} => found_version,
RepositoryProtocolStatus::Current => {
return Ok(MigrationReport {
from_version: CURRENT_FORMAT_VERSION,
to_version: CURRENT_FORMAT_VERSION,
changes_rewritten: 0,
commit_members_rewritten: 0,
});
}
RepositoryProtocolStatus::MigrationRequired { found_version } => {
return Err(migration_error(format!(
"no registered migration from repository v{found_version}"
)));
}
RepositoryProtocolStatus::TooNew { found_version } => {
return Err(migration_error(format!(
"repository v{found_version} is newer than this engine"
)));
}
RepositoryProtocolStatus::Malformed | RepositoryProtocolStatus::Missing => {
return Err(migration_error(
"repository has no valid versioned protocol marker",
));
}
};
read.finish().map_err(storage_error)?;
if from_version <= 78 {
super::deterministic_witness::backfill(&adapter, options, false).await?;
}
let commit_records_rewritten = if from_version <= 74 {
rewrite_commit_records_to_v6(&adapter, &storage, options, from_version).await?
} else {
0
};
if from_version <= 72 {
migrate_v72_account_profile_uri(&adapter, &storage, options).await?;
}
if from_version <= 73 {
migrate_v73_row_pk_indexes(&adapter, &storage, options).await?;
}
let commit_members_rewritten = if from_version <= 74 {
migrate_v74_complete_snapshot_commits(&adapter, &storage, options).await?
} else {
0
};
if from_version == 75 {
let read = super::MigrationPlanningRead::new(&adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
drop(read);
crate::migration::publish::publish(
&adapter,
expected_revision,
crate::init::REPOSITORY_PROTOCOL_V75,
crate::init::REPOSITORY_PROTOCOL_V76,
crate::migration::publish::PublicationPlan::bounded(0, 0),
)
.await?;
}
if from_version <= 76 {
let read = super::MigrationPlanningRead::new(&adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
drop(read);
crate::migration::publish::publish(
&adapter,
expected_revision,
crate::init::REPOSITORY_PROTOCOL_V76,
crate::init::REPOSITORY_PROTOCOL_V77,
crate::migration::publish::PublicationPlan::bounded(0, 0),
)
.await?;
}
let checkpoint_records_rewritten = if from_version <= 77 {
super::checkpoint_metadata::migrate(&adapter, options).await?
} else {
0
};
if from_version <= 78 {
backfill_missing_row_pk_indexes(
&adapter,
&storage,
options,
78,
crate::init::REPOSITORY_PROTOCOL_V78,
crate::init::REPOSITORY_PROTOCOL_V78,
"v79 complete row-PK catalog repair",
false,
true,
)
.await?;
super::deterministic_witness::backfill(&adapter, options, true).await?;
}
if from_version <= 79 {
super::incorporation::migrate(&adapter, options, false).await?;
}
super::runtime_epoch::migrate(&adapter, false).await?;
Ok(MigrationReport {
from_version,
to_version: CURRENT_FORMAT_VERSION,
changes_rewritten: commit_records_rewritten + checkpoint_records_rewritten,
commit_members_rewritten,
})
}
async fn migrate_v72_account_profile_uri<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
storage: &S,
options: MigrationOptions,
) -> Result<(), LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let marker = load_repository_protocol_marker(adapter).await?;
let backfill_source = match marker.as_deref() {
Some(value) if value == REPOSITORY_PROTOCOL_V72 => Some(REPOSITORY_PROTOCOL_V72),
Some(value) if value == crate::init::REPOSITORY_PROTOCOL_V72_COMMIT_REWRITE => {
Some(crate::init::REPOSITORY_PROTOCOL_V72_COMMIT_REWRITE)
}
Some(value) if value == crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP => {
Some(crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP)
}
_ => {
return Err(migration_error(
"v72 row-PK-index bootstrap observed an unexpected protocol marker",
));
}
};
if let Some(source_protocol) = backfill_source {
backfill_missing_row_pk_indexes(
adapter,
storage,
options,
72,
source_protocol,
crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP,
"v72 row-PK-index bootstrap",
false,
true,
)
.await?;
}
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
match crate::init::repository_protocol_status(&read).await? {
RepositoryProtocolStatus::MigrationRequired { found_version: 72 } => {}
status => {
return Err(migration_error(format!(
"v72 account-schema migration observed unexpected repository status {status:?}"
)));
}
}
let branch_ids = BranchHeadControlContext::new()
.reader(read.clone())
.scan()
.await?
.into_iter()
.map(|(branch_id, _)| branch_id)
.collect::<Vec<_>>();
if branch_ids.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"v72 account-schema migration exceeds configured branch bound: {} branches",
branch_ids.len()
),
));
}
read.finish().map_err(storage_error)?;
let target = crate::schema::seed_schema_definition(ACCOUNT_SCHEMA_KEY)
.expect("lix_account is a built-in schema")
.clone();
let target_schema = lix_schema::from_value(target.clone()).map_err(|error| {
migration_error(format!("bundled lix_account schema is invalid: {error}"))
})?;
let engine = crate::engine::Engine::new_for_migration_with_adapter(adapter.clone(), 72).await?;
for branch_id in branch_ids {
let session = engine.open_session_at_for_migration(&branch_id);
let result = session
.execute(
"SELECT value FROM lix_registered_schema \
WHERE schema_key = 'lix_account'",
&[],
)
.await?;
let rows = result.rows();
let [row] = rows else {
return Err(migration_error(format!(
"branch '{branch_id}' must expose exactly one lix_account schema row, found {}",
rows.len()
)));
};
let [crate::Value::Jsonb(stored)] = row.values() else {
return Err(migration_error(format!(
"branch '{branch_id}' has a non-JSON lix_account schema row"
)));
};
let stored = stored.to_value();
if stored != target {
let stored_schema = lix_schema::from_value(stored).map_err(|error| {
migration_error(format!(
"branch '{branch_id}' has an invalid persisted lix_account schema: {error}"
))
})?;
if let Err(error) = lix_schema::validate_amendment(&stored_schema, &target_schema) {
if lix_schema::validate_amendment(&target_schema, &stored_schema).is_err() {
return Err(migration_error(format!(
"branch '{branch_id}' has a divergent lix_account schema: {error}"
)));
}
} else {
let updated = session
.execute(
"UPDATE lix_registered_schema SET value = $1 \
WHERE schema_key = 'lix_account'",
&[crate::Value::Jsonb(target.clone().into())],
)
.await?;
if updated.rows_affected() != 1 {
return Err(migration_error(format!(
"branch '{branch_id}' updated {} lix_account schema rows instead of one",
updated.rows_affected()
)));
}
}
}
session.close().await?;
}
drop(engine);
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
read.finish().map_err(storage_error)?;
crate::migration::publish::publish(
adapter,
expected_revision,
crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP,
REPOSITORY_PROTOCOL_V73,
crate::migration::publish::PublicationPlan::bounded(0, 0),
)
.await
}
pub(super) async fn load_repository_protocol_marker<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
) -> Result<Option<Bytes>, LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
let values = crate::storage_adapter::PointReadPlan::new(
REPOSITORY_PROTOCOL_SPACE,
&[Key(Bytes::from_static(REPOSITORY_PROTOCOL_KEY))],
)
.materialize(&read, GetOptions::default())
.await
.map_err(storage_error)?;
Ok(values
.value
.into_iter()
.next()
.flatten()
.and_then(|value| match value {
ProjectedValue::FullValue(value) => Some(value),
ProjectedValue::KeyOnly => None,
}))
}
async fn migrate_v73_row_pk_indexes<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
storage: &S,
options: MigrationOptions,
) -> Result<(), LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let marker = load_repository_protocol_marker(adapter).await?;
let source_protocol = match marker.as_deref() {
Some(value) if value == REPOSITORY_PROTOCOL_V73 => REPOSITORY_PROTOCOL_V73,
Some(value) if value == crate::init::REPOSITORY_PROTOCOL_V73_COMMIT_REWRITE => {
crate::init::REPOSITORY_PROTOCOL_V73_COMMIT_REWRITE
}
_ => {
return Err(migration_error(
"v73 row-PK-index migration observed an unexpected protocol marker",
));
}
};
backfill_missing_row_pk_indexes(
adapter,
storage,
options,
73,
source_protocol,
REPOSITORY_PROTOCOL_V74,
"v73 row-PK-index migration",
true,
true,
)
.await
}
async fn migrate_v74_complete_snapshot_commits<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
storage: &S,
options: MigrationOptions,
) -> Result<u64, LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let marker = load_repository_protocol_marker(adapter).await?;
let source_protocol = match marker.as_deref() {
Some(value) if value == REPOSITORY_PROTOCOL_V74 => REPOSITORY_PROTOCOL_V74,
Some(value) if value == crate::init::REPOSITORY_PROTOCOL_V74_COMMIT_REWRITE => {
crate::init::REPOSITORY_PROTOCOL_V74_COMMIT_REWRITE
}
_ => {
return Err(migration_error(
"v74 complete-snapshot migration observed an unexpected protocol marker",
));
}
};
let (members_injected, repaired_manifests) =
repair_filesystem_closure(adapter, storage, options, source_protocol).await?;
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
read.finish().map_err(storage_error)?;
let mut publication = crate::migration::publish::PublicationPlan::bounded(
repaired_manifests.len().saturating_mul(2),
options.max_preflight_bytes,
);
for manifest in repaired_manifests {
for (space, key, value) in
encode_commit_state_manifest_replacement_for_migration(&manifest)?
{
publication.replace_immutable(space, vec![(key, value)])?;
}
}
crate::migration::publish::publish(
adapter,
expected_revision,
source_protocol,
crate::init::REPOSITORY_PROTOCOL_V76,
publication,
)
.await?;
Ok(members_injected)
}
fn migration_snapshot_text(
row: &crate::tracked_state::MaterializedTrackedStateRow,
field: &str,
) -> Option<String> {
let snapshot = row.decoded_snapshot.as_ref()?;
match snapshot.row.get(field) {
Some(lix_schema::Value::Text(value)) => Some(value.clone()),
Some(lix_schema::Value::Uuid(value)) => Some(value.to_string()),
_ => None,
}
}
struct CommitGraphNode {
parents: Vec<crate::changelog::CommitId>,
base: Option<crate::changelog::CommitId>,
}
fn commit_ancestors(
commit_id: crate::changelog::CommitId,
commit_graph: &BTreeMap<crate::changelog::CommitId, CommitGraphNode>,
) -> BTreeSet<crate::changelog::CommitId> {
let mut seen = BTreeSet::new();
let mut worklist = vec![commit_id];
while let Some(current) = worklist.pop() {
if !seen.insert(current) {
continue;
}
if let Some(node) = commit_graph.get(¤t) {
worklist.extend(node.parents.iter().copied());
if let Some(base) = node.base {
worklist.push(base);
}
}
}
seen
}
fn resolve_missing_directory_closure(
present: &BTreeSet<String>,
referenced: &BTreeSet<String>,
ancestors: &BTreeSet<crate::changelog::CommitId>,
candidates: &BTreeMap<String, Vec<crate::tracked_state::MaterializedTrackedStateRow>>,
commit_graph: &BTreeMap<crate::changelog::CommitId, CommitGraphNode>,
) -> Result<
BTreeMap<String, crate::tracked_state::MaterializedTrackedStateRow>,
(String, &'static str),
> {
let mut resolved = BTreeMap::new();
let mut worklist = referenced.difference(present).cloned().collect::<Vec<_>>();
while let Some(id) = worklist.pop() {
if resolved.contains_key(&id) {
continue;
}
let admissible = candidates
.get(&id)
.into_iter()
.flatten()
.filter(|row| ancestors.contains(&row.commit_id))
.collect::<Vec<_>>();
if admissible.is_empty() {
return Err((id, "no tree owned by that commit's ancestry carries it"));
}
let mut owner_ancestry = BTreeMap::new();
for row in &admissible {
owner_ancestry
.entry(row.commit_id)
.or_insert_with(|| commit_ancestors(row.commit_id, commit_graph));
}
let mut winner: Option<&crate::tracked_state::MaterializedTrackedStateRow> = None;
let mut ambiguous = false;
for row in &admissible {
let strictly_superseded = admissible.iter().any(|other| {
owner_ancestry[&other.commit_id].contains(&row.commit_id)
&& !owner_ancestry[&row.commit_id].contains(&other.commit_id)
});
if strictly_superseded {
continue;
}
match winner {
None => winner = Some(row),
Some(current) if current.change_id == row.change_id => {}
Some(_) => ambiguous = true,
}
}
if ambiguous {
return Err((
id,
"its surviving versions are causally incomparable merge ancestors",
));
}
let row = winner
.expect("a non-empty admissible set has a maximal element")
.clone();
if let Some(parent) = migration_snapshot_text(&row, "parent_id")
&& !present.contains(&parent)
{
worklist.push(parent);
}
resolved.insert(id, row);
}
Ok(resolved)
}
async fn repair_filesystem_closure<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
_storage: &S,
options: MigrationOptions,
source_protocol: &'static [u8],
) -> Result<(u64, Vec<crate::tracked_state::CommitStateManifest>), LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
const FILE_DESCRIPTOR_SCHEMA_KEY: &str = "lix_file_descriptor";
const DIRECTORY_DESCRIPTOR_SCHEMA_KEY: &str = "lix_directory_descriptor";
fn index_value_ref(
row: &crate::tracked_state::MaterializedTrackedStateRow,
) -> Result<crate::tracked_state::TrackedStateIndexValueRef, LixError> {
Ok(crate::tracked_state::TrackedStateIndexValueRef {
change_id: row.change_id,
commit_id: row.commit_id,
deleted: false,
created_at: crate::common::LixTimestamp::parse(&row.created_at)
.map_err(|error| migration_error(format!("repair created_at: {error}")))?,
updated_at: crate::common::LixTimestamp::parse(&row.updated_at)
.map_err(|error| migration_error(format!("repair updated_at: {error}")))?,
})
}
fn push_missing_directories(
builder: &mut crate::tracked_state::TrackedStateMutationBatchBuilder,
missing: &[String],
source: &BTreeMap<String, crate::tracked_state::MaterializedTrackedStateRow>,
) -> Result<(), LixError> {
for id in missing {
let row = source.get(id).expect("resolved above");
builder.push(
TrackedStateKeyRef {
schema_key: DIRECTORY_DESCRIPTOR_SCHEMA_KEY,
file_id: row.file_id.as_deref(),
row_pk: &row.row_pk,
},
index_value_ref(row)?,
);
}
Ok(())
}
let operation = "v74 filesystem-closure repair";
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
let commit_ids = crate::tracked_state::scan_commit_state_manifest_commit_ids(&read).await?;
let mut commit_graph = BTreeMap::new();
let mut cursor = read
.begin_scan(
crate::changelog::COMMIT_SPACE,
StorageKeyRange {
lower: Bound::Unbounded,
upper: Bound::Unbounded,
},
StorageBeginScanOptions {
projection: CoreProjection::FullValue,
..StorageBeginScanOptions::default()
},
)
.await
.map_err(storage_error)?;
while let Some(entries) = cursor.next_chunk().await.map_err(storage_error)? {
for entry in entries {
if commit_graph.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} exceeds configured commit bound: {} commits",
commit_graph.len()
),
));
}
let ProjectedValue::FullValue(value) = entry.value else {
return Err(migration_error(format!(
"{operation} commit scan omitted a value"
)));
};
let record = crate::storage_codec::decode::<crate::changelog::CommitRecord>(
"commit record",
&value,
)
.map_err(|error| {
migration_error(format!(
"{operation} could not decode a commit record: {error}"
))
})?;
commit_graph.insert(
record.commit_id,
CommitGraphNode {
parents: record.parent_commit_ids.clone(),
base: record.base_commit_id,
},
);
}
}
drop(cursor);
let mut candidates =
BTreeMap::<String, Vec<crate::tracked_state::MaterializedTrackedStateRow>>::new();
let mut reader = TrackedStateContext::new().reader(read.clone());
let mut audits = Vec::new();
let mut visited_rows = 0_usize;
for commit_id in commit_ids {
let commit_key = commit_id.to_string();
let budget = options.max_changes.saturating_sub(visited_rows);
let rows = reader
.scan_batch_at_commit(
&commit_key,
&TrackedStateScanRequest {
filter: TrackedStateFilter {
schema_keys: vec![
FILE_DESCRIPTOR_SCHEMA_KEY.to_owned(),
DIRECTORY_DESCRIPTOR_SCHEMA_KEY.to_owned(),
],
..TrackedStateFilter::default()
},
read_columns: TrackedStateReadColumns::default(),
limit: Some(budget.saturating_add(1)),
},
)
.await?
.into_rows();
if rows.len() > budget {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!("{operation} exceeds configured row bound"),
));
}
visited_rows = visited_rows.saturating_add(rows.len());
let mut present = BTreeSet::new();
let mut referenced = BTreeSet::new();
for row in &rows {
match row.schema_key.as_str() {
DIRECTORY_DESCRIPTOR_SCHEMA_KEY => {
if let Some(id) = migration_snapshot_text(row, "id") {
present.insert(id);
}
if let Some(parent) = migration_snapshot_text(row, "parent_id") {
referenced.insert(parent);
}
}
FILE_DESCRIPTOR_SCHEMA_KEY => {
if let Some(directory) = migration_snapshot_text(row, "directory_id") {
referenced.insert(directory);
}
}
_ => {}
}
}
for row in rows {
if row.schema_key == DIRECTORY_DESCRIPTOR_SCHEMA_KEY
&& let Some(id) = migration_snapshot_text(&row, "id")
{
let versions = candidates.entry(id).or_default();
if !versions
.iter()
.any(|existing| existing.change_id == row.change_id)
{
versions.push(row);
}
}
}
if referenced.difference(&present).next().is_some() {
audits.push((commit_id, present, referenced));
}
}
let mut chunk_writes = adapter.new_write_set();
let mut replacements = Vec::new();
let mut injected_total = 0_u64;
for (commit_id, present, referenced) in audits {
let commit_key = commit_id.to_string();
let ancestors = commit_ancestors(commit_id, &commit_graph);
let missing_rows = resolve_missing_directory_closure(
&present,
&referenced,
&ancestors,
&candidates,
&commit_graph,
)
.map_err(|(id, reason)| {
migration_error(format!(
"{operation} cannot resolve directory '{id}' referenced by commit \
'{commit_id}': {reason}"
))
})?;
let missing = missing_rows.keys().cloned().collect::<Vec<_>>();
let mut manifest = crate::tracked_state::load_commit_state_manifest(&read, commit_id)
.await?
.ok_or_else(|| {
migration_error(format!(
"{operation} commit '{commit_id}' has no commit-state manifest"
))
})?;
let mut injected =
crate::tracked_state::TrackedStateMutationBatchBuilder::with_row_capacity(
missing.len(),
);
push_missing_directories(&mut injected, &missing, &missing_rows)?;
let (primary, secondary) =
crate::tracked_state::with_row_pk_index_mutations(injected.finish())?;
let mut overlay = crate::tracked_state::TrackedStateChunkOverlay::new();
let tree = crate::tracked_state::TrackedStateTree::new();
if let Some(existing_root) = manifest.snapshot_root.as_deref() {
let result = tree
.apply_mutations_with_overlay(
&read,
&mut chunk_writes,
&mut overlay,
Some(&existing_root.root_id),
primary,
Some(&commit_key),
)
.await?;
let mut repaired_root = existing_root.clone();
repaired_root.root_id = result.root_id;
repaired_root.changed_key_count = repaired_root
.changed_key_count
.saturating_add(missing.len() as u64);
repaired_root.row_count_estimate = repaired_root
.row_count_estimate
.saturating_add(missing.len() as u64);
repaired_root.tree_height = result.tree_height as u32;
manifest.snapshot_root = Some(Box::new(repaired_root));
} else {
let budget = options.max_changes.saturating_sub(visited_rows);
let full_rows = reader
.scan_batch_at_commit(
&commit_key,
&TrackedStateScanRequest {
limit: Some(budget.saturating_add(1)),
..TrackedStateScanRequest::default()
},
)
.await?
.into_rows();
if full_rows.len() > budget {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} rootless commit '{commit_id}' exceeds configured row bound"
),
));
}
visited_rows = visited_rows.saturating_add(full_rows.len());
let mut full =
crate::tracked_state::TrackedStateMutationBatchBuilder::with_row_capacity(
full_rows.len() + missing.len(),
);
for row in &full_rows {
full.push(
TrackedStateKeyRef {
schema_key: &row.schema_key,
file_id: row.file_id.as_deref(),
row_pk: &row.row_pk,
},
index_value_ref(row)?,
);
}
push_missing_directories(&mut full, &missing, &missing_rows)?;
let result = tree
.apply_mutations_with_overlay(
&read,
&mut chunk_writes,
&mut overlay,
None,
full.finish(),
Some(&commit_key),
)
.await?;
manifest.snapshot_root = Some(Box::new(crate::tracked_state::TrackedStateCommitRoot {
commit_id,
root_id: result.root_id,
parent_roots: Vec::new(),
changed_key_count: result.row_count as u64,
row_count_estimate: result.row_count as u64,
tree_height: result.tree_height as u32,
complete_state_fence: true,
}));
manifest.replay_debt = CommitStateReplayDebt::default();
}
let row_pk_base = manifest.row_pk_index_root_id.clone();
let result = tree
.apply_mutations_with_overlay(
&read,
&mut chunk_writes,
&mut overlay,
row_pk_base.as_ref(),
secondary,
Some(&commit_key),
)
.await?;
manifest.row_pk_index_root_id = Some(result.root_id);
injected_total = injected_total.saturating_add(missing.len() as u64);
replacements.push(manifest);
}
drop(reader);
read.finish().map_err(storage_error)?;
let chunk_stats = chunk_writes.stats();
if chunk_stats.written_bytes > options.max_preflight_bytes as u64 {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} exceeds configured byte bound: {} bytes",
chunk_stats.written_bytes
),
));
}
if chunk_stats.staged_puts != 0 || chunk_stats.staged_deletes != 0 {
let mut write = adapter
.begin_migration_write(WriteOptions {
await_durable: true,
preconditions: vec![
Precondition::KeyValueEquals {
space: REPOSITORY_PROTOCOL_SPACE,
key: Key(Bytes::from_static(REPOSITORY_PROTOCOL_KEY)),
expected: Bytes::from_static(source_protocol),
},
crate::storage_adapter::repository_mutation_revision_precondition(
expected_revision,
),
],
..WriteOptions::default()
})
.await
.map_err(storage_error)?;
if let Err(error) = chunk_writes.lower_into(&mut write).await {
let _ = write.rollback().await;
return Err(error.into());
}
write.commit().await.map_err(storage_error)?;
}
Ok((injected_total, replacements))
}
async fn rewrite_commit_records_to_v6<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
_storage: &S,
options: MigrationOptions,
expected_version: u32,
) -> Result<u64, LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
#[derive(Debug, musli::Decode)]
#[musli(packed)]
struct CommitRecordV5 {
format_version: u32,
commit_id: crate::changelog::CommitId,
generation: u64,
parent_commit_ids: Vec<crate::changelog::CommitId>,
first_parent_jump_commit_id: crate::changelog::CommitId,
first_parent_jump_span: u64,
account_id: String,
created_at: crate::common::LixTimestamp,
touched_scope_digest: crate::changelog::CommitTouchedScopeDigest,
}
struct CommitChronology {
created_at: crate::common::LixTimestamp,
generation: u64,
first_parent: Option<crate::changelog::CommitId>,
}
let operation = "commit-record v6 rewrite";
let marker = load_repository_protocol_marker(adapter)
.await?
.ok_or_else(|| migration_error(format!("{operation} found no protocol marker")))?;
match parse_repository_protocol(&marker) {
RepositoryProtocolStatus::MigrationRequired { found_version }
if found_version == expected_version => {}
status => {
return Err(migration_error(format!(
"{operation} observed unexpected repository status {status:?}"
)));
}
}
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
let mut cursor = read
.begin_scan(
crate::changelog::COMMIT_SPACE,
StorageKeyRange {
lower: Bound::Unbounded,
upper: Bound::Unbounded,
},
StorageBeginScanOptions {
projection: CoreProjection::FullValue,
..StorageBeginScanOptions::default()
},
)
.await
.map_err(storage_error)?;
let mut pending = Vec::new();
let mut pending_v6 = Vec::new();
let mut chronology = BTreeMap::new();
while let Some(entries) = cursor.next_chunk().await.map_err(storage_error)? {
for entry in entries {
if chronology.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} exceeds configured commit bound: {} commits",
chronology.len()
),
));
}
let ProjectedValue::FullValue(value) = entry.value else {
return Err(migration_error(format!(
"{operation} commit scan omitted a value"
)));
};
if let Ok(record) = crate::storage_codec::decode::<crate::changelog::CommitRecord>(
"commit record",
&value,
) && record.format_version == crate::changelog::COMMIT_RECORD_FORMAT_VERSION
{
chronology.insert(
record.commit_id,
CommitChronology {
created_at: record.created_at,
generation: record.generation,
first_parent: record.parent_commit_ids.first().copied(),
},
);
continue;
}
if let Some(record) = super::checkpoint_metadata::decode_v6(&value) {
chronology.insert(
record.commit_id,
CommitChronology {
created_at: record.created_at,
generation: record.generation,
first_parent: record.parent_commit_ids.first().copied(),
},
);
pending_v6.push((entry.key, record));
continue;
}
let record = crate::storage_codec::decode::<CommitRecordV5>("v5 commit record", &value)
.map_err(|error| {
migration_error(format!(
"{operation} could not decode a commit record as v5: {error}"
))
})?;
if record.format_version != 5 {
return Err(migration_error(format!(
"{operation} commit '{}' has unsupported record format v{}",
record.commit_id, record.format_version
)));
}
chronology.insert(
record.commit_id,
CommitChronology {
created_at: record.created_at,
generation: record.generation,
first_parent: record.parent_commit_ids.first().copied(),
},
);
pending.push((entry.key, record));
}
}
drop(cursor);
if pending.is_empty() && pending_v6.is_empty() {
read.finish().map_err(storage_error)?;
return Ok(0);
}
let global_head = BranchHeadControlContext::new()
.reader(read.clone())
.load(crate::GLOBAL_BRANCH_ID)
.await?
.ok_or_else(|| migration_error(format!("{operation} found no global branch")))?
.head_commit_id;
let mut global_chronology = Vec::new();
let mut global_commits = BTreeSet::new();
let mut next_global = Some(global_head);
while let Some(commit_id) = next_global {
if !global_commits.insert(commit_id) {
return Err(migration_error(format!(
"{operation} global lineage revisits commit '{commit_id}'"
)));
}
let node = chronology.get(&commit_id).ok_or_else(|| {
migration_error(format!(
"{operation} global lineage commit '{commit_id}' is missing"
))
})?;
global_chronology.push((node.created_at, node.generation, commit_id));
next_global = node.first_parent;
}
global_chronology.sort();
read.finish().map_err(storage_error)?;
let fence = match expected_version {
72 => crate::init::REPOSITORY_PROTOCOL_V72_COMMIT_REWRITE,
73 => crate::init::REPOSITORY_PROTOCOL_V73_COMMIT_REWRITE,
74 => crate::init::REPOSITORY_PROTOCOL_V74_COMMIT_REWRITE,
other => {
return Err(migration_error(format!(
"{operation} has no rewrite fence for source version v{other}"
)));
}
};
let mut writes = adapter.new_write_set();
writes.put(REPOSITORY_PROTOCOL_SPACE, REPOSITORY_PROTOCOL_KEY, fence);
let mut rewritten = 0_u64;
for (key, record) in pending_v6 {
writes.put(
crate::changelog::COMMIT_SPACE,
key.0.to_vec(),
crate::storage_codec::encode("commit record", &record)?,
);
rewritten += 1;
}
for (key, record) in pending {
let base_commit_id = if global_commits.contains(&record.commit_id) {
None
} else {
let newest_not_after = global_chronology
.partition_point(|(created_at, _, _)| *created_at <= record.created_at);
let Some((_, _, base)) = newest_not_after
.checked_sub(1)
.and_then(|index| global_chronology.get(index))
else {
return Err(migration_error(format!(
"{operation} local commit '{}' predates every global-lineage commit",
record.commit_id
)));
};
Some(*base)
};
let upgraded = crate::changelog::CommitRecord {
is_checkpoint: false,
format_version: crate::changelog::COMMIT_RECORD_FORMAT_VERSION,
commit_id: record.commit_id,
generation: record.generation,
parent_commit_ids: record.parent_commit_ids,
base_commit_id,
first_parent_jump_commit_id: record.first_parent_jump_commit_id,
first_parent_jump_span: record.first_parent_jump_span,
account_id: record.account_id,
created_at: record.created_at,
touched_scope_digest: record.touched_scope_digest,
};
let encoded = crate::storage_codec::encode("commit record", &upgraded)?;
writes.put(crate::changelog::COMMIT_SPACE, key.0.to_vec(), encoded);
rewritten += 1;
}
let mut write = adapter
.begin_migration_write(WriteOptions {
await_durable: true,
preconditions: vec![
Precondition::KeyValueEquals {
space: REPOSITORY_PROTOCOL_SPACE,
key: Key(Bytes::from_static(REPOSITORY_PROTOCOL_KEY)),
expected: marker,
},
crate::storage_adapter::repository_mutation_revision_precondition(
expected_revision,
),
],
..WriteOptions::default()
})
.await
.map_err(storage_error)?;
if let Err(error) = writes.lower_into(&mut write).await {
let _ = write.rollback().await;
return Err(error.into());
}
write.commit().await.map_err(storage_error)?;
Ok(rewritten)
}
#[allow(clippy::too_many_arguments)]
async fn backfill_missing_row_pk_indexes<S>(
adapter: &crate::storage_adapter::StorageAdapter<S>,
_storage: &S,
options: MigrationOptions,
expected_version: u32,
expected_protocol: &'static [u8],
target_protocol: &'static [u8],
operation: &'static str,
preflight_reserved_columns: bool,
rebuild_existing: bool,
) -> Result<(), LixError>
where
S: Storage + Clone + Send + Sync + 'static,
{
let read = super::MigrationPlanningRead::new(adapter)
.await
.map_err(storage_error)?;
match crate::init::repository_protocol_status(&read).await? {
RepositoryProtocolStatus::MigrationRequired { found_version }
if found_version == expected_version => {}
status => {
return Err(migration_error(format!(
"{operation} observed unexpected repository status {status:?}"
)));
}
}
if preflight_reserved_columns {
preflight_v74_registered_schemas(&read, options).await?;
}
let expected_revision = crate::storage_adapter::load_repository_mutation_revision(&read)
.await
.map_err(storage_error)?;
let commit_ids = crate::tracked_state::scan_commit_state_manifest_commit_ids(&read).await?;
if commit_ids.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} exceeds configured commit bound: {} commits",
commit_ids.len()
),
));
}
let mut global_commits = BTreeSet::new();
if expected_version < 74 {
let global_head = BranchHeadControlContext::new()
.reader(read.clone())
.load(crate::GLOBAL_BRANCH_ID)
.await?
.ok_or_else(|| migration_error("row-PK-index migration found no global branch"))?
.head_commit_id;
let mut next_global = Some(global_head);
while let Some(commit_id) = next_global {
if !global_commits.insert(commit_id) {
return Err(migration_error(format!(
"{operation} found a cycle in the global first-parent lineage at '{commit_id}'"
)));
}
if global_commits.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!("{operation} exceeds configured global-lineage commit bound"),
));
}
let ids = [commit_id];
let record = ChangelogContext::new()
.reader(read.clone())
.load_commits(CommitLoadRequest { commit_ids: &ids })
.await?
.into_iter()
.next()
.and_then(|(_, record)| record)
.ok_or_else(|| {
migration_error(format!(
"{operation} global lineage commit '{commit_id}' is missing"
))
})?;
next_global = record.parent_commit_ids.first().copied();
}
}
let mut chunk_writes = adapter.new_write_set();
let mut replacements = Vec::new();
let mut visited_rows = 0usize;
for commit_id in commit_ids {
let mut manifest = crate::tracked_state::load_commit_state_manifest(&read, commit_id)
.await?
.ok_or_else(|| {
migration_error(format!(
"{operation} commit '{commit_id}' has no commit-state manifest"
))
})?;
if !rebuild_existing && manifest.row_pk_index_root_id.is_some() {
continue;
}
let (root, row_count) = backfill_row_pk_index_for_commit(
&read,
&mut chunk_writes,
&manifest,
options.max_changes.saturating_sub(visited_rows),
)
.await?;
visited_rows = visited_rows
.checked_add(row_count)
.ok_or_else(|| migration_error(format!("{operation} row count exceeds usize")))?;
if expected_version < 74 {
manifest.global_scope = global_commits.contains(&commit_id);
}
manifest.row_pk_index_root_id = root;
replacements.push(manifest);
}
let chunk_stats = chunk_writes.stats();
if chunk_stats.written_bytes > options.max_preflight_bytes as u64 {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"{operation} exceeds configured byte bound: {} bytes",
chunk_stats.written_bytes
),
));
}
if chunk_stats.staged_puts != 0 || chunk_stats.staged_deletes != 0 {
let mut write = adapter
.begin_migration_write(WriteOptions {
await_durable: true,
preconditions: vec![
Precondition::KeyValueEquals {
space: REPOSITORY_PROTOCOL_SPACE,
key: Key(Bytes::from_static(REPOSITORY_PROTOCOL_KEY)),
expected: Bytes::from_static(expected_protocol),
},
crate::storage_adapter::repository_mutation_revision_precondition(
expected_revision.clone(),
),
],
..WriteOptions::default()
})
.await
.map_err(storage_error)?;
if let Err(error) = chunk_writes.lower_into(&mut write).await {
let _ = write.rollback().await;
return Err(error.into());
}
write.commit().await.map_err(storage_error)?;
}
read.finish().map_err(storage_error)?;
let mut publication = crate::migration::publish::PublicationPlan::bounded(
replacements.len().saturating_mul(2),
options.max_preflight_bytes,
);
for manifest in replacements {
for (space, key, value) in
encode_commit_state_manifest_replacement_for_migration(&manifest)?
{
publication.replace_immutable(space, vec![(key, value)])?;
}
}
crate::migration::publish::publish(
adapter,
expected_revision,
expected_protocol,
target_protocol,
publication,
)
.await
}
async fn preflight_v74_registered_schemas(
read: &(impl crate::storage_adapter::StorageAdapterRead + Clone),
options: MigrationOptions,
) -> Result<(), LixError> {
let heads = BranchHeadControlContext::new()
.reader(read.clone())
.scan()
.await?
.into_iter()
.map(|(_, control)| control.head_commit_id)
.collect::<BTreeSet<_>>();
if heads.len() > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"v74 schema preflight exceeds configured branch-head bound: {} heads",
heads.len()
),
));
}
let mut inspected = 0usize;
let mut inspected_bytes = 0usize;
for head in heads {
let remaining = options.max_changes.saturating_sub(inspected);
let rows = TrackedStateContext::new()
.reader(read.clone())
.scan_batch_at_commit(
&head.to_string(),
&TrackedStateScanRequest {
filter: TrackedStateFilter {
schema_keys: vec!["lix_registered_schema".to_owned()],
file_ids: vec![crate::NullableKeyFilter::Null],
..TrackedStateFilter::default()
},
read_columns: TrackedStateReadColumns {
columns: vec!["snapshot".to_owned()],
},
limit: Some(remaining.saturating_add(1)),
},
)
.await?;
inspected = inspected
.checked_add(rows.len())
.ok_or_else(|| migration_error("v74 schema preflight row count exceeds usize"))?;
if inspected > options.max_changes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!("v74 schema preflight exceeds configured row bound: {inspected} rows"),
));
}
for row in rows.iter() {
let schema_key = row.row_pk().as_single_string_owned().map_err(|error| {
migration_error(format!(
"v74 schema preflight found an invalid registered-schema identity: {error}"
))
})?;
let schema = match row
.decoded_snapshot()
.and_then(|typed| typed.row.get("value"))
{
Some(lix_schema::Value::Jsonb(value)) => value.clone().into_value(),
_ => {
return Err(migration_error(format!(
"v74 schema preflight could not decode registered schema '{schema_key}' at commit '{head}'"
)));
}
};
inspected_bytes = inspected_bytes
.checked_add(serde_json::to_vec(&schema).map_err(|error| {
migration_error(format!(
"v74 schema preflight could not size registered schema '{schema_key}': {error}"
))
})?.len())
.ok_or_else(|| migration_error("v74 schema preflight byte count exceeds usize"))?;
if inspected_bytes > options.max_preflight_bytes {
return Err(LixError::new(
"LIX_ERROR_MIGRATION_LIMIT_EXCEEDED",
format!(
"v74 schema preflight exceeds configured byte bound: {inspected_bytes} bytes"
),
));
}
validate_v74_registered_schema(&schema_key, &schema)?;
}
}
Ok(())
}
fn validate_v74_registered_schema(
schema_key: &str,
schema: &serde_json::Value,
) -> Result<(), LixError> {
crate::schema::parse_lix_schema(schema).map(|_| ()).map_err(|error| {
LixError::new(
"LIX_ERROR_MIGRATION_SCHEMA_INCOMPATIBLE",
format!(
"repository cannot migrate to v74 while registered schema '{schema_key}' uses a reserved column: {}. With a v73-capable engine, register a replacement schema under a new key, migrate its rows, remove the old schema, and then retry",
error.message
),
)
})
}
fn storage_error(error: StorageError) -> LixError {
if matches!(error, StorageError::ReadExpired) {
return LixError::from(error);
}
LixError::new(
LixError::CODE_INTERNAL_ERROR,
format!("repository migration storage error: {error}"),
)
}
fn migration_error(message: impl Into<String>) -> LixError {
LixError::new("LIX_ERROR_MIGRATION_FAILED", message.into())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::changelog::CommitId;
use crate::storage::{Memory, StorageWrite, WriteOptions};
use crate::storage_adapter::{PutBatch, PutEntry, StorageValue};
use crate::tracked_state::TrackedStateRootId;
#[test]
fn v74_preflight_teaches_legacy_reserved_column_rename() {
let schema = serde_json::json!({
"$schema": "https://lix.dev/schema-v1.json",
"key": "legacy_reserved_column",
"columns": [
{ "name": "id", "type": "text", "nullable": false },
{ "name": "lixcol_user_value", "type": "text", "nullable": true }
],
"primary_key": ["id"]
});
let error = validate_v74_registered_schema("legacy_reserved_column", &schema)
.expect_err("v74 must reject a formerly valid reserved column name");
assert_eq!(error.code, "LIX_ERROR_MIGRATION_SCHEMA_INCOMPATIBLE");
assert!(error.message.contains("lixcol_user_value"));
assert!(error.message.contains("replacement schema under a new key"));
}
#[tokio::test]
async fn v69_marker_is_rejected_by_the_v75_hard_cut() {
let storage = Memory::new();
let adapter = crate::storage_adapter::StorageAdapter::new(storage.clone());
let mut seed = adapter.new_write_set();
seed.put(
REPOSITORY_PROTOCOL_SPACE,
REPOSITORY_PROTOCOL_KEY,
&b"tracked-default-branch.v69"[..],
);
adapter
.commit_write_set(seed, WriteOptions::default())
.await
.unwrap();
let error = migrate_lix_with_adapter(
storage.clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect_err("v75 intentionally has no in-place migration from v69");
assert_eq!(error.code, "LIX_ERROR_MIGRATION_FAILED");
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Required {
from_version: 69,
to_version: CURRENT_FORMAT_VERSION,
}
);
}
#[tokio::test]
async fn v76_repository_advances_through_the_builtin_catalog_hard_cut() {
let storage = Memory::new();
let adapter = crate::storage_adapter::StorageAdapter::new(storage.clone());
crate::engine::Engine::initialize_with_adapter(adapter.clone(), None)
.await
.expect("repository should initialize");
let mut marker = adapter.new_write_set();
marker.put(
REPOSITORY_PROTOCOL_SPACE,
REPOSITORY_PROTOCOL_KEY,
crate::init::REPOSITORY_PROTOCOL_V76,
);
adapter
.commit_write_set(marker, WriteOptions::default())
.await
.expect("v76 marker should stage");
let report = migrate_lix_with_adapter(
storage.clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect("v76 repository should migrate to engine-authoritative built-ins");
assert_eq!(report.from_version, 76);
assert_eq!(report.to_version, CURRENT_FORMAT_VERSION);
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Current {
version: CURRENT_FORMAT_VERSION,
}
);
let engine =
crate::engine::Engine::new_with_adapter(adapter, crate::engine::EngineOptions::new())
.await
.expect("migrated repository should open");
let session = engine.open_session().await.expect("session should open");
session
.execute(
"INSERT INTO lix_file (path, content) VALUES ('/after-v77.txt', CAST('ok' AS BYTEA))",
&[],
)
.await
.expect("migrated repository should accept file descriptor writes");
}
#[tokio::test]
async fn v73_repository_migrates_through_the_v74_backfill() {
let storage = Memory::new();
let (adapter, _branch_id, _head_commit_id, rootless_commit_id, _) =
seed_rooted_head_with_rootless_checkpoint_cursor(&storage, REPOSITORY_PROTOCOL_V73)
.await;
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Required {
from_version: 73,
to_version: CURRENT_FORMAT_VERSION,
}
);
let pre_migration = SharedStorageAdapterRead::new(
adapter.begin_read(ReadOptions::default()).await.unwrap(),
);
let pre_migration_commit_ids =
crate::tracked_state::scan_commit_state_manifest_commit_ids(&pre_migration)
.await
.unwrap();
for commit_id in pre_migration_commit_ids {
let manifest =
crate::tracked_state::load_commit_state_manifest(&pre_migration, commit_id)
.await
.unwrap()
.unwrap();
assert!(
manifest.row_pk_index_root_id.is_none(),
"v73 fixture commit '{commit_id}' must genuinely lack the new index"
);
}
assert!(
crate::tracked_state::load_commit_state_manifest(&pre_migration, rootless_commit_id,)
.await
.unwrap()
.unwrap()
.snapshot_root
.is_none()
);
pre_migration.finish().unwrap();
let report = migrate_lix_with_adapter(
adapter.storage().clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect("a v73 repository migrates through the v74 backfill to v75");
assert_eq!(report.from_version, 73);
assert_eq!(report.to_version, CURRENT_FORMAT_VERSION);
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Current {
version: CURRENT_FORMAT_VERSION,
}
);
}
#[tokio::test]
async fn v72_index_bootstrap_marker_resumes_through_the_chain() {
let storage = Memory::new();
let (adapter, ..) =
seed_rooted_head_with_rootless_checkpoint_cursor(&storage, REPOSITORY_PROTOCOL_V72)
.await;
backfill_missing_row_pk_indexes(
&adapter,
adapter.storage(),
MigrationOptions::default(),
72,
REPOSITORY_PROTOCOL_V72,
crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP,
"test v72 bootstrap",
false,
false,
)
.await
.unwrap();
assert_eq!(
load_repository_protocol_marker(&adapter)
.await
.unwrap()
.as_deref(),
Some(crate::init::REPOSITORY_PROTOCOL_V72_ROW_PK_BOOTSTRAP)
);
assert!(
crate::engine::Engine::new(storage.clone()).await.is_err(),
"normal engine open must reject an interrupted bootstrap"
);
let report = migrate_lix_with_adapter(
adapter.storage().clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect("an interrupted v72 bootstrap resumes through the chain to v75");
assert_eq!(report.from_version, 72);
assert_eq!(report.to_version, CURRENT_FORMAT_VERSION);
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Current {
version: CURRENT_FORMAT_VERSION,
}
);
}
#[tokio::test]
async fn v74_commit_rewrite_fence_resumes_to_v75() {
let storage = Memory::new();
let (adapter, ..) = seed_rooted_head_with_rootless_checkpoint_cursor(
&storage,
crate::init::REPOSITORY_PROTOCOL_V74_COMMIT_REWRITE,
)
.await;
assert!(
crate::engine::Engine::new(storage.clone()).await.is_err(),
"normal engine open must reject an interrupted rewrite"
);
let report = migrate_lix_with_adapter(
adapter.storage().clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect("an interrupted v74 rewrite resumes to v75");
assert_eq!(report.from_version, 74);
assert_eq!(report.to_version, CURRENT_FORMAT_VERSION);
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Current {
version: CURRENT_FORMAT_VERSION,
}
);
}
#[tokio::test]
async fn v72_commit_rewrite_fence_resumes_through_the_chain() {
let storage = Memory::new();
let (adapter, ..) = seed_rooted_head_with_rootless_checkpoint_cursor(
&storage,
crate::init::REPOSITORY_PROTOCOL_V72_COMMIT_REWRITE,
)
.await;
let report = migrate_lix_with_adapter(
adapter.storage().clone(),
adapter.clone(),
MigrationOptions::default(),
)
.await
.expect("an interrupted v72 rewrite resumes through the chain to v75");
assert_eq!(report.from_version, 72);
assert_eq!(report.to_version, CURRENT_FORMAT_VERSION);
assert_eq!(
inspect_lix_with_adapter(&adapter).await.unwrap(),
MigrationStatus::Current {
version: CURRENT_FORMAT_VERSION,
}
);
}
fn repair_candidate_row(
directory_id: &str,
created_at: &str,
change_label: &str,
commit_id: CommitId,
) -> crate::tracked_state::MaterializedTrackedStateRow {
crate::tracked_state::MaterializedTrackedStateRow {
row_pk: crate::row_pk::RowPk::single(directory_id),
schema_key: "lix_directory_descriptor".to_string(),
file_id: None,
snapshot_content: None,
decoded_snapshot: None,
metadata: None,
deleted: false,
created_at: created_at.to_string(),
updated_at: created_at.to_string(),
change_id: crate::changelog::ChangeId::for_test_label(change_label),
commit_id,
}
}
fn repair_test_commit(label: u128) -> CommitId {
CommitId::with_change_address_space(uuid::Uuid::from_u128(label))
}
#[test]
fn repair_resolution_ranks_versions_causally_not_by_wall_clock() {
let a = repair_test_commit(0xA1);
let b = repair_test_commit(0xB2);
let c = repair_test_commit(0xC3);
let mut graph = BTreeMap::new();
graph.insert(
a,
CommitGraphNode {
parents: vec![],
base: None,
},
);
graph.insert(
b,
CommitGraphNode {
parents: vec![a],
base: None,
},
);
graph.insert(
c,
CommitGraphNode {
parents: vec![b],
base: None,
},
);
let mut candidates = BTreeMap::new();
candidates.insert(
"d1".to_string(),
vec![
repair_candidate_row("d1", "2026-01-05T00:00:00.000Z", "stale-late-clock", a),
repair_candidate_row("d1", "2026-01-01T00:00:00.000Z", "current-early-clock", b),
],
);
let present = BTreeSet::new();
let referenced = BTreeSet::from(["d1".to_string()]);
let ancestors = commit_ancestors(c, &graph);
let resolved = resolve_missing_directory_closure(
&present,
&referenced,
&ancestors,
&candidates,
&graph,
)
.expect("linear versions resolve");
assert_eq!(
resolved["d1"].change_id,
crate::changelog::ChangeId::for_test_label("current-early-clock"),
"the causally newest admissible version wins regardless of clocks"
);
}
#[test]
fn repair_resolution_fails_closed_on_incomparable_versions() {
let a = repair_test_commit(0xA1);
let b = repair_test_commit(0xB2);
let d = repair_test_commit(0xD4);
let m = repair_test_commit(0xE5);
let mut graph = BTreeMap::new();
graph.insert(
a,
CommitGraphNode {
parents: vec![],
base: None,
},
);
graph.insert(
b,
CommitGraphNode {
parents: vec![a],
base: None,
},
);
graph.insert(
d,
CommitGraphNode {
parents: vec![a],
base: None,
},
);
graph.insert(
m,
CommitGraphNode {
parents: vec![b, d],
base: None,
},
);
let mut candidates = BTreeMap::new();
candidates.insert(
"d1".to_string(),
vec![
repair_candidate_row("d1", "2026-01-01T00:00:00.000Z", "left-version", b),
repair_candidate_row("d1", "2026-01-02T00:00:00.000Z", "right-version", d),
],
);
let present = BTreeSet::new();
let referenced = BTreeSet::from(["d1".to_string()]);
let ancestors = commit_ancestors(m, &graph);
let (id, reason) = resolve_missing_directory_closure(
&present,
&referenced,
&ancestors,
&candidates,
&graph,
)
.expect_err("incomparable merge-ancestor versions must fail closed");
assert_eq!(id, "d1");
assert!(reason.contains("incomparable"), "got reason: {reason}");
}
#[tokio::test]
async fn filesystem_repair_surfaces_its_row_bound_instead_of_truncating() {
let storage = Memory::new();
let (adapter, ..) =
seed_rooted_head_with_rootless_checkpoint_cursor(&storage, REPOSITORY_PROTOCOL_V74)
.await;
let error = repair_filesystem_closure(
&adapter,
adapter.storage(),
MigrationOptions {
max_changes: 1,
..MigrationOptions::default()
},
REPOSITORY_PROTOCOL_V74,
)
.await
.expect_err("a bound below the tracked filesystem rows must error, not truncate");
assert_eq!(error.code, "LIX_ERROR_MIGRATION_LIMIT_EXCEEDED");
}
async fn seed_rooted_head_with_rootless_checkpoint_cursor(
storage: &Memory,
protocol: &'static [u8],
) -> (
crate::storage_adapter::StorageAdapter<crate::storage::StorageSession<Memory>>,
String,
CommitId,
CommitId,
TrackedStateRootId,
) {
let lix = crate::open_lix()
.with_storage(storage.clone())
.await
.unwrap();
lix.execute(
"INSERT INTO lix_key_value (key, value) VALUES ('migration-chronology-head', 'rooted')",
&[],
)
.await
.unwrap();
lix.execute(
"INSERT INTO lix_file (id, path) VALUES \
('01920000-0000-7000-8000-0000000000a1', '/migration-a'), \
('01920000-0000-7000-8000-0000000000a2', '/migration-b')",
&[],
)
.await
.unwrap();
lix.execute(
"INSERT INTO lix_key_value (key, value, lixcol_file_id) VALUES \
('migration-shared-pk', 'a', '01920000-0000-7000-8000-0000000000a1'), \
('migration-shared-pk', 'b', '01920000-0000-7000-8000-0000000000a2')",
&[],
)
.await
.unwrap();
let adapter = lix.storage_adapter().clone();
lix.close().await.unwrap();
let read = SharedStorageAdapterRead::new(
adapter.begin_read(ReadOptions::default()).await.unwrap(),
);
let controls = BranchHeadControlContext::new()
.reader(read.clone())
.scan()
.await
.unwrap();
let mut candidates = controls.into_iter().filter_map(|(branch_id, control)| {
control
.working_diff_checkpoint_commit_id
.filter(|checkpoint_commit_id| *checkpoint_commit_id != control.head_commit_id)
.map(|_| (branch_id, control))
});
let (branch_id, control) = candidates
.next()
.expect("fixture should have one branch with a distinct checkpoint cursor");
assert!(
candidates.next().is_none(),
"fixture should have only one branch with a distinct checkpoint cursor"
);
let head_commit_id = control.head_commit_id;
let checkpoint_commit_id = control
.working_diff_checkpoint_commit_id
.expect("fixture branch should retain its initial checkpoint cursor");
assert_ne!(head_commit_id, checkpoint_commit_id);
let head_manifest = crate::tracked_state::load_commit_state_manifest(&read, head_commit_id)
.await
.unwrap()
.unwrap();
let head_root_id = head_manifest
.snapshot_root
.as_ref()
.expect("fixture head should already be rooted")
.root_id
.clone();
let commit_ids = crate::tracked_state::scan_commit_state_manifest_commit_ids(&read)
.await
.unwrap();
let mut manifests = Vec::with_capacity(commit_ids.len());
for commit_id in commit_ids {
manifests.push(
crate::tracked_state::load_commit_state_manifest(&read, commit_id)
.await
.unwrap()
.unwrap(),
);
}
let checkpoint_manifest = manifests
.iter_mut()
.find(|manifest| manifest.commit_id == checkpoint_commit_id)
.expect("fixture checkpoint manifest should be retained");
assert!(checkpoint_manifest.snapshot_root.is_some());
checkpoint_manifest.snapshot_root = None;
checkpoint_manifest.replay_debt = CommitStateReplayDebt {
depth: 1,
rows: u64::from(checkpoint_manifest.mutations.member_count),
bytes: 1,
};
let replacements = manifests
.into_iter()
.flat_map(|mut manifest| {
manifest.row_pk_index_root_id = None;
encode_commit_state_manifest_replacement_for_migration(&manifest).unwrap()
})
.collect::<Vec<_>>();
read.finish().unwrap();
let mut fixture = adapter
.begin_migration_write(WriteOptions::default())
.await
.unwrap();
for (space, key, value) in replacements {
fixture
.replace_many(
space,
PutBatch {
entries: vec![PutEntry {
key: Key(Bytes::from(key)),
value: StorageValue {
bytes: Bytes::from(value),
},
}],
},
)
.await
.unwrap();
}
fixture
.put_many(
REPOSITORY_PROTOCOL_SPACE,
PutBatch {
entries: vec![PutEntry {
key: Key(Bytes::from_static(REPOSITORY_PROTOCOL_KEY)),
value: StorageValue {
bytes: Bytes::from_static(protocol),
},
}],
},
)
.await
.unwrap();
fixture.commit().await.unwrap();
(
adapter,
branch_id,
head_commit_id,
checkpoint_commit_id,
head_root_id,
)
}
}