#![forbid(unsafe_code)]
use std::sync::Arc;
use crate::core::extent::ChunkId;
use crate::core::materialize::materialize_to_vec;
use crate::optimizer::policy::OptimizeOptions;
use crate::optimizer::search::{GuidedContext, SearchMode, encode_guided};
use crate::store::transaction::CrashHooks;
use crate::store::{Store, StoreError};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct BackgroundStats {
pub scanned: u64,
pub rewritten: u64,
pub saved_bytes: u64,
pub stale_skips: u64,
pub no_gain: u64,
pub errors: u64,
}
pub fn current_persisted_bytes(
store: &Store,
desc: &crate::core::representation::Representation,
) -> u64 {
let mut total = desc.encoded_size();
let mut object_ids: Vec<&ChunkId> = Vec::new();
match desc {
crate::core::representation::Representation::Raw { obj, .. } => object_ids.push(obj),
crate::core::representation::Representation::Rans { model, enc_obj, .. } => {
object_ids.push(model);
object_ids.push(enc_obj);
}
_ => {}
}
for id in object_ids {
if let Some(loc) = store.object_index().get(id) {
total = total.saturating_add(loc.stored_len);
}
}
total
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PassCursor {
pub ino_index: usize,
pub offset: u64,
}
pub fn optimize_pass(
store: &mut Store,
options: OptimizeOptions,
max_extents: Option<u64>,
mut cursor: Option<&mut PassCursor>,
) -> Result<BackgroundStats, StoreError> {
let limits = *store.limits();
let mut stats = BackgroundStats::default();
let inos = store.all_inodes()?;
let start_idx = cursor.as_ref().map(|c| c.ino_index).unwrap_or(0);
let mut resume_offset = cursor.as_ref().map(|c| c.offset).unwrap_or(0);
let mut idx = 0usize;
let mut truncated = false;
'ino: for ino in &inos {
let at = *ino;
if idx < start_idx {
idx += 1;
continue;
}
let inode = match store.get_inode(at)? {
Some(i) => i,
None => {
idx += 1;
continue;
}
};
let extent_root = match inode.data {
crate::store::inode::InodeData::File { extent_root } => extent_root,
_ => {
idx += 1;
continue;
}
};
if extent_root.is_zero() {
idx += 1;
continue;
}
let entries = crate::store::extent_tree::scan_all(
extent_root,
crate::store::BTREE_ORDER,
limits.max_fanout,
store,
)?;
for (start, desc_bytes) in entries {
if start < resume_offset {
continue; }
if let Some(max) = max_extents {
if stats.scanned >= max {
if let Some(c) = cursor.as_mut() {
c.ino_index = idx;
c.offset = start;
}
truncated = true;
break 'ino;
}
}
stats.scanned += 1;
resume_offset = 0; let desc = match crate::format::descriptor::decode(
&desc_bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
) {
Ok(d) => d,
Err(_) => {
stats.errors += 1;
continue;
}
};
let bytes = match materialize_to_vec(&desc, store, &limits) {
Ok(b) => b,
Err(_) => {
stats.errors += 1;
continue;
}
};
let cid = ChunkId::of(&bytes);
let rebased =
crate::optimizer::rebase::flatten_if_deep(store, start, &desc, &bytes, &cid)?;
let ctx = GuidedContext {
ino: at,
offset: start,
target: &bytes,
prev_version: None,
mode: SearchMode::Background,
};
let searched = match encode_guided(store, &ctx, options) {
Ok(o) => Some(o),
Err(_) => {
stats.errors += 1;
None
}
};
let current_bytes = current_persisted_bytes(store, &desc);
let mut best: Option<crate::store::ExtentUpdate> = None;
let mut best_bytes = u64::MAX;
if let Some(outcome) = &searched {
if outcome.update.descriptor != desc {
let b = update_persisted_bytes(&outcome.update);
if b < best_bytes {
best_bytes = b;
best = Some(outcome.update.clone());
}
}
}
if let Some(u) = rebased {
if u.descriptor != desc {
let b = update_persisted_bytes(&u);
if b < best_bytes {
best_bytes = b;
best = Some(u);
}
}
}
let Some(update) = best else {
stats.no_gain += 1;
continue;
};
if best_bytes >= current_bytes {
stats.no_gain += 1;
continue;
}
let current_desc = store.extent_descriptor(at, start)?;
let stale = match current_desc {
Some(cur) => cur != desc_bytes,
None => true,
};
if stale {
stats.stale_skips += 1;
continue;
}
if update.content_id != cid {
stats.errors += 1;
continue;
}
store.commit_file_extents(at, vec![update], None, &CrashHooks::none())?;
stats.rewritten += 1;
stats.saved_bytes = stats.saved_bytes.saturating_add(current_bytes - best_bytes);
}
idx += 1;
}
if !truncated {
if let Some(c) = cursor {
*c = PassCursor::default();
}
}
Ok(stats)
}
fn update_persisted_bytes(update: &crate::store::ExtentUpdate) -> u64 {
let mut total = update.descriptor.encoded_size();
for o in &update.objects {
total = total.saturating_add(o.payload.len() as u64);
}
total.saturating_add(4) }
const WORKER_CYCLE_EXTENTS: u64 = 64;
const WORKER_IDLE_SECS: u64 = 3;
pub fn spawn_background_worker(
store: Arc<std::sync::Mutex<Store>>,
ops: Arc<std::sync::atomic::AtomicU64>,
stop: Arc<std::sync::atomic::AtomicBool>,
options: OptimizeOptions,
) -> std::thread::JoinHandle<()> {
use std::sync::atomic::Ordering;
std::thread::Builder::new()
.name("entropyfs-optimizer".into())
.spawn(move || {
let mut cursor = PassCursor::default();
let mut last_ops = ops.load(Ordering::Relaxed);
while !stop.load(Ordering::SeqCst) {
std::thread::sleep(std::time::Duration::from_secs(WORKER_IDLE_SECS));
if stop.load(Ordering::SeqCst) {
break;
}
let now = ops.load(Ordering::Relaxed);
if now != last_ops {
last_ops = now;
continue;
}
if let Ok(mut s) = store.try_lock() {
let _ = optimize_pass(
&mut s,
options,
Some(WORKER_CYCLE_EXTENTS),
Some(&mut cursor),
);
}
last_ops = ops.load(Ordering::Relaxed);
}
})
.expect("spawn background worker")
}