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