#![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();
for id in crate::store::transaction::descriptor_objects(desc, &store.config().limits) {
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: &Store,
options: OptimizeOptions,
max_extents: Option<u64>,
mut cursor: Option<&mut PassCursor>,
) -> Result<BackgroundStats, StoreError> {
store.ensure_epoch_flushed(&crate::store::transaction::CrashHooks::none())?;
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 limits = *store.limits();
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; evaluate_commit_extent(store, at, start, desc_bytes, &[], options, &mut stats)?;
}
idx += 1;
}
if !truncated {
if let Some(c) = cursor {
*c = PassCursor::default();
}
}
let rebased = store.rebase_overdepth_extents(&crate::store::transaction::CrashHooks::none())?;
stats.rewritten = stats.rewritten.saturating_add(rebased);
Ok(stats)
}
fn evaluate_commit_extent(
store: &Store,
ino: u64,
start: u64,
desc_bytes: Vec<u8>,
shared_pool: &[crate::core::candidate::BaseChunk],
options: OptimizeOptions,
stats: &mut BackgroundStats,
) -> Result<(), StoreError> {
let limits = *store.limits();
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;
return Ok(());
}
};
let bytes = match materialize_to_vec(&desc, store, &limits) {
Ok(b) => b,
Err(_) => {
stats.errors += 1;
return Ok(());
}
};
let cid = ChunkId::of(&bytes);
let rebased = crate::optimizer::rebase::flatten_if_deep(store, start, &desc, &bytes, &cid)?;
let mut searched: Vec<(crate::store::ExtentUpdate, u64)> = Vec::new();
if shared_pool.is_empty() {
let ctx = GuidedContext {
ino,
offset: start,
target: &bytes,
prev_version: None,
dictionary: if start >= limits.chunk_class {
store.base_chunk_at(ino, start - limits.chunk_class, bytes.len())?
} else {
None
},
shared: None,
pending: None,
mode: SearchMode::Background,
};
match encode_guided(
store,
&ctx,
options,
crate::optimizer::foreground::ForegroundPolicy::full(),
) {
Ok(o) => searched.push((o.update, 0)),
Err(_) => stats.errors += 1,
}
} else {
for anchor in shared_pool {
let ctx = GuidedContext {
ino,
offset: start,
target: &bytes,
prev_version: None,
dictionary: if start >= limits.chunk_class {
store.base_chunk_at(ino, start - limits.chunk_class, bytes.len())?
} else {
None
},
shared: Some(anchor.clone()),
pending: None,
mode: SearchMode::Background,
};
match encode_guided(
store,
&ctx,
options,
crate::optimizer::foreground::ForegroundPolicy::full(),
) {
Ok(o) => searched.push((o.update, 0)),
Err(_) => stats.errors += 1,
}
}
}
let current_bytes = current_persisted_bytes(store, &desc);
let mut best: Option<crate::store::ExtentUpdate> = None;
let mut best_bytes = u64::MAX;
for (update, _) in &searched {
if update.descriptor != desc {
let b = update_persisted_bytes(update);
if b < best_bytes {
best_bytes = b;
best = Some(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;
return Ok(());
};
if best_bytes >= current_bytes {
stats.no_gain += 1;
return Ok(());
}
let _lock = store.inode_lock(ino);
let current_desc = store.extent_descriptor(ino, start)?;
let stale = match current_desc {
Some(cur) => cur != desc_bytes,
None => true,
};
if stale {
stats.stale_skips += 1;
return Ok(());
}
if update.content_id != cid {
stats.errors += 1;
return Ok(());
}
store.commit_file_extents(ino, vec![update], None, &CrashHooks::none())?;
stats.rewritten += 1;
stats.saved_bytes = stats.saved_bytes.saturating_add(current_bytes - best_bytes);
Ok(())
}
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) }
pub fn shared_dict_pass(
store: &Store,
options: OptimizeOptions,
max_extents: Option<u64>,
) -> Result<BackgroundStats, StoreError> {
shared_dict_pass_pool(store, options, max_extents, MAX_ANCHOR_POOL)
}
pub(crate) fn shared_dict_pass_pool(
store: &Store,
options: OptimizeOptions,
max_extents: Option<u64>,
max_pool: usize,
) -> Result<BackgroundStats, StoreError> {
let mut stats = BackgroundStats::default();
if !options.allow_shared_dict {
return Ok(stats);
}
let dir_of = build_dir_map(store)?;
let mut by_dir: std::collections::BTreeMap<u64, Vec<MemberChunk>> =
std::collections::BTreeMap::new();
for ino in store.all_inodes()? {
let Some(inode) = store.get_inode(ino)? else {
continue;
};
let extent_root = match inode.data {
crate::store::inode::InodeData::File { extent_root } => extent_root,
_ => continue,
};
if extent_root.is_zero() {
continue;
}
let limits = *store.limits();
let Some(first_bytes) = store.extent_descriptor(ino, 0)? else {
continue;
};
let desc = match crate::format::descriptor::decode(
&first_bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
) {
Ok(d) => d,
Err(_) => continue,
};
if desc.len() > limits.chunk_class {
continue;
}
let bytes = match materialize_to_vec(&desc, store, &limits) {
Ok(b) => b,
Err(_) => continue,
};
let dir = dir_of.get(&ino).copied().unwrap_or(1);
let incumbent = current_persisted_bytes(store, &desc);
by_dir
.entry(dir)
.or_default()
.push(MemberChunk { bytes, incumbent });
}
let mut pools: std::collections::BTreeMap<u64, Vec<crate::core::candidate::BaseChunk>> =
std::collections::BTreeMap::new();
for (dir, members) in &by_dir {
if members.len() < 2 {
continue;
}
let pool = select_anchor_pool(store, members, max_pool)?;
if !pool.is_empty() {
pools.insert(*dir, pool);
}
}
let inos = store.all_inodes()?;
let mut scanned = 0u64;
'ino: for ino in &inos {
let at = *ino;
let Some(inode) = store.get_inode(at)? else {
continue;
};
let extent_root = match inode.data {
crate::store::inode::InodeData::File { extent_root } => extent_root,
_ => continue,
};
if extent_root.is_zero() {
continue;
}
let Some(dir) = dir_of.get(&at).copied() else {
continue;
};
let Some(pool) = pools.get(&dir) else {
continue;
};
let limits = *store.limits();
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 let Some(max) = max_extents {
if scanned >= max {
break 'ino;
}
}
scanned += 1;
stats.scanned += 1;
evaluate_commit_extent(store, at, start, desc_bytes, pool, options, &mut stats)?;
}
}
let rebased = store.rebase_overdepth_extents(&crate::store::transaction::CrashHooks::none())?;
stats.rewritten = stats.rewritten.saturating_add(rebased);
Ok(stats)
}
fn build_dir_map(store: &Store) -> Result<std::collections::HashMap<u64, u64>, StoreError> {
let fanout = store.config().limits.max_fanout;
let mut map: std::collections::HashMap<u64, u64> = std::collections::HashMap::new();
let root_dir = store.current_root().root_dir_ino;
for dir_ino in store.all_inodes()? {
let Some(inode) = store.get_inode(dir_ino)? else {
continue;
};
let dir_root = match inode.data {
crate::store::inode::InodeData::Directory { dir_root } => dir_root,
_ => continue,
};
if dir_root.is_zero() {
continue;
}
let entries =
crate::store::index::scan_all(dir_root, crate::store::BTREE_ORDER, fanout, store)?;
for (_, v) in entries {
if let Ok(e) = crate::store::directory::DirEntry::decode(&v) {
map.entry(e.ino).or_insert(dir_ino);
}
}
}
map.entry(root_dir).or_insert(root_dir);
Ok(map)
}
struct MemberChunk {
bytes: Vec<u8>,
incumbent: u64,
}
fn select_anchor_pool(
store: &Store,
members: &[MemberChunk],
max_pool: usize,
) -> Result<Vec<crate::core::candidate::BaseChunk>, StoreError> {
use crate::core::candidate::{CandidateContext, Encoder};
let limits = *store.limits();
let policy = *store.policy();
let mut seen = std::collections::HashSet::new();
let mut cands: Vec<&[u8]> = Vec::new();
let mut sorted: Vec<&[u8]> = members.iter().map(|m| m.bytes.as_slice()).collect();
sorted.sort_by_key(|b| std::cmp::Reverse(b.len()));
for b in sorted {
if b.len() < MIN_ANCHOR_BYTES {
continue;
}
let id = ChunkId::of(b);
if seen.insert(id) {
cands.push(b);
}
if cands.len() >= MAX_ANCHOR_CANDIDATES {
break;
}
}
struct Cand {
cid: ChunkId,
bytes: Vec<u8>,
saved: Vec<u64>, }
let mut pool_cands: Vec<Cand> = Vec::new();
for c in cands {
let Some(desc_bytes) = store.chunk_descriptor(&ChunkId::of(c))? else {
continue;
};
let Ok(desc) = crate::format::descriptor::decode(
&desc_bytes,
limits.max_descriptor_bytes,
limits.max_inline_bytes,
limits.max_palette,
limits.max_period,
limits.max_chunk_size,
) else {
continue;
};
if crate::core::cost::reference_depth(&desc) != 0 {
continue;
}
let cid = ChunkId::of(c);
let mut saved = Vec::with_capacity(members.len());
for m in members {
if m.bytes == c {
saved.push(0); continue;
}
let ctx = CandidateContext {
limits: &limits,
policy: &policy,
content_id: ChunkId::of(&m.bytes),
bases: &[],
dedup: None,
};
let enc = crate::rans::sequence::SequenceSharedDictEncoder {
dictionary: crate::core::extent::ChunkId::ZERO,
dict_bytes: Vec::new(),
dict_depth: 0,
shared: cid,
shared_bytes: c.to_vec(),
shared_depth: 0,
};
let cost = enc
.encode(&m.bytes, &ctx)
.into_iter()
.map(|cand| cand.cost.persisted_bytes())
.min()
.unwrap_or(m.incumbent);
saved.push(m.incumbent.saturating_sub(cost));
}
pool_cands.push(Cand {
cid,
bytes: c.to_vec(),
saved,
});
}
let mut pool: Vec<crate::core::candidate::BaseChunk> = Vec::new();
let mut covered: Vec<u64> = vec![0; members.len()];
let mut selected: Vec<bool> = vec![false; pool_cands.len()];
for _ in 0..max_pool {
let mut best_idx: Option<usize> = None;
let mut best_gain: u64 = 0;
let mut best_id: ChunkId = ChunkId::ZERO;
for (i, cand) in pool_cands.iter().enumerate() {
if selected[i] {
continue;
}
let mut gain = 0u64;
for (m, &s) in cand.saved.iter().enumerate() {
gain = gain.saturating_add(s.saturating_sub(covered[m]));
}
let better = match best_idx {
None => gain > 0,
Some(_) => {
gain > best_gain
|| (gain == best_gain && cand.cid.as_bytes() < best_id.as_bytes())
}
};
if better && gain > 0 {
best_idx = Some(i);
best_gain = gain;
best_id = cand.cid;
}
}
let Some(idx) = best_idx else {
break; };
let cand = &pool_cands[idx];
for (m, &s) in cand.saved.iter().enumerate() {
covered[m] = covered[m].max(s);
}
selected[idx] = true;
pool.push(crate::core::candidate::BaseChunk {
id: cand.cid,
bytes: cand.bytes.clone(),
depth: 0,
});
}
Ok(pool)
}
const MAX_ANCHOR_CANDIDATES: usize = 12;
const MAX_ANCHOR_POOL: usize = 4;
const MIN_ANCHOR_BYTES: usize = 512;
const WORKER_CYCLE_EXTENTS: u64 = 64;
const WORKER_IDLE_SECS: u64 = 3;
pub fn spawn_background_worker(
store: Arc<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;
}
let _ = optimize_pass(
&store,
options,
Some(WORKER_CYCLE_EXTENTS),
Some(&mut cursor),
);
let _ = shared_dict_pass(&store, options, Some(WORKER_CYCLE_EXTENTS));
let _ = model_bundle_pass(&store, options, Some(WORKER_CYCLE_EXTENTS));
last_ops = ops.load(Ordering::Relaxed);
}
})
.expect("spawn background worker")
}
const MAX_MODEL_POOL: usize = 4;
struct ModelMember {
ino: u64,
start: u64,
desc_bytes: Vec<u8>,
desc: crate::core::representation::Representation,
streams: Vec<Vec<u8>>,
types: Vec<u8>,
pinned: u64,
}
type ModelBundle = std::collections::BTreeMap<u8, crate::rans::model::RansModel>;
struct BundleEncode {
descriptor: crate::core::representation::Representation,
enc_payload: Vec<u8>,
model_payload: Vec<u8>,
pinned: u64,
}
fn bundle_model_bytes(b: &ModelBundle) -> u64 {
b.values()
.map(|m| crate::rans::metadata::encode_model(m).len() as u64)
.sum()
}
fn bundle_key(b: &ModelBundle) -> Vec<u8> {
b.iter()
.flat_map(|(_, m)| crate::rans::metadata::encode_model(m))
.collect()
}
fn collect_model_member(
store: &Store,
ino: u64,
start: u64,
desc_bytes: Vec<u8>,
limits: &crate::core::limits::Limits,
) -> Result<Option<ModelMember>, StoreError> {
use crate::core::representation::Representation;
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(_) => return Ok(None),
};
let (scale_bits, codec) = crate::rans::sequence::sequence_scale_codec();
let (streams, types, enc_id): (Vec<Vec<u8>>, Vec<u8>, ChunkId) = match &desc {
Representation::SequenceRans {
model,
enc_obj,
scale_bits: sb,
codec: c,
seq_len,
lit_len,
off_len,
cmds,
lit_out,
..
} => {
if *sb != scale_bits || *c != codec {
return Ok(None);
}
let v = match crate::rans::sequence::decode_streams_n(
store,
limits,
crate::rans::sequence::StreamRefs {
model: *model,
enc_obj: *enc_obj,
scale_bits: *sb,
codec: *c,
},
&[*seq_len, *lit_len, *off_len],
*cmds as u64,
*lit_out as u64,
None,
2,
) {
Ok(v) => v,
Err(_) => return Ok(None),
};
(v, vec![0, 1, 2], *enc_obj)
}
Representation::SequenceDeep {
model,
enc_obj,
scale_bits: sb,
codec: c,
seq_len,
lit_len,
off_len,
len_len,
cmds,
lit_out,
..
} => {
if *sb != scale_bits || *c != codec {
return Ok(None);
}
let d = match crate::rans::sequence::decode_deep_streams(
store,
limits,
crate::rans::sequence::StreamRefs {
model: *model,
enc_obj: *enc_obj,
scale_bits: *sb,
codec: *c,
},
crate::rans::sequence::DeepLens {
seq_len: *seq_len,
lit_len: *lit_len,
off_len: *off_len,
len_len: *len_len,
cmds: *cmds,
lit_out: *lit_out,
},
) {
Ok(d) => d,
Err(_) => return Ok(None),
};
(
vec![d.commands, d.literals, d.offsets, d.lengths],
vec![0, 1, 2, 4],
*enc_obj,
)
}
Representation::SequenceDict {
dictionary: _,
dictionary_len: _,
model,
enc_obj,
scale_bits: sb,
codec: c,
seq_len,
lit_len,
off_len,
src_len,
cmds,
lit_out,
..
}
| Representation::SequenceSharedDict {
dictionary: _,
dictionary_len: _,
shared: _,
shared_len: _,
model,
enc_obj,
scale_bits: sb,
codec: c,
seq_len,
lit_len,
off_len,
src_len,
cmds,
lit_out,
..
} => {
if *sb != scale_bits || *c != codec {
return Ok(None);
}
let d = match crate::rans::sequence::decode_four_streams(
store,
limits,
crate::rans::sequence::StreamRefs {
model: *model,
enc_obj: *enc_obj,
scale_bits: *sb,
codec: *c,
},
crate::rans::sequence::FourStreams {
seq_len: *seq_len,
lit_len: *lit_len,
off_len: *off_len,
src_len: *src_len,
cmds: *cmds,
lit_out: *lit_out,
},
) {
Ok(d) => d,
Err(_) => return Ok(None),
};
(
vec![d.commands, d.literals, d.offsets, d.sources],
vec![0, 1, 2, 3],
*enc_obj,
)
}
_ => return Ok(None),
};
let pinned = desc.encoded_size().saturating_add(
store
.object_index()
.get(&enc_id)
.map(|l| l.stored_len)
.unwrap_or(0),
);
Ok(Some(ModelMember {
ino,
start,
desc_bytes,
desc,
streams,
types,
pinned,
}))
}
fn rebuild_descriptor(
incumbent: &crate::core::representation::Representation,
streams: &[Vec<u8>],
enc: &crate::rans::sequence::EncodedStreams,
model_id: ChunkId,
enc_id: ChunkId,
) -> Option<crate::core::representation::Representation> {
use crate::core::representation::Representation;
let (scale_bits, codec) = crate::rans::sequence::sequence_scale_codec();
match incumbent {
Representation::SequenceRans { len, .. } => Some(Representation::SequenceRans {
model: model_id,
enc_obj: enc_id,
scale_bits,
codec,
seq_len: enc.lens[0],
lit_len: enc.lens[1],
off_len: enc.lens[2],
cmds: streams[0].len() as u32,
lit_out: streams[1].len() as u32,
len: *len,
}),
Representation::SequenceDeep { len, .. } => Some(Representation::SequenceDeep {
model: model_id,
enc_obj: enc_id,
scale_bits,
codec,
seq_len: enc.lens[0],
lit_len: enc.lens[1],
off_len: enc.lens[2],
len_len: enc.lens[3],
cmds: streams[0].len() as u32,
lit_out: streams[1].len() as u32,
len: *len,
}),
Representation::SequenceDict {
dictionary,
dictionary_len,
len,
..
} => Some(Representation::SequenceDict {
dictionary: *dictionary,
dictionary_len: *dictionary_len,
model: model_id,
enc_obj: enc_id,
scale_bits,
codec,
seq_len: enc.lens[0],
lit_len: enc.lens[1],
off_len: enc.lens[2],
src_len: enc.lens[3],
cmds: streams[0].len() as u32,
lit_out: streams[1].len() as u32,
len: *len,
}),
Representation::SequenceSharedDict {
dictionary,
dictionary_len,
shared,
shared_len,
len,
..
} => Some(Representation::SequenceSharedDict {
dictionary: *dictionary,
dictionary_len: *dictionary_len,
shared: *shared,
shared_len: *shared_len,
model: model_id,
enc_obj: enc_id,
scale_bits,
codec,
seq_len: enc.lens[0],
lit_len: enc.lens[1],
off_len: enc.lens[2],
src_len: enc.lens[3],
cmds: streams[0].len() as u32,
lit_out: streams[1].len() as u32,
len: *len,
}),
_ => None,
}
}
fn encode_member_against(member: &ModelMember, bundle: &ModelBundle) -> Option<BundleEncode> {
let mut forced: Vec<Option<&crate::rans::model::RansModel>> = vec![None; member.types.len()];
for (si, &t) in member.types.iter().enumerate() {
forced[si] = bundle.get(&t);
}
let enc = crate::rans::sequence::encode_streams_n_with_models(&member.streams, &forced)?;
let enc_obj = crate::core::candidate::ObjectRecord::data(enc.enc_obj.clone());
let model_obj = crate::core::candidate::ObjectRecord::model(enc.model_obj.clone());
let descriptor = rebuild_descriptor(
&member.desc,
&member.streams,
&enc,
model_obj.id,
enc_obj.id,
)?;
let pinned = descriptor
.encoded_size()
.saturating_add(enc_obj.payload.len() as u64)
.saturating_add(4); if pinned >= member.pinned {
return None;
}
Some(BundleEncode {
descriptor,
enc_payload: enc_obj.payload,
model_payload: model_obj.payload,
pinned,
})
}
pub fn model_bundle_pass(
store: &Store,
options: OptimizeOptions,
max_extents: Option<u64>,
) -> Result<BackgroundStats, StoreError> {
let mut stats = BackgroundStats::default();
if !(options.allow_sequence_rans
|| options.allow_sequence_rans_deep
|| options.allow_sequence_dict
|| options.allow_shared_dict)
{
return Ok(stats);
}
let dir_of = build_dir_map(store)?;
let limits = *store.limits();
let mut by_dir: std::collections::BTreeMap<u64, Vec<ModelMember>> =
std::collections::BTreeMap::new();
let mut collected = 0u64;
'ino: for ino in store.all_inodes()? {
let Some(inode) = store.get_inode(ino)? else {
continue;
};
let extent_root = match inode.data {
crate::store::inode::InodeData::File { extent_root } => extent_root,
_ => continue,
};
if extent_root.is_zero() {
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 let Some(max) = max_extents {
if collected >= max {
break 'ino;
}
}
let Some(dir) = dir_of.get(&ino).copied() else {
continue;
};
let Some(member) = collect_model_member(store, ino, start, desc_bytes, &limits)? else {
continue;
};
collected += 1;
stats.scanned = stats.scanned.saturating_add(1);
by_dir.entry(dir).or_default().push(member);
}
}
for (_dir, members) in &by_dir {
if members.len() < 2 {
continue;
}
let mut agg_hist: std::collections::BTreeMap<u8, [u32; 256]> = Default::default();
for m in members {
for (si, &t) in m.types.iter().enumerate() {
let h = agg_hist.entry(t).or_insert([0u32; 256]);
for &b in &m.streams[si] {
h[b as usize] = h[b as usize].saturating_add(1);
}
}
}
let mut agg_bundle: ModelBundle = Default::default();
for (t, h) in &agg_hist {
if let Some(model) = crate::rans::sequence::aggregate_model(h) {
agg_bundle.insert(*t, model);
}
}
if agg_bundle.is_empty() {
continue;
}
let mut cands: Vec<ModelBundle> = Vec::new();
let mut seen: std::collections::HashSet<Vec<u8>> = std::collections::HashSet::new();
seen.insert(bundle_key(&agg_bundle));
cands.push(agg_bundle);
for m in members {
let mut b: ModelBundle = Default::default();
for (si, &t) in m.types.iter().enumerate() {
if b.contains_key(&t) {
continue;
}
let mut h = [0u32; 256];
for &x in &m.streams[si] {
h[x as usize] += 1;
}
if let Some(model) = crate::rans::sequence::aggregate_model(&h) {
b.insert(t, model);
}
}
if seen.insert(bundle_key(&b)) {
cands.push(b);
}
}
let bundle_costs: Vec<u64> = cands.iter().map(bundle_model_bytes).collect();
let mut gains: Vec<Vec<u64>> = Vec::with_capacity(cands.len());
let mut encodes: Vec<Vec<Option<BundleEncode>>> = Vec::with_capacity(cands.len());
for bundle in &cands {
let mut g = Vec::with_capacity(members.len());
let mut es = Vec::with_capacity(members.len());
for m in members {
let trial = encode_member_against(m, bundle);
match &trial {
Some(e) => {
g.push(m.pinned - e.pinned);
es.push(trial);
}
_ => {
g.push(0);
es.push(None);
}
}
}
gains.push(g);
encodes.push(es);
}
let mut selected: Vec<usize> = Vec::new();
let mut covered: Vec<u64> = vec![0; members.len()];
for _ in 0..MAX_MODEL_POOL {
let mut best_idx: Option<usize> = None;
let mut best_marginal = 0u64;
for i in 0..cands.len() {
if selected.contains(&i) {
continue;
}
let mut marginal = 0u64;
for m in 0..members.len() {
marginal = marginal.saturating_add(gains[i][m].saturating_sub(covered[m]));
}
if marginal > bundle_costs[i] && marginal - bundle_costs[i] > best_marginal {
best_marginal = marginal - bundle_costs[i];
best_idx = Some(i);
}
}
let Some(idx) = best_idx else {
break;
};
selected.push(idx);
for m in 0..members.len() {
covered[m] = covered[m].max(gains[idx][m]);
}
}
if selected.is_empty() {
continue;
}
let mut rewrite: Vec<(usize, usize)> = Vec::new(); let mut model_payloads: Vec<Vec<u8>> = Vec::new();
for m in 0..members.len() {
if covered[m] == 0 {
continue;
}
let mut best_i = selected[0];
let mut best_gain = gains[selected[0]][m];
for &i in &selected[1..] {
if gains[i][m] > best_gain {
best_gain = gains[i][m];
best_i = i;
}
}
let e = encodes[best_i][m]
.as_ref()
.expect("gain > 0 implies an encode");
let mp = e.model_payload.clone();
if !model_payloads.contains(&mp) {
model_payloads.push(mp);
}
rewrite.push((m, best_i));
}
if rewrite.is_empty() {
continue;
}
let mut model_cost = 0u64;
for p in &model_payloads {
if !store.object_index().contains(&ChunkId::of(p)) {
model_cost = model_cost.saturating_add(p.len() as u64);
}
}
let group_savings: u64 = covered.iter().sum();
if group_savings <= model_cost {
stats.no_gain = stats.no_gain.saturating_add(members.len() as u64);
continue;
}
let group_gain = group_savings - model_cost;
for (m, i) in rewrite {
let member = &members[m];
let e = encodes[i][m]
.as_ref()
.expect("rewrite set implies an encode");
let bytes = match materialize_to_vec(&member.desc, store, &limits) {
Ok(b) => b,
Err(_) => {
stats.errors += 1;
continue;
}
};
let cid = ChunkId::of(&bytes);
let objects = vec![
crate::core::candidate::ObjectRecord::data(e.enc_payload.clone()),
crate::core::candidate::ObjectRecord::model(e.model_payload.clone()),
];
let candidate = crate::core::candidate::Candidate {
representation: e.descriptor.clone(),
objects: objects.clone(),
cost: Default::default(),
content_id: cid,
};
let resolver = crate::optimizer::search::CandidateResolver::new(
store,
objects.iter().map(|o| (o.id, o.payload.clone())).collect(),
None,
);
if crate::core::candidate::validate_candidate(&candidate, &bytes, &resolver, &limits)
.is_err()
{
stats.errors += 1;
continue;
}
let _lock = store.inode_lock(member.ino);
let current = store.extent_descriptor(member.ino, member.start)?;
if current.as_deref() != Some(member.desc_bytes.as_slice()) {
stats.stale_skips += 1;
continue;
}
store.commit_file_extents(
member.ino,
vec![crate::store::ExtentUpdate {
offset: member.start,
descriptor: e.descriptor.clone(),
content_id: cid,
objects,
}],
None,
&CrashHooks::none(),
)?;
stats.rewritten = stats.rewritten.saturating_add(1);
}
stats.saved_bytes = stats.saved_bytes.saturating_add(group_gain);
}
let rebased = store.rebase_overdepth_extents(&CrashHooks::none())?;
stats.rewritten = stats.rewritten.saturating_add(rebased);
Ok(stats)
}