use super::block_fetch::load_segment_index_for_reorganization;
use super::build::{
build_manifest_tables_from_rows, debug_assert_manifest_table_segments_do_not_overlap,
MetadataTableSegmentation,
};
use super::error::ManifestLoadError;
use super::flush::{ensure_metadata_publication_budget, next_manifest_id_after};
use super::load::{
load_namespace_manifest_envelope_if_present, load_verified_manifest_tables,
validate_direntry_child_bind_index, validate_revision_by_inode_desc_index,
};
use super::publish::{publish_metadata_root, write_namespace_manifest, ManifestPublicationOutcome};
use super::runs::{
flatten_manifest_tables, l0_run_count, MetadataLsmPolicy, MetadataRunManifest,
CHECKPOINT_BASE_RUN_LEVEL, CHECKPOINT_L0_RUN_LEVEL,
};
use super::scan::VerifiedMetadataTables;
use crate::context::MutationContext;
use crate::error::{CoreError, MetadataProjectionLoadError, Result};
use crate::limits::CONTENTION_RETRY_LIMIT;
use crate::namespace::basis::resolve_retention_floor_seq;
use crate::namespace::control::{read_head_object, read_metadata_root_object_if_present};
use crate::timing::{MonotonicTimer, StdMonotonicTimer};
use loonfs_api::wire::manifest::{
ActiveDeletionRowAction, MetadataFileRef, MetadataRow, MetadataTableFamily,
NamespaceManifestEnvelope, NamespaceManifestPayload,
};
use loonfs_api::{ChangeSeq, InodeId, ManifestId, ManifestObjectId, NamespaceId};
use loonfs_objectstore::keys::metadata_manifest_object;
use loonfs_objectstore::ObjectStore;
use std::collections::{BTreeMap, BTreeSet};
const REORGANIZE_FAMILY_GROUPS: [&[MetadataTableFamily]; 6] = [
&[
MetadataTableFamily::DirentryBinds,
MetadataTableFamily::DirentryChildBinds,
MetadataTableFamily::DirentryUnbinds,
],
&[
MetadataTableFamily::Revisions,
MetadataTableFamily::RevisionsByInodeDesc,
],
&[MetadataTableFamily::Inodes],
&[MetadataTableFamily::Tombstones],
&[MetadataTableFamily::ActiveDeletions],
&[MetadataTableFamily::CommitReceipts],
];
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum MetadataReorganizeOutcome {
NotNeeded { l0_runs: usize },
UnitPublished {
families: Vec<MetadataTableFamily>,
folded_l0_rows: u64,
input_runs: usize,
decoded_input_rows: u64,
decoded_input_bytes: u64,
manifest_id: ManifestId,
},
BudgetExhausted {
families: Vec<MetadataTableFamily>,
l0_runs: usize,
},
Superseded,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MetadataReorganizeReport {
pub namespace_id: NamespaceId,
pub outcome: MetadataReorganizeOutcome,
}
#[tracing::instrument(
level = "info",
name = "loonfs.phase",
err,
skip_all,
fields(phase = "reorganize_metadata", key_class = "manifest")
)]
pub(crate) async fn reorganize_metadata_step<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
context: &MutationContext,
policy: MetadataLsmPolicy,
) -> Result<MetadataReorganizeReport> {
let timer = StdMonotonicTimer::default();
reorganize_metadata_step_with_timer(store, namespace_id, context, policy, &timer).await
}
pub(super) async fn reorganize_metadata_step_with_timer<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
context: &MutationContext,
policy: MetadataLsmPolicy,
timer: &dyn MonotonicTimer,
) -> Result<MetadataReorganizeReport> {
let publication_started_ms = timer.monotonic_now_ms();
let Some(root) = read_metadata_root_object_if_present(store, namespace_id)
.await
.map_err(CoreError::load_head)?
.map(|loaded| loaded.envelope.state)
else {
return Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::NotNeeded { l0_runs: 0 },
});
};
let tables = load_verified_manifest_tables(store, namespace_id, &root.manifest_object_id)
.await
.map_err(|error| {
CoreError::MetadataProjection(MetadataProjectionLoadError::ManifestLoad(error))
})?;
let previous = tables.manifest();
let l0_runs = l0_run_count(&previous.payload);
if l0_runs < policy.max_l0_runs.get()
&& !manifest_has_partial_reorganization(tables.scan_runs.as_ref())
{
return Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::NotNeeded { l0_runs },
});
}
let Some(group) = select_family_group(&previous.payload) else {
return Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::NotNeeded { l0_runs },
});
};
let Some(input) = select_reorganization_input(&tables, group, policy)
.await
.map_err(|error| {
CoreError::MetadataProjection(MetadataProjectionLoadError::ManifestLoad(error))
})?
else {
return Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::BudgetExhausted {
families: group.to_vec(),
l0_runs,
},
});
};
let mut rows_by_family = BTreeMap::<MetadataTableFamily, Vec<MetadataRow>>::new();
for family in group {
let rows = tables
.scan_prefix_in_runs(&input.runs, *family, "")
.await
.map_err(|error| {
CoreError::MetadataProjection(MetadataProjectionLoadError::ManifestLoad(error))
})?;
rows_by_family.insert(*family, rows);
}
if group.contains(&MetadataTableFamily::DirentryBinds) {
validate_direntry_child_bind_index(
root.manifest_object_id.as_ref(),
rows_by_family
.get(&MetadataTableFamily::DirentryBinds)
.map_or(&[], Vec::as_slice),
rows_by_family
.get(&MetadataTableFamily::DirentryChildBinds)
.map_or(&[], Vec::as_slice),
)
.map_err(|error| {
CoreError::MetadataProjection(MetadataProjectionLoadError::ManifestLoad(error))
})?;
}
if group.contains(&MetadataTableFamily::Revisions) {
validate_revision_by_inode_desc_index(
root.manifest_object_id.as_ref(),
rows_by_family
.get(&MetadataTableFamily::Revisions)
.map_or(&[], Vec::as_slice),
rows_by_family
.get(&MetadataTableFamily::RevisionsByInodeDesc)
.map_or(&[], Vec::as_slice),
)
.map_err(|error| {
CoreError::MetadataProjection(MetadataProjectionLoadError::ManifestLoad(error))
})?;
}
let head = read_head_object(store, namespace_id)
.await
.map_err(CoreError::load_head)?
.envelope
.state;
let floor_seq = resolve_retention_floor_seq(store, &head)
.await
.map_err(CoreError::load_head)?;
drop_rows_below_retention_floor(&mut rows_by_family, floor_seq)?;
let run_tables = build_manifest_tables_from_rows(
store,
namespace_id,
previous.payload.head_seq,
CHECKPOINT_BASE_RUN_LEVEL,
|family| rows_by_family.remove(&family).unwrap_or_default(),
MetadataTableSegmentation::Base {
max_rows_per_segment: policy.max_rows_per_segment,
},
)
.await?;
debug_assert_manifest_table_segments_do_not_overlap(&run_tables);
let mut metadata_files: Vec<_> = previous
.payload
.metadata_files
.iter()
.filter(|descriptor| {
!group.contains(&descriptor.family)
|| !input
.run_ids
.contains(&(descriptor.run_seq, descriptor.level))
})
.cloned()
.collect();
metadata_files.extend(flatten_manifest_tables(run_tables));
let base_seq = metadata_files
.iter()
.map(|descriptor| descriptor.run_seq)
.min()
.unwrap_or(previous.payload.base_seq);
let manifest =
write_reorganized_manifest(store, namespace_id, previous, metadata_files, base_seq, {
let mut payload_floor = previous.payload.retention_floor_seq;
if floor_seq > payload_floor {
payload_floor = floor_seq;
}
payload_floor
})
.await?;
ensure_metadata_publication_budget(timer, publication_started_ms, namespace_id)?;
match publish_metadata_root(
store,
namespace_id,
&manifest,
Some(root.manifest_object_id.clone()),
context.now_ms,
)
.await?
{
ManifestPublicationOutcome::Published(_) => Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::UnitPublished {
families: group.to_vec(),
folded_l0_rows: input.folded_l0_rows,
input_runs: input.runs.len(),
decoded_input_rows: input.decoded_rows,
decoded_input_bytes: input.decoded_bytes,
manifest_id: manifest.payload.manifest_id,
},
}),
ManifestPublicationOutcome::Superseded(_) | ManifestPublicationOutcome::RootCasRaceLost => {
Ok(MetadataReorganizeReport {
namespace_id: namespace_id.clone(),
outcome: MetadataReorganizeOutcome::Superseded,
})
}
}
}
struct ReorganizationInput {
runs: Vec<MetadataRunManifest>,
run_ids: BTreeSet<(ChangeSeq, u32)>,
folded_l0_rows: u64,
decoded_rows: u64,
decoded_bytes: u64,
}
async fn select_reorganization_input<S: ObjectStore + ?Sized>(
tables: &VerifiedMetadataTables<'_, S>,
group: &[MetadataTableFamily],
policy: MetadataLsmPolicy,
) -> std::result::Result<Option<ReorganizationInput>, ManifestLoadError> {
let mut candidates = tables
.scan_runs
.iter()
.filter(|run| run_has_group_rows(run, group))
.collect::<Vec<_>>();
candidates.sort_by(|left, right| {
let left_is_l0 = left.level == CHECKPOINT_L0_RUN_LEVEL;
let right_is_l0 = right.level == CHECKPOINT_L0_RUN_LEVEL;
left_is_l0
.cmp(&right_is_l0)
.then(left.run_seq.cmp(&right.run_seq))
.then(right.level.cmp(&left.level))
});
let candidate_count = candidates.len();
let row_budget =
u64::try_from(policy.max_decoded_input_rows_per_step.get()).unwrap_or(u64::MAX);
let byte_budget =
u64::try_from(policy.max_decoded_input_bytes_per_step.get()).unwrap_or(u64::MAX);
let mut runs = Vec::new();
let mut decoded_rows = 0u64;
let mut decoded_bytes = 0u64;
let mut folded_l0_rows = 0u64;
for run in candidates
.into_iter()
.take(policy.max_input_runs_per_step.get())
{
let run_rows = group_run_descriptors(run, group)
.map(|descriptor| descriptor.row_count)
.sum::<u64>();
if decoded_rows.saturating_add(run_rows) > row_budget {
break;
}
let run_bytes = decoded_group_run_bytes(tables, run, group).await?;
if decoded_bytes.saturating_add(run_bytes) > byte_budget {
break;
}
if run.level == CHECKPOINT_L0_RUN_LEVEL {
folded_l0_rows = folded_l0_rows.saturating_add(run_rows);
}
decoded_rows = decoded_rows.saturating_add(run_rows);
decoded_bytes = decoded_bytes.saturating_add(run_bytes);
runs.push(run.clone());
}
let selected_l0 = runs.iter().any(|run| run.level == CHECKPOINT_L0_RUN_LEVEL);
let makes_progress = selected_l0 && (runs.len() > 1 || candidate_count == 1);
if !makes_progress {
return Ok(None);
}
let run_ids = runs.iter().map(|run| (run.run_seq, run.level)).collect();
Ok(Some(ReorganizationInput {
runs,
run_ids,
folded_l0_rows,
decoded_rows,
decoded_bytes,
}))
}
fn manifest_has_partial_reorganization(runs: &[MetadataRunManifest]) -> bool {
let Some(oldest_l0_seq) = runs
.iter()
.filter(|run| run.level == CHECKPOINT_L0_RUN_LEVEL)
.map(|run| run.run_seq)
.min()
else {
return false;
};
runs.iter()
.any(|run| run.level != CHECKPOINT_L0_RUN_LEVEL && run.run_seq >= oldest_l0_seq)
}
fn run_has_group_rows(run: &MetadataRunManifest, group: &[MetadataTableFamily]) -> bool {
group_run_descriptors(run, group).next().is_some()
}
fn group_run_descriptors<'a>(
run: &'a MetadataRunManifest,
group: &'a [MetadataTableFamily],
) -> impl Iterator<Item = &'a MetadataFileRef> {
run.tables
.iter()
.filter(|table| group.contains(&table.family))
.flat_map(|table| &table.segments)
}
async fn decoded_group_run_bytes<S: ObjectStore + ?Sized>(
tables: &VerifiedMetadataTables<'_, S>,
run: &MetadataRunManifest,
group: &[MetadataTableFamily],
) -> std::result::Result<u64, ManifestLoadError> {
let mut decoded_bytes = 0u64;
for descriptor in group_run_descriptors(run, group) {
let index = load_segment_index_for_reorganization(
tables.store,
tables.table_cache,
&tables.block_memo,
descriptor,
)
.await?;
for entry in index.iter() {
decoded_bytes = decoded_bytes.saturating_add(u64::from(entry.block.decoded_len));
}
}
Ok(decoded_bytes)
}
fn select_family_group(
payload: &NamespaceManifestPayload,
) -> Option<&'static [MetadataTableFamily]> {
REORGANIZE_FAMILY_GROUPS
.into_iter()
.map(|group| (group_l0_rows(payload, group), group))
.filter(|(rows, _)| *rows > 0)
.max_by(|(left_rows, left), (right_rows, right)| {
left_rows.cmp(right_rows).then_with(|| {
position_of(right).cmp(&position_of(left))
})
})
.map(|(_, group)| group)
}
fn position_of(group: &[MetadataTableFamily]) -> usize {
REORGANIZE_FAMILY_GROUPS
.iter()
.position(|candidate| candidate.as_ptr() == group.as_ptr())
.unwrap_or(usize::MAX)
}
fn group_l0_rows(payload: &NamespaceManifestPayload, group: &[MetadataTableFamily]) -> u64 {
payload
.metadata_files
.iter()
.filter(|descriptor| {
descriptor.level == CHECKPOINT_L0_RUN_LEVEL && group.contains(&descriptor.family)
})
.map(|descriptor| descriptor.row_count)
.sum()
}
async fn write_reorganized_manifest<S: ObjectStore + ?Sized>(
store: &S,
namespace_id: &NamespaceId,
previous: &NamespaceManifestEnvelope,
metadata_files: Vec<MetadataFileRef>,
base_seq: ChangeSeq,
retention_floor_seq: ChangeSeq,
) -> Result<NamespaceManifestEnvelope> {
let manifest_id = next_manifest_id_after(previous.payload.manifest_id)?;
for _allocation_attempt in 0..CONTENTION_RETRY_LIMIT {
let manifest_object_id = ManifestObjectId::generate(manifest_id);
let manifest_key = metadata_manifest_object(namespace_id.as_str(), &manifest_object_id);
match load_namespace_manifest_envelope_if_present(
store,
namespace_id,
&manifest_object_id,
&manifest_key,
)
.await
{
Ok(Some(_existing)) => continue,
Ok(None) => {}
Err(error) => {
return Err(CoreError::MetadataProjection(
MetadataProjectionLoadError::ManifestLoad(error),
))
}
}
let manifest = NamespaceManifestEnvelope::from_payload(NamespaceManifestPayload {
namespace_id: namespace_id.clone(),
manifest_id,
manifest_object_id,
head_seq: previous.payload.head_seq,
head_commit_id: previous.payload.head_commit_id.clone(),
base_seq,
writer_epoch: previous.payload.writer_epoch,
next_inode_id: previous.payload.next_inode_id,
retention_floor_seq,
metadata_files: metadata_files.clone(),
})
.map_err(|err| {
CoreError::Internal(format!("failed to build reorganized manifest: {err}"))
})?;
match write_namespace_manifest(store, &manifest).await {
Ok(()) => return Ok(manifest),
Err(MetadataProjectionLoadError::ManifestLoad(
ManifestLoadError::ManifestConflict { .. },
)) => continue,
Err(error) => return Err(CoreError::MetadataProjection(error)),
}
}
Err(CoreError::Internal(
"reorganized manifest allocation retry exhausted".to_owned(),
))
}
pub(super) fn drop_rows_below_retention_floor(
rows_by_family: &mut BTreeMap<MetadataTableFamily, Vec<MetadataRow>>,
retention_floor_seq: ChangeSeq,
) -> Result<()> {
let mut unbound_at_floor = BTreeSet::new();
for row in rows_by_family
.get(&MetadataTableFamily::DirentryUnbinds)
.into_iter()
.flatten()
{
if let MetadataRow::DirentryUnbind {
parent_inode_id,
name_key,
bind_seq,
bind_delta_index,
unbind_seq,
..
} = row
{
if *unbind_seq <= retention_floor_seq {
unbound_at_floor.insert((
*parent_inode_id,
name_key.clone(),
*bind_seq,
*bind_delta_index,
));
}
}
}
let mut latest_bind_at_floor = BTreeMap::new();
for row in rows_by_family
.get(&MetadataTableFamily::DirentryBinds)
.into_iter()
.flatten()
{
if let MetadataRow::DirentryBind {
parent_inode_id,
name_key,
bind_seq,
bind_delta_index,
..
} = row
{
if *bind_seq <= retention_floor_seq {
let candidate = (*bind_seq, *bind_delta_index);
let latest = latest_bind_at_floor
.entry((*parent_inode_id, name_key.clone()))
.or_insert(candidate);
if candidate > *latest {
*latest = candidate;
}
}
}
}
for row in rows_by_family
.get(&MetadataTableFamily::DirentryBinds)
.into_iter()
.flatten()
{
if let MetadataRow::DirentryBind {
parent_inode_id,
name_key,
bind_seq,
bind_delta_index,
..
} = row
{
if *bind_seq <= retention_floor_seq
&& latest_bind_at_floor.get(&(*parent_inode_id, name_key.clone()))
!= Some(&(*bind_seq, *bind_delta_index))
&& !unbound_at_floor.contains(&(
*parent_inode_id,
name_key.clone(),
*bind_seq,
*bind_delta_index,
))
{
return Err(CoreError::NamespaceCorrupt(format!(
"bind at seq `{bind_seq}` delta {bind_delta_index} for parent `{parent_inode_id}` is superseded at or below the retention floor without an unbind; refusing to drop rows"
)));
}
}
}
let retain_bind = |row: &MetadataRow| match row {
MetadataRow::DirentryBind {
parent_inode_id,
name_key,
bind_seq,
bind_delta_index,
..
} => {
*bind_seq > retention_floor_seq
|| (latest_bind_at_floor.get(&(*parent_inode_id, name_key.clone()))
== Some(&(*bind_seq, *bind_delta_index))
&& !unbound_at_floor.contains(&(
*parent_inode_id,
name_key.clone(),
*bind_seq,
*bind_delta_index,
)))
}
_ => true,
};
for family in [
MetadataTableFamily::DirentryBinds,
MetadataTableFamily::DirentryChildBinds,
] {
if let Some(rows) = rows_by_family.get_mut(&family) {
rows.retain(retain_bind);
}
}
if let Some(rows) = rows_by_family.get_mut(&MetadataTableFamily::DirentryUnbinds) {
rows.retain(|row| match row {
MetadataRow::DirentryUnbind { unbind_seq, .. } => *unbind_seq > retention_floor_seq,
_ => true,
});
}
if let Some(rows) = rows_by_family.get_mut(&MetadataTableFamily::ActiveDeletions) {
let revoked: BTreeSet<(ChangeSeq, InodeId)> = rows
.iter()
.filter_map(|row| match row {
MetadataRow::ActiveDeletion {
root_inode_id,
deleted_at_seq,
action: ActiveDeletionRowAction::Removed { .. },
} => Some((*deleted_at_seq, *root_inode_id)),
_ => None,
})
.collect();
rows.retain(|row| match row {
MetadataRow::ActiveDeletion {
root_inode_id,
deleted_at_seq,
action,
} => match action {
ActiveDeletionRowAction::Removed { .. } => false,
ActiveDeletionRowAction::Listed { .. } => {
!revoked.contains(&(*deleted_at_seq, *root_inode_id))
}
},
_ => true,
});
}
if let Some(rows) = rows_by_family.get_mut(&MetadataTableFamily::CommitReceipts) {
rows.retain(|row| match row {
MetadataRow::CommitReceipt { committed_seq, .. } => {
*committed_seq >= retention_floor_seq
}
_ => true,
});
}
Ok(())
}