Skip to main content

sim_lib_journal/
host.rs

1use crate::native_codec::*;
2use crate::{
3    Admission, JournalBackend, JournalEntry, JournalError, JournalHead, JournalObject, Lease,
4    StoredDatumRef, StoredState,
5};
6use sha2::{Digest, Sha256};
7use sim_kernel::{ContentId, Datum};
8use sim_storage_port::{HostDirErrorKind, HostDirPort, NeverCancel};
9use std::{
10    collections::{BTreeMap, BTreeSet},
11    sync::{
12        Arc,
13        atomic::{AtomicU64, Ordering},
14    },
15    time::{SystemTime, UNIX_EPOCH},
16};
17
18type FailHook = Arc<dyn Fn(Failpoint) -> bool + Send + Sync>;
19static NAMESPACE_COUNTER: AtomicU64 = AtomicU64::new(0);
20
21/// Crash-durable journal composition over the canonical host-backed Table/Dir port.
22///
23/// V1 storage is a read-only compatibility prefix. The first v2 lease verifies
24/// that prefix and atomically selects its descriptor, a higher fence, and a
25/// fresh namespace. Every later append publishes immutable content-addressed
26/// leaves before comparing the complete state envelope.
27pub struct HostDirJournalBackend {
28    port: Arc<dyn HostDirPort>,
29    capabilities: BackendCapabilities,
30    work_bound: usize,
31    fail: Option<FailHook>,
32}
33
34impl HostDirJournalBackend {
35    /// Opens a binding and verifies the complete committed closure.
36    pub fn open(
37        port: Arc<dyn HostDirPort>,
38        capabilities: BackendCapabilities,
39        work_bound: usize,
40    ) -> Result<Self, JournalError> {
41        if work_bound == 0 {
42            return Err(JournalError::WorkBoundExceeded);
43        }
44        let backend = Self {
45            port,
46            capabilities,
47            work_bound,
48            fail: None,
49        };
50        backend.read_state()?;
51        Ok(backend)
52    }
53
54    /// Installs a deterministic crash hook for conformance models.
55    pub fn with_failpoint_hook(
56        mut self,
57        hook: impl Fn(Failpoint) -> bool + Send + Sync + 'static,
58    ) -> Self {
59        self.fail = Some(Arc::new(hook));
60        self
61    }
62
63    /// Reports the capability evidence used to admit or refuse writes.
64    pub fn capabilities(&self) -> BackendCapabilities {
65        self.capabilities
66    }
67
68    /// Returns the selected v2 state envelope, if migration has occurred.
69    pub fn state_envelope(&self) -> Result<Option<NativeStateEnvelope>, JournalError> {
70        let Some(bytes) = self.state_bytes()? else {
71            return Ok(None);
72        };
73        if bytes.starts_with(b"SIMJSTATE1") {
74            return Ok(None);
75        }
76        Ok(Some(decode_envelope(&bytes)?))
77    }
78
79    fn trip(&self, point: Failpoint) -> Result<(), JournalError> {
80        if self.fail.as_ref().is_some_and(|hook| hook(point)) {
81            Err(JournalError::InjectedCrash(point.label()))
82        } else {
83            Ok(())
84        }
85    }
86
87    fn write_capable(&self) -> Result<(), JournalError> {
88        if !self.capabilities.linearizable_cas {
89            return Err(JournalError::WriteRefused(
90                "linearizable table/cas unavailable",
91            ));
92        }
93        if !self.capabilities.durable_publish {
94            return Err(JournalError::WriteRefused(
95                "storage durability receipt unavailable",
96            ));
97        }
98        Ok(())
99    }
100
101    fn ensure_layout(&self) -> Result<(), JournalError> {
102        for path in [
103            vec!["objects-v2".into()],
104            vec!["namespaces-v2".into()],
105            vec!["compat-v1".into()],
106            vec!["temporary".into()],
107        ] {
108            self.port.create_dir(&path).map_err(port_error)?;
109        }
110        Ok(())
111    }
112
113    fn state_bytes(&self) -> Result<Option<Vec<u8>>, JournalError> {
114        match self.port.read(&["state".into()]) {
115            Ok(bytes) => Ok(Some(bytes)),
116            Err(error) if error.kind == HostDirErrorKind::NotFound => Ok(None),
117            Err(error) => Err(port_error(error)),
118        }
119    }
120
121    fn put_immutable(&self, path: &[String], bytes: &[u8]) -> Result<(), JournalError> {
122        let outcome = self
123            .port
124            .compare_exchange(path, None, Some(bytes), &NeverCancel)
125            .map_err(port_error)?;
126        if outcome.exchanged || outcome.observed.as_deref() == Some(bytes) {
127            Ok(())
128        } else {
129            Err(JournalError::ConflictingObject)
130        }
131    }
132
133    fn load(&self, path: &[String]) -> Result<Vec<u8>, JournalError> {
134        self.port.read(path).map_err(port_error)
135    }
136
137    fn reserve_namespace(&self, seed: &[u8]) -> Result<EntryNamespace, JournalError> {
138        for _ in 0..128 {
139            let nonce = NAMESPACE_COUNTER.fetch_add(1, Ordering::Relaxed);
140            let now = SystemTime::now()
141                .duration_since(UNIX_EPOCH)
142                .map_err(|_| JournalError::Backend("system time before epoch".into()))?
143                .as_nanos();
144            let mut hash = Sha256::new();
145            hash.update(b"sim-journal-namespace-v2\0");
146            hash.update(seed);
147            hash.update(std::process::id().to_be_bytes());
148            hash.update(now.to_be_bytes());
149            hash.update(nonce.to_be_bytes());
150            let token = hex(&hash.finalize());
151            let namespace = EntryNamespace(token);
152            let root = namespace_root(&namespace);
153            self.port.create_dir(&root).map_err(port_error)?;
154            let marker = [root.clone(), vec!["format".into()]].concat();
155            let result = self
156                .port
157                .compare_exchange(&marker, None, Some(b"SIMJNAMESPACE2"), &NeverCancel)
158                .map_err(port_error)?;
159            if result.exchanged {
160                self.port
161                    .create_dir(&[root, vec!["entries".into()]].concat())
162                    .map_err(port_error)?;
163                return Ok(namespace);
164            }
165        }
166        Err(JournalError::Backend(
167            "could not reserve a fresh journal namespace".into(),
168        ))
169    }
170
171    fn put_object(&self, object: &JournalObject) -> Result<StoredDatumRef, JournalError> {
172        object.verify()?;
173        let bytes = object.storage_bytes()?;
174        let storage = crate::object::storage_id(&bytes);
175        let meaning_dir = vec!["objects-v2".into(), id_key(&object.id)];
176        self.port.create_dir(&meaning_dir).map_err(port_error)?;
177        let path = [meaning_dir, vec![id_key(&storage)]].concat();
178        self.put_immutable(&path, &bytes)?;
179        Ok(StoredDatumRef {
180            meaning: object.id.clone(),
181            storage,
182        })
183    }
184
185    fn find_object(
186        &self,
187        meaning: &ContentId,
188    ) -> Result<(StoredDatumRef, JournalObject), JournalError> {
189        let dir = vec!["objects-v2".into(), id_key(meaning)];
190        let entries = self.port.list(&dir).map_err(port_error)?;
191        let files: Vec<_> = entries
192            .into_iter()
193            .filter(|entry| entry.kind == sim_storage_port::HostEntryKind::File)
194            .collect();
195        if files.is_empty() {
196            return Err(JournalError::MissingSemanticObject(meaning.clone()));
197        }
198        if files.len() != 1 {
199            return Err(JournalError::CorruptState("ambiguous semantic object"));
200        }
201        let storage = parse_id_key(&files[0].name)?;
202        let bytes = self.load(&[dir, vec![files[0].name.clone()]].concat())?;
203        if crate::object::storage_id(&bytes) != storage {
204            return Err(JournalError::CorruptState("object storage id"));
205        }
206        let object = JournalObject::from_storage_bytes(&bytes)?;
207        if object.id != *meaning {
208            return Err(JournalError::CorruptObject(meaning.clone()));
209        }
210        Ok((
211            StoredDatumRef {
212                meaning: meaning.clone(),
213                storage,
214            },
215            object,
216        ))
217    }
218
219    fn load_object_ref(&self, reference: &StoredDatumRef) -> Result<JournalObject, JournalError> {
220        let bytes = self.load(&object_path(reference))?;
221        if crate::object::storage_id(&bytes) != reference.storage {
222            return Err(JournalError::CorruptState("object storage id"));
223        }
224        let object = JournalObject::from_storage_bytes(&bytes)?;
225        if object.id != reference.meaning {
226            return Err(JournalError::CorruptObject(reference.meaning.clone()));
227        }
228        Ok(object)
229    }
230
231    fn read_v1(&self, state_bytes: &[u8]) -> Result<(StoredState, JournalHead), JournalError> {
232        let (_, head) = decode_v1_state(state_bytes)?;
233        let Some(physical_head) = head else {
234            return Err(JournalError::CorruptState("empty v1 prefix"));
235        };
236        let needed = usize::try_from(physical_head.sequence)
237            .ok()
238            .and_then(|value| value.checked_add(1))
239            .ok_or(JournalError::WorkBoundExceeded)?;
240        if needed > self.work_bound {
241            return Err(JournalError::WorkBoundExceeded);
242        }
243        let mut entries = BTreeMap::new();
244        let mut objects = BTreeMap::new();
245        let mut datums = BTreeMap::new();
246        let mut physical_previous = None;
247        let mut canonical_previous = None;
248        for sequence in 0..=physical_head.sequence {
249            let bytes = self.load(&v1_entry_path(sequence))?;
250            let old = decode_v1_entry(&bytes)?;
251            if old.sequence != sequence || old.previous != physical_previous {
252                return Err(JournalError::CorruptState("v1 chain"));
253            }
254            if old.id != old.canonical_id() {
255                return Err(JournalError::CorruptEntry);
256            }
257            let mut payloads = Vec::with_capacity(old.payloads.len());
258            for old_id in &old.payloads {
259                let payload = self.load(&v1_object_path(old_id))?;
260                if v1_object_id(&payload) != *old_id {
261                    return Err(JournalError::CorruptObject(old_id.clone()));
262                }
263                let object = JournalObject::from_bytes(payload);
264                payloads.push(object.id.clone());
265                objects.insert(object.id.clone(), object.bytes.clone());
266                datums.insert(object.id.clone(), object.datum().clone());
267            }
268            let entry = JournalEntry::new(sequence, canonical_previous.clone(), old.kind, payloads);
269            physical_previous = Some(old.id);
270            canonical_previous = Some(entry.id.clone());
271            entries.insert(sequence, entry);
272        }
273        if physical_previous.as_ref() != Some(&physical_head.entry) {
274            return Err(JournalError::CorruptState("v1 head"));
275        }
276        let canonical_head = entries
277            .last_key_value()
278            .map(|(_, entry)| JournalHead {
279                sequence: entry.sequence,
280                entry: entry.id.clone(),
281            })
282            .ok_or(JournalError::CorruptState("empty v1 prefix"))?;
283        let state = StoredState {
284            objects,
285            datums,
286            entries,
287            head: Some(canonical_head.clone()),
288        };
289        crate::verify::verify_state(&state)?;
290        Ok((state, canonical_head))
291    }
292
293    fn read_prefix(&self, prefix: &VerifiedNativePrefixRef) -> Result<StoredState, JournalError> {
294        let bytes = self.load(&descriptor_path(&prefix.descriptor))?;
295        if crate::object::storage_id(&bytes) != prefix.descriptor {
296            return Err(JournalError::CorruptState("prefix descriptor id"));
297        }
298        let (old_state, described_canonical_head) = decode_descriptor(&bytes)?;
299        let (_, physical) = decode_v1_state(&old_state)?;
300        if physical.as_ref() != Some(&prefix.physical_head) {
301            return Err(JournalError::CorruptState("prefix physical head"));
302        }
303        let (state, canonical) = self.read_v1(&old_state)?;
304        if described_canonical_head != prefix.canonical_head
305            || canonical != prefix.canonical_head
306            || prefix.entries != canonical.sequence.saturating_add(1)
307        {
308            return Err(JournalError::CorruptState("prefix canonical head"));
309        }
310        Ok(state)
311    }
312
313    fn read_v2(&self, envelope: &NativeStateEnvelope) -> Result<StoredState, JournalError> {
314        if envelope.format != NativeFormatId::V2 {
315            return Err(JournalError::CorruptState("native format"));
316        }
317        let marker =
318            self.load(&[namespace_root(&envelope.namespace), vec!["format".into()]].concat())?;
319        if marker != b"SIMJNAMESPACE2" {
320            return Err(JournalError::CorruptState("namespace format"));
321        }
322        let mut state = match &envelope.prefix {
323            Some(prefix) => self.read_prefix(prefix)?,
324            None => StoredState::default(),
325        };
326        let prefix_len = state.entries.len();
327        let mut location = envelope.head_location.clone();
328        let mut suffix = Vec::new();
329        let mut seen = BTreeSet::new();
330        while let Some(current) = location {
331            if current.namespace != envelope.namespace || !seen.insert(current.clone()) {
332                return Err(JournalError::CorruptState("entry locator chain"));
333            }
334            if prefix_len + suffix.len() >= self.work_bound {
335                return Err(JournalError::WorkBoundExceeded);
336            }
337            let bytes = self.load(&entry_path(&current))?;
338            if crate::object::storage_id(&bytes) != current.storage {
339                return Err(JournalError::CorruptState("entry storage id"));
340            }
341            let physical = decode_v2_entry(&bytes)?;
342            if physical.entry.id != current.entry
343                || physical.entry.sequence != current.sequence
344                || physical.entry.canonical_id()? != physical.entry.id
345            {
346                return Err(JournalError::CorruptEntry);
347            }
348            for (meaning, reference) in physical.entry.payloads.iter().zip(&physical.payloads) {
349                if meaning != &reference.meaning {
350                    return Err(JournalError::CorruptState("payload locator"));
351                }
352                let object = self.load_object_ref(reference)?;
353                state
354                    .objects
355                    .insert(object.id.clone(), object.bytes.clone());
356                state
357                    .datums
358                    .insert(object.id.clone(), object.datum().clone());
359            }
360            if physical.entry.payloads.len() != physical.payloads.len() {
361                return Err(JournalError::CorruptState("payload locator count"));
362            }
363            location = physical.previous_location.clone();
364            suffix.push(physical);
365        }
366        suffix.reverse();
367        for physical in suffix {
368            if state
369                .entries
370                .insert(physical.entry.sequence, physical.entry)
371                .is_some()
372            {
373                return Err(JournalError::CorruptState("overlapping suffix"));
374            }
375        }
376        state.head = envelope.head.clone();
377        crate::verify::verify_state(&state)?;
378        let expected_location = state.entries.len() > prefix_len;
379        if expected_location != envelope.head_location.is_some() {
380            return Err(JournalError::CorruptState("head locator"));
381        }
382        Ok(state)
383    }
384
385    fn install_v2(&self, observed: Option<&[u8]>) -> Result<Option<Lease>, JournalError> {
386        let (old_fence, prefix, canonical_head) = match observed {
387            Some(bytes) => {
388                let (fence, physical_head) = decode_v1_state(bytes)?;
389                match physical_head {
390                    Some(physical_head) => {
391                        let (_, canonical_head) = self.read_v1(bytes)?;
392                        self.trip(Failpoint::BeforePrefixDescriptor)?;
393                        let descriptor_bytes = encode_descriptor(bytes, &canonical_head);
394                        let descriptor = crate::object::storage_id(&descriptor_bytes);
395                        self.put_immutable(&descriptor_path(&descriptor), &descriptor_bytes)?;
396                        self.trip(Failpoint::AfterPrefixDescriptor)?;
397                        (
398                            fence,
399                            Some(VerifiedNativePrefixRef {
400                                entries: canonical_head.sequence + 1,
401                                physical_head,
402                                canonical_head: canonical_head.clone(),
403                                descriptor,
404                            }),
405                            Some(canonical_head),
406                        )
407                    }
408                    None => (fence, None, None),
409                }
410            }
411            None => (0, None, None),
412        };
413        let fence = old_fence
414            .checked_add(1)
415            .ok_or_else(|| JournalError::Backend("fence exhausted".into()))?;
416        let namespace = self.reserve_namespace(observed.unwrap_or_default())?;
417        self.trip(Failpoint::AfterNamespaceReservation)?;
418        let envelope = NativeStateEnvelope {
419            format: NativeFormatId::V2,
420            fence,
421            namespace,
422            prefix,
423            head: canonical_head,
424            head_location: None,
425        };
426        let replacement = encode_envelope(&envelope);
427        self.trip(Failpoint::BeforeFormatCas)?;
428        let result = self
429            .port
430            .compare_exchange(
431                &["state".into()],
432                observed,
433                Some(&replacement),
434                &NeverCancel,
435            )
436            .map_err(port_error)?;
437        if !result.exchanged {
438            return Ok(None);
439        }
440        self.trip(Failpoint::AfterFormatCas)?;
441        Ok(Some(Lease { fence }))
442    }
443}
444
445impl JournalBackend for HostDirJournalBackend {
446    fn acquire_lease(&self) -> Result<Lease, JournalError> {
447        self.write_capable()?;
448        self.ensure_layout()?;
449        loop {
450            let observed = self.state_bytes()?;
451            match observed.as_deref() {
452                None => {
453                    if let Some(lease) = self.install_v2(observed.as_deref())? {
454                        return Ok(lease);
455                    }
456                }
457                Some(bytes) if bytes.starts_with(b"SIMJSTATE1") => {
458                    if let Some(lease) = self.install_v2(Some(bytes))? {
459                        return Ok(lease);
460                    }
461                }
462                Some(bytes) => {
463                    let mut envelope = decode_envelope(bytes)?;
464                    self.read_v2(&envelope)?;
465                    envelope.fence = envelope
466                        .fence
467                        .checked_add(1)
468                        .ok_or_else(|| JournalError::Backend("fence exhausted".into()))?;
469                    let replacement = encode_envelope(&envelope);
470                    let result = self
471                        .port
472                        .compare_exchange(
473                            &["state".into()],
474                            Some(bytes),
475                            Some(&replacement),
476                            &NeverCancel,
477                        )
478                        .map_err(port_error)?;
479                    if result.exchanged {
480                        return Ok(Lease {
481                            fence: envelope.fence,
482                        });
483                    }
484                }
485            }
486        }
487    }
488
489    fn read_state(&self) -> Result<StoredState, JournalError> {
490        let Some(bytes) = self.state_bytes()? else {
491            return Ok(StoredState::default());
492        };
493        if bytes.starts_with(b"SIMJSTATE1") {
494            let (_, head) = decode_v1_state(&bytes)?;
495            return match head {
496                Some(_) => self.read_v1(&bytes).map(|(state, _)| state),
497                None => Ok(StoredState::default()),
498            };
499        }
500        let envelope = decode_envelope(&bytes)?;
501        self.read_v2(&envelope)
502    }
503
504    fn admit(&self, admission: Admission) -> Result<JournalHead, JournalError> {
505        self.write_capable()?;
506        self.ensure_layout()?;
507        let observed = self.state_bytes()?.ok_or(JournalError::WriteRefused(
508            "acquire a v2 lease before append",
509        ))?;
510        if observed.starts_with(b"SIMJSTATE1") {
511            return Err(JournalError::WriteRefused(
512                "acquire a v2 lease before append",
513            ));
514        }
515        let mut envelope = decode_envelope(&observed)?;
516        if envelope.fence != admission.fence {
517            return Err(JournalError::StaleLease);
518        }
519        if envelope.head != admission.expected {
520            if admission.entries.last().is_some_and(|entry| {
521                envelope
522                    .head
523                    .as_ref()
524                    .is_some_and(|head| head.entry == entry.id)
525            }) {
526                return envelope.head.ok_or(JournalError::WrongHead);
527            }
528            return Err(JournalError::WrongHead);
529        }
530        self.trip(Failpoint::BeforeObjectPublish)?;
531        let mut references = BTreeMap::new();
532        for object in &admission.objects {
533            let reference = self.put_object(object)?;
534            references.insert(reference.meaning.clone(), reference);
535        }
536        self.trip(Failpoint::AfterObjectPublish)?;
537        let mut previous_location = envelope.head_location.clone();
538        let mut last_location = None;
539        for entry in &admission.entries {
540            let mut payloads = Vec::with_capacity(entry.payloads.len());
541            for meaning in &entry.payloads {
542                let reference = match references.get(meaning) {
543                    Some(reference) => reference.clone(),
544                    None => self.find_object(meaning)?.0,
545                };
546                payloads.push(reference);
547            }
548            let physical = PhysicalEntry {
549                entry: entry.clone(),
550                previous_location: previous_location.clone(),
551                payloads,
552            };
553            let bytes = encode_v2_entry(&physical);
554            let storage = crate::object::storage_id(&bytes);
555            let location = EntryLocation {
556                namespace: envelope.namespace.clone(),
557                sequence: entry.sequence,
558                entry: entry.id.clone(),
559                storage,
560            };
561            ensure_entry_dirs(self.port.as_ref(), &location)?;
562            self.put_immutable(&entry_path(&location), &bytes)?;
563            previous_location = Some(location.clone());
564            last_location = Some(location);
565        }
566        self.trip(Failpoint::AfterDurabilityReceipt)?;
567        let last = admission.entries.last().ok_or(JournalError::EmptyBatch)?;
568        let new_head = JournalHead {
569            sequence: last.sequence,
570            entry: last.id.clone(),
571        };
572        envelope.head = Some(new_head.clone());
573        envelope.head_location = last_location;
574        self.trip(Failpoint::BeforeCas)?;
575        let replacement = encode_envelope(&envelope);
576        let result = self
577            .port
578            .compare_exchange(
579                &["state".into()],
580                Some(&observed),
581                Some(&replacement),
582                &NeverCancel,
583            )
584            .map_err(port_error)?;
585        if !result.exchanged {
586            return Err(JournalError::WrongHead);
587        }
588        self.trip(Failpoint::AfterCas)?;
589        self.trip(Failpoint::BeforeAcknowledgement)?;
590        Ok(new_head)
591    }
592
593    fn put_datum(&self, object: JournalObject) -> Result<StoredDatumRef, JournalError> {
594        self.write_capable()?;
595        self.ensure_layout()?;
596        self.put_object(&object)
597    }
598
599    fn get_datum(&self, meaning: &ContentId) -> Result<Datum, JournalError> {
600        Ok(self.find_object(meaning)?.1.datum().clone())
601    }
602
603    fn rebuild_datum_index(&self) -> Result<Vec<StoredDatumRef>, JournalError> {
604        let mut rebuilt = Vec::new();
605        let entries = match self.port.list(&["objects-v2".into()]) {
606            Ok(entries) => entries,
607            Err(error) if error.kind == HostDirErrorKind::NotFound => return Ok(rebuilt),
608            Err(error) => return Err(port_error(error)),
609        };
610        for entry in entries {
611            if entry.kind != sim_storage_port::HostEntryKind::Directory {
612                return Err(JournalError::CorruptState("object index entry"));
613            }
614            if rebuilt.len() >= self.work_bound {
615                return Err(JournalError::WorkBoundExceeded);
616            }
617            let meaning = parse_id_key(&entry.name)?;
618            rebuilt.push(self.find_object(&meaning)?.0);
619        }
620        rebuilt.sort_by(|left, right| left.meaning.cmp(&right.meaning));
621        Ok(rebuilt)
622    }
623}