1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
//! Container-backed object store (RFC 102 Stage 3). Public API (`FileObjectStore`, `ObjectReader`,
//! `ObjectWriter`) is unchanged from the loose-file implementation this replaces -- every other call
//! site in the workspace uses only that trait interface, so none of them needed to change. Only the
//! internals moved: reads and writes now go through `index.rs`'s lookup/write-protocol functions,
//! which target `container.rs`'s per-type container files instead of one file per object.
use prikk_error::{PrikkError, Result};
use prikk_object::{ObjectEnvelope, ObjectId, ObjectType};
use crate::index::{
self, IndexEntry, WriteDecision, append_object_to_container, decide_write_outcome,
lookup_object_location, read_object_envelope_at,
};
use crate::layout::RepositoryLayout;
/// Read-only object access boundary.
pub trait ObjectReader {
/// Read an object by ID.
fn read_object(&self, id: ObjectId) -> Result<Option<ObjectEnvelope>>;
/// Read and require a specific object type. Default-implemented in terms of `read_object` alone
/// (RFC 111 §6.1: every implementor -- `FileObjectStore`, `ObjectReadSnapshot`,
/// `ObjectWriteSession`, `MemoryObjectStore` -- gets this for free, and any function generic over
/// `impl ObjectReader` can call it without depending on a concrete type). One body, not one per
/// implementor (Stage 1 review v1 §3): every concrete type's own `read_typed` delegates here too.
fn read_typed(&self, id: ObjectId, object_type: ObjectType) -> Result<Option<ObjectEnvelope>> {
let Some(envelope) = self.read_object(id)? else {
return Ok(None);
};
if envelope.object_type != object_type {
return Err(PrikkError::ObjectTypeMismatch {
expected: object_type.to_string(),
actual: envelope.object_type.to_string(),
});
}
Ok(Some(envelope))
}
}
/// Write object boundary.
pub trait ObjectWriter {
/// Write an object envelope after validation.
fn write_object(&mut self, envelope: &ObjectEnvelope) -> Result<ObjectId>;
}
/// File-backed object store.
#[derive(Debug, Clone)]
pub struct FileObjectStore {
layout: RepositoryLayout,
}
impl FileObjectStore {
/// Create a file object store for a repository layout.
#[must_use]
pub fn new(layout: RepositoryLayout) -> Self {
Self { layout }
}
/// Return the repository layout.
#[must_use]
pub fn layout(&self) -> &RepositoryLayout {
&self.layout
}
/// Return true if an object with this id and type is indexed.
#[must_use]
pub fn contains_object(&self, object_type: ObjectType, id: ObjectId) -> bool {
if object_type == ObjectType::RefUpdate {
return false;
}
matches!(
lookup_object_location(&self.layout, id),
Ok(Some(entry)) if entry.object_type == object_type
)
}
}
impl ObjectReader for FileObjectStore {
fn read_object(&self, id: ObjectId) -> Result<Option<ObjectEnvelope>> {
let Some(entry) = lookup_object_location(&self.layout, id)? else {
return Ok(None);
};
read_object_at_entry(&self.layout, &entry, id)
}
}
impl ObjectWriter for FileObjectStore {
fn write_object(&mut self, envelope: &ObjectEnvelope) -> Result<ObjectId> {
if envelope.object_type == ObjectType::RefUpdate {
return Err(PrikkError::UnsupportedObjectType(
"RefUpdate is stored inline in ref logs for v1".to_string(),
));
}
self.layout.validate_format()?;
crate::format::validate_object_envelope(self.layout.format(), envelope)?;
// The write protocol (design §5, handoff §3) lives in `index.rs`, not here: append the
// object record to its container and make it durable, then and only then append the index
// entry. Stated at that call site too, not only here. The idempotency decision itself (RFC
// 111 §6.1 addendum, C2) is `index::decide_write_outcome`, shared verbatim with
// `ObjectWriteSession` below -- only where its `existing` lookup comes from differs: this
// type always re-decodes the whole index (unchanged cost, a safe default for any call site
// not migrated to a snapshot-backed type).
let existing = lookup_object_location(&self.layout, envelope.object_id())?;
match decide_write_outcome(
&self.layout,
envelope.object_type,
envelope,
existing.as_ref(),
)? {
WriteDecision::AlreadyPresent(id) => Ok(id),
WriteDecision::New => {
append_object_to_container(&self.layout, envelope.object_type, envelope)
.map(|entry| entry.object_id)
}
}
}
}
/// Read validation shared by every reader below (`FileObjectStore`, `ObjectReadSnapshot`,
/// `ObjectWriteSession`): the index is trusted for *location*, but the bytes found there are always
/// checked against the id actually asked for by recomputing it from the decoded content -- free,
/// since decoding already happened. A mismatch is reported, never silently accepted and never a
/// fallback to scanning ("one seek", design §12/§10.3).
fn read_object_at_entry(
layout: &RepositoryLayout,
entry: &IndexEntry,
id: ObjectId,
) -> Result<Option<ObjectEnvelope>> {
let envelope = read_object_envelope_at(layout, entry)?;
let computed = envelope.object_id();
if computed != id {
return Err(PrikkError::Integrity(format!(
"index entry for {id} resolves to an envelope with computed id {computed}"
)));
}
if envelope.object_type != entry.object_type {
return Err(PrikkError::Integrity(format!(
"index entry for {id} names type {}, envelope decoded as {}",
entry.object_type, envelope.object_type
)));
}
crate::format::validate_read_schema(layout.format(), &envelope)?;
Ok(Some(envelope))
}
/// A decoded object-index snapshot, taken once. Backs both `ObjectReadSnapshot` and
/// `ObjectWriteSession` (RFC 111 §6.1) -- the read logic (lookup, then decode at a known offset) is
/// identical between them, so it exists here once rather than twice. Has no public API of its own;
/// both public types below wrap it.
struct IndexSnapshot {
entries: Vec<IndexEntry>,
/// The object index's byte length as of the last time `entries` was known-current --
/// `bytes.len() - trailing_partial_bytes` from whichever decode produced `entries`, never the
/// raw stat size (RFC 111 §6.1 addendum §3.2: a torn trailing write must not be counted as
/// decoded).
known_length: u64,
}
impl IndexSnapshot {
fn open(layout: &RepositoryLayout) -> Result<Self> {
let (replay, known_length) = index::replay_index_with_extent(layout)?;
if replay.has_item_failure() {
return Err(PrikkError::Integrity(
"object index has a damaged entry; run doctor before reading".to_string(),
));
}
Ok(Self {
entries: replay.entries,
known_length,
})
}
/// Same last-entry-wins semantics `lookup_object_location` already has, preserved verbatim.
fn lookup(&self, id: ObjectId) -> Option<&IndexEntry> {
self.entries
.iter()
.rev()
.find(|entry| entry.object_id == id)
}
/// Re-stat only; decode only if the stat disagrees with what this snapshot already knows. Every
/// write decision calls this first (RFC 111 §6.1 addendum, C1) -- it is what makes a stale
/// idempotency decision structurally impossible regardless of *what* wrote the new bytes: a
/// nested unmediated writer in the same process (`refs/publication.rs`'s current shape), a call
/// site not yet migrated to a snapshot-backed type, or a genuinely separate process. All three
/// grow the index file, and this catches every one the same way, because it checks the one fact
/// that is true regardless of cause. The common case -- nothing else wrote -- costs one stat, no
/// decode. `ObjectReadSnapshot` never calls this: a reader's staleness is already accepted and
/// bounded (RFC 111 Q3/Q4), so charging every read a stat here would buy nothing.
fn ensure_current(&mut self, layout: &RepositoryLayout) -> Result<()> {
let relative = layout.repository_relative(&layout.container_index_path())?;
let current_length =
crate::fsutil::stat_file_state_if_exists(layout.repository_mutation_root(), &relative)?
.map_or(0, |stat| stat.size);
if current_length == self.known_length {
return Ok(());
}
if current_length < self.known_length {
// The object index is append-only and must never shrink (it is not one of the four
// compactable containers). A shorter file than this snapshot last knew means either the
// file was rebuilt out from under an open session or something is badly wrong -- fail
// closed rather than decode from an offset past the new end (RFC 111 §6.1 addendum §3.1).
return Err(PrikkError::Integrity(format!(
"object index shrank from {} to {current_length} bytes since it was last read; \
the object index is append-only and must never shrink -- run doctor",
self.known_length
)));
}
let (tail, new_extent) = index::replay_index_tail_with_extent(layout, self.known_length)?;
if tail.has_item_failure() {
return Err(PrikkError::Integrity(
"object index has a damaged entry; run doctor before reading".to_string(),
));
}
self.entries.extend(tail.entries);
self.known_length = new_extent;
Ok(())
}
}
/// Read-only object access for one operation's lifetime (RFC 111 §6.1). Takes one decoded index
/// snapshot at construction and never re-decodes -- correct because a reader never writes, so it can
/// never observe its *own* write as missing the way a writer holding a stale snapshot could (RFC 111
/// Q3). A snapshot taken here may miss an object a concurrent writer appends after construction; that
/// is `verify`'s own already-documented point-in-time semantics, unchanged by this type (RFC 111 Q4).
pub struct ObjectReadSnapshot {
layout: RepositoryLayout,
snapshot: IndexSnapshot,
}
impl ObjectReadSnapshot {
/// Open a read-only snapshot of `layout`'s object index, decoding it exactly once.
pub fn open(layout: &RepositoryLayout) -> Result<Self> {
Ok(Self {
layout: layout.clone(),
snapshot: IndexSnapshot::open(layout)?,
})
}
/// Return true if an object with this id and type is indexed, as of when this snapshot was
/// taken.
#[must_use]
pub fn contains_object(&self, object_type: ObjectType, id: ObjectId) -> bool {
if object_type == ObjectType::RefUpdate {
return false;
}
matches!(self.snapshot.lookup(id), Some(entry) if entry.object_type == object_type)
}
}
impl ObjectReader for ObjectReadSnapshot {
fn read_object(&self, id: ObjectId) -> Result<Option<ObjectEnvelope>> {
let Some(entry) = self.snapshot.lookup(id) else {
return Ok(None);
};
read_object_at_entry(&self.layout, entry, id)
}
}
/// Read-write object access for one writing operation's lifetime (RFC 111 §6.1). Holds the same kind
/// of in-memory index snapshot `ObjectReadSnapshot` does, but every write decision first calls
/// `IndexSnapshot::ensure_current` (see its own doc), and every successful write calls it again
/// afterward instead of computing what changed itself (Stage 1 review v1, B1) -- `ensure_current`'s
/// own tail-decode is the only thing ever allowed to grow `entries`/`known_length`, whether what it
/// finds is this session's own write, a concurrent one, or both.
pub struct ObjectWriteSession {
layout: RepositoryLayout,
snapshot: IndexSnapshot,
}
impl ObjectWriteSession {
/// Open a read-write session over `layout`'s object index, decoding it exactly once.
pub fn open(layout: &RepositoryLayout) -> Result<Self> {
Ok(Self {
layout: layout.clone(),
snapshot: IndexSnapshot::open(layout)?,
})
}
/// Return true if an object with this id and type is indexed, refreshing the snapshot first if
/// something else has grown the index since it was last known-current.
pub fn contains_object(&mut self, object_type: ObjectType, id: ObjectId) -> Result<bool> {
if object_type == ObjectType::RefUpdate {
return Ok(false);
}
self.snapshot.ensure_current(&self.layout)?;
Ok(matches!(self.snapshot.lookup(id), Some(entry) if entry.object_type == object_type))
}
}
impl ObjectReader for ObjectWriteSession {
fn read_object(&self, id: ObjectId) -> Result<Option<ObjectEnvelope>> {
let Some(entry) = self.snapshot.lookup(id) else {
return Ok(None);
};
read_object_at_entry(&self.layout, entry, id)
}
}
impl ObjectWriter for ObjectWriteSession {
fn write_object(&mut self, envelope: &ObjectEnvelope) -> Result<ObjectId> {
if envelope.object_type == ObjectType::RefUpdate {
return Err(PrikkError::UnsupportedObjectType(
"RefUpdate is stored inline in ref logs for v1".to_string(),
));
}
self.layout.validate_format()?;
crate::format::validate_object_envelope(self.layout.format(), envelope)?;
self.snapshot.ensure_current(&self.layout)?;
let existing = self.snapshot.lookup(envelope.object_id());
match decide_write_outcome(&self.layout, envelope.object_type, envelope, existing)? {
WriteDecision::AlreadyPresent(id) => Ok(id),
WriteDecision::New => {
let object_id = envelope.object_id();
append_object_to_container(&self.layout, envelope.object_type, envelope)?;
// Do not trust the append's own return value for what changed (RFC 111 Stage 1
// review v1, B1): re-derive by re-checking freshness through the same primitive
// every other read decision uses. `ensure_current`'s tail-decode is frame-aligned
// by construction and correctly picks up this write, a concurrent writer's, or
// both -- a stat taken immediately after only this call's own append cannot tell
// those apart.
self.snapshot.ensure_current(&self.layout)?;
Ok(object_id)
}
}
}
}
// DC-97 correction of the comment this replaced: the Linux/macOS-only reasoning was true when
// written (DC-71/DC-81, before DC-87 made Windows a mutating platform) and nobody revisited it once
// Windows mutation shipped -- found only by DC-97's own G5 investigation, back when this module's
// now-deleted `tests::immutable` still made the claimed Windows evidence for G5. `publish_immutable`
// and its tests are gone entirely as of DC-98 (G5 retired, zero production callers). What remains
// here is gated the same way regardless: `RepositoryLayout::init` and real repository mutation are
// not Linux/macOS-only, so what is still unix-only inside this module (failpoints, symlinks, FIFOs)
// is gated per-test/per-file instead of by one blanket gate.
#[cfg(test)]
impl ObjectWriteSession {
/// The session's own current view of the object index's byte extent -- exposed only so tests can
/// assert it lands on the true file length rather than a value this type accumulated itself (RFC
/// 111 §6.1 addendum §3.3, the load-bearing assertion the design review named explicitly).
pub(crate) fn known_index_length_for_test(&self) -> u64 {
self.snapshot.known_length
}
/// The session's own current view of how many index entries it holds -- exposed only for tests.
pub(crate) fn entry_count_for_test(&self) -> usize {
self.snapshot.entries.len()
}
}
#[cfg(all(
test,
any(target_os = "linux", target_os = "macos", target_os = "windows")
))]
mod tests;