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
21pub struct HostDirJournalBackend {
28 port: Arc<dyn HostDirPort>,
29 capabilities: BackendCapabilities,
30 work_bound: usize,
31 fail: Option<FailHook>,
32}
33
34impl HostDirJournalBackend {
35 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 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 pub fn capabilities(&self) -> BackendCapabilities {
65 self.capabilities
66 }
67
68 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(¤t))?;
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}