Skip to main content

objects/store/fs/
fs_impl.rs

1// SPDX-License-Identifier: Apache-2.0
2//! ObjectStore implementation for FsStore.
3
4use 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
42/// Bytes we read off disk to recover a blob's uncompressed size.
43/// Must cover the 9-byte modern header **plus** the 4-byte ZSTD
44/// magic that `header_uncompressed_size` uses to disambiguate
45/// modern from legacy (5-byte) headers — without the magic in the
46/// peek buffer the lookup silently returns the on-disk byte length
47/// instead of the recorded uncompressed size, which left `stat`
48/// reporting the compressed size of every loose blob.
49const 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
64/// A streamable body for a tree held only in memory: the lean anchor, or the
65/// full HTR4 body when the tree records a Git source layout HLR1 cannot carry.
66fn 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
74/// A V4 salted (HSR1) tree cannot be streamed through the paging reader — its
75/// Merkle-root id is not verifiable by the reader's incremental single-pass
76/// hasher. Surface a loud `Err` (never a silent `Ok(None)` / NotFound) so the
77/// caller falls back to an eager `get_tree` decode instead of mistaking the
78/// tree for missing. Mirrors how `InMemoryStore::open_tree` surfaces an Err
79/// (via `encode_lean` refusing V4).
80fn 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    /// Whether this store already holds `id` loose, packed or in its recent
229    /// cache. Unlike [`ObjectStore::has_state`] this never consults the
230    /// external source: a packed copy already shadows that source.
231    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    /// The incoming pack's States that need a loose copy: those the store
239    /// already holds. Call this before the pack itself is installed.
240    ///
241    /// Two copies under one StateId can differ in fields outside the hash
242    /// (`git_lossy`, sub-second `created_at`), and readers take a loose copy
243    /// before any pack. The loose copy keeps the later install's bytes in
244    /// front of the copy the store held. A State the store does not hold yet
245    /// has no other copy to shadow, so it is served from its pack, exactly as
246    /// after [`prune_loose_objects`](Self::prune_loose_objects). This keeps
247    /// the install's file and fsync count independent of how many States a
248    /// pack carries (HeddleCo/heddle#2023).
249    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    /// Publish loose copies of packed States the store already held behind
271    /// one parent-directory durability barrier (see
272    /// [`states_needing_loose_copies`](Self::states_needing_loose_copies)).
273    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    /// Publish the attachment index for objects already made durable in a
362    /// snapshot pack. This sidecar is only a materialized view: if a crash
363    /// loses it, [`rebuild_state_attachment_index`](Self::rebuild_state_attachment_index)
364    /// reconstructs it by scanning authoritative loose objects and packs.
365    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
393/// Validate all logical objects before the buffered install returns its IDs.
394/// The streaming seam uses the same validator with a disk metadata visitor.
395fn 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    /// Insert into the recent-blob cache when the payload fits the size gate.
539    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    /// Single-pass blob lookup. The wrapper in `ObjectStore::get_blob`
582    /// retries this once after a stale-reload on miss.
583    fn try_get_blob_once(&self, hash: &ContentHash) -> Result<Option<Blob>> {
584        // Cache first — avoid `path.exists()` / pack probes on warm hits.
585        // Access bits are atomic, so hits remain concurrent under a read lock.
586        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                // Loose blobs are bare bytes on disk: a half-written
613                // file or bit-rot inside the payload would slip past
614                // the path-is-the-hash invariant. Keep the verify on
615                // this path. Pack-resident reads above skip it because
616                // pack entries are framed with offset + length records
617                // that fail to parse if the pack is corrupt.
618                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    /// Shared body for `try_has_{blob,tree,state}_once`: object is
632    /// present iff the loose path exists or the pack manager
633    /// resolves it. Callers pass the loose path and the
634    /// pack-manager probe; the helper handles the lock.
635    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        // This is the native-ownership probe used by `has_blob_locally`.
651        // Recent-object entries may be read-through values from an external
652        // Git overlay, so cache presence cannot establish local durability.
653        let path = hash_path(&blobs_dir(&self.root), hash);
654        self.loose_or_packed(&path, |m| m.has_object(hash))
655    }
656
657    /// Header-only size lookup for a single attempt. Tries:
658    /// 1. The recent-blob cache (we already have the bytes in
659    ///    memory — `len()` is free).
660    /// 2. The loose blob: peek the 9-byte compression header. For a
661    ///    compressed blob the recorded uncompressed size lives in the
662    ///    header. For an uncompressed blob (no recognised header) the
663    ///    on-disk file length IS the blob size.
664    /// 3. Any loaded pack: the pack format records the uncompressed
665    ///    size as a varint right after the tagged id, so we can decode
666    ///    it without touching the body.
667    ///
668    /// Cost: one short read (typically 9 bytes) for loose blobs, or a
669    /// pure in-memory varint decode for packed blobs. *No*
670    /// decompression.
671    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            // No recognised compression header — the file is raw
684            // blob bytes. The on-disk length is the blob size.
685            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        // Cache first. The recent-object cache only ever holds trees we
892        // wrote or read this process, so a hit is authoritative for a
893        // read. Atomic second-chance marking keeps the map under a shared lock.
894        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        // Loose trees may be migration-promoted V2 shadows of an older packed
902        // V1 encoding at the same semantic tree hash. Prefer the loose copy
903        // when it exists, then fall through to pack lookup.
904        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                // This can only be a read-through external tree: native bodies
1022                // were checked above so a cached delta cannot hide a chain.
1023                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        // This is the native-ownership probe used by `has_tree_locally`.
1116        // Recent-object entries may be read-through values from an external
1117        // Git overlay, so cache presence cannot establish local durability.
1118        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        // Cache first — avoid `path.exists()` / pack probes on warm hits.
1130        // Atomic second-chance marking keeps hits under a shared lock. Put
1131        // paths and successful reads below keep the cache coherent for the
1132        // process.
1133        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        // Read-lock `contains`: an existence check needs no clock
1169        // promotion, so it must not serialize on the write lock.
1170        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    /// Lightweight repository-open seam for authoritative snapshot recovery.
1237    #[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    /// Internal repository seam for the local authoritative snapshot artifact.
1248    /// Kept off [`ObjectStore`] so other stores do not acquire a filesystem
1249    /// recovery contract.
1250    #[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    /// O(1) lookup for the authoritative snapshot pack associated with a
1261    /// pushed state.
1262    #[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    /// Drop process-local decoded-object caches for benchmarks and
1276    /// diagnostics. This is intentionally an FsStore operation rather than a
1277    /// durable [`ObjectStore`] requirement.
1278    pub fn clear_recent_caches(&self) {
1279        self.clear_recent_object_caches();
1280    }
1281
1282    /// Repack loose native objects in this filesystem-backed store.
1283    #[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    /// Remove loose objects already represented by an installed native pack.
1289    #[instrument(skip(self))]
1290    pub fn prune_loose_objects(&self) -> Result<(u64, u64)> {
1291        self.prune_loose_objects_impl()
1292    }
1293
1294    /// Remove only pack/index pairs that fail validation so clone repair can
1295    /// advertise their objects as missing on the next pull.
1296    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    /// Zero-copy pack fast path. When the blob lives in a packfile
1373    /// and is non-delta + uncompressed, returns a `Bytes::slice`
1374    /// view of the pack's mmap — no decompression, no allocation,
1375    /// no memcpy. Compressed pack entries, delta entries, and
1376    /// loose blobs fall back to `get_blob` and wrap the result in a
1377    /// `Bytes` (the `Vec` → `Bytes` conversion is itself zero-copy).
1378    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        // Miss path: a sibling FsStore (e.g. the worktree's repo
1400        // backing the same `.heddle/`) may have installed a new pack
1401        // since we loaded ours. Cheap disk-count check first; full
1402        // reload only when the count grew.
1403        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    /// Loose blob path safe for clonefile/copy materialization.
1501    ///
1502    /// Returns `Some(path)` only when the loose file exists, is
1503    /// stored uncompressed, *and* its bytes hash to the expected
1504    /// content hash. Compressed blobs and pack-only blobs fall
1505    /// through to `None`; so do *torn* cache-mirror files (the
1506    /// `AtomicWriteMode::NoSync` write side may leave one if the
1507    /// host crashed during a previous promote). On the torn case
1508    /// the caller re-promotes from the authoritative pack copy.
1509    ///
1510    /// Verification is amortised: a hash that passes the check once
1511    /// in this process is recorded in `verified_loose_blobs` and
1512    /// subsequent calls skip the read+hash. So the cost on the
1513    /// materialize hot path is at most one BLAKE3 over each unique
1514    /// blob per process lifetime — negligible for tiny blobs,
1515    /// bounded by working-set size for huge ones.
1516    fn loose_blob_path(&self, hash: &ContentHash) -> Option<PathBuf> {
1517        let path = hash_path(&blobs_dir(&self.root), hash);
1518        // Fast path: this process already verified (or wrote) this
1519        // hash's loose mirror in `promote_to_loose_uncompressed`.
1520        // Trust without re-hashing — `path.exists()` is the only
1521        // I/O we need.
1522        if let Ok(verified) = self.verified_loose_blobs.read()
1523            && verified.contains(hash)
1524            && path.exists()
1525        {
1526            return Some(path);
1527        }
1528
1529        // First-time-this-process check: peek the header to filter
1530        // out compressed-loose files cheaply, then verify the
1531        // body's hash matches what the caller expects. A torn-write
1532        // (post-crash) cache mirror fails this and the caller
1533        // re-promotes from the pack.
1534        //
1535        // Header peek must cover the 9-byte modern header **plus**
1536        // the 4-byte ZSTD magic that `is_compressed` checks —
1537        // peeking only 9 bytes makes `is_compressed` falsely
1538        // return `false` on a properly-compressed blob, and we'd
1539        // hand the caller the compressed file path. Same off-by-4
1540        // we fixed in `BLOB_HEADER_PEEK`.
1541        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            // Torn write or unrelated corruption. Leave the file on
1549            // disk; the caller's `promote_to_loose_uncompressed`
1550            // will overwrite it via the standard temp+rename path.
1551            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    /// Promote a blob to its uncompressed-loose canonical path so
1560    /// `loose_blob_path` returns `Some(path)` and hardlink-first
1561    /// materialization fires.
1562    ///
1563    /// Three cases:
1564    /// 1. Already loose+uncompressed: peek the header, no-op.
1565    /// 2. Loose but compressed: read+decompress, atomically rewrite
1566    ///    the canonical path with raw bytes.
1567    /// 3. Pack-only: read out of the pack via `get_blob`, atomically
1568    ///    write to the canonical loose path. Pack copy is left in
1569    ///    place — the next prune cycle will discard the loose mirror
1570    ///    and a future materialize will re-promote.
1571    #[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        // External-only overlay blobs stay external. If Heddle also owns a
1576        // native copy, promotion is a native storage optimization and does not
1577        // cross the source-authority boundary.
1578        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        // Idempotent fast path: already loose AND uncompressed.
1586        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        // Either compressed-loose or pack-only. Reading via
1594        // `get_blob` covers both: compressed-loose decompresses on
1595        // the way out, pack-only reads from the loaded pack manager.
1596        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        // Install the uncompressed bytes at the canonical loose path
1604        // via the cache-mirror atomic-write variant: no fsync, just
1605        // temp+rename. The fsync skip is what makes promotion fast
1606        // (measured: ~5 ms/blob with `sync_data` vs ~0.2 ms without
1607        // on macOS APFS); the safety comes from the read-side hash
1608        // check in `loose_blob_path`. A torn write after a crash
1609        // produces a file whose content hash doesn't match, so the
1610        // next reader rejects it and re-promotes from the pack.
1611        //
1612        // Record the hash in this process's verified-blobs cache:
1613        // we just wrote the bytes ourselves, so the subsequent read
1614        // path can trust them without re-hashing.
1615        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        // Sibling-store recovery, mirroring the read path: if a
1632        // concurrent writer just installed a pack we don't know about,
1633        // reload and retry once before reporting a miss.
1634        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        // `put_tree` is an ownership boundary: a native state that references
1762        // this tree must survive loss or pruning of an overlay read-through
1763        // source. Descriptor-only states do not call this method; they retain
1764        // their explicit external-source semantics.
1765        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        // A full canonical tree has landed for this hash — an explicit full
1777        // fetch backfilling what a partial clone withheld. Drop any lingering
1778        // redacted projection so the DERIVED partial marker
1779        // (`list_partial_trees`) clears and the full tree is authoritative.
1780        // Idempotent when no partial slot is held.
1781        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        // An HRT1 redacted projection is not a full tree: route it to the
1789        // partial slot (monotone) rather than through the full-tree decoder.
1790        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        // Full tree backfilled: drop any lingering redacted projection for this
1810        // hash (no auto-backfill; the DERIVED partial marker clears). Idempotent.
1811        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                // Raw canonical storage body: skips a full tree decode +
2166                // re-encode on the pack-building path (the receiver installs
2167                // through `put_tree_serialized`).
2168                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        // Mirrors and attachments are also staged on disk. Validate every
2262        // object before publishing either immutable files or derived metadata.
2263        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}