Skip to main content

objects/store/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Backend-neutral object storage abstractions and concrete implementations.
3
4use std::path::PathBuf;
5
6use crate::object::{
7    Action, ActionId, AnnotatedTag, Blob, ContentHash, OpenedTreeBody, PartialTree, State,
8    StateAttachment, StateAttachmentId, StateId, Tree, TreeEntry, TreeEntryReader,
9    TreeResumeCursor, is_redacted_tree, is_streamable_tree,
10};
11
12pub mod codec;
13#[cfg(feature = "fs")]
14mod delta_source;
15#[cfg(feature = "fs")]
16pub mod fs;
17pub mod liveness;
18#[cfg(any(test, feature = "memory-backend"))]
19pub mod memory;
20#[cfg(test)]
21mod partial_tree_tests;
22pub use heddle_pack::store::pack;
23#[cfg(feature = "fs")]
24pub mod shallow;
25#[cfg(feature = "fs")]
26mod snapshot_commit;
27pub mod source;
28pub mod store_compliance;
29#[cfg(feature = "fs")]
30pub mod writer_lease;
31
32#[cfg(feature = "fs")]
33pub use fs::{
34    DEFAULT_PACK_INSTALL_INTENT_TTL_SECS, FsRepackOperation, FsStore, PackInstallIntent,
35    PackInstallMetricsSnapshot, PackInstallPhase, PackInstallRecoverReport,
36    install_pack_bytes_journaled, pack_install_metrics_reset, pack_install_metrics_snapshot,
37    recover_pack_install_intents, recover_pack_install_intents_with_ttl,
38};
39pub use heddle_format::compression::{CompressionConfig, CompressionError, compress, decompress};
40pub use liveness::{
41    AGENT_LEASE_DURATION, Liveness, current_boot_id, process_alive, process_birth,
42    reservation_liveness_at,
43};
44#[cfg(any(test, feature = "memory-backend"))]
45pub use memory::InMemoryStore;
46pub use pack::{
47    CancellationToken as RepackCancellationToken, LoadMonitor as RepackLoadMonitor, PackBuilder,
48    PackObjectId, PackReader, PackStats, RepackContext, RepackError, RepackHandle, RepackInventory,
49    RepackOperation, RepackOutcome, RepackPolicy, RepackReason, RepackReport, RepackResourceLimits,
50    RepackSchedule, RepackScheduler, StreamingPackBuilder, SyncData,
51};
52#[cfg(feature = "fs")]
53pub use shallow::ShallowInfo;
54#[cfg(feature = "fs")]
55#[doc(hidden)]
56pub use snapshot_commit::{
57    SNAPSHOT_COMMIT_ARTIFACT_SCHEMA, SnapshotCommitArtifact, SnapshotCommitDescriptor,
58    SnapshotPackManager,
59};
60#[cfg(feature = "async-source")]
61pub use source::AsyncObjectSource;
62pub use source::ObjectSource;
63#[cfg(feature = "fs")]
64pub use writer_lease::{
65    WriterLease, WriterLeaseAuthOutcome, WriterLeaseDraft, WriterLeaseGrant,
66    WriterLeaseReserveOutcome, WriterLeaseStatus, WriterLeaseStore, checkout_writer_lock,
67    generate_writer_lease_id, generate_writer_lease_token,
68};
69
70/// A newly-authored tree plus its immediate parent, when capture already knows
71/// that relationship. Stores may use the hint for bounded HDC1 encoding; it
72/// never changes the tree's semantic content hash.
73#[derive(Clone, Debug)]
74pub struct TreeWrite {
75    pub tree: Tree,
76    pub parent: Option<ContentHash>,
77}
78
79impl TreeWrite {
80    pub fn anchor(tree: Tree) -> Self {
81        Self { tree, parent: None }
82    }
83
84    pub fn descendant(tree: Tree, parent: ContentHash) -> Self {
85        Self {
86            tree,
87            parent: Some(parent),
88        }
89    }
90}
91
92/// Read-only objects whose authoritative representation lives outside the
93/// native Heddle object directory. Git-overlay repositories use this seam to
94/// translate objects directly from `.git` without importing a second copy.
95pub trait ExternalObjectSource: Send + Sync {
96    fn get_blob(&self, hash: &ContentHash) -> Result<Option<Blob>>;
97    fn get_tree(&self, hash: &ContentHash) -> Result<Option<Tree>>;
98    fn get_state(&self, id: &StateId) -> Result<Option<State>>;
99    fn list_states(&self) -> Result<Vec<StateId>>;
100}
101
102/// Explicit cache control for benchmarks and diagnostic tools.
103///
104/// Cache invalidation is not part of durable object storage semantics. Keeping
105/// it separate prevents remote stores such as Weft's backend from having to
106/// pretend they expose process-local cache controls merely to implement
107/// [`ObjectStore`].
108pub trait ObjectCacheControl: Send + Sync {
109    /// Drop process-local decoded-object caches, if this implementation has
110    /// any. The next read should observe the implementation's cold path.
111    fn clear_recent_caches(&self);
112}
113
114pub use crate::error::{HeddleError as StoreError, HeddleError, Result};
115
116/// Sidecar records that live outside the content-addressed object graph —
117/// signed redactions and state-visibility tiers. They never ride native packs
118/// and are transferred out-of-band. Backends that do not model them can use
119/// the default methods, while native stores override the relevant operations.
120pub trait SidecarStore: Send + Sync {
121    /// Whether the store holds any redaction record for the given blob.
122    ///
123    /// Redactions live in a sidecar (`<heddle_dir>/redactions/`) that is
124    /// structurally outside the content-addressed object graph so GC
125    /// can't reach them. The wire layer needs a cheap probe to decide
126    /// whether to ship a redaction for a blob in the closure, so this
127    /// is a separate method rather than a `get_*` + null check.
128    ///
129    /// Default impl returns `Ok(false)` — stores that don't model
130    /// redactions silently report "no redactions," which is the
131    /// correct behaviour for purely in-memory or remote-shim stores.
132    fn has_redactions_for_blob(&self, _blob: &ContentHash) -> Result<bool> {
133        Ok(false)
134    }
135
136    /// Return the raw rmp-encoded `RedactionsBlob` bytes for the given
137    /// blob, or `Ok(None)` if no redaction record exists. The bytes
138    /// are byte-identical to what was written by `put_redactions_bytes_for_blob`
139    /// (or by `Repository::put_redaction`); this is the wire-transfer
140    /// payload, not a re-serialized view.
141    ///
142    /// Default impl returns `Ok(None)`.
143    fn get_redactions_bytes_for_blob(&self, _blob: &ContentHash) -> Result<Option<Vec<u8>>> {
144        Ok(None)
145    }
146
147    /// Persist the rmp-encoded `RedactionsBlob` bytes for the given
148    /// blob. Receiver-side replay calls this after signature
149    /// verification so the bytes land in the same sidecar that the
150    /// sender's `Repository::put_redaction` writes to.
151    ///
152    /// Default impl returns an "unsupported" error — stores that don't
153    /// model redactions (e.g. read-only shims) refuse rather than
154    /// silently dropping the record.
155    fn put_redactions_bytes_for_blob(&self, _blob: &ContentHash, _bytes: &[u8]) -> Result<()> {
156        Err(HeddleError::InvalidObject(
157            "this object store does not support persisting redactions".to_string(),
158        ))
159    }
160
161    /// List every blob that has at least one redaction record. Used by
162    /// the GC pin guard and by sync to enumerate redactions for the
163    /// state closure. Order is unspecified; callers that need stable
164    /// ordering should sort.
165    ///
166    /// Default impl returns `Ok(vec![])`.
167    fn list_blobs_with_redactions(&self) -> Result<Vec<ContentHash>> {
168        Ok(Vec::new())
169    }
170
171    /// Whether the store holds any state-visibility record for `state`.
172    ///
173    /// Like redactions, state-visibility records live in a sidecar outside
174    /// the content-addressed object graph and cannot ride native packs.
175    /// Sync uses this probe while enumerating a state closure so a non-public
176    /// state can advertise the sidecar that must travel out-of-pack.
177    ///
178    /// Default impl returns `Ok(false)` for stores that do not model this
179    /// sidecar.
180    fn has_state_visibility_for_state(&self, _state: &StateId) -> Result<bool> {
181        Ok(false)
182    }
183
184    /// Return the raw rmp-encoded `StateVisibilityBlob` bytes for `state`,
185    /// or `Ok(None)` if no sidecar exists. The bytes are the wire-transfer
186    /// payload for state visibility.
187    ///
188    /// Default impl returns `Ok(None)`.
189    fn get_state_visibility_bytes_for_state(&self, _state: &StateId) -> Result<Option<Vec<u8>>> {
190        Ok(None)
191    }
192
193    /// Persist raw `StateVisibilityBlob` bytes for `state`.
194    ///
195    /// Default impl returns an "unsupported" error so stores that do not
196    /// model the sidecar refuse instead of dropping it.
197    fn put_state_visibility_bytes_for_state(&self, _state: &StateId, _bytes: &[u8]) -> Result<()> {
198        Err(HeddleError::InvalidObject(
199            "this object store does not support persisting state visibility".to_string(),
200        ))
201    }
202
203    /// List every state with at least one state-visibility record.
204    ///
205    /// Default impl returns `Ok(vec![])`.
206    fn list_states_with_visibility(&self) -> Result<Vec<StateId>> {
207        Ok(Vec::new())
208    }
209}
210
211/// The result of resolving a tree hash that may be held either as the full
212/// canonical object or only as a redacted partial projection (HRT1).
213///
214/// A full tree SUPERSEDES a partial for the same hash. `Partial` is distinct
215/// from `Absent` on purpose: a partial clone that withholds an entry must be
216/// distinguishable from a repository that is missing or corrupt, so callers can
217/// present "withheld" rather than "gone".
218#[derive(Clone, Debug)]
219pub enum TreeRead {
220    /// The full canonical tree is held.
221    Full(Tree),
222    /// Only a redacted partial projection is held: its visible entries carry
223    /// their preimage and its withheld entries are marked opaque. Verifies
224    /// against the declared `State.tree` via `PartialTree::reconstruct_root`.
225    Partial(PartialTree),
226    /// Neither a full tree nor a partial projection is held for this hash.
227    Absent,
228}
229
230/// The outcome of storing a redacted partial projection under the monotone
231/// discipline enforced by [`ObjectStore::put_partial_tree`].
232#[derive(Clone, Copy, Debug, PartialEq, Eq)]
233pub enum PartialTreeWrite {
234    /// The projection was written to the partial slot.
235    Stored {
236        /// Number of withheld (redacted) leaves.
237        redacted: usize,
238        /// Number of visible leaves.
239        visible: usize,
240    },
241    /// A full canonical tree for this hash is already held, so the projection
242    /// was dropped rather than stored — a full tree supersedes a partial, and a
243    /// partial never overwrites a full (monotone).
244    SupersededByFull,
245}
246
247/// Trait for object storage backends.
248///
249/// Sidecars remain a separate implementation seam, but every object store
250/// exposes that seam. This preserves object-safe `dyn ObjectStore` consumers
251/// such as Weft's local filesystem backend without coupling its S3 backend to
252/// the native store implementation.
253pub trait ObjectStore: SidecarStore + Send + Sync {
254    /// Resolve a recorded tree hash without inventing an empty baseline.
255    fn require_tree(&self, hash: &ContentHash) -> Result<Tree> {
256        self.get_tree(hash)?
257            .ok_or_else(|| HeddleError::MissingObject {
258                object_type: "tree".to_string(),
259                id: hash.to_hex(),
260            })
261    }
262
263    /// Resolve a recorded blob hash without inventing content.
264    fn require_blob(&self, hash: &ContentHash) -> Result<Blob> {
265        self.get_blob(hash)?
266            .ok_or_else(|| HeddleError::MissingObject {
267                object_type: "blob".to_string(),
268                id: hash.to_hex(),
269            })
270    }
271
272    fn get_annotated_tag(&self, _hash: &ContentHash) -> Result<Option<AnnotatedTag>> {
273        Ok(None)
274    }
275    fn put_annotated_tag(&self, _tag: &AnnotatedTag) -> Result<ContentHash> {
276        Err(HeddleError::InvalidObject(
277            "object store does not support annotated tags".to_string(),
278        ))
279    }
280    fn list_annotated_tags(&self) -> Result<Vec<ContentHash>> {
281        Ok(Vec::new())
282    }
283    fn get_blob(&self, hash: &ContentHash) -> Result<Option<Blob>>;
284    fn put_blob(&self, blob: &Blob) -> Result<ContentHash>;
285
286    /// Zero-copy variant of `get_blob`. Returns a [`bytes::Bytes`]
287    /// view of the blob's content, which for `FsStore` reads is a
288    /// slice into the pack file's mmap when the entry is non-delta
289    /// and uncompressed — no allocation, no memcpy.
290    ///
291    /// Default impl wraps `get_blob`'s `Vec<u8>` in a `Bytes` (one
292    /// Arc allocation, no body copy) so backends without a native
293    /// fast path still satisfy the contract. The mount's hot read
294    /// path goes through this method instead of `get_blob` so the
295    /// pack-mmap fast path lights up automatically.
296    fn get_blob_bytes(&self, hash: &ContentHash) -> Result<Option<bytes::Bytes>> {
297        Ok(self
298            .get_blob(hash)?
299            .map(|blob| bytes::Bytes::from(blob.into_content())))
300    }
301
302    /// Return the *uncompressed* byte length of the blob identified by
303    /// `hash`, or `Ok(None)` when the blob is not in the store.
304    ///
305    /// The contract is "size without paying for content": backends are
306    /// expected to honour this with a header read or index lookup
307    /// rather than a full decompression. This is the hot path for
308    /// directory listings (`ls -l` over a thread mount) where loading
309    /// every blob just to learn its size would dominate.
310    ///
311    /// The default implementation falls back to `get_blob` so backends
312    /// without a cheap size accessor still satisfy the contract; native
313    /// stores (`FsStore`, `InMemoryStore`) override this with a
314    /// header- or hashmap-only path.
315    fn blob_size(&self, hash: &ContentHash) -> Result<Option<u64>> {
316        Ok(self.get_blob(hash)?.map(|blob| blob.content().len() as u64))
317    }
318
319    /// Filesystem path of the loose blob whose on-disk bytes are
320    /// byte-identical to the blob's *uncompressed* content, suitable
321    /// for `hard_link`/`clonefile` materialization without going
322    /// through `get_blob`.
323    ///
324    /// Returns `None` when the blob is missing, is only available via
325    /// a packfile, is stored compressed (the on-disk bytes wouldn't
326    /// match what a worktree consumer needs to read), or the backend
327    /// doesn't expose stable filesystem paths (e.g. `InMemoryStore`). The
328    /// default impl returns `None` so non-`FsStore` backends silently fall
329    /// through to the bytes path.
330    fn loose_blob_path(&self, _hash: &ContentHash) -> Option<PathBuf> {
331        None
332    }
333
334    /// Ensure the blob identified by `hash` is materialized as an
335    /// uncompressed loose file at the canonical loose path so that
336    /// `loose_blob_path` returns `Some(path)` on a subsequent call.
337    ///
338    /// This is the "warm canonical store" path that lets the
339    /// hardlink-first materializer keep its 5–10× wall-clock and
340    /// storage-allocation wins after `pack_objects + prune_loose_objects`
341    /// has moved everything into a packfile. Without this, the lazy
342    /// hardlink path silently degrades to `fs::write(decompressed)` on
343    /// every materialize, because `loose_blob_path` returns `None` for
344    /// pack-only and compressed-loose blobs.
345    ///
346    /// Cost-amortization: the first promotion of a blob pays
347    /// `decompress + atomic write`. Every subsequent materialize of
348    /// the same blob — into the same worktree on `goto`, or into a
349    /// sibling worktree on `delegate` — is a single `link(2)`. Net
350    /// win for any N > 1 materializations; break-even at N == 1.
351    ///
352    /// Pack invariants are preserved: this method does not remove the
353    /// pack-resident copy. The blob lives in both pack and loose-
354    /// uncompressed until the next `prune_loose_objects` cycle, at
355    /// which point the loose mirror is discarded and a future
356    /// materialize re-promotes on demand.
357    ///
358    /// Idempotent: a blob that's already loose-and-uncompressed is a
359    /// no-op fast path. A blob that's loose-but-compressed is
360    /// rewritten in place (atomically) with the uncompressed bytes.
361    /// A blob that's pack-resident is decompressed out of the pack
362    /// and written loose without touching the pack.
363    ///
364    /// Returns `Ok(true)` when the call did real work (a write
365    /// happened), `Ok(false)` when it was a no-op (blob was already
366    /// loose+uncompressed), and `Err` when the blob isn't in the
367    /// store at all. The default impl returns `Ok(false)` for
368    /// backends that don't expose loose paths (`InMemoryStore`), since the
369    /// hardlink path is fundamentally inapplicable there.
370    fn promote_to_loose_uncompressed(&self, _hash: &ContentHash) -> Result<bool> {
371        Ok(false)
372    }
373
374    fn put_blob_with_hash(&self, blob: &Blob, hash: ContentHash) -> Result<ContentHash> {
375        if blob.hash() != hash {
376            return Err(HeddleError::InvalidObject("blob hash mismatch".to_string()));
377        }
378        self.put_blob(blob)
379    }
380
381    fn has_blob(&self, hash: &ContentHash) -> Result<bool>;
382    /// Return whether the blob is owned by this store, excluding any configured
383    /// read-through source. Snapshot builders use this to ensure a new native
384    /// state owns its complete object closure.
385    fn has_blob_locally(&self, hash: &ContentHash) -> Result<bool> {
386        self.has_blob(hash)
387    }
388    fn get_tree(&self, hash: &ContentHash) -> Result<Option<Tree>>;
389    /// Resolve one named tree entry. Pack-capable stores override this so a
390    /// lookup can use a restartable packed record instead of materializing the
391    /// complete tree.
392    fn get_tree_entry(&self, hash: &ContentHash, name: &str) -> Result<Option<TreeEntry>> {
393        Ok(self
394            .get_tree(hash)?
395            .and_then(|tree| tree.get(name).cloned()))
396    }
397    fn put_tree(&self, tree: &Tree) -> Result<ContentHash>;
398    fn has_tree(&self, hash: &ContentHash) -> Result<bool>;
399    /// Return whether the tree is owned by this store, excluding any configured
400    /// read-through source.
401    fn has_tree_locally(&self, hash: &ContentHash) -> Result<bool> {
402        self.has_tree(hash)
403    }
404
405    // ── Redacted partial-tree projections (HRT1) ──────────────────────
406    //
407    // A partial projection lives in a slot keyed by the canonical tree hash it
408    // projects, DISTINCT from the full-tree object slot. The full-tree methods
409    // above (`get_tree`/`has_tree`) never see or return a projection; the
410    // methods below are the only way in and out of the partial slot. Backends
411    // that do not model projections inherit the default impls (report "none" /
412    // refuse the write), mirroring the sidecar seam.
413
414    /// Whether the store holds a redacted partial projection keyed by `hash`.
415    /// Independent of [`ObjectStore::has_tree`], which reports the full object.
416    fn has_partial_tree(&self, _hash: &ContentHash) -> Result<bool> {
417        Ok(false)
418    }
419
420    /// Raw HRT1 partial-projection bytes for `hash`, or `Ok(None)`. These are
421    /// byte-identical to what [`ObjectStore::put_partial_tree_bytes`] wrote —
422    /// the wire-transfer payload, not a re-serialized view.
423    fn get_partial_tree_bytes(&self, _hash: &ContentHash) -> Result<Option<Vec<u8>>> {
424        Ok(None)
425    }
426
427    /// Persist raw HRT1 partial-projection bytes keyed by `hash`. This is the
428    /// low-level slot; callers should prefer [`ObjectStore::put_partial_tree`],
429    /// which verifies the projection and enforces the monotone discipline.
430    fn put_partial_tree_bytes(&self, _hash: &ContentHash, _bytes: &[u8]) -> Result<()> {
431        Err(HeddleError::InvalidObject(
432            "this object store does not support partial tree projections".to_string(),
433        ))
434    }
435
436    /// List every tree hash for which a partial projection is held. Order is
437    /// unspecified; callers that need stable ordering should sort.
438    fn list_partial_trees(&self) -> Result<Vec<ContentHash>> {
439        Ok(Vec::new())
440    }
441
442    /// Drop any partial projection held for `hash`. Idempotent — a no-op when
443    /// none is held. Used to reclaim the partial slot once the full tree is
444    /// backfilled.
445    fn remove_partial_tree(&self, _hash: &ContentHash) -> Result<()> {
446        Ok(())
447    }
448
449    /// Store a redacted partial projection under monotone discipline.
450    ///
451    /// `hrt1` must be an HRT1 body whose visible preimages + withheld leaf
452    /// hashes reconstruct `expected` (verified through
453    /// [`codec::decode_partial_tree`], i.e. Leg 1's `reconstruct_root`).
454    /// Monotone: a partial NEVER overwrites a full tree. If the full canonical
455    /// tree for `expected` is already held, the projection is dropped and
456    /// [`PartialTreeWrite::SupersededByFull`] is returned; otherwise the raw
457    /// bytes land in the partial slot.
458    fn put_partial_tree(&self, expected: &ContentHash, hrt1: &[u8]) -> Result<PartialTreeWrite> {
459        let partial = codec::decode_partial_tree(hrt1, *expected)?;
460        if self.has_tree(expected)? {
461            return Ok(PartialTreeWrite::SupersededByFull);
462        }
463        self.put_partial_tree_bytes(expected, hrt1)?;
464        let redacted = partial.redacted_count();
465        Ok(PartialTreeWrite::Stored {
466            redacted,
467            visible: partial.leaves().len() - redacted,
468        })
469    }
470
471    /// Read the tree for `hash` as the full canonical tree, a redacted partial
472    /// projection, or absent.
473    ///
474    /// A full tree SUPERSEDES a partial: when the full object is held it is
475    /// returned even if a partial projection also exists for the same hash.
476    /// [`ObjectStore::get_tree`] by contrast returns the full tree or `None`
477    /// and NEVER the partial, so a caller that needs the complete tree cannot
478    /// be silently handed a projection.
479    fn read_tree(&self, hash: &ContentHash) -> Result<TreeRead> {
480        if let Some(full) = self.get_tree(hash)? {
481            return Ok(TreeRead::Full(full));
482        }
483        match self.get_partial_tree_bytes(hash)? {
484            Some(bytes) => Ok(TreeRead::Partial(codec::decode_partial_tree(
485                &bytes, *hash,
486            )?)),
487            None => Ok(TreeRead::Absent),
488        }
489    }
490    /// Open a streamable HTR4 tree body. Store backends use sequential
491    /// verify: resume at ordinal > 0 is refused until the bytes are hashed.
492    fn open_tree(
493        &self,
494        tree_id: &ContentHash,
495        cursor: Option<&TreeResumeCursor>,
496    ) -> Result<Option<TreeEntryReader<OpenedTreeBody>>> {
497        let Some(body) = self.get_tree_serialized(tree_id)? else {
498            return Ok(None);
499        };
500        let body = if is_streamable_tree(&body) {
501            body
502        } else {
503            let tree = self
504                .get_tree(tree_id)?
505                .ok_or_else(|| HeddleError::NotFound(format!("tree {tree_id}")))?;
506            tree.encode_lean()?
507        };
508        Ok(Some(TreeEntryReader::open(
509            OpenedTreeBody::Bytes(crate::object::BytesTreeSource::sequential_verify(body)),
510            *tree_id,
511            cursor,
512        )?))
513    }
514    fn get_state(&self, id: &StateId) -> Result<Option<State>>;
515    fn put_state(&self, state: &State) -> Result<()>;
516    fn has_state(&self, id: &StateId) -> Result<bool>;
517    fn list_states(&self) -> Result<Vec<StateId>>;
518    fn get_state_attachment(
519        &self,
520        _state: &StateId,
521        _id: &StateAttachmentId,
522    ) -> Result<Option<StateAttachment>> {
523        Ok(None)
524    }
525    fn put_state_attachment(&self, _attachment: &StateAttachment) -> Result<StateAttachmentId> {
526        Err(HeddleError::InvalidObject(
527            "object store does not support state attachments".to_string(),
528        ))
529    }
530    fn list_state_attachments(&self, _state: &StateId) -> Result<Vec<StateAttachment>> {
531        Ok(Vec::new())
532    }
533    fn get_action(&self, id: &ActionId) -> Result<Option<Action>>;
534    fn put_action(&self, action: &mut Action) -> Result<ActionId>;
535    fn list_actions(&self) -> Result<Vec<ActionId>>;
536    fn list_blobs(&self) -> Result<Vec<ContentHash>>;
537    fn list_trees(&self) -> Result<Vec<ContentHash>>;
538
539    fn put_blob_bytes_with_hash(&self, data: &[u8], hash: ContentHash) -> Result<ContentHash> {
540        self.put_blob_with_hash(&Blob::from_slice(data), hash)
541    }
542
543    /// Return the stored tree body for `hash`, without requiring HTR4.
544    ///
545    /// This is a migration seam, not a runtime compatibility reader: callers
546    /// that need current tree semantics should use [`ObjectStore::get_tree`].
547    /// Loose and packed backends must return the raw stored bytes so one-shot
548    /// migrations can canonicalize older encodings without a current-decoder
549    /// gate. Default impls that only have `get_tree` re-encode current trees.
550    fn get_tree_serialized(&self, hash: &ContentHash) -> Result<Option<Vec<u8>>> {
551        self.get_tree(hash)?
552            .map(|tree| tree.encode_canonical().map_err(HeddleError::from))
553            .transpose()
554    }
555
556    fn put_tree_serialized(&self, data: &[u8], hash: ContentHash) -> Result<ContentHash> {
557        // An HRT1 redacted projection is not a full tree: route it to the
558        // partial slot (monotone) instead of hard-refusing in the full-tree
559        // decoder.
560        if is_redacted_tree(data) {
561            self.put_partial_tree(&hash, data)?;
562            return Ok(hash);
563        }
564        let tree = codec::decode_tree_serialized_with_key(data, hash, None)?;
565        self.put_tree(&tree)
566    }
567
568    fn put_state_serialized(&self, data: &[u8], id: StateId) -> Result<()> {
569        let state = State::decode_current_msgpack(data)?;
570        if !state.accepts_stored_id(&id) {
571            return Err(HeddleError::InvalidObject(format!(
572                "state id mismatch: expected {id}, computed {}",
573                state.id()
574            )));
575        }
576        self.put_state(&state)
577    }
578
579    fn put_action_serialized(&self, data: &[u8], id: ActionId) -> Result<()> {
580        let mut action: Action = rmp_serde::from_slice(data)?;
581        let found_id = action.compute_id();
582        if found_id != id {
583            return Err(HeddleError::InvalidObject(format!(
584                "action id mismatch: expected {}, found {}",
585                id, found_id
586            )));
587        }
588        let stored_id = self.put_action(&mut action)?;
589        if stored_id != id {
590            return Err(HeddleError::InvalidObject(format!(
591                "action id mismatch after write: expected {}, found {}",
592                id, stored_id
593            )));
594        }
595        Ok(())
596    }
597
598    fn get_pack_object(
599        &self,
600        id: &pack::PackObjectId,
601    ) -> Result<Option<(pack::ObjectType, Vec<u8>)>> {
602        match id {
603            pack::PackObjectId::AnnotatedTag(hash) => Ok(self
604                .get_annotated_tag(hash)?
605                .map(|tag| (pack::ObjectType::AnnotatedTag, tag.encode_current_msgpack()))),
606            pack::PackObjectId::Hash(hash) => {
607                if let Some(blob) = self.get_blob(hash)? {
608                    return Ok(Some((pack::ObjectType::Blob, blob.content().to_vec())));
609                }
610                if let Some(tree) = self.get_tree(hash)? {
611                    return Ok(Some((pack::ObjectType::Tree, tree.encode_canonical()?)));
612                }
613                if let Some(action) = self.get_action(&ActionId::from_hash(*hash))? {
614                    return Ok(Some((
615                        pack::ObjectType::Action,
616                        rmp_serde::to_vec_named(&action)?,
617                    )));
618                }
619                Ok(None)
620            }
621            pack::PackObjectId::StateId(change_id) => {
622                if let Some(state) = self.get_state(change_id)? {
623                    Ok(Some((
624                        pack::ObjectType::State,
625                        state.encode_current_msgpack()?,
626                    )))
627                } else {
628                    Ok(None)
629                }
630            }
631        }
632    }
633
634    /// Bulk-write a batch of blobs as a single durable unit. The default
635    /// implementation falls back to per-blob writes; backends that
636    /// support packfiles (i.e. `FsStore`) override this to install one
637    /// packfile + index — two fsyncs total instead of N. Used by the
638    /// snapshot hot path so writing 1000 small files takes ~one fsync,
639    /// not 1000.
640    ///
641    /// Blobs already present in the store are skipped on the way in
642    /// (the caller would otherwise duplicate them in the pack).
643    fn put_blobs_packed(&self, blobs: Vec<(ContentHash, Vec<u8>)>) -> Result<()> {
644        for (hash, data) in blobs {
645            if !self.has_blob(&hash)? {
646                self.put_blob_bytes_with_hash(&data, hash)?;
647            }
648        }
649        Ok(())
650    }
651
652    /// Durably install a snapshot's newly-authored immutable object closure as
653    /// one storage batch. Pack-capable backends override this to share one pack
654    /// installation across blobs, the root tree, and the state; other backends
655    /// preserve the same ordering through their ordinary object methods.
656    fn put_snapshot_objects_packed(
657        &self,
658        blobs: Vec<(ContentHash, Vec<u8>)>,
659        tree: &Tree,
660        state: &State,
661    ) -> Result<()> {
662        self.put_blobs_packed(blobs)?;
663        self.put_tree(tree)?;
664        self.put_state(state)
665    }
666
667    /// Snapshot closure variant that also durably installs immutable authored
668    /// attachments. The separate method preserves the existing backend API;
669    /// pack-capable stores override it to share the snapshot pack barrier.
670    fn put_snapshot_objects_and_attachments_packed(
671        &self,
672        blobs: Vec<(ContentHash, Vec<u8>)>,
673        tree: &Tree,
674        state: &State,
675        attachments: Vec<StateAttachment>,
676    ) -> Result<()> {
677        self.put_snapshot_objects_packed(blobs, tree, state)?;
678        for attachment in attachments {
679            self.put_state_attachment(&attachment)?;
680        }
681        Ok(())
682    }
683    fn install_pack(&self, pack_data: &[u8], index_data: &[u8]) -> Result<Vec<pack::PackObjectId>> {
684        let reader = pack::PackReader::from_slice_in_memory(pack_data, index_data)?;
685        let ids = reader.list_ids()?;
686        for id in &ids {
687            let Some((obj_type, data)) = reader.get_object(id)? else {
688                continue;
689            };
690            install_pack_entry(self, *id, obj_type, &data)?;
691        }
692        Ok(ids)
693    }
694
695    /// Install staged files with a bounded reader. The returned disk inventory
696    /// owns its identities independently of the consumed staging files.
697    fn install_pack_streaming(
698        &self,
699        pack_path: &std::path::Path,
700        index_path: &std::path::Path,
701    ) -> Result<pack::PackInventory> {
702        let scratch = index_path
703            .parent()
704            .ok_or_else(|| StoreError::InvalidObject("pack index path has no parent".into()))?;
705        let inventory = pack::PackInventory::copy_from_index(index_path, scratch)?;
706        let reader = pack::PackReader::open(pack_path, index_path, scratch)?;
707        reader.visit_objects(|id, kind, data| install_pack_entry(self, id, kind, data))?;
708        drop(reader);
709        std::fs::remove_file(pack_path)?;
710        std::fs::remove_file(index_path)?;
711        Ok(inventory)
712    }
713
714    fn begin_snapshot_write_batch(&self) -> Result<()> {
715        Ok(())
716    }
717
718    fn flush_snapshot_write_batch(&self) -> Result<()> {
719        Ok(())
720    }
721
722    fn abort_snapshot_write_batch(&self) {}
723}
724
725fn install_pack_entry(
726    store: &(impl ObjectStore + ?Sized),
727    id: pack::PackObjectId,
728    obj_type: pack::ObjectType,
729    data: &[u8],
730) -> Result<()> {
731    match (&id, obj_type) {
732        (pack::PackObjectId::Hash(hash), pack::ObjectType::Blob) => {
733            store.put_blob_bytes_with_hash(data, *hash)?;
734        }
735        (pack::PackObjectId::AnnotatedTag(hash), pack::ObjectType::AnnotatedTag) => {
736            let tag = AnnotatedTag::decode_current_msgpack(data)
737                .map_err(|error| HeddleError::InvalidObject(error.to_string()))?;
738            if tag.hash() != *hash {
739                return Err(HeddleError::InvalidObject(
740                    "annotated tag hash mismatch".to_string(),
741                ));
742            }
743            store.put_annotated_tag(&tag)?;
744        }
745        (pack::PackObjectId::Hash(hash), pack::ObjectType::Tree) => {
746            store.put_tree_serialized(data, *hash)?;
747        }
748        (pack::PackObjectId::Hash(hash), pack::ObjectType::Action) => {
749            store.put_action_serialized(data, ActionId::from_hash(*hash))?;
750        }
751        (pack::PackObjectId::StateId(change_id), pack::ObjectType::State) => {
752            store.put_state_serialized(data, *change_id)?;
753        }
754        (_, pack::ObjectType::TimelineOperation) => {
755            return Err(HeddleError::InvalidObject(
756                "timeline operations belong in the timeline pack store".to_string(),
757            ));
758        }
759        _ => {
760            return Err(HeddleError::InvalidObject(format!(
761                "unsupported native pack object: {:?} {:?}",
762                id, obj_type
763            )));
764        }
765    }
766    Ok(())
767}