1use std::{
5 collections::HashSet,
6 fs::{self, File, OpenOptions},
7 path::{Path, PathBuf},
8};
9
10use fs2::FileExt;
11use heddle_format::compression::{header_uncompressed_size, is_compressed};
12use tracing::{debug, instrument, trace};
13
14use super::{
15 FsStore,
16 fs_io::{list_hashes_from_dir, read_file_bytes, read_file_header},
17 fs_paths::{
18 action_path, actions_dir, annotated_tags_dir, blobs_dir, hash_path, partial_tree_path,
19 partial_trees_dir, redaction_path, redactions_dir, state_attachment_index_lock_path,
20 state_attachment_index_path, state_attachment_path, state_attachments_dir, state_path,
21 state_visibility_dir, state_visibility_path, states_dir, tree_lineage_path, trees_dir,
22 },
23};
24use crate::{
25 object::{
26 Action, ActionId, AnnotatedTag, Blob, BytesTreeSource, ContentHash, FileTreeSource,
27 OpenedTreeBody, State, StateAttachment, StateAttachmentId, StateId, TREE_CANONICAL_MAGIC,
28 TREE_DELTA_HEADER_LEN, TREE_DELTA_MAGIC, TREE_LEAN_MAGIC, TREE_SALTED_MAGIC, Tree,
29 TreeByteSource, TreeEntry, TreeEntryReader, TreeResumeCursor, decode_tree_delta_header,
30 decode_tree_delta_header_prefix, is_delta_tree, is_redacted_tree, is_salted_tree,
31 is_streamable_tree,
32 },
33 store::{
34 HeddleError, ObjectCacheControl, ObjectStore, Result, SidecarStore,
35 SnapshotCommitDescriptor, TreeWrite, codec,
36 codec::{EncodedTree, TreeDeltaBase, TreeEncodingKind, TreeLineage},
37 delta_source::DeltaTreeSource,
38 pack::{ObjectType, PackManager, PackObjectId},
39 },
40};
41
42const BLOB_HEADER_PEEK: usize = 13;
50
51fn validate_loaded_tree(tree: Tree) -> Result<Tree> {
52 tree.validate()?;
53 Ok(tree)
54}
55
56fn streamable_body(tree: &Tree) -> Result<Vec<u8>> {
59 if tree.has_git_layout() {
60 Ok(tree.encode_canonical()?)
61 } else {
62 Ok(tree.encode_lean()?)
63 }
64}
65
66fn salted_tree_not_streamable(tree_id: &ContentHash) -> HeddleError {
73 HeddleError::InvalidObject(format!(
74 "tree {tree_id} is an HSR1 salted (v4) tree; the paging reader cannot \
75 stream it — decode it eagerly via get_tree"
76 ))
77}
78
79fn validate_blob_bytes(data: &[u8], hash: ContentHash) -> Result<()> {
80 let mut hasher = ContentHash::typed_hasher("blob", data.len() as u64);
81 hasher.update(data);
82 let found = ContentHash::from_bytes(hasher.finalize().into());
83 if found != hash {
84 return Err(HeddleError::Corruption {
85 expected: hash,
86 found,
87 });
88 }
89
90 Ok(())
91}
92
93fn validate_tree_serialized(data: &[u8], hash: ContentHash) -> Result<Tree> {
94 let tree = codec::decode_tree_serialized_with_key(data, hash, None)?;
95 let tree = validate_loaded_tree(tree)?;
96 let found = tree.hash();
97 if found != hash {
98 return Err(HeddleError::Corruption {
99 expected: hash,
100 found,
101 });
102 }
103
104 Ok(tree)
105}
106
107fn validate_annotated_tag(data: &[u8], hash: ContentHash) -> Result<AnnotatedTag> {
108 let tag = AnnotatedTag::decode_current_msgpack(data)
109 .map_err(|error| HeddleError::InvalidObject(error.to_string()))?;
110 if tag.hash() != hash {
111 return Err(HeddleError::Corruption {
112 expected: hash,
113 found: tag.hash(),
114 });
115 }
116 Ok(tag)
117}
118
119fn validate_loaded_state(requested_id: &StateId, mut state: State) -> Result<State> {
120 if !state.accepts_stored_id(requested_id) {
121 return Err(HeddleError::InvalidObject(format!(
122 "state id mismatch: requested {requested_id}, computed {}",
123 state.id()
124 )));
125 }
126 state.state_id = *requested_id;
127 Ok(state)
128}
129
130pub(super) fn validate_state_serialized(data: &[u8], id: StateId) -> Result<State> {
131 let state = State::decode_current_msgpack(data)?;
132 validate_loaded_state(&id, state)
133}
134
135fn validate_loaded_action(requested_id: &ActionId, action: Action) -> Result<Action> {
136 let found_id = action.compute_id();
137 if found_id != *requested_id {
138 return Err(HeddleError::InvalidObject(format!(
139 "action id mismatch: requested {}, found {}",
140 requested_id, found_id
141 )));
142 }
143
144 Ok(action)
145}
146
147fn validate_action_serialized(data: &[u8], id: ActionId) -> Result<Action> {
148 let action: Action = rmp_serde::from_slice(data)?;
149 validate_loaded_action(&id, action)
150}
151
152trait EnumerationCounter {
153 fn membership_check(&mut self);
154 fn header_read(&mut self);
155}
156
157struct NoopEnumerationCounter;
158
159impl EnumerationCounter for NoopEnumerationCounter {
160 fn membership_check(&mut self) {}
161 fn header_read(&mut self) {}
162}
163
164fn append_packed_hashes_with_counter(
165 hashes: &mut Vec<ContentHash>,
166 manager: &PackManager,
167 expected_type: ObjectType,
168 counter: &mut impl EnumerationCounter,
169) -> Result<()> {
170 let mut known: HashSet<_> = hashes.iter().copied().collect();
171 for id in manager.list_all_ids()? {
172 let hash = match id {
173 PackObjectId::Hash(hash) if expected_type != ObjectType::AnnotatedTag => hash,
174 PackObjectId::AnnotatedTag(hash) if expected_type == ObjectType::AnnotatedTag => hash,
175 PackObjectId::Hash(_) | PackObjectId::StateId(_) | PackObjectId::AnnotatedTag(_) => {
176 continue;
177 }
178 };
179 counter.membership_check();
180 if known.contains(&hash) {
181 continue;
182 }
183 counter.header_read();
184 let found_type = if expected_type == ObjectType::AnnotatedTag {
185 manager
186 .get_object(&PackObjectId::AnnotatedTag(hash))?
187 .map(|(object_type, _)| object_type)
188 } else {
189 manager.get_hashed_object_type(&hash)?
190 };
191 if found_type == Some(expected_type) {
192 known.insert(hash);
193 hashes.push(hash);
194 }
195 }
196 Ok(())
197}
198
199fn append_packed_hashes(
200 hashes: &mut Vec<ContentHash>,
201 manager: &PackManager,
202 expected_type: ObjectType,
203) -> Result<()> {
204 append_packed_hashes_with_counter(hashes, manager, expected_type, &mut NoopEnumerationCounter)
205}
206
207fn append_unique_states(
208 states: &mut Vec<StateId>,
209 known: &mut HashSet<StateId>,
210 incoming: impl IntoIterator<Item = StateId>,
211) {
212 for id in incoming {
213 if known.insert(id) {
214 states.push(id);
215 }
216 }
217}
218
219impl FsStore {
220 fn holds_state(&self, id: &StateId) -> Result<bool> {
224 if self.try_has_state_once(id)? {
225 return Ok(true);
226 }
227 Ok(self.reload_packs_if_stale()? && self.try_has_state_once(id)?)
228 }
229
230 fn states_needing_loose_copies(
242 &self,
243 reader: &crate::store::pack::PackReader,
244 ids: &[PackObjectId],
245 ) -> Result<Vec<(StateId, Vec<u8>)>> {
246 let mut held = HashSet::new();
247 for id in ids {
248 if let PackObjectId::StateId(state) = id
249 && self.holds_state(state)?
250 {
251 held.insert(*state);
252 }
253 }
254 if held.is_empty() {
255 return Ok(Vec::new());
256 }
257 let mut states = state_entries_from_pack(reader, ids)?;
258 states.retain(|(id, _)| held.contains(id));
259 Ok(states)
260 }
261
262 fn write_packed_state_mirrors_batch(&self, states: Vec<(StateId, Vec<u8>)>) -> Result<()> {
266 if states.is_empty() {
267 return Ok(());
268 }
269
270 self.begin_snapshot_write_batch_impl()?;
271 for (id, data) in states {
272 if let Err(error) = ObjectStore::put_state_serialized(self, &data, id) {
273 self.abort_snapshot_write_batch_impl();
274 return Err(error);
275 }
276 }
277 if let Err(error) = self.flush_snapshot_write_batch_impl() {
278 self.abort_snapshot_write_batch_impl();
279 return Err(error);
280 }
281 Ok(())
282 }
283
284 fn with_state_attachment_index_lock<T>(
285 &self,
286 state: &StateId,
287 operation: impl FnOnce() -> Result<T>,
288 ) -> Result<T> {
289 let path = state_attachment_index_lock_path(&self.root, state);
290 if let Some(parent) = path.parent() {
291 fs::create_dir_all(parent)?;
292 }
293 let file = OpenOptions::new()
294 .create(true)
295 .truncate(false)
296 .read(true)
297 .write(true)
298 .open(path)?;
299 file.lock_exclusive()?;
300 let result = operation();
301 file.unlock()?;
302 result
303 }
304
305 fn collect_state_attachment_ids(&self, state: &StateId) -> Result<Vec<StateAttachmentId>> {
306 let mut ids = Vec::new();
307 let dir = state_attachments_dir(&self.root, state);
308 if let Ok(entries) = fs::read_dir(dir) {
309 for entry in entries {
310 let attachment =
311 StateAttachment::decode_current_msgpack(&fs::read(entry?.path())?)?;
312 if attachment.state_id != *state {
313 return Err(HeddleError::InvalidObject(
314 "state attachment stored under wrong state".to_string(),
315 ));
316 }
317 ids.push(attachment.id());
318 }
319 }
320 if let Ok(manager) = self.pack_manager().read() {
321 for pack_id in manager.list_all_ids()? {
322 let PackObjectId::Hash(hash) = pack_id else {
323 continue;
324 };
325 let Some((ObjectType::StateAttachment, bytes)) =
326 manager.get_hashed_object(&hash)?
327 else {
328 continue;
329 };
330 let attachment = StateAttachment::decode_current_msgpack(&bytes)?;
331 if attachment.state_id == *state {
332 ids.push(attachment.id());
333 }
334 }
335 }
336 ids.sort();
337 ids.dedup();
338 Ok(ids)
339 }
340
341 fn rebuild_state_attachment_index(&self, state: &StateId) -> Result<Vec<StateAttachmentId>> {
342 #[cfg(test)]
343 fs::write(
344 state_attachment_index_path(&self.root, state).with_extension("rebuild-marker"),
345 b"rebuilt",
346 )?;
347 let ids = self.collect_state_attachment_ids(state)?;
348 let path = state_attachment_index_path(&self.root, state);
349 self.write_loose_object_atomic(&path, &rmp_serde::to_vec_named(&ids)?)?;
350 Ok(ids)
351 }
352
353 pub(super) fn materialize_packed_attachment_index(
358 &self,
359 state: &StateId,
360 packed_ids: &[StateAttachmentId],
361 state_was_present: bool,
362 ) -> Result<()> {
363 if packed_ids.is_empty() {
364 return Ok(());
365 }
366 self.with_state_attachment_index_lock(state, || {
367 let path = state_attachment_index_path(&self.root, state);
368 let mut ids = if state_was_present {
369 match read_file_bytes(&path)? {
370 Some(bytes) => rmp_serde::from_slice(bytes.as_slice())?,
371 None => self.collect_state_attachment_ids(state)?,
372 }
373 } else {
374 Vec::new()
375 };
376 ids.extend_from_slice(packed_ids);
377 ids.sort();
378 ids.dedup();
379 self.write_reconstructible_cache(&path, &rmp_serde::to_vec_named(&ids)?)?;
380 Ok(())
381 })
382 }
383}
384
385fn validate_and_list_pack(
393 store: &FsStore,
394 reader: &crate::store::pack::PackReader,
395) -> Result<Vec<PackObjectId>> {
396 let ids = reader.list_ids()?;
397 reader.visit_objects(|id, object_type, data| {
398 if let (PackObjectId::Hash(hash), ObjectType::Tree) = (id, object_type)
399 && is_delta_tree(data)
400 {
401 let header = decode_tree_delta_header(data)?;
402 if header.anchor == hash {
403 return Err(HeddleError::InvalidObject(
404 "HDC1 result id must differ from its anchor id".to_string(),
405 ));
406 }
407 let anchor_body = match reader.get_object(&PackObjectId::Hash(header.anchor))? {
408 Some((ObjectType::Tree, body)) => Some(body),
409 Some((kind, _)) => {
410 return Err(HeddleError::InvalidObject(format!(
411 "HDC1 anchor {} is indexed as {kind:?}, expected Tree",
412 header.anchor
413 )));
414 }
415 None => store.try_get_tree_serialized_once(&header.anchor)?,
416 }
417 .ok_or_else(|| HeddleError::NotFound(format!("tree delta anchor {}", header.anchor)))?;
418 if is_delta_tree(&anchor_body) {
419 return Err(HeddleError::InvalidObject(
420 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
421 ));
422 }
423 let anchor = codec::decode_tree_serialized_with_key(&anchor_body, header.anchor, None)?;
424 codec::decode_tree_serialized_with_key(data, hash, Some(&anchor))?;
425 return Ok(());
426 }
427 validate_pack_entry(&id, object_type, data)
428 })?;
429 Ok(ids)
430}
431
432fn state_entries_from_pack(
433 reader: &crate::store::pack::PackReader,
434 ids: &[PackObjectId],
435) -> Result<Vec<(StateId, Vec<u8>)>> {
436 let mut states = Vec::new();
437 let expected = ids.iter().copied().collect::<HashSet<_>>();
438 reader.visit_objects(|id, object_type, data| {
439 if !expected.contains(&id) {
440 return Err(HeddleError::InvalidObject(
441 "pack visitor yielded an unindexed object".into(),
442 ));
443 }
444 if let PackObjectId::StateId(state_id) = id {
445 if object_type != ObjectType::State {
446 return Err(HeddleError::InvalidObject(format!(
447 "pack id {} is indexed as {object_type:?}, expected State",
448 state_id.to_string_full()
449 )));
450 }
451 validate_state_serialized(data, state_id)?;
452 states.push((state_id, data.to_vec()));
453 }
454 Ok(())
455 })?;
456 Ok(states)
457}
458
459fn attachment_entries_from_pack(
460 reader: &crate::store::pack::PackReader,
461 ids: &[PackObjectId],
462) -> Result<Vec<StateAttachment>> {
463 let mut attachments = Vec::new();
464 let expected = ids.iter().copied().collect::<HashSet<_>>();
465 reader.visit_objects(|id, object_type, data| {
466 if expected.contains(&id) && object_type == ObjectType::StateAttachment {
467 attachments.push(StateAttachment::decode_current_msgpack(data)?);
468 }
469 Ok(())
470 })?;
471 Ok(attachments)
472}
473
474pub(super) fn validate_pack_entry(
475 id: &PackObjectId,
476 obj_type: ObjectType,
477 data: &[u8],
478) -> Result<()> {
479 match (id, obj_type) {
480 (PackObjectId::Hash(hash), ObjectType::Blob) => validate_blob_bytes(data, *hash),
481 (PackObjectId::AnnotatedTag(hash), ObjectType::AnnotatedTag) => {
482 validate_annotated_tag(data, *hash).map(|_| ())
483 }
484 (PackObjectId::Hash(hash), ObjectType::Tree) => {
485 validate_tree_serialized(data, *hash).map(|_| ())
486 }
487 (PackObjectId::Hash(hash), ObjectType::Action) => {
488 validate_action_serialized(data, ActionId::from_hash(*hash)).map(|_| ())
489 }
490 (PackObjectId::StateId(change_id), ObjectType::State) => {
491 validate_state_serialized(data, *change_id).map(|_| ())
492 }
493 (PackObjectId::Hash(hash), ObjectType::StateAttachment) => {
494 let attachment = StateAttachment::decode_current_msgpack(data)?;
495 if attachment.id().as_hash() != hash {
496 return Err(HeddleError::InvalidObject(
497 "state attachment pack id mismatch".to_string(),
498 ));
499 }
500 Ok(())
501 }
502 (PackObjectId::Hash(hash), ObjectType::SnapshotCommit) => {
503 let artifact: crate::store::SnapshotCommitArtifact = rmp_serde::from_slice(data)?;
504 artifact.validate()?;
505 if artifact.id() != *hash {
506 return Err(HeddleError::InvalidObject(
507 "snapshot commit artifact pack id mismatch".to_string(),
508 ));
509 }
510 Ok(())
511 }
512 (_, ObjectType::TimelineOperation) => Err(HeddleError::InvalidObject(
513 "timeline operations belong in the timeline pack store".to_string(),
514 )),
515 _ => Err(HeddleError::InvalidObject(format!(
516 "unsupported native pack object: {:?} {:?}",
517 id, obj_type
518 ))),
519 }
520}
521
522impl FsStore {
523 fn cache_recent_blob(&self, hash: ContentHash, blob: &Blob) {
525 if blob.content().len() > super::fs_store::RECENT_BLOB_CACHE_MAX_BYTES {
526 return;
527 }
528 if let Ok(mut cache) = self.recent_blobs.write() {
529 cache.insert(hash, blob.clone());
530 }
531 }
532
533 fn cache_recent_tree(&self, hash: ContentHash, tree: &Tree) {
534 if let Ok(mut cache) = self.recent_trees.write() {
535 cache.insert(hash, tree.clone());
536 }
537 }
538
539 fn cache_recent_state(&self, id: StateId, state: &State) {
540 if let Ok(mut cache) = self.recent_states.write() {
541 cache.insert(id, state.clone());
542 }
543 }
544
545 fn recent_blob(&self, hash: &ContentHash) -> Option<Blob> {
546 self.recent_blobs
547 .read()
548 .ok()
549 .and_then(|cache| cache.get(hash).cloned())
550 }
551
552 fn recent_tree(&self, hash: &ContentHash) -> Option<Tree> {
553 self.recent_trees
554 .read()
555 .ok()
556 .and_then(|cache| cache.get(hash).cloned())
557 }
558
559 fn recent_state(&self, id: &StateId) -> Option<State> {
560 self.recent_states
561 .read()
562 .ok()
563 .and_then(|cache| cache.get(id).cloned())
564 }
565
566 fn try_get_blob_once(&self, hash: &ContentHash) -> Result<Option<Blob>> {
569 if let Ok(cache) = self.recent_blobs.read()
572 && let Some(blob) = cache.get(hash)
573 {
574 trace!("Found blob in recent object cache");
575 return Ok(Some(blob.clone()));
576 }
577
578 if let Ok(manager) = self.pack_manager().read()
579 && let Some((obj_type, data)) = manager.get_hashed_object(hash)?
580 && obj_type == ObjectType::Blob
581 {
582 trace!("Found blob in packfile");
583 validate_blob_bytes(&data, *hash)?;
584 let blob = Blob::new(data);
585 heddle_perf_contract::record_object_decode();
586 self.cache_recent_blob(*hash, &blob);
587 return Ok(Some(blob));
588 }
589
590 let path = hash_path(&blobs_dir(&self.root), hash);
591 match read_file_bytes(&path)? {
592 Some(data) => {
593 trace!(size = data.as_slice().len(), "Blob data read");
594 let content = codec::decode_blob_content(data.as_slice())?;
595 let blob = Blob::new(content);
596 heddle_perf_contract::record_object_decode();
597 if blob.hash() != *hash {
604 return Err(HeddleError::Corruption {
605 expected: *hash,
606 found: blob.hash(),
607 });
608 }
609 self.cache_recent_blob(*hash, &blob);
610 Ok(Some(blob))
611 }
612 None => Ok(None),
613 }
614 }
615
616 fn loose_or_packed(
621 &self,
622 loose_path: &Path,
623 in_pack: impl FnOnce(&PackManager) -> bool,
624 ) -> Result<bool> {
625 if loose_path.exists() {
626 return Ok(true);
627 }
628 if let Ok(manager) = self.pack_manager().read() {
629 return Ok(in_pack(&manager));
630 }
631 Ok(false)
632 }
633
634 fn try_has_blob_once(&self, hash: &ContentHash) -> Result<bool> {
635 let path = hash_path(&blobs_dir(&self.root), hash);
639 self.loose_or_packed(&path, |m| m.has_object(hash))
640 }
641
642 fn try_get_blob_size_once(&self, hash: &ContentHash) -> Result<Option<u64>> {
657 if let Ok(cache) = self.recent_blobs.read()
658 && let Some(blob) = cache.get(hash)
659 {
660 return Ok(Some(blob.content().len() as u64));
661 }
662
663 let path = hash_path(&blobs_dir(&self.root), hash);
664 if let Some((header, file_len)) = read_file_header(&path, BLOB_HEADER_PEEK)? {
665 if let Some(size) = header_uncompressed_size(&header) {
666 return Ok(Some(size));
667 }
668 return Ok(Some(file_len));
671 }
672
673 if let Ok(manager) = self.pack_manager().read()
674 && let Some(size) = manager.get_hashed_object_size(hash)?
675 {
676 return Ok(Some(size));
677 }
678 Ok(None)
679 }
680
681 fn try_open_tree_once(
682 &self,
683 tree_id: &ContentHash,
684 cursor: Option<&TreeResumeCursor>,
685 ) -> Result<Option<TreeEntryReader<OpenedTreeBody>>> {
686 let path = hash_path(&trees_dir(&self.root), tree_id);
687 if path.exists()
688 && let Some((header, len)) = read_file_header(&path, TREE_CANONICAL_MAGIC.len())?
689 {
690 if header.starts_with(TREE_CANONICAL_MAGIC) || header.starts_with(TREE_LEAN_MAGIC) {
691 let file = File::open(&path)?;
692 return Ok(Some(TreeEntryReader::open(
693 OpenedTreeBody::File(FileTreeSource::sequential_verify(file, len)),
694 *tree_id,
695 cursor,
696 )?));
697 }
698 if header.starts_with(TREE_DELTA_MAGIC) {
699 let file = File::open(&path)?;
700 return self.open_delta_tree_source(
701 *tree_id,
702 cursor,
703 OpenedTreeBody::File(FileTreeSource::sequential_verify(file, len)),
704 );
705 }
706 if header.starts_with(TREE_SALTED_MAGIC) {
707 return Err(salted_tree_not_streamable(tree_id));
708 }
709 }
710 if path.exists()
711 && let Some(data) = read_file_bytes(&path)?
712 {
713 let body = codec::decode_tree_body(data.as_slice())?;
714 if is_streamable_tree(&body) {
715 return Ok(Some(TreeEntryReader::open(
716 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(body)),
717 *tree_id,
718 cursor,
719 )?));
720 }
721 if is_delta_tree(&body) {
722 return self.open_delta_tree_source(
723 *tree_id,
724 cursor,
725 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(body)),
726 );
727 }
728 if is_salted_tree(&body) {
729 return Err(salted_tree_not_streamable(tree_id));
730 }
731 }
732 let packed = if let Ok(manager) = self.pack_manager().read() {
733 manager.get_hashed_object(tree_id)?
734 } else {
735 None
736 };
737 if let Some((ObjectType::Tree, data)) = packed {
738 if is_streamable_tree(&data) {
739 return Ok(Some(TreeEntryReader::open(
740 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(data)),
741 *tree_id,
742 cursor,
743 )?));
744 }
745 if is_delta_tree(&data) {
746 return self.open_delta_tree_source(
747 *tree_id,
748 cursor,
749 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(data)),
750 );
751 }
752 if is_salted_tree(&data) {
753 return Err(salted_tree_not_streamable(tree_id));
754 }
755 }
756 let npk_tree = if let Ok(manager) = self.npk1_manager().read() {
757 manager.get_tree(tree_id)?
758 } else {
759 None
760 };
761 if let Some(tree) = npk_tree {
762 return Ok(Some(TreeEntryReader::open(
763 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(streamable_body(&tree)?)),
764 *tree_id,
765 cursor,
766 )?));
767 }
768 Ok(None)
769 }
770
771 fn open_delta_tree_source(
772 &self,
773 tree_id: ContentHash,
774 cursor: Option<&TreeResumeCursor>,
775 mut delta: OpenedTreeBody,
776 ) -> Result<Option<TreeEntryReader<OpenedTreeBody>>> {
777 let object_len = usize::try_from(delta.len())
778 .map_err(|_| HeddleError::InvalidObject("HDC1 body exceeds usize".to_string()))?;
779 let mut header_bytes = [0u8; TREE_DELTA_HEADER_LEN];
780 delta.read_exact_at(0, &mut header_bytes)?;
781 let header = decode_tree_delta_header_prefix(&header_bytes, object_len)?;
782 let anchor = self
783 .try_open_materialized_tree_once(&header.anchor)?
784 .ok_or_else(|| HeddleError::NotFound(format!("tree delta anchor {}", header.anchor)))?;
785 let source = DeltaTreeSource::open(delta, anchor)?;
786 Ok(Some(TreeEntryReader::open(
787 OpenedTreeBody::Dynamic(Box::new(source)),
788 tree_id,
789 cursor,
790 )?))
791 }
792
793 fn try_open_materialized_tree_once(
794 &self,
795 tree_id: &ContentHash,
796 ) -> Result<Option<TreeEntryReader<OpenedTreeBody>>> {
797 let path = hash_path(&trees_dir(&self.root), tree_id);
798 if path.exists()
799 && let Some((header, len)) = read_file_header(&path, TREE_CANONICAL_MAGIC.len())?
800 {
801 if header.starts_with(TREE_DELTA_MAGIC) {
802 return Err(HeddleError::InvalidObject(
803 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
804 ));
805 }
806 if header.starts_with(TREE_CANONICAL_MAGIC) || header.starts_with(TREE_LEAN_MAGIC) {
807 let file = File::open(&path)?;
808 return Ok(Some(TreeEntryReader::open(
809 OpenedTreeBody::File(FileTreeSource::sequential_verify(file, len)),
810 *tree_id,
811 None,
812 )?));
813 }
814 }
815 if path.exists()
816 && let Some(data) = read_file_bytes(&path)?
817 {
818 let body = codec::decode_tree_body(data.as_slice())?;
819 if is_delta_tree(&body) {
820 return Err(HeddleError::InvalidObject(
821 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
822 ));
823 }
824 if is_streamable_tree(&body) {
825 return Ok(Some(TreeEntryReader::open(
826 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(body)),
827 *tree_id,
828 None,
829 )?));
830 }
831 }
832 let packed = if let Ok(manager) = self.pack_manager().read() {
833 manager.get_hashed_object(tree_id)?
834 } else {
835 None
836 };
837 if let Some((ObjectType::Tree, data)) = packed {
838 if is_delta_tree(&data) {
839 return Err(HeddleError::InvalidObject(
840 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
841 ));
842 }
843 if is_streamable_tree(&data) {
844 return Ok(Some(TreeEntryReader::open(
845 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(data)),
846 *tree_id,
847 None,
848 )?));
849 }
850 }
851 let npk_tree = if let Ok(manager) = self.npk1_manager().read() {
852 manager.get_tree(tree_id)?
853 } else {
854 None
855 };
856 if let Some(tree) = npk_tree {
857 return Ok(Some(TreeEntryReader::open(
858 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(streamable_body(&tree)?)),
859 *tree_id,
860 None,
861 )?));
862 }
863 if let Some(source) = &self.external_source
864 && let Some(tree) = source.get_tree(tree_id)?
865 {
866 return Ok(Some(TreeEntryReader::open(
867 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(streamable_body(&tree)?)),
868 *tree_id,
869 None,
870 )?));
871 }
872 Ok(None)
873 }
874
875 fn try_get_tree_once(&self, hash: &ContentHash) -> Result<Option<Tree>> {
876 if let Ok(cache) = self.recent_trees.read()
880 && let Some(tree) = cache.get(hash)
881 {
882 trace!("Found tree in recent object cache");
883 return Ok(Some(tree.clone()));
884 }
885
886 let path = hash_path(&trees_dir(&self.root), hash);
890 if path.exists()
891 && let Some(data) = read_file_bytes(&path)?
892 {
893 trace!(size = data.as_slice().len(), "Tree data read");
894 let body = codec::decode_tree_body(data.as_slice())?;
895 let tree = validate_loaded_tree(self.decode_tree_storage_body(*hash, &body)?)?;
896 heddle_perf_contract::record_object_decode();
897 if tree.hash() != *hash {
898 return Err(HeddleError::Corruption {
899 expected: *hash,
900 found: tree.hash(),
901 });
902 }
903 if let Ok(mut cache) = self.recent_trees.write() {
904 cache.insert(*hash, tree.clone());
905 }
906 return Ok(Some(tree));
907 }
908
909 if let Ok(manager) = self.npk1_manager().read()
910 && let Some(tree) = manager.get_tree(hash)?
911 {
912 trace!("Found tree in NPK1 pack");
913 heddle_perf_contract::record_object_decode();
914 self.cache_recent_tree(*hash, &tree);
915 return Ok(Some(tree));
916 }
917 if let Ok(manager) = self.pack_manager().read()
918 && let Some((obj_type, data)) = manager.get_hashed_object(hash)?
919 && obj_type == ObjectType::Tree
920 {
921 trace!("Found tree in packfile");
922 let tree = validate_loaded_tree(self.decode_tree_storage_body(*hash, &data)?)?;
923 heddle_perf_contract::record_object_decode();
924 if tree.hash() != *hash {
925 return Err(HeddleError::Corruption {
926 expected: *hash,
927 found: tree.hash(),
928 });
929 }
930 if let Ok(mut cache) = self.recent_trees.write() {
931 cache.insert(*hash, tree.clone());
932 }
933 return Ok(Some(tree));
934 }
935 Ok(None)
936 }
937
938 fn try_get_tree_entry_once(&self, hash: &ContentHash, name: &str) -> Result<Option<TreeEntry>> {
939 if let Some(tree) = self.recent_tree(hash) {
940 return Ok(tree.get(name).cloned());
941 }
942 let path = hash_path(&trees_dir(&self.root), hash);
943 if path.exists() {
944 return Ok(self
945 .try_get_tree_once(hash)?
946 .and_then(|tree| tree.get(name).cloned()));
947 }
948 if let Ok(manager) = self.npk1_manager().read()
949 && manager.has_tree(hash)?
950 {
951 return manager.get_entry(hash, name);
952 }
953 if let Ok(manager) = self.pack_manager().read()
954 && manager.has_object(hash)
955 {
956 return Ok(self
957 .try_get_tree_once(hash)?
958 .and_then(|tree| tree.get(name).cloned()));
959 }
960 Ok(None)
961 }
962
963 pub(super) fn try_get_tree_serialized_once(
964 &self,
965 hash: &ContentHash,
966 ) -> Result<Option<Vec<u8>>> {
967 let path = hash_path(&trees_dir(&self.root), hash);
968 if path.exists()
969 && let Some(data) = read_file_bytes(&path)?
970 {
971 return Ok(Some(codec::decode_tree_body(data.as_slice())?));
972 }
973
974 if let Ok(manager) = self.npk1_manager().read()
975 && let Some(tree) = manager.get_tree(hash)?
976 {
977 return tree.encode_lean().map(Some).map_err(HeddleError::from);
978 }
979
980 if let Ok(manager) = self.pack_manager().read()
981 && let Some((obj_type, data)) = manager.get_hashed_object(hash)?
982 && obj_type == ObjectType::Tree
983 {
984 return Ok(Some(data));
985 }
986
987 Ok(None)
988 }
989
990 pub(super) fn decode_tree_storage_body(&self, hash: ContentHash, data: &[u8]) -> Result<Tree> {
991 let anchor = if is_delta_tree(data) {
992 let header = decode_tree_delta_header(data)?;
993 let anchor = if let Some(anchor_body) =
994 self.try_get_tree_serialized_once(&header.anchor)?
995 {
996 if is_delta_tree(&anchor_body) {
997 return Err(HeddleError::InvalidObject(
998 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
999 ));
1000 }
1001 let tree =
1002 codec::decode_tree_serialized_with_key(&anchor_body, header.anchor, None)?;
1003 self.cache_recent_tree(header.anchor, &tree);
1004 Some(tree)
1005 } else if let Some(tree) = self.recent_tree(&header.anchor) {
1006 Some(tree)
1009 } else if let Some(source) = &self.external_source {
1010 source.get_tree(&header.anchor)?
1011 } else {
1012 None
1013 };
1014 Some(anchor.ok_or_else(|| {
1015 HeddleError::NotFound(format!("tree delta anchor {}", header.anchor))
1016 })?)
1017 } else {
1018 None
1019 };
1020 codec::decode_tree_serialized_with_key(data, hash, anchor.as_ref())
1021 }
1022
1023 pub(super) fn encode_tree_write(&self, write: &TreeWrite) -> Result<EncodedTree> {
1024 let Some(parent) = write.parent else {
1025 return codec::encode_tree_hot(&write.tree, None);
1026 };
1027 let Some(parent_body) = self.try_get_tree_serialized_once(&parent)? else {
1028 return codec::encode_tree_hot(&write.tree, None);
1029 };
1030 let base = if is_delta_tree(&parent_body) {
1031 let header = decode_tree_delta_header(&parent_body)?;
1032 let Some(lineage) = self.read_tree_lineage(&parent)? else {
1033 return codec::encode_tree_hot(&write.tree, None);
1034 };
1035 if lineage.anchor != header.anchor || lineage.depth == 0 {
1036 return codec::encode_tree_hot(&write.tree, None);
1037 }
1038 let Some(anchor_body) = self.try_get_tree_serialized_once(&lineage.anchor)? else {
1039 return codec::encode_tree_hot(&write.tree, None);
1040 };
1041 if is_delta_tree(&anchor_body) {
1042 return Err(HeddleError::InvalidObject(
1043 "HDC1 lineage points to another delta".to_string(),
1044 ));
1045 }
1046 let anchor =
1047 codec::decode_tree_serialized_with_key(&anchor_body, lineage.anchor, None)?;
1048 Some((lineage.anchor, anchor, lineage.depth))
1049 } else {
1050 let anchor = codec::decode_tree_serialized_with_key(&parent_body, parent, None)?;
1051 Some((parent, anchor, 0))
1052 };
1053 let Some((anchor_id, anchor, parent_depth)) = base else {
1054 return codec::encode_tree_hot(&write.tree, None);
1055 };
1056 codec::encode_tree_hot(
1057 &write.tree,
1058 Some(TreeDeltaBase {
1059 anchor_id,
1060 anchor: &anchor,
1061 parent_depth,
1062 }),
1063 )
1064 }
1065
1066 fn read_tree_lineage(&self, hash: &ContentHash) -> Result<Option<TreeLineage>> {
1067 let Some(bytes) = read_file_bytes(&tree_lineage_path(&self.root, hash))? else {
1068 return Ok(None);
1069 };
1070 let data = bytes.as_slice();
1071 if data.len() != 33 {
1072 return Ok(None);
1073 }
1074 let anchor = match data[..32].try_into() {
1075 Ok(bytes) => ContentHash::from_bytes(bytes),
1076 Err(_) => return Ok(None),
1077 };
1078 let depth = data[32];
1079 if depth == 0 || depth >= crate::object::TREE_DELTA_ANCHOR_INTERVAL {
1080 return Ok(None);
1081 }
1082 Ok(Some(TreeLineage { anchor, depth }))
1083 }
1084
1085 pub(super) fn remember_tree_encoding(
1086 &self,
1087 hash: ContentHash,
1088 kind: TreeEncodingKind,
1089 ) -> Result<()> {
1090 let TreeEncodingKind::Delta { anchor, depth, .. } = kind else {
1091 return Ok(());
1092 };
1093 let mut bytes = Vec::with_capacity(33);
1094 bytes.extend_from_slice(anchor.as_bytes());
1095 bytes.push(depth);
1096 self.write_reconstructible_cache(&tree_lineage_path(&self.root, &hash), &bytes)
1097 }
1098
1099 fn try_has_tree_once(&self, hash: &ContentHash) -> Result<bool> {
1100 let path = hash_path(&trees_dir(&self.root), hash);
1104 if self.loose_or_packed(&path, |m| m.has_object(hash))? {
1105 return Ok(true);
1106 }
1107 if let Ok(manager) = self.npk1_manager().read() {
1108 return manager.has_tree(hash);
1109 }
1110 Ok(false)
1111 }
1112
1113 fn try_get_state_once(&self, id: &StateId) -> Result<Option<State>> {
1114 if let Ok(cache) = self.recent_states.read()
1119 && let Some(state) = cache.get(id)
1120 {
1121 trace!("Found state in recent object cache");
1122 return Ok(Some(state.clone()));
1123 }
1124
1125 let path = state_path(&self.root, id);
1126 if let Some(data) = read_file_bytes(&path)? {
1127 trace!(size = data.as_slice().len(), "State read from loose object");
1128 let state = validate_loaded_state(id, codec::decode_state(data.as_slice())?)?;
1129 heddle_perf_contract::record_object_decode();
1130 if let Ok(mut cache) = self.recent_states.write() {
1131 cache.insert(*id, state.clone());
1132 }
1133 return Ok(Some(state));
1134 }
1135
1136 if let Ok(manager) = self.pack_manager().read()
1137 && let Some((obj_type, data)) = manager.get_object(&PackObjectId::StateId(*id))?
1138 && obj_type == ObjectType::State
1139 {
1140 trace!("Found state in packfile");
1141 let state = validate_loaded_state(id, State::decode_current_msgpack(&data)?)?;
1142 heddle_perf_contract::record_object_decode();
1143 if let Ok(mut cache) = self.recent_states.write() {
1144 cache.insert(*id, state.clone());
1145 }
1146 return Ok(Some(state));
1147 }
1148
1149 Ok(None)
1150 }
1151
1152 fn try_has_state_once(&self, id: &StateId) -> Result<bool> {
1153 if let Ok(cache) = self.recent_states.read()
1156 && cache.contains(id)
1157 {
1158 return Ok(true);
1159 }
1160 let path = state_path(&self.root, id);
1161 self.loose_or_packed(&path, |m| m.has_object_id(&PackObjectId::StateId(*id)))
1162 }
1163
1164 fn try_get_action_once(&self, id: &ActionId) -> Result<Option<Action>> {
1165 let path = action_path(&self.root, id);
1166 if let Some(data) = read_file_bytes(&path)? {
1167 trace!(size = data.as_slice().len(), "Action data read");
1168 return Ok(Some(validate_loaded_action(
1169 id,
1170 codec::decode_action(data.as_slice())?,
1171 )?));
1172 }
1173 if let Ok(manager) = self.pack_manager().read()
1174 && let Some((ObjectType::Action, data)) = manager.get_hashed_object(id.as_hash())?
1175 {
1176 trace!("Found action in packfile");
1177 return Ok(Some(validate_loaded_action(
1178 id,
1179 rmp_serde::from_slice(&data)?,
1180 )?));
1181 }
1182 Ok(None)
1183 }
1184
1185 fn try_get_state_attachment_once(
1186 &self,
1187 state: &StateId,
1188 id: &StateAttachmentId,
1189 ) -> Result<Option<StateAttachment>> {
1190 let path = state_attachment_path(&self.root, state, id);
1191 let file_bytes = read_file_bytes(&path)?;
1192 if let Some(bytes) = file_bytes.as_ref() {
1193 let attachment = StateAttachment::decode_current_msgpack(bytes.as_slice())?;
1194 return Self::validate_state_attachment(attachment, state, id).map(Some);
1195 }
1196 if let Ok(manager) = self.pack_manager().read()
1197 && let Some((ObjectType::StateAttachment, pack_bytes)) =
1198 manager.get_hashed_object(id.as_hash())?
1199 {
1200 let attachment = StateAttachment::decode_current_msgpack(&pack_bytes)?;
1201 return Self::validate_state_attachment(attachment, state, id).map(Some);
1202 }
1203 Ok(None)
1204 }
1205
1206 fn validate_state_attachment(
1207 attachment: StateAttachment,
1208 state: &StateId,
1209 id: &StateAttachmentId,
1210 ) -> Result<StateAttachment> {
1211 if attachment.state_id != *state || attachment.id() != *id {
1212 return Err(HeddleError::InvalidObject(
1213 "state attachment address does not match content".to_string(),
1214 ));
1215 }
1216 Ok(attachment)
1217 }
1218}
1219
1220impl FsStore {
1221 #[doc(hidden)]
1223 pub fn snapshot_commit_recovery_descriptors(&self) -> Result<Vec<SnapshotCommitDescriptor>> {
1224 self.reload_packs_if_stale()?;
1225 let manager = self
1226 .pack_manager()
1227 .read()
1228 .map_err(|_| HeddleError::Config("Failed to acquire pack manager lock".to_string()))?;
1229 manager.snapshot_commit_recovery_descriptors()
1230 }
1231
1232 #[doc(hidden)]
1236 pub fn snapshot_commit_descriptors(&self) -> Result<Vec<SnapshotCommitDescriptor>> {
1237 self.reload_packs_if_stale()?;
1238 let manager = self
1239 .pack_manager()
1240 .read()
1241 .map_err(|_| HeddleError::Config("Failed to acquire pack manager lock".to_string()))?;
1242 manager.snapshot_commit_descriptors()
1243 }
1244
1245 #[doc(hidden)]
1248 pub fn snapshot_commit_descriptor_for_state(
1249 &self,
1250 state: &StateId,
1251 ) -> Result<Option<SnapshotCommitDescriptor>> {
1252 self.reload_packs_if_stale()?;
1253 let manager = self
1254 .pack_manager()
1255 .read()
1256 .map_err(|_| HeddleError::Config("Failed to acquire pack manager lock".to_string()))?;
1257 manager.snapshot_commit_descriptor_for_state(state)
1258 }
1259
1260 pub fn clear_recent_caches(&self) {
1264 self.clear_recent_object_caches();
1265 }
1266
1267 #[instrument(skip(self))]
1269 pub fn pack_objects(&self, delta_search: bool) -> Result<(u64, u64)> {
1270 self.pack_objects_impl(delta_search)
1271 }
1272
1273 #[instrument(skip(self))]
1275 pub fn prune_loose_objects(&self) -> Result<(u64, u64)> {
1276 self.prune_loose_objects_impl()
1277 }
1278
1279 pub fn discard_corrupt_clone_packs(&self) -> Result<usize> {
1282 let packs = super::fs_paths::packs_dir(&self.root);
1283 let mut removed = 0;
1284 for entry in match fs::read_dir(&packs) {
1285 Ok(entries) => entries,
1286 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(0),
1287 Err(error) => return Err(error.into()),
1288 } {
1289 let path = entry?.path();
1290 match path.extension().and_then(|value| value.to_str()) {
1291 Some("pack") => {
1292 let index = path.with_extension("idx");
1293 let valid = crate::store::pack::PackReader::open(&path, &index)
1294 .and_then(|reader| validate_and_list_pack(self, &reader).map(|_| ()))
1295 .is_ok();
1296 if !valid {
1297 let _ = fs::remove_file(&path);
1298 let _ = fs::remove_file(&index);
1299 removed += 1;
1300 }
1301 }
1302 Some("npk") if super::npk1::Npk1Pack::open(&path).is_err() => {
1303 let _ = fs::remove_file(&path);
1304 removed += 1;
1305 }
1306 _ => {}
1307 }
1308 }
1309 if removed > 0 {
1310 self.reload_packs()?;
1311 self.clear_recent_object_caches();
1312 }
1313 Ok(removed)
1314 }
1315}
1316
1317impl ObjectCacheControl for FsStore {
1318 fn clear_recent_caches(&self) {
1319 FsStore::clear_recent_caches(self);
1320 }
1321}
1322
1323impl ObjectStore for FsStore {
1324 fn get_annotated_tag(&self, hash: &ContentHash) -> Result<Option<AnnotatedTag>> {
1325 let path = hash_path(&annotated_tags_dir(&self.root), hash);
1326 if let Some(data) = read_file_bytes(&path)? {
1327 return validate_annotated_tag(data.as_slice(), *hash).map(Some);
1328 }
1329 self.reload_packs_if_stale()?;
1330 if let Ok(manager) = self.pack_manager().read()
1331 && let Some((ObjectType::AnnotatedTag, data)) =
1332 manager.get_object(&PackObjectId::AnnotatedTag(*hash))?
1333 {
1334 return validate_annotated_tag(&data, *hash).map(Some);
1335 }
1336 Ok(None)
1337 }
1338
1339 fn put_annotated_tag(&self, tag: &AnnotatedTag) -> Result<ContentHash> {
1340 let hash = tag.hash();
1341 let path = hash_path(&annotated_tags_dir(&self.root), &hash);
1342 if !path.exists() {
1343 self.write_loose_object_atomic(&path, &tag.encode_current_msgpack())?;
1344 }
1345 Ok(hash)
1346 }
1347
1348 fn list_annotated_tags(&self) -> Result<Vec<ContentHash>> {
1349 self.reload_packs_if_stale()?;
1350 let mut hashes = list_hashes_from_dir(&annotated_tags_dir(&self.root))?;
1351 if let Ok(manager) = self.pack_manager().read() {
1352 append_packed_hashes(&mut hashes, &manager, ObjectType::AnnotatedTag)?;
1353 }
1354 Ok(hashes)
1355 }
1356
1357 fn get_blob_bytes(&self, hash: &ContentHash) -> Result<Option<bytes::Bytes>> {
1364 if let Ok(manager) = self.pack_manager().read()
1365 && let Some((obj_type, data)) = manager.get_hashed_object_bytes(hash)?
1366 && obj_type == crate::store::pack::ObjectType::Blob
1367 {
1368 validate_blob_bytes(data.as_ref(), *hash)?;
1369 return Ok(Some(data));
1370 }
1371 Ok(self
1372 .get_blob(hash)?
1373 .map(|blob| bytes::Bytes::from(blob.into_content())))
1374 }
1375
1376 #[instrument(skip(self), fields(hash = %hash.short()))]
1377 fn get_blob(&self, hash: &ContentHash) -> Result<Option<Blob>> {
1378 if let Some(blob) = self.recent_blob(hash) {
1379 return Ok(Some(blob));
1380 }
1381 if let Some(blob) = self.try_get_blob_once(hash)? {
1382 return Ok(Some(blob));
1383 }
1384 if self.reload_packs_if_stale()?
1389 && let Some(blob) = self.try_get_blob_once(hash)?
1390 {
1391 return Ok(Some(blob));
1392 }
1393 if let Some(source) = &self.external_source
1394 && let Some(blob) = source.get_blob(hash)?
1395 {
1396 self.cache_recent_blob(*hash, &blob);
1397 return Ok(Some(blob));
1398 }
1399 trace!("Blob not found");
1400 Ok(None)
1401 }
1402
1403 #[instrument(skip(self, blob), fields(size = blob.content().len()))]
1404 fn put_blob(&self, blob: &Blob) -> Result<ContentHash> {
1405 let hash = blob.hash();
1406 let path = hash_path(&blobs_dir(&self.root), &hash);
1407
1408 if !path.exists() {
1409 let data = codec::encode_blob_content(blob.content(), &self.compression)?;
1410 trace!(compressed_size = data.len(), "Writing blob");
1411 self.write_loose_object_atomic(&path, &data)?;
1412 } else {
1413 trace!("Blob already exists, skipping write");
1414 }
1415 self.cache_recent_blob(hash, blob);
1416
1417 Ok(hash)
1418 }
1419
1420 #[instrument(skip(self, blob), fields(hash = %hash.short()))]
1421 fn put_blob_with_hash(&self, blob: &Blob, hash: ContentHash) -> Result<ContentHash> {
1422 if blob.hash() != hash {
1423 return Err(HeddleError::Corruption {
1424 expected: hash,
1425 found: blob.hash(),
1426 });
1427 }
1428
1429 let path = hash_path(&blobs_dir(&self.root), &hash);
1430
1431 if !path.exists() {
1432 let data = codec::encode_blob_content(blob.content(), &self.compression)?;
1433 trace!(
1434 compressed_size = data.len(),
1435 "Writing blob with precomputed hash"
1436 );
1437 self.write_loose_object_atomic(&path, &data)?;
1438 }
1439 self.cache_recent_blob(hash, blob);
1440
1441 Ok(hash)
1442 }
1443
1444 #[instrument(skip(self, data), fields(hash = %hash.short(), size = data.len()))]
1445 fn put_blob_bytes_with_hash(&self, data: &[u8], hash: ContentHash) -> Result<ContentHash> {
1446 validate_blob_bytes(data, hash)?;
1447
1448 let path = hash_path(&blobs_dir(&self.root), &hash);
1449 if !path.exists() {
1450 trace!(
1451 size = data.len(),
1452 "Writing raw blob bytes with precomputed hash"
1453 );
1454 self.write_loose_object_atomic(&path, data)?;
1455 }
1456 self.cache_recent_blob(hash, &Blob::from_slice(data));
1457
1458 Ok(hash)
1459 }
1460
1461 #[instrument(skip(self), fields(hash = %hash.short()))]
1462 fn has_blob(&self, hash: &ContentHash) -> Result<bool> {
1463 if ObjectStore::has_blob_locally(self, hash)? {
1464 return Ok(true);
1465 }
1466 if let Some(source) = &self.external_source {
1467 if self.recent_blob(hash).is_some() {
1468 return Ok(true);
1469 }
1470 if let Some(blob) = source.get_blob(hash)? {
1471 self.cache_recent_blob(*hash, &blob);
1472 return Ok(true);
1473 }
1474 }
1475 Ok(false)
1476 }
1477
1478 fn has_blob_locally(&self, hash: &ContentHash) -> Result<bool> {
1479 if self.try_has_blob_once(hash)? {
1480 return Ok(true);
1481 }
1482 Ok(self.reload_packs_if_stale()? && self.try_has_blob_once(hash)?)
1483 }
1484
1485 fn loose_blob_path(&self, hash: &ContentHash) -> Option<PathBuf> {
1502 let path = hash_path(&blobs_dir(&self.root), hash);
1503 if let Ok(verified) = self.verified_loose_blobs.read()
1508 && verified.contains(hash)
1509 && path.exists()
1510 {
1511 return Some(path);
1512 }
1513
1514 let (header, _) = read_file_header(&path, BLOB_HEADER_PEEK).ok().flatten()?;
1527 if is_compressed(&header) {
1528 return None;
1529 }
1530 let bytes = read_file_bytes(&path).ok().flatten()?;
1531 let actual = ContentHash::compute_typed("blob", bytes.as_slice());
1532 if actual != *hash {
1533 return None;
1537 }
1538 if let Ok(mut verified) = self.verified_loose_blobs.write() {
1539 verified.insert(*hash, ());
1540 }
1541 Some(path)
1542 }
1543
1544 #[instrument(skip(self), fields(hash = %hash.short()))]
1557 fn promote_to_loose_uncompressed(&self, hash: &ContentHash) -> Result<bool> {
1558 let path = hash_path(&blobs_dir(&self.root), hash);
1559
1560 if !ObjectStore::has_blob_locally(self, hash)?
1564 && let Some(source) = &self.external_source
1565 && (self.recent_blob(hash).is_some() || source.get_blob(hash)?.is_some())
1566 {
1567 return Ok(false);
1568 }
1569
1570 if let Some((header, _)) = read_file_header(&path, 9)?
1572 && !is_compressed(&header)
1573 {
1574 trace!("Blob already loose+uncompressed; skipping promotion");
1575 return Ok(false);
1576 }
1577
1578 let blob = self.get_blob(hash)?.ok_or_else(|| {
1582 HeddleError::NotFound(format!(
1583 "blob {} not found in store; cannot promote to loose-uncompressed",
1584 hash
1585 ))
1586 })?;
1587
1588 debug!(
1601 size = blob.content().len(),
1602 "Promoting blob to loose-uncompressed canonical store"
1603 );
1604 self.write_loose_object_cache(&path, blob.content())?;
1605 if let Ok(mut verified) = self.verified_loose_blobs.write() {
1606 verified.insert(*hash, ());
1607 }
1608 Ok(true)
1609 }
1610
1611 #[instrument(skip(self), fields(hash = %hash.short()))]
1612 fn blob_size(&self, hash: &ContentHash) -> Result<Option<u64>> {
1613 if let Some(size) = self.try_get_blob_size_once(hash)? {
1614 return Ok(Some(size));
1615 }
1616 if self.reload_packs_if_stale()?
1620 && let Some(size) = self.try_get_blob_size_once(hash)?
1621 {
1622 return Ok(Some(size));
1623 }
1624 if let Some(source) = &self.external_source {
1625 if let Some(blob) = self.recent_blob(hash) {
1626 return Ok(Some(blob.content().len() as u64));
1627 }
1628 if let Some(blob) = source.get_blob(hash)? {
1629 let size = blob.content().len() as u64;
1630 self.cache_recent_blob(*hash, &blob);
1631 return Ok(Some(size));
1632 }
1633 }
1634 Ok(None)
1635 }
1636
1637 #[instrument(skip(self), fields(hash = %hash.short()))]
1638 fn get_tree(&self, hash: &ContentHash) -> Result<Option<Tree>> {
1639 if let Some(tree) = self.recent_tree(hash) {
1640 return Ok(Some(tree));
1641 }
1642 if let Some(tree) = self.try_get_tree_once(hash)? {
1643 return Ok(Some(tree));
1644 }
1645 if self.reload_packs_if_stale()?
1646 && let Some(tree) = self.try_get_tree_once(hash)?
1647 {
1648 return Ok(Some(tree));
1649 }
1650 if let Some(source) = &self.external_source
1651 && let Some(tree) = source.get_tree(hash)?
1652 {
1653 self.cache_recent_tree(*hash, &tree);
1654 return Ok(Some(tree));
1655 }
1656 trace!("Tree not found");
1657 Ok(None)
1658 }
1659
1660 #[instrument(skip(self), fields(hash = %hash.short(), name))]
1661 fn get_tree_entry(&self, hash: &ContentHash, name: &str) -> Result<Option<TreeEntry>> {
1662 if let Some(entry) = self.try_get_tree_entry_once(hash, name)? {
1663 return Ok(Some(entry));
1664 }
1665 if self.reload_packs_if_stale()?
1666 && let Some(entry) = self.try_get_tree_entry_once(hash, name)?
1667 {
1668 return Ok(Some(entry));
1669 }
1670 if let Some(source) = &self.external_source
1671 && let Some(tree) = source.get_tree(hash)?
1672 {
1673 let entry = tree.get(name).cloned();
1674 self.cache_recent_tree(*hash, &tree);
1675 return Ok(entry);
1676 }
1677 Ok(None)
1678 }
1679
1680 #[instrument(skip(self), fields(hash = %hash.short()))]
1681 fn get_tree_serialized(&self, hash: &ContentHash) -> Result<Option<Vec<u8>>> {
1682 if let Some(data) = self.try_get_tree_serialized_once(hash)? {
1683 return Ok(Some(data));
1684 }
1685 if self.reload_packs_if_stale()?
1686 && let Some(data) = self.try_get_tree_serialized_once(hash)?
1687 {
1688 return Ok(Some(data));
1689 }
1690 let external_tree = if let Some(tree) = self.recent_tree(hash) {
1691 Some(tree)
1692 } else if let Some(source) = &self.external_source {
1693 let tree = source.get_tree(hash)?;
1694 if let Some(tree) = &tree {
1695 self.cache_recent_tree(*hash, tree);
1696 }
1697 tree
1698 } else {
1699 None
1700 };
1701 if let Some(tree) = external_tree {
1702 return tree.encode_canonical().map(Some).map_err(HeddleError::from);
1703 }
1704 Ok(None)
1705 }
1706
1707 fn open_tree(
1708 &self,
1709 tree_id: &ContentHash,
1710 cursor: Option<&TreeResumeCursor>,
1711 ) -> Result<Option<TreeEntryReader<OpenedTreeBody>>> {
1712 if let Some(reader) = self.try_open_tree_once(tree_id, cursor)? {
1713 return Ok(Some(reader));
1714 }
1715 if self.reload_packs_if_stale()?
1716 && let Some(reader) = self.try_open_tree_once(tree_id, cursor)?
1717 {
1718 return Ok(Some(reader));
1719 }
1720 if let Some(data) = ObjectStore::get_tree_serialized(self, tree_id)? {
1721 let body = if is_streamable_tree(&data) {
1722 data
1723 } else if data.starts_with(TREE_DELTA_MAGIC) {
1724 self.get_tree(tree_id)?
1725 .ok_or_else(|| HeddleError::NotFound(format!("tree {tree_id}")))?
1726 .encode_lean()?
1727 } else if is_salted_tree(&data) {
1728 return Err(salted_tree_not_streamable(tree_id));
1729 } else {
1730 return Ok(None);
1731 };
1732 return Ok(Some(TreeEntryReader::open(
1733 OpenedTreeBody::Bytes(BytesTreeSource::sequential_verify(body)),
1734 *tree_id,
1735 cursor,
1736 )?));
1737 }
1738 Ok(None)
1739 }
1740
1741 #[instrument(skip(self, tree), fields(entry_count = tree.entries().len()))]
1742 fn put_tree(&self, tree: &Tree) -> Result<ContentHash> {
1743 let hash = tree.hash();
1744 let path = hash_path(&trees_dir(&self.root), &hash);
1745
1746 if !ObjectStore::has_tree_locally(self, &hash)? {
1751 let (_, data) = codec::encode_tree(tree, &self.compression)?;
1752 trace!(compressed_size = data.len(), "Writing tree");
1753 self.write_loose_object_atomic(&path, &data)?;
1754 } else {
1755 trace!("Tree already exists, skipping write");
1756 }
1757 if let Ok(mut cache) = self.recent_trees.write() {
1758 cache.insert(hash, tree.clone());
1759 }
1760
1761 self.remove_partial_tree(&hash)?;
1767
1768 Ok(hash)
1769 }
1770
1771 #[instrument(skip(self, data), fields(hash = %hash.short(), size = data.len()))]
1772 fn put_tree_serialized(&self, data: &[u8], hash: ContentHash) -> Result<ContentHash> {
1773 if is_redacted_tree(data) {
1776 self.put_partial_tree(&hash, data)?;
1777 return Ok(hash);
1778 }
1779 let tree = validate_loaded_tree(self.decode_tree_storage_body(hash, data)?)?;
1780
1781 let path = hash_path(&trees_dir(&self.root), &hash);
1782 let should_write = match read_file_bytes(&path)? {
1783 Some(existing) => codec::decode_tree_body(existing.as_slice())? != data,
1784 None => true,
1785 };
1786 if should_write {
1787 trace!(size = data.len(), "Writing raw serialized tree");
1788 self.write_loose_object_atomic(&path, data)?;
1789 }
1790 if let Ok(mut cache) = self.recent_trees.write() {
1791 cache.insert(hash, tree);
1792 }
1793
1794 self.remove_partial_tree(&hash)?;
1797
1798 Ok(hash)
1799 }
1800
1801 #[instrument(skip(self), fields(hash = %hash.short()))]
1802 fn has_tree(&self, hash: &ContentHash) -> Result<bool> {
1803 if ObjectStore::has_tree_locally(self, hash)? {
1804 return Ok(true);
1805 }
1806 if let Some(source) = &self.external_source {
1807 if self.recent_tree(hash).is_some() {
1808 return Ok(true);
1809 }
1810 if let Some(tree) = source.get_tree(hash)? {
1811 self.cache_recent_tree(*hash, &tree);
1812 return Ok(true);
1813 }
1814 }
1815 Ok(false)
1816 }
1817
1818 fn has_tree_locally(&self, hash: &ContentHash) -> Result<bool> {
1819 if self.try_has_tree_once(hash)? {
1820 return Ok(true);
1821 }
1822 Ok(self.reload_packs_if_stale()? && self.try_has_tree_once(hash)?)
1823 }
1824
1825 fn has_partial_tree(&self, hash: &ContentHash) -> Result<bool> {
1826 Ok(partial_tree_path(&self.root, hash).exists())
1827 }
1828
1829 fn get_partial_tree_bytes(&self, hash: &ContentHash) -> Result<Option<Vec<u8>>> {
1830 let path = partial_tree_path(&self.root, hash);
1831 match fs::read(&path) {
1832 Ok(bytes) => Ok(Some(bytes)),
1833 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(None),
1834 Err(err) => Err(HeddleError::Io(err)),
1835 }
1836 }
1837
1838 fn put_partial_tree_bytes(&self, hash: &ContentHash, bytes: &[u8]) -> Result<()> {
1839 let dir = partial_trees_dir(&self.root);
1840 if !dir.exists() {
1841 crate::fs_atomic::create_dir_all_durable(&dir)?;
1842 }
1843 let path = partial_tree_path(&self.root, hash);
1844 crate::fs_atomic::write_file_atomic(&path, bytes)?;
1845 Ok(())
1846 }
1847
1848 fn list_partial_trees(&self) -> Result<Vec<ContentHash>> {
1849 let dir = partial_trees_dir(&self.root);
1850 if !dir.exists() {
1851 return Ok(Vec::new());
1852 }
1853 let mut out = Vec::new();
1854 for entry in fs::read_dir(&dir)? {
1855 let entry = entry?;
1856 let path = entry.path();
1857 if path.extension().and_then(|e| e.to_str()) != Some("bin") {
1858 continue;
1859 }
1860 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
1861 continue;
1862 };
1863 if let Ok(hash) = ContentHash::from_hex(stem) {
1864 out.push(hash);
1865 }
1866 }
1867 Ok(out)
1868 }
1869
1870 fn remove_partial_tree(&self, hash: &ContentHash) -> Result<()> {
1871 let path = partial_tree_path(&self.root, hash);
1872 match fs::remove_file(&path) {
1873 Ok(()) => Ok(()),
1874 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()),
1875 Err(err) => Err(HeddleError::Io(err)),
1876 }
1877 }
1878
1879 #[instrument(skip(self), fields(id = %id.short()))]
1880 fn get_state(&self, id: &StateId) -> Result<Option<State>> {
1881 if let Some(state) = self.recent_state(id) {
1882 return Ok(Some(state));
1883 }
1884 if let Some(state) = self.try_get_state_once(id)? {
1885 return Ok(Some(state));
1886 }
1887 if self.reload_packs_if_stale()?
1888 && let Some(state) = self.try_get_state_once(id)?
1889 {
1890 return Ok(Some(state));
1891 }
1892 if let Some(source) = &self.external_source
1893 && let Some(state) = source.get_state(id)?
1894 {
1895 self.cache_recent_state(*id, &state);
1896 return Ok(Some(state));
1897 }
1898 trace!("State not found");
1899 Ok(None)
1900 }
1901
1902 #[instrument(skip(self, state), fields(id = %state.id().short()))]
1903 fn put_state(&self, state: &State) -> Result<()> {
1904 let state_id = state.id();
1905 let path = state_path(&self.root, &state_id);
1906 let data = codec::encode_state(state, &self.compression)?;
1907 trace!(compressed_size = data.len(), "Writing state");
1908 self.write_loose_object_atomic(&path, &data)?;
1909 if let Ok(mut cache) = self.recent_states.write() {
1910 let mut cached = state.clone();
1911 cached.state_id = state_id;
1912 cache.insert(state_id, cached);
1913 }
1914 Ok(())
1915 }
1916
1917 #[instrument(skip(self, data), fields(id = %id.short(), size = data.len()))]
1918 fn put_state_serialized(&self, data: &[u8], id: StateId) -> Result<()> {
1919 let state = validate_state_serialized(data, id)?;
1920 let path = state_path(&self.root, &id);
1921 trace!(size = data.len(), "Writing raw serialized state");
1922 self.write_loose_object_atomic(&path, data)?;
1923 if let Ok(mut cache) = self.recent_states.write() {
1924 cache.insert(id, state);
1925 }
1926 Ok(())
1927 }
1928
1929 #[instrument(skip(self), fields(id = %id.short()))]
1930 fn has_state(&self, id: &StateId) -> Result<bool> {
1931 if self.try_has_state_once(id)? {
1932 return Ok(true);
1933 }
1934 if self.reload_packs_if_stale()? && self.try_has_state_once(id)? {
1935 return Ok(true);
1936 }
1937 if let Some(source) = &self.external_source {
1938 if self.recent_state(id).is_some() {
1939 return Ok(true);
1940 }
1941 if let Some(state) = source.get_state(id)? {
1942 self.cache_recent_state(*id, &state);
1943 return Ok(true);
1944 }
1945 }
1946 Ok(false)
1947 }
1948
1949 #[instrument(skip(self))]
1950 fn list_states(&self) -> Result<Vec<StateId>> {
1951 self.reload_packs_if_stale()?;
1952
1953 let mut states = Vec::new();
1954 let mut known = HashSet::new();
1955 let dir = states_dir(&self.root);
1956 if dir.exists() {
1957 for entry in fs::read_dir(&dir)? {
1958 let entry = entry?;
1959 let path = entry.path();
1960 if let Some(name) = path.file_stem()
1961 && let Some(name_str) = name.to_str()
1962 && let Ok(id) = StateId::parse(name_str)
1963 && known.insert(id)
1964 {
1965 states.push(id);
1966 }
1967 }
1968 }
1969 if let Ok(manager) = self.pack_manager().read() {
1970 append_unique_states(
1971 &mut states,
1972 &mut known,
1973 manager
1974 .list_all_ids()?
1975 .into_iter()
1976 .filter_map(|id| match id {
1977 PackObjectId::StateId(state) => Some(state),
1978 PackObjectId::Hash(_) | PackObjectId::AnnotatedTag(_) => None,
1979 }),
1980 );
1981 }
1982 if let Some(source) = &self.external_source {
1983 append_unique_states(&mut states, &mut known, source.list_states()?);
1984 }
1985 debug!(count = states.len(), "Listed states");
1986 Ok(states)
1987 }
1988
1989 fn get_state_attachment(
1990 &self,
1991 state: &StateId,
1992 id: &StateAttachmentId,
1993 ) -> Result<Option<StateAttachment>> {
1994 if let Some(attachment) = self.try_get_state_attachment_once(state, id)? {
1995 return Ok(Some(attachment));
1996 }
1997 if self.reload_packs_if_stale()? {
1998 return self.try_get_state_attachment_once(state, id);
1999 }
2000 Ok(None)
2001 }
2002
2003 fn put_state_attachment(&self, attachment: &StateAttachment) -> Result<StateAttachmentId> {
2004 let id = attachment.id();
2005 self.with_state_attachment_index_lock(&attachment.state_id, || {
2006 let index_path = state_attachment_index_path(&self.root, &attachment.state_id);
2007 let mut ids: Vec<StateAttachmentId> = match read_file_bytes(&index_path)? {
2008 Some(bytes) => rmp_serde::from_slice(bytes.as_slice())?,
2009 None => self.rebuild_state_attachment_index(&attachment.state_id)?,
2010 };
2011 if !ids.contains(&id) {
2012 ids.push(id);
2013 ids.sort();
2014 self.write_loose_object_atomic(&index_path, &rmp_serde::to_vec_named(&ids)?)?;
2015 }
2016 let path = state_attachment_path(&self.root, &attachment.state_id, &id);
2017 self.write_loose_object_atomic(&path, &attachment.encode_current_msgpack()?)?;
2018 Ok(id)
2019 })
2020 }
2021
2022 fn list_state_attachments(&self, state: &StateId) -> Result<Vec<StateAttachment>> {
2023 self.with_state_attachment_index_lock(state, || {
2024 let index_path = state_attachment_index_path(&self.root, state);
2025 let mut ids: Vec<StateAttachmentId> = match read_file_bytes(&index_path)? {
2026 Some(bytes) => rmp_serde::from_slice(bytes.as_slice())?,
2027 None => self.rebuild_state_attachment_index(state)?,
2028 };
2029 let mut attachments = Vec::new();
2030 let mut stale = false;
2031 for id in &ids {
2032 match self.get_state_attachment(state, id)? {
2033 Some(attachment) => attachments.push(attachment),
2034 None => stale = true,
2035 }
2036 }
2037 if stale {
2038 ids = self.rebuild_state_attachment_index(state)?;
2039 attachments.clear();
2040 for id in ids {
2041 let attachment = self.get_state_attachment(state, &id)?.ok_or_else(|| {
2042 HeddleError::InvalidObject(format!(
2043 "rebuilt state attachment index references missing {id}"
2044 ))
2045 })?;
2046 attachments.push(attachment);
2047 }
2048 }
2049 Ok(attachments)
2050 })
2051 }
2052
2053 #[instrument(skip(self), fields(id = %id))]
2054 fn get_action(&self, id: &ActionId) -> Result<Option<Action>> {
2055 if let Some(action) = self.try_get_action_once(id)? {
2056 return Ok(Some(action));
2057 }
2058 if self.reload_packs_if_stale()? {
2059 return self.try_get_action_once(id);
2060 }
2061 trace!("Action not found");
2062 Ok(None)
2063 }
2064
2065 #[instrument(skip(self, action))]
2066 fn put_action(&self, action: &mut Action) -> Result<ActionId> {
2067 let id = action.id();
2068 let path = action_path(&self.root, &id);
2069
2070 if !path.exists() {
2071 let (_, data) = codec::encode_action(action, &self.compression)?;
2072 trace!(id = %id, compressed_size = data.len(), "Writing action");
2073 self.write_loose_object_atomic(&path, &data)?;
2074 }
2075
2076 Ok(id)
2077 }
2078
2079 #[instrument(skip(self))]
2080 fn list_actions(&self) -> Result<Vec<ActionId>> {
2081 self.reload_packs_if_stale()?;
2082 let dir = actions_dir(&self.root);
2083 let mut action_hashes = Vec::new();
2084 if dir.exists() {
2085 for entry in fs::read_dir(&dir)? {
2086 let entry = entry?;
2087 let path = entry.path();
2088 if let Some(name) = path.file_stem()
2089 && let Some(name_str) = name.to_str()
2090 && let Ok(hash) = ContentHash::from_hex(name_str)
2091 {
2092 action_hashes.push(hash);
2093 }
2094 }
2095 }
2096 if let Ok(manager) = self.pack_manager().read() {
2097 append_packed_hashes(&mut action_hashes, &manager, ObjectType::Action)?;
2098 }
2099 let actions = action_hashes
2100 .into_iter()
2101 .map(ActionId::from_hash)
2102 .collect::<Vec<_>>();
2103 debug!(count = actions.len(), "Listed actions");
2104 Ok(actions)
2105 }
2106
2107 #[instrument(skip(self))]
2108 fn list_blobs(&self) -> Result<Vec<ContentHash>> {
2109 self.reload_packs_if_stale()?;
2110 let dir = blobs_dir(&self.root);
2111 let mut blobs = list_hashes_from_dir(&dir)?;
2112 if let Ok(manager) = self.pack_manager().read() {
2113 append_packed_hashes(&mut blobs, &manager, ObjectType::Blob)?;
2114 }
2115 Ok(blobs)
2116 }
2117
2118 #[instrument(skip(self))]
2119 fn list_trees(&self) -> Result<Vec<ContentHash>> {
2120 self.reload_packs_if_stale()?;
2121 let dir = trees_dir(&self.root);
2122 let mut trees = list_hashes_from_dir(&dir)?;
2123 if let Ok(manager) = self.pack_manager().read() {
2124 append_packed_hashes(&mut trees, &manager, ObjectType::Tree)?;
2125 }
2126 if let Ok(manager) = self.npk1_manager().read() {
2127 trees.extend(manager.list_ids()?);
2128 }
2129 trees.sort();
2130 trees.dedup();
2131 Ok(trees)
2132 }
2133
2134 #[instrument(skip(self), fields(id = ?id))]
2135 fn get_pack_object(&self, id: &PackObjectId) -> Result<Option<(ObjectType, Vec<u8>)>> {
2136 if let Ok(manager) = self.pack_manager().read()
2137 && let Some((obj_type, data)) = manager.get_object(id)?
2138 {
2139 return Ok(Some((obj_type, data)));
2140 }
2141
2142 match id {
2143 PackObjectId::AnnotatedTag(hash) => Ok(self
2144 .get_annotated_tag(hash)?
2145 .map(|tag| (ObjectType::AnnotatedTag, tag.encode_current_msgpack()))),
2146 PackObjectId::Hash(hash) => {
2147 if let Some(blob) = self.get_blob(hash)? {
2148 return Ok(Some((ObjectType::Blob, blob.into_content())));
2149 }
2150 if let Some(tree_data) = self.get_tree_serialized(hash)? {
2154 return Ok(Some((ObjectType::Tree, tree_data)));
2155 }
2156 if let Some(action) = self.get_action(&ActionId::from_hash(*hash))? {
2157 return Ok(Some((
2158 ObjectType::Action,
2159 rmp_serde::to_vec_named(&action)?,
2160 )));
2161 }
2162 Ok(None)
2163 }
2164 PackObjectId::StateId(change_id) => {
2165 if let Some(state) = self.get_state(change_id)? {
2166 Ok(Some((ObjectType::State, state.encode_current_msgpack()?)))
2167 } else {
2168 Ok(None)
2169 }
2170 }
2171 }
2172 }
2173
2174 #[instrument(skip(self, pack_data, index_data))]
2175 fn install_pack(&self, pack_data: &[u8], index_data: &[u8]) -> Result<Vec<PackObjectId>> {
2176 let reader = crate::store::pack::PackReader::from_slice(pack_data, index_data)?;
2177 let ids = validate_and_list_pack(self, &reader)?;
2178 let state_entries = self.states_needing_loose_copies(&reader, &ids)?;
2179 let attachment_entries = attachment_entries_from_pack(&reader, &ids)?;
2180 self.install_pack_files(pack_data, index_data)?;
2181 self.write_packed_state_mirrors_batch(state_entries)?;
2182 for attachment in attachment_entries {
2183 self.put_state_attachment(&attachment)?;
2184 }
2185 self.clear_recent_object_caches();
2186 Ok(ids)
2187 }
2188
2189 #[instrument(skip(self, blobs), fields(count = blobs.len()))]
2190 fn put_blobs_packed(&self, blobs: Vec<(crate::object::ContentHash, Vec<u8>)>) -> Result<()> {
2191 self.put_blobs_packed_impl(blobs)
2192 }
2193
2194 #[instrument(skip(self, blobs, tree, state), fields(blob_count = blobs.len()))]
2195 fn put_snapshot_objects_packed(
2196 &self,
2197 blobs: Vec<(ContentHash, Vec<u8>)>,
2198 tree: &Tree,
2199 state: &State,
2200 ) -> Result<()> {
2201 self.put_snapshot_objects_packed_impl(
2202 blobs,
2203 Vec::new(),
2204 &TreeWrite::anchor(tree.clone()),
2205 state,
2206 Vec::new(),
2207 None,
2208 )
2209 .map(|_| ())
2210 }
2211
2212 fn put_snapshot_objects_and_attachments_packed(
2213 &self,
2214 blobs: Vec<(ContentHash, Vec<u8>)>,
2215 tree: &Tree,
2216 state: &State,
2217 attachments: Vec<StateAttachment>,
2218 ) -> Result<()> {
2219 self.put_snapshot_objects_packed_impl(
2220 blobs,
2221 Vec::new(),
2222 &TreeWrite::anchor(tree.clone()),
2223 state,
2224 attachments,
2225 None,
2226 )
2227 .map(|_| ())
2228 }
2229
2230 #[instrument(skip(self))]
2231 fn install_pack_streaming(
2232 &self,
2233 pack_path: &std::path::Path,
2234 index_path: &std::path::Path,
2235 ) -> Result<Vec<PackObjectId>> {
2236 let ids = {
2242 let reader = crate::store::pack::PackReader::open(pack_path, index_path)?;
2243 validate_and_list_pack(self, &reader)?
2244 };
2245 let state_entries = {
2246 let reader = crate::store::pack::PackReader::open(pack_path, index_path)?;
2247 self.states_needing_loose_copies(&reader, &ids)?
2248 };
2249 let attachment_entries = {
2250 let reader = crate::store::pack::PackReader::open(pack_path, index_path)?;
2251 attachment_entries_from_pack(&reader, &ids)?
2252 };
2253 self.install_pack_files_streaming(pack_path, index_path)?;
2254 self.write_packed_state_mirrors_batch(state_entries)?;
2255 for attachment in attachment_entries {
2256 self.put_state_attachment(&attachment)?;
2257 }
2258 Ok(ids)
2259 }
2260
2261 #[instrument(skip(self))]
2262 fn begin_snapshot_write_batch(&self) -> Result<()> {
2263 self.begin_snapshot_write_batch_impl()
2264 }
2265
2266 #[instrument(skip(self))]
2267 fn flush_snapshot_write_batch(&self) -> Result<()> {
2268 self.flush_snapshot_write_batch_impl()
2269 }
2270
2271 #[instrument(skip(self))]
2272 fn abort_snapshot_write_batch(&self) {
2273 self.abort_snapshot_write_batch_impl();
2274 }
2275}
2276
2277impl SidecarStore for FsStore {
2278 fn has_redactions_for_blob(&self, blob: &ContentHash) -> Result<bool> {
2279 Ok(redaction_path(&self.root, blob).exists())
2280 }
2281
2282 fn get_redactions_bytes_for_blob(&self, blob: &ContentHash) -> Result<Option<Vec<u8>>> {
2283 let path = redaction_path(&self.root, blob);
2284 match fs::read(&path) {
2285 Ok(bytes) => Ok(Some(bytes)),
2286 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(None),
2287 Err(err) => Err(HeddleError::Io(err)),
2288 }
2289 }
2290
2291 fn put_redactions_bytes_for_blob(&self, blob: &ContentHash, bytes: &[u8]) -> Result<()> {
2292 let dir = redactions_dir(&self.root);
2293 if !dir.exists() {
2294 crate::fs_atomic::create_dir_all_durable(&dir)?;
2295 }
2296 let path = redaction_path(&self.root, blob);
2297 crate::fs_atomic::write_file_atomic(&path, bytes)?;
2298 Ok(())
2299 }
2300
2301 fn list_blobs_with_redactions(&self) -> Result<Vec<ContentHash>> {
2302 let dir = redactions_dir(&self.root);
2303 if !dir.exists() {
2304 return Ok(Vec::new());
2305 }
2306 let mut out = Vec::new();
2307 for entry in fs::read_dir(&dir)? {
2308 let entry = entry?;
2309 let path = entry.path();
2310 if path.extension().and_then(|e| e.to_str()) != Some("bin") {
2311 continue;
2312 }
2313 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
2314 continue;
2315 };
2316 if let Ok(hash) = ContentHash::from_hex(stem) {
2317 out.push(hash);
2318 }
2319 }
2320 Ok(out)
2321 }
2322
2323 fn has_state_visibility_for_state(&self, state: &StateId) -> Result<bool> {
2324 Ok(state_visibility_path(&self.root, state).exists())
2325 }
2326
2327 fn get_state_visibility_bytes_for_state(&self, state: &StateId) -> Result<Option<Vec<u8>>> {
2328 let path = state_visibility_path(&self.root, state);
2329 match fs::read(&path) {
2330 Ok(bytes) => Ok(Some(bytes)),
2331 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(None),
2332 Err(err) => Err(HeddleError::Io(err)),
2333 }
2334 }
2335
2336 fn put_state_visibility_bytes_for_state(&self, state: &StateId, bytes: &[u8]) -> Result<()> {
2337 let dir = state_visibility_dir(&self.root);
2338 if !dir.exists() {
2339 crate::fs_atomic::create_dir_all_durable(&dir)?;
2340 }
2341 let path = state_visibility_path(&self.root, state);
2342 crate::fs_atomic::write_file_atomic(&path, bytes)?;
2343 Ok(())
2344 }
2345
2346 fn list_states_with_visibility(&self) -> Result<Vec<StateId>> {
2347 let dir = state_visibility_dir(&self.root);
2348 if !dir.exists() {
2349 return Ok(Vec::new());
2350 }
2351 let mut out = Vec::new();
2352 for entry in fs::read_dir(&dir)? {
2353 let entry = entry?;
2354 let path = entry.path();
2355 if path.extension().and_then(|e| e.to_str()) != Some("bin") {
2356 continue;
2357 }
2358 let Some(stem) = path.file_stem().and_then(|s| s.to_str()) else {
2359 continue;
2360 };
2361 if let Ok(state) = StateId::parse(stem) {
2362 out.push(state);
2363 }
2364 }
2365 Ok(out)
2366 }
2367}
2368
2369#[cfg(test)]
2370mod state_attachment_tests {
2371 use std::sync::Arc;
2372
2373 use chrono::Utc;
2374
2375 use super::*;
2376 use crate::{
2377 object::{Attribution, Principal, StateAttachmentBody},
2378 store::{CompressionConfig, pack::PackBuilder},
2379 };
2380
2381 fn fixture(store: &FsStore) -> (State, StateAttachment) {
2382 let tree = store.put_tree(&Tree::new()).unwrap();
2383 let attribution = Attribution::human(Principal::new("Test", "test@example.com"));
2384 let state = State::new(tree, vec![], attribution.clone());
2385 store.put_state(&state).unwrap();
2386 let attachment = StateAttachment {
2387 state_id: state.id(),
2388 body: StateAttachmentBody::Context(ContentHash::compute(b"context")),
2389 attribution,
2390 created_at: Utc::now(),
2391 supersedes: None,
2392 };
2393 (state, attachment)
2394 }
2395
2396 #[test]
2397 fn concurrent_attachment_writes_keep_every_index_entry() {
2398 let temp = tempfile::TempDir::new().unwrap();
2399 let store = Arc::new(FsStore::new(temp.path()));
2400 let (state, base) = fixture(&store);
2401 let mut threads = Vec::new();
2402 for byte in 0..16u8 {
2403 let store = Arc::clone(&store);
2404 let mut attachment = base.clone();
2405 attachment.body = StateAttachmentBody::Context(ContentHash::compute(&[byte]));
2406 threads.push(std::thread::spawn(move || {
2407 store.put_state_attachment(&attachment).unwrap();
2408 }));
2409 }
2410 for thread in threads {
2411 thread.join().unwrap();
2412 }
2413 assert_eq!(store.list_state_attachments(&state.id()).unwrap().len(), 16);
2414 }
2415
2416 #[test]
2417 fn missing_index_rebuilds_from_loose_objects() {
2418 let temp = tempfile::TempDir::new().unwrap();
2419 let store = FsStore::new(temp.path());
2420 let (state, attachment) = fixture(&store);
2421 store.put_state_attachment(&attachment).unwrap();
2422 fs::remove_file(state_attachment_index_path(&store.root, &state.id())).unwrap();
2423 assert_eq!(
2424 store.list_state_attachments(&state.id()).unwrap(),
2425 vec![attachment]
2426 );
2427 }
2428
2429 #[test]
2430 fn packed_attachment_uses_state_index_for_lookup() {
2431 let temp = tempfile::TempDir::new().unwrap();
2432 let store = FsStore::new(temp.path());
2433 let (state, attachment) = fixture(&store);
2434 let mut builder = PackBuilder::new(CompressionConfig::default());
2435 builder.add(
2436 *attachment.id().as_hash(),
2437 ObjectType::StateAttachment,
2438 rmp_serde::to_vec_named(&attachment).unwrap(),
2439 );
2440 let (pack, index, _) = builder.build().unwrap();
2441 store.install_pack(&pack, &index).unwrap();
2442 fs::remove_file(state_attachment_path(
2443 &store.root,
2444 &state.id(),
2445 &attachment.id(),
2446 ))
2447 .unwrap();
2448 let rebuild_marker =
2449 state_attachment_index_path(&store.root, &state.id()).with_extension("rebuild-marker");
2450 let _ = fs::remove_file(&rebuild_marker);
2451 assert_eq!(
2452 store.list_state_attachments(&state.id()).unwrap(),
2453 vec![attachment.clone()]
2454 );
2455 assert_eq!(
2456 store.list_state_attachments(&state.id()).unwrap(),
2457 vec![attachment]
2458 );
2459 assert!(!rebuild_marker.exists());
2460 }
2461}
2462
2463#[cfg(test)]
2464mod enumeration_tests {
2465 use heddle_format::{compression::CompressionConfig, delta::DeltaEncoder};
2466 use tempfile::TempDir;
2467
2468 use super::*;
2469 use crate::store::pack::{
2470 PackBuilder, PackContainerSpec, PackIndex, append_container_checksum,
2471 encode_tagged_entry_parts, write_container_header,
2472 };
2473
2474 #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
2475 struct TestEnumerationMetrics {
2476 membership_checks: u64,
2477 header_reads: u64,
2478 full_object_decodes: u64,
2479 }
2480
2481 impl EnumerationCounter for TestEnumerationMetrics {
2482 fn membership_check(&mut self) {
2483 self.membership_checks += 1;
2484 }
2485
2486 fn header_read(&mut self) {
2487 self.header_reads += 1;
2488 }
2489 }
2490
2491 fn install_pack_files(
2492 dir: &TempDir,
2493 name: &str,
2494 pack_data: &[u8],
2495 index_data: &[u8],
2496 ) -> PackManager {
2497 fs::write(dir.path().join(format!("{name}.pack")), pack_data).unwrap();
2498 fs::write(dir.path().join(format!("{name}.idx")), index_data).unwrap();
2499 PackManager::new(dir.path().to_path_buf())
2500 }
2501
2502 fn raw_mixed_manager() -> (TempDir, PackManager, Vec<(ContentHash, ObjectType)>) {
2503 let dir = TempDir::new().unwrap();
2504 let objects = [
2505 (ObjectType::Blob, b"packed blob".as_slice()),
2506 (ObjectType::Tree, b"packed tree".as_slice()),
2507 (ObjectType::Action, b"packed action".as_slice()),
2508 ];
2509 let mut builder = PackBuilder::new(CompressionConfig::disabled());
2510 let mut classified = Vec::new();
2511 for (obj_type, data) in objects {
2512 let hash = ContentHash::compute_typed("enumeration-test", data);
2513 builder.add(hash, obj_type, data.to_vec());
2514 classified.push((hash, obj_type));
2515 }
2516 let (pack, index, _) = builder.build().unwrap();
2517 let manager = install_pack_files(&dir, "mixed", &pack, &index);
2518 (dir, manager, classified)
2519 }
2520
2521 fn delta_chain_manager() -> (TempDir, PackManager, Vec<ContentHash>) {
2522 const SPEC: PackContainerSpec = PackContainerSpec {
2523 magic: b"LMPK",
2524 version: 4,
2525 };
2526 let dir = TempDir::new().unwrap();
2527 let base = b"delta-chain base payload ".repeat(64);
2528 let mut middle = base.clone();
2529 middle[200..208].copy_from_slice(b"middle!!");
2530 let mut tip = middle.clone();
2531 tip[900..908].copy_from_slice(b"tip!!!!!");
2532 let bodies = [&base, &middle, &tip];
2533 let hashes = bodies
2534 .iter()
2535 .map(|body| ContentHash::compute_typed("blob", body))
2536 .collect::<Vec<_>>();
2537 let middle_delta = DeltaEncoder::encode(&base, &middle);
2538 let tip_delta = DeltaEncoder::encode(&middle, &tip);
2539
2540 let mut pack = Vec::new();
2541 let mut index = PackIndex::new();
2542 write_container_header(&mut pack, SPEC, 3);
2543 for (position, payload) in [
2544 base.as_slice(),
2545 middle_delta.as_slice(),
2546 tip_delta.as_slice(),
2547 ]
2548 .into_iter()
2549 .enumerate()
2550 {
2551 index.add(PackObjectId::Hash(hashes[position]), pack.len() as u64);
2552 let (stored_type, base_id) = if position == 0 {
2553 (ObjectType::Blob, None)
2554 } else {
2555 (
2556 ObjectType::Delta,
2557 Some(PackObjectId::Hash(hashes[position - 1])),
2558 )
2559 };
2560 encode_tagged_entry_parts(
2561 &mut pack,
2562 PackObjectId::Hash(hashes[position]),
2563 stored_type,
2564 bodies[position].len(),
2565 base_id,
2566 payload,
2567 )
2568 .unwrap();
2569 }
2570 index.sort();
2571 append_container_checksum(&mut pack);
2572 let manager = install_pack_files(&dir, "delta-chain", &pack, &index.to_bytes());
2573 (dir, manager, hashes)
2574 }
2575
2576 fn legacy_append_packed_hashes(
2577 hashes: &mut Vec<ContentHash>,
2578 manager: &PackManager,
2579 expected_type: ObjectType,
2580 ) -> Result<TestEnumerationMetrics> {
2581 let mut metrics = TestEnumerationMetrics::default();
2582 for id in manager.list_all_ids()? {
2583 let PackObjectId::Hash(hash) = id else {
2584 continue;
2585 };
2586 let mut already_listed = false;
2587 for listed in hashes.iter() {
2588 metrics.membership_checks += 1;
2589 if listed == &hash {
2590 already_listed = true;
2591 break;
2592 }
2593 }
2594 if already_listed {
2595 continue;
2596 }
2597 metrics.full_object_decodes += 1;
2598 if let Some((obj_type, _)) = manager.get_hashed_object(&hash)?
2599 && obj_type == expected_type
2600 {
2601 hashes.push(hash);
2602 }
2603 }
2604 Ok(metrics)
2605 }
2606
2607 fn assert_new_matches_legacy(
2608 label: &str,
2609 manager: &PackManager,
2610 loose: Vec<ContentHash>,
2611 expected_type: ObjectType,
2612 ) {
2613 let mut new = loose.clone();
2614 let mut new_metrics = TestEnumerationMetrics::default();
2615 append_packed_hashes_with_counter(&mut new, manager, expected_type, &mut new_metrics)
2616 .unwrap();
2617 let mut legacy = loose;
2618 legacy_append_packed_hashes(&mut legacy, manager, expected_type).unwrap();
2619 assert_eq!(new, legacy, "fixture {label} changed output or ordering");
2620 assert_eq!(new_metrics.full_object_decodes, 0, "fixture {label}");
2621 }
2622
2623 #[test]
2624 fn type_only_enumeration_matches_full_decode_across_fixture_set() {
2625 let empty_dir = TempDir::new().unwrap();
2626 let empty = PackManager::new(empty_dir.path().to_path_buf());
2627 let loose_hash = ContentHash::compute(b"loose only");
2628 assert_new_matches_legacy("loose-only", &empty, vec![loose_hash], ObjectType::Blob);
2629
2630 let (_raw_dir, raw, classified) = raw_mixed_manager();
2631 for expected_type in [ObjectType::Blob, ObjectType::Tree, ObjectType::Action] {
2632 assert_new_matches_legacy("packed-only", &raw, Vec::new(), expected_type);
2633 assert_new_matches_legacy(
2634 "mixed",
2635 &raw,
2636 vec![ContentHash::compute_typed("loose", &[expected_type as u8])],
2637 expected_type,
2638 );
2639 let duplicate = classified
2640 .iter()
2641 .find_map(|(hash, obj_type)| (*obj_type == expected_type).then_some(*hash))
2642 .unwrap();
2643 assert_new_matches_legacy(
2644 "duplicate-loose-packed",
2645 &raw,
2646 vec![duplicate],
2647 expected_type,
2648 );
2649 }
2650 for (hash, expected_type) in classified {
2651 assert_eq!(
2652 raw.get_hashed_object_type(&hash).unwrap(),
2653 raw.get_hashed_object(&hash)
2654 .unwrap()
2655 .map(|(obj_type, _)| obj_type)
2656 );
2657 assert_eq!(
2658 raw.get_hashed_object_type(&hash).unwrap(),
2659 Some(expected_type)
2660 );
2661 }
2662
2663 let (_delta_dir, delta, delta_hashes) = delta_chain_manager();
2664 assert_new_matches_legacy("two-link-delta-chain", &delta, Vec::new(), ObjectType::Blob);
2665 for hash in delta_hashes {
2666 assert_eq!(
2667 delta.get_hashed_object_type(&hash).unwrap(),
2668 Some(ObjectType::Blob)
2669 );
2670 assert_eq!(
2671 delta.get_hashed_object_type(&hash).unwrap(),
2672 delta
2673 .get_hashed_object(&hash)
2674 .unwrap()
2675 .map(|(obj_type, _)| obj_type)
2676 );
2677 }
2678 }
2679
2680 #[test]
2681 fn structural_counter_rejects_vec_scan_and_full_decode_negative_control() {
2682 let (_dir, manager, _) = raw_mixed_manager();
2683 let loose = (0..64u8)
2684 .map(|byte| ContentHash::compute_typed("loose", &[byte]))
2685 .collect::<Vec<_>>();
2686 let packed_hashes = manager
2687 .list_all_ids()
2688 .unwrap()
2689 .into_iter()
2690 .filter(|id| matches!(id, PackObjectId::Hash(_)))
2691 .count() as u64;
2692
2693 for expected_type in [ObjectType::Blob, ObjectType::Tree, ObjectType::Action] {
2694 let mut optimized = loose.clone();
2695 let mut optimized_metrics = TestEnumerationMetrics::default();
2696 append_packed_hashes_with_counter(
2697 &mut optimized,
2698 &manager,
2699 expected_type,
2700 &mut optimized_metrics,
2701 )
2702 .unwrap();
2703 assert_eq!(optimized_metrics.membership_checks, packed_hashes);
2704 assert_eq!(optimized_metrics.header_reads, packed_hashes);
2705 assert_eq!(optimized_metrics.full_object_decodes, 0);
2706
2707 let mut legacy = loose.clone();
2708 let legacy_metrics =
2709 legacy_append_packed_hashes(&mut legacy, &manager, expected_type).unwrap();
2710 assert!(legacy_metrics.membership_checks >= loose.len() as u64 * packed_hashes);
2711 assert_eq!(legacy_metrics.full_object_decodes, packed_hashes);
2712 assert!(
2713 !(legacy_metrics.membership_checks <= packed_hashes
2714 && legacy_metrics.full_object_decodes == 0),
2715 "negative control unexpectedly passed the structural contract: {legacy_metrics:?}"
2716 );
2717 }
2718 }
2719
2720 #[test]
2721 fn state_union_preserves_first_seen_order_at_scale() {
2722 let first = (0..20_000u32)
2723 .map(|value| {
2724 StateId::from_bytes(*ContentHash::compute(&value.to_le_bytes()).as_bytes())
2725 })
2726 .collect::<Vec<_>>();
2727 let second = first[10_000..]
2728 .iter()
2729 .copied()
2730 .chain((20_000..30_000u32).map(|value| {
2731 StateId::from_bytes(*ContentHash::compute(&value.to_le_bytes()).as_bytes())
2732 }))
2733 .collect::<Vec<_>>();
2734 let mut states = Vec::new();
2735 let mut known = HashSet::new();
2736
2737 append_unique_states(&mut states, &mut known, first.iter().copied());
2738 append_unique_states(&mut states, &mut known, second);
2739
2740 assert_eq!(states.len(), 30_000);
2741 assert_eq!(&states[..20_000], first.as_slice());
2742 }
2743}