use std::sync::Arc;
use object_store::ObjectStore;
use tracing::{info, instrument, warn};
use uni_common::api::error::UniError;
use uni_common::core::fork::{ForkInfo, ForkStatus};
use super::registry::ForkRegistryHandle;
use crate::backend::branching::ForkBranching;
#[instrument(
skip(registry, storage_store, candidate_datasets, branching),
level = "info"
)]
pub async fn recover_forks(
registry: &ForkRegistryHandle,
storage_store: &Arc<dyn ObjectStore>,
candidate_datasets: &[String],
branching: Option<&dyn ForkBranching>,
) -> Result<usize, UniError> {
let mut reconciled = 0usize;
let snapshot = registry.snapshot().await;
let pending: Vec<ForkInfo> = snapshot
.forks
.values()
.filter(|f| f.status == ForkStatus::Pending)
.cloned()
.collect();
for info in pending {
info!(fork_name = %info.name, fork_id = %info.id, "rolling back Pending create");
rollback_branches(&info, candidate_datasets, branching).await;
registry.rollback_create(&info.name).await?;
reconciled += 1;
}
let snapshot = registry.snapshot().await;
let tombstoned: Vec<ForkInfo> = snapshot
.forks
.values()
.filter(|f| f.status == ForkStatus::Tombstoned)
.cloned()
.collect();
for info in tombstoned {
info!(fork_name = %info.name, fork_id = %info.id, "completing tombstoned drop");
delete_all_branches(&info, branching).await;
registry.finish_drop(&info).await?;
super::delete_fork_artifacts(storage_store, &info.id).await;
reconciled += 1;
}
let orphans = registry.list_tombstones().await?;
for info in orphans {
info!(
fork_name = %info.name,
fork_id = %info.id,
"sweeping orphan tombstone"
);
delete_all_branches(&info, branching).await;
registry.finish_drop(&info).await?;
super::delete_fork_artifacts(storage_store, &info.id).await;
reconciled += 1;
}
Ok(reconciled)
}
async fn delete_all_branches(info: &ForkInfo, branching: Option<&dyn ForkBranching>) {
let Some(branching) = branching else {
return;
};
for (dataset, branch) in &info.datasets {
if let Err(e) = branching.delete_branch(dataset, branch).await {
warn!(
dataset = %dataset,
branch = %branch,
"delete_branch during recovery failed: {e}"
);
}
}
}
async fn rollback_branches(
info: &ForkInfo,
candidate_datasets: &[String],
branching: Option<&dyn ForkBranching>,
) {
if !info.datasets.is_empty() {
delete_all_branches(info, branching).await;
}
let Some(branching) = branching else {
return;
};
for dataset in candidate_datasets {
let branch = format!("fork_{}_{}", info.id, dataset);
if let Err(e) = branching.delete_branch(dataset, &branch).await {
warn!(
dataset = %dataset,
branch = %branch,
"zombie-branch reclamation during recovery failed: {e}"
);
}
}
}