1use std::{
5 collections::HashMap,
6 fs,
7 path::{Path, PathBuf},
8 sync::{OnceLock, RwLock},
9 time::SystemTime,
10};
11
12use tracing::{debug, instrument, trace};
13
14use crate::{
15 object::ContentHash,
16 store::{
17 Result,
18 pack::{ObjectType, PackObjectId, PackReadTier, PackReader},
19 },
20};
21
22pub struct PackManager {
26 packs_dir: PathBuf,
27 packs: Vec<CachedPack>,
28 scratch_root: PathBuf,
29 object_locations: RwLock<ObjectLocationIndex>,
30 eager_object_locations: bool,
31}
32
33#[derive(Default)]
34struct ObjectLocationIndex {
35 locations: HashMap<PackObjectId, ObjectLocation>,
36 complete: bool,
37}
38
39#[derive(Clone, Copy)]
40struct ObjectLocation {
41 pack_index: usize,
42 tier: PackReadTier,
43}
44
45struct CachedPack {
46 pack_path: PathBuf,
47 index_path: PathBuf,
48 reader: OnceLock<Option<PackReader<'static>>>,
49 scratch_root: PathBuf,
50}
51
52impl CachedPack {
53 fn discovered(pack_path: PathBuf, index_path: PathBuf, scratch_root: PathBuf) -> Self {
54 Self {
55 pack_path,
56 index_path,
57 reader: OnceLock::new(),
58 scratch_root,
59 }
60 }
61
62 fn validated(
63 pack_path: PathBuf,
64 index_path: PathBuf,
65 reader: PackReader<'static>,
66 scratch_root: PathBuf,
67 ) -> Self {
68 Self {
69 pack_path,
70 index_path,
71 reader: OnceLock::from(Some(reader)),
72 scratch_root,
73 }
74 }
75
76 fn reader(&self) -> Option<&PackReader<'static>> {
77 self.reader
78 .get_or_init(|| {
79 match PackReader::open_lazy(&self.pack_path, &self.index_path, &self.scratch_root) {
80 Ok(reader) => Some(reader),
81 Err(error) => {
82 debug!(pack = ?self.pack_path, %error, "Failed to open pack");
83 None
84 }
85 }
86 })
87 .as_ref()
88 }
89
90 fn verified_reader(&self) -> Option<&PackReader<'static>> {
91 self.reader
92 .get_or_init(|| {
93 match PackReader::open(&self.pack_path, &self.index_path, &self.scratch_root) {
94 Ok(reader) => Some(reader),
95 Err(error) => {
96 debug!(pack = ?self.pack_path, %error, "Failed to open pack");
97 None
98 }
99 }
100 })
101 .as_ref()
102 }
103}
104
105impl PackManager {
106 pub fn new(packs_dir: PathBuf, scratch_root: PathBuf) -> Self {
107 Self::new_with_scratch(packs_dir, scratch_root, force_eager_pack_index())
108 }
109
110 #[cfg(test)]
111 fn new_with_index_mode(packs_dir: PathBuf, eager_object_locations: bool) -> Self {
112 let scratch_root = packs_dir.join("tmp");
113 Self::new_with_scratch(packs_dir, scratch_root, eager_object_locations)
114 }
115
116 fn new_with_scratch(
117 packs_dir: PathBuf,
118 scratch_root: PathBuf,
119 eager_object_locations: bool,
120 ) -> Self {
121 let packs = Self::load_packs(&packs_dir, &scratch_root).unwrap_or_default();
122 let object_locations = Self::initial_object_locations(&packs, eager_object_locations);
123 Self {
124 packs_dir,
125 packs,
126 scratch_root,
127 object_locations: RwLock::new(object_locations),
128 eager_object_locations,
129 }
130 }
131
132 fn discover_pack_paths(packs_dir: &Path) -> Result<Vec<(PathBuf, PathBuf)>> {
133 let mut packs = Vec::new();
134
135 if !packs_dir.exists() {
136 return Ok(packs);
137 }
138
139 for entry in fs::read_dir(packs_dir)? {
140 let entry = entry?;
141 let path = entry.path();
142
143 if path.extension().map(|e| e == "pack").unwrap_or(false) {
144 let index_path = path.with_extension("idx");
145 if index_path.exists() {
146 packs.push((path, index_path));
147 }
148 }
149 }
150
151 packs.sort_by(|left, right| {
157 pack_modified(&left.0)
158 .cmp(&pack_modified(&right.0))
159 .then_with(|| left.0.cmp(&right.0))
160 });
161
162 debug!(count = packs.len(), "Discovered pack files");
163 Ok(packs)
164 }
165
166 fn load_packs(packs_dir: &Path, scratch_root: &Path) -> Result<Vec<CachedPack>> {
167 Ok(Self::discover_pack_paths(packs_dir)?
168 .into_iter()
169 .map(|(pack_path, index_path)| {
170 CachedPack::discovered(pack_path, index_path, scratch_root.to_path_buf())
171 })
172 .collect())
173 }
174
175 pub fn reload(&mut self) -> Result<()> {
176 self.packs = Self::load_packs(&self.packs_dir, &self.scratch_root)?;
177 self.reset_object_locations();
178 Ok(())
179 }
180
181 fn initial_object_locations(packs: &[CachedPack], eager: bool) -> ObjectLocationIndex {
182 if !eager {
183 return ObjectLocationIndex::default();
184 }
185 let mut locations = HashMap::new();
186 for (pack_index, pack) in packs.iter().enumerate() {
187 let Some(reader) = pack.verified_reader() else {
188 continue;
189 };
190 let Ok(objects) = reader.indexed_read_tiers() else {
191 continue;
192 };
193 for (id, tier) in objects {
194 remember_location(&mut locations, id, pack_index, tier);
195 }
196 }
197 ObjectLocationIndex {
198 locations,
199 complete: true,
200 }
201 }
202
203 fn reset_object_locations(&mut self) {
204 self.object_locations = RwLock::new(Self::initial_object_locations(
205 &self.packs,
206 self.eager_object_locations,
207 ));
208 }
209
210 fn object_location(&self, id: &PackObjectId) -> Result<Option<usize>> {
211 {
212 let index = self
213 .object_locations
214 .read()
215 .unwrap_or_else(std::sync::PoisonError::into_inner);
216 if let Some(location) = index.locations.get(id) {
217 return Ok(Some(location.pack_index));
218 }
219 if index.complete {
220 return Ok(None);
221 }
222 }
223
224 let mut index = self
225 .object_locations
226 .write()
227 .unwrap_or_else(std::sync::PoisonError::into_inner);
228 if !index.complete {
229 for (pack_index, pack) in self.packs.iter().enumerate() {
230 let Some(reader) = pack.reader() else {
231 continue;
232 };
233 let Ok(objects) = reader.indexed_read_tiers() else {
234 continue;
235 };
236 for (object_id, tier) in objects {
237 remember_location(&mut index.locations, object_id, pack_index, tier);
238 }
239 }
240 index.complete = true;
241 }
242 Ok(index.locations.get(id).map(|location| location.pack_index))
243 }
244
245 fn point_object_location(&self, id: &PackObjectId) -> Result<Option<usize>> {
249 {
250 let index = self
251 .object_locations
252 .read()
253 .unwrap_or_else(std::sync::PoisonError::into_inner);
254 if let Some(location) = index.locations.get(id) {
255 return Ok(Some(location.pack_index));
256 }
257 if index.complete {
258 return Ok(None);
259 }
260 }
261 for (pack_index, pack) in self.packs.iter().enumerate().rev() {
262 let Some(reader) = pack.reader() else {
263 continue;
264 };
265 if reader.contains_object(id)? {
266 return Ok(Some(pack_index));
267 }
268 }
269 Ok(None)
270 }
271
272 pub fn add_pack(&mut self, pack_path: PathBuf, index_path: PathBuf) -> Result<()> {
274 if self.packs.iter().any(|pack| pack.pack_path == pack_path) {
275 return Ok(());
276 }
277 let reader = PackReader::open(&pack_path, &index_path, &self.scratch_root)?;
278 let pack_index = self.packs.len();
279 let cached =
280 CachedPack::validated(pack_path, index_path, reader, self.scratch_root.clone());
281 self.packs.push(cached);
282 let mut index = self
283 .object_locations
284 .write()
285 .unwrap_or_else(std::sync::PoisonError::into_inner);
286 if index.complete {
287 let objects = self.packs[pack_index]
288 .reader()
289 .ok_or_else(|| {
290 crate::store::StoreError::InvalidObject("new pack reader unavailable".into())
291 })?
292 .indexed_read_tiers()?;
293 for (id, tier) in objects {
294 remember_location(&mut index.locations, id, pack_index, tier);
295 }
296 }
297 Ok(())
298 }
299
300 pub fn needs_reload(&self) -> Result<bool> {
307 let discovered = Self::discover_pack_paths(&self.packs_dir)?;
308 Ok(discovered.len() != self.packs.len()
309 || discovered
310 .iter()
311 .zip(&self.packs)
312 .any(|((pack, index), cached)| {
313 *pack != cached.pack_path || *index != cached.index_path
314 }))
315 }
316
317 pub fn reload_if_stale(&mut self) -> Result<bool> {
327 if !self.needs_reload()? {
328 return Ok(false);
329 }
330 debug!("PackManager: pack set changed under us, reloading");
331 self.reload()?;
332 Ok(true)
333 }
334
335 pub fn get_object(&self, id: &PackObjectId) -> Result<Option<(ObjectType, Vec<u8>)>> {
336 let Some(pack_index) = self.point_object_location(id)? else {
337 trace!("Object not found in any pack");
338 return Ok(None);
339 };
340 let Some(reader) = self.packs[pack_index].reader() else {
341 return Ok(None);
342 };
343 let object = reader.get_object(id)?;
344 if object.is_some() {
345 trace!("Found object in pack");
346 }
347 Ok(object)
348 }
349
350 pub fn object_read_tier(&self, id: &PackObjectId) -> Result<Option<PackReadTier>> {
355 let _ = self.object_location(id)?;
356 let index = self
357 .object_locations
358 .read()
359 .unwrap_or_else(std::sync::PoisonError::into_inner);
360 Ok(index.locations.get(id).map(|location| location.tier))
361 }
362
363 pub fn get_object_from_pack(
367 &self,
368 pack_path: &Path,
369 id: &PackObjectId,
370 ) -> Result<Option<(ObjectType, Vec<u8>)>> {
371 let Some(pack) = self.packs.iter().find(|pack| pack.pack_path == pack_path) else {
372 return Ok(None);
373 };
374 let Some(reader) = pack.reader() else {
375 return Ok(None);
376 };
377 reader.get_object(id)
378 }
379
380 pub fn list_ids_from_pack(&self, pack_path: &Path) -> Result<Vec<PackObjectId>> {
383 let Some(pack) = self.packs.iter().find(|pack| pack.pack_path == pack_path) else {
384 return Ok(Vec::new());
385 };
386 let Some(reader) = pack.reader() else {
387 return Ok(Vec::new());
388 };
389 reader.list_ids()
390 }
391
392 #[instrument(skip(self), fields(hash = %hash.short()))]
393 pub fn get_hashed_object(&self, hash: &ContentHash) -> Result<Option<(ObjectType, Vec<u8>)>> {
394 self.get_object(&PackObjectId::Hash(*hash))
395 }
396
397 pub fn get_hashed_object_type(&self, hash: &ContentHash) -> Result<Option<ObjectType>> {
399 let id = PackObjectId::Hash(*hash);
400 let Some(pack_index) = self.point_object_location(&id)? else {
401 return Ok(None);
402 };
403 let Some(reader) = self.packs[pack_index].reader() else {
404 return Ok(None);
405 };
406 reader.get_hashed_object_type(hash)
407 }
408
409 pub fn get_hashed_object_bytes(
414 &self,
415 hash: &ContentHash,
416 ) -> Result<Option<(ObjectType, bytes::Bytes)>> {
417 let id = PackObjectId::Hash(*hash);
418 let Some(pack_index) = self.point_object_location(&id)? else {
419 return Ok(None);
420 };
421 let Some(reader) = self.packs[pack_index].reader() else {
422 return Ok(None);
423 };
424 reader.get_object_bytes(&id)
425 }
426
427 pub fn has_object(&self, hash: &ContentHash) -> bool {
428 self.point_object_location(&PackObjectId::Hash(*hash))
429 .is_ok_and(|location| location.is_some())
430 }
431
432 pub fn get_hashed_object_size(&self, hash: &ContentHash) -> Result<Option<u64>> {
436 let id = PackObjectId::Hash(*hash);
437 let Some(pack_index) = self.point_object_location(&id)? else {
438 return Ok(None);
439 };
440 let Some(reader) = self.packs[pack_index].reader() else {
441 return Ok(None);
442 };
443 reader.get_hashed_object_size(hash)
444 }
445
446 pub fn has_object_id(&self, id: &PackObjectId) -> bool {
447 self.point_object_location(id)
448 .is_ok_and(|location| location.is_some())
449 }
450
451 pub fn list_all_hashes(&self) -> Result<Vec<ContentHash>> {
453 let mut hashes = Vec::new();
454 for pack in &self.packs {
455 if let Some(reader) = pack.reader() {
456 hashes.extend(reader.list_hashes()?);
457 }
458 }
459 Ok(hashes)
460 }
461
462 pub fn list_all_ids(&self) -> Result<Vec<PackObjectId>> {
463 let mut ids = Vec::new();
464 for pack in &self.packs {
465 if let Some(reader) = pack.reader() {
466 ids.extend(reader.list_ids()?);
467 }
468 }
469 Ok(ids)
470 }
471
472 pub fn pack_file_paths(&self) -> Vec<(&Path, &Path)> {
474 self.packs
475 .iter()
476 .map(|pack| (pack.pack_path.as_path(), pack.index_path.as_path()))
477 .collect()
478 }
479
480 pub fn pack_count(&self) -> usize {
481 self.packs.len()
482 }
483
484 pub fn packs_dir(&self) -> &Path {
485 &self.packs_dir
486 }
487}
488
489fn pack_modified(path: &Path) -> SystemTime {
490 fs::metadata(path)
491 .and_then(|metadata| metadata.modified())
492 .unwrap_or(SystemTime::UNIX_EPOCH)
493}
494
495fn remember_location(
496 locations: &mut HashMap<PackObjectId, ObjectLocation>,
497 id: PackObjectId,
498 pack_index: usize,
499 tier: PackReadTier,
500) {
501 let candidate = ObjectLocation { pack_index, tier };
502 match locations.get_mut(&id) {
503 Some(existing)
504 if existing.tier == PackReadTier::SolidFrame && tier == PackReadTier::Hot =>
505 {
506 *existing = candidate;
507 }
508 Some(_) => {}
509 None => {
510 locations.insert(id, candidate);
511 }
512 }
513}
514
515fn force_eager_pack_index() -> bool {
516 std::env::var("HEDDLE_PERF_FORCE_EAGER_PACK_INDEX")
517 .is_ok_and(|value| matches!(value.as_str(), "1" | "true" | "yes"))
518}
519
520#[cfg(test)]
521#[path = "manager_tests.rs"]
522mod tests;