From 2fec6a39e01d449be216d01dae7ed186375d289e Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Fri, 4 Sep 2026 00:22:10 -0400 Subject: [PATCH] feat(pds)!: put record bodies in a heap and the journal on its slots MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A record's canonical DAG-CBOR is appended to `records/pds.heap` as a framed body and the journal carries `RecordStored` — the three names, the content identifier and the slot — instead of the value. Bytes reach the medium before the fact that names them, so a crash between the two costs bodies nothing names, which the next boot reclaims; a slot past the end of the heap is a torn tail and the record is dropped, and a slot back inside it after one past the end is refused. Change-Id: I6ec51e0e8672aecb933dbfa80e52be585044fc14 (cherry picked from commit ca0f8521f7efdc497e3b83847d0961628c12ff6c) --- crates/didbot-pds/src/blobs.rs | 17 +- crates/didbot-pds/src/durable.rs | 194 +++++-- crates/didbot-pds/src/heap.rs | 760 ++++++++++++++++++++++++++ crates/didbot-pds/src/journal.rs | 33 +- crates/didbot-pds/src/lib.rs | 8 +- crates/didbot-pds/src/records.rs | 697 ++++++++++++++++++++++- crates/didbot-pds/src/wal/mod.rs | 22 + crates/didbot-pds/tests/durability.rs | 2 +- crates/didbot-pds/tests/upgrade.rs | 6 + 9 files changed, 1666 insertions(+), 73 deletions(-) create mode 100644 crates/didbot-pds/src/heap.rs diff --git a/crates/didbot-pds/src/blobs.rs b/crates/didbot-pds/src/blobs.rs index a9ac0fdc..2b415a21 100644 --- a/crates/didbot-pds/src/blobs.rs +++ b/crates/didbot-pds/src/blobs.rs @@ -38,7 +38,7 @@ //! # The seam //! //! [`BlobStore`] is a trait beside [`AccountStore`](crate::AccountStore) and -//! [`RecordStore`], with the same two implementations +//! [`RecordStore`](crate::records::RecordStore), with the same two implementations //! each of those has: [`MemoryBlobStore`] for tests and an ephemeral run, and //! [`FileBlobStore`](crate::FileBlobStore) for a deployment given `--data`. //! The trait is shaped for a third that does not exist yet — object storage — which is what this @@ -52,7 +52,6 @@ use std::time::Duration; use didbot_data::{Cid, RawHasher}; use time::OffsetDateTime; -use crate::records::RecordStore; use crate::wal::Entry; /// Longest a single blob may be, when a deployment does not say. @@ -515,7 +514,7 @@ pub trait BlobStore: Send + Sync { /// Drops every blob belonging to an account, and the bytes with them. /// /// Called when the account is deleted, for the reason - /// [`RecordStore::remove_repo`] gives: + /// [`RecordStore::remove_repo`](crate::records::RecordStore::remove_repo) gives: /// bytes that outlive the identity that uploaded them are unattributable /// and unreachable, and here they also go on occupying a disk. fn remove_repo(&self, did: &str); @@ -1337,8 +1336,8 @@ impl crate::journal::Journaled for crate::FileBlobStore { } impl crate::journal::Bodied for crate::FileBlobStore { - fn bodies(&self) -> &std::path::Path { - self.root() + fn bodies(&self) -> crate::journal::Bodies<'_> { + crate::journal::Bodies::Tree(self.root()) } } @@ -1631,10 +1630,10 @@ mod tests { let dir = std::env::temp_dir().join(format!("didbot-blob-bodies-{}", std::process::id())); let _ = std::fs::remove_dir_all(&dir); let store = file_store(&dir); - assert_eq!( - crate::journal::Bodied::bodies(&store), - dir.join(crate::BLOB_DIR) - ); + let crate::journal::Bodies::Tree(root) = crate::journal::Bodied::bodies(&store) else { + panic!("the blob store files one value per name, so its bodies are a tree"); + }; + assert_eq!(root, dir.join(crate::BLOB_DIR)); let _ = std::fs::remove_dir_all(&dir); } } diff --git a/crates/didbot-pds/src/durable.rs b/crates/didbot-pds/src/durable.rs index 303f4500..2482c5b2 100644 --- a/crates/didbot-pds/src/durable.rs +++ b/crates/didbot-pds/src/durable.rs @@ -56,6 +56,7 @@ use crate::blobs::{ BlobTally, BlobUpload, CollectedBlob, Fetch, }; use crate::credential::{AgentTokenStore, IssuedToken, MemoryAgentTokenStore, TokenError}; +use crate::heap::Heap; use crate::history::{CommitStore, Head, HistoryStats, MemoryCommitStore}; use crate::journal::{Derived, Journaled, Stores}; use crate::ledger::{AgentLedger, LedgerEntry, LedgerEvent, LedgerStore, MemoryLedger}; @@ -64,9 +65,9 @@ use crate::oauth::{ Grant, GrantRequest, Lifetimes, MemoryOAuthGrantStore, MintedGrant, OAuthGrantStore, Rotation, }; use crate::records::{ - content_id, plan_batch, resolve_key, validate, validate_removal, validate_shape, BatchOp, - BatchOutcome, ListParams, MemoryRecordStore, Precondition, RecordError, RecordStats, - RecordStore, ResolvedKey, Stance, Written, + encode_record, heap_backend, plan_batch, resolve_key, validate, validate_removal, + validate_shape, BatchOp, BatchOutcome, HeapRecordStore, Held, ListParams, Precondition, + RecordError, RecordStats, RecordStore, ResolvedKey, Stance, Written, }; use crate::sequence::SequenceFloor; use crate::wal::{Durability, Entry, Replay, Wal, WalError}; @@ -182,7 +183,12 @@ impl Durable { let wal = Arc::new(wal); let accounts = MemoryAccountStore::new(); - let records = MemoryRecordStore::new(); + // Before the replay, because a replayed slot is judged against the + // length the heap already has: bodies reached the medium in front of + // the journal entries naming them, so the heap as it stands now is + // what decides which of those entries are covered. + let heap = Arc::new(Heap::open(&dir.join(crate::wal::RECORD_DIR)).map_err(heap_open)?); + let records = HeapRecordStore::new(heap.clone()); let blobs = FileBlobStore::open(dir, wal.clone(), limits)?; let ledger = MemoryLedger::new(); let commits = MemoryCommitStore::new(); @@ -208,6 +214,14 @@ impl Durable { entry, ); } + // After the replay and before anything is served, for the same + // reason the blob sweep below runs there: a body is an orphan exactly + // when the journal does not name it, and the journal is only fully + // read now. This is also where a journal that names bodies the heap + // does not hold is resolved — dropped if it is the tail, refused if + // the two files disagree about their order. See + // `HeapRecordStore::settle`. + let reclaimed_bodies = records.settle().map_err(heap_open)?; // After the replay and before anything is served: a file is an orphan // exactly when the log does not name it, and the log is only fully // read now. Only when there was a log to read: with no log the index @@ -272,6 +286,8 @@ impl Durable { credentials = credential_count, oauth_grants = grant_count, history, + record_bytes = heap.len(), + reclaimed_bodies, swept = reconciled.swept.bytes, missing_blobs = reconciled.missing.count, "restored the deployment from its write-ahead log" @@ -309,6 +325,7 @@ impl Durable { records: Arc::new(FileRecordStore { guard: Mutex::new(()), inner: records, + heap, wal: wal.clone(), }), blobs: Arc::new(blobs), @@ -452,7 +469,7 @@ impl Durable { #[allow(clippy::too_many_arguments)] fn apply( accounts: &MemoryAccountStore, - records: &MemoryRecordStore, + records: &HeapRecordStore, blobs: &FileBlobStore, ledger: &MemoryLedger, commits: &MemoryCommitStore, @@ -468,7 +485,32 @@ fn apply( // about to replace, and that record is only readable until the record // store applies the entry. See the [`Derived`] implementation in // [`crate::blobs`]. - blobs.observe(&entry, Stores { records }); + // A stored record reaches the blob index as the put it is. The index + // counts the blobs a record names by reading the record, and a + // `RecordStored` carries a slot instead — so the body is resolved out of + // the heap here and the index is handed the fact in the terms it holds + // counts in. Its own move to a carried count is the pull after this one. + match &entry { + Entry::RecordStored { + did, + collection, + rkey, + .. + } => { + if let Some(record) = records.replayed_body(&entry) { + blobs.observe( + &Entry::RecordPut { + did: did.clone(), + collection: collection.clone(), + rkey: rkey.clone(), + record, + }, + Stores { records }, + ); + } + } + _ => blobs.observe(&entry, Stores { records }), + } // One deletion of one repository is one fact, and three stores each hold // a piece of it: a repository is its records, its blobs and its commit // head. The fan-out is here rather than in any of them because the fact @@ -486,7 +528,7 @@ fn apply( // store's file and leaves this list where it is. if MemoryAccountStore::owns(&entry) { accounts.apply(entry); - } else if MemoryRecordStore::owns(&entry) { + } else if HeapRecordStore::owns(&entry) { records.apply(entry); } else if FileBlobStore::owns(&entry) { blobs.apply(entry); @@ -506,10 +548,9 @@ fn apply( floor.apply(entry); } else if !matches!( entry, - Entry::RecordStored { .. } | Entry::BlobReferenced { .. } | Entry::CheckpointSealed { .. } + Entry::BlobReferenced { .. } | Entry::CheckpointSealed { .. } ) { - // `RecordStored`, `BlobReferenced` and - // `CheckpointSealed` are the destination format + // `BlobReferenced` and `CheckpointSealed` are the destination format // [`crate::layout::ENTRY_SHAPE`] declares: a log carrying one is read // here, and the state each describes belongs to the store that will // write it. Anything else arriving here is an entry the log carries @@ -525,7 +566,7 @@ fn apply( #[allow(clippy::too_many_arguments)] fn snapshot( accounts: &MemoryAccountStore, - records: &MemoryRecordStore, + records: &HeapRecordStore, blobs: &FileBlobStore, ledger: &MemoryLedger, commits: &MemoryCommitStore, @@ -734,6 +775,15 @@ fn store_backend(error: WalError) -> StoreError { } } +/// A heap failure at startup, as the refusal [`Durable::open`] hands back. +/// +/// The heap is a second file under the same directory, so a directory whose +/// bodies will not open is a directory this binary cannot serve — the same +/// answer `pds.wal` gives for its own file. +fn heap_open(error: crate::heap::HeapError) -> WalError { + WalError::Bodies(error) +} + fn backend(error: WalError) -> String { let mut message = error.to_string(); let mut source = std::error::Error::source(&error); @@ -872,7 +922,8 @@ impl AccountStore for FileAccountStore { #[derive(Debug)] pub struct FileRecordStore { guard: Mutex<()>, - inner: MemoryRecordStore, + inner: HeapRecordStore, + heap: Arc, wal: Arc, } @@ -883,9 +934,42 @@ impl FileRecordStore { .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) } + + /// The heap the record bodies live in. + pub fn heap(&self) -> &Arc { + &self.heap + } + + /// Appends a body, and answers where it landed. + /// + /// The half of the write that goes first. **Bytes reach the medium before + /// the fact that names them does** — see [`crate::heap`] — so this runs + /// under the write lock, before the journal append, and the heap is + /// synced ahead of any journal append that will sync. + fn store_body( + &self, + body: &[u8], + durability: Durability, + ) -> Result { + let slot = self.heap.append(body, durability).map_err(heap_backend)?; + // The journal is about to sync, so the bytes go first. `sync_due` + // is the journal's own answer about its deferred window, which is + // what makes this exact rather than a guess about timers. + if durability == Durability::Sync || self.wal.sync_due() { + self.heap.sync().map_err(heap_backend)?; + } + Ok(slot) + } } impl RecordStore for FileRecordStore { + /// Check, append the body, append the fact, apply — under one lock. + /// + /// [`RecordStore::put`]'s order with one more append inside it, and the + /// order of those two is the whole of the crash argument. The heap frame + /// reaches the medium first, so the crash between them costs bytes + /// nothing names — dead space a boot reclaims — rather than a fact naming + /// bytes nothing has. fn put( &self, did: &str, @@ -895,10 +979,11 @@ impl RecordStore for FileRecordStore { expect: &Precondition, ) -> Result { validate(collection, &record)?; - // Named before the lock for the same reason the key is resolved - // before it: the CID is a pure function of the record the caller - // handed over, and computing it is work that does not need the lock. - let cid = content_id(&record)?; + // Encoded before the lock for the same reason the key is resolved + // before it: the body and the identifier over it are pure functions + // of the record the caller handed over, and both are work that does + // not need the lock. + let (body, cid) = encode_record(&record)?; // Resolved before the lock, because it reads only the lexicon and // what the caller asked for, and a refusal should not have queued // behind another write to get there. @@ -913,28 +998,34 @@ impl RecordStore for FileRecordStore { ResolvedKey::Mint => self.inner.mint_key(), ResolvedKey::Fixed(fixed) => fixed, }; - // Checked before the log entry is appended, under the write lock, so - // a refused write leaves nothing in the log to replay. The read is of - // the same in-memory state the write is about to change. - expect.check( - self.inner - .get(did, collection, &rkey) - .map(|record| content_id(&record)) - .transpose()? - .as_ref(), - )?; + // Checked before anything is appended, under the write lock, so a + // refused write leaves neither a body in the heap nor a fact in the + // log. The comparison is against the slot's stored identifier, which + // is what makes it a map lookup rather than a read of the body. + expect.check(self.inner.cid_at(did, collection, &rkey).as_ref())?; + let slot = self.store_body(&body, Durability::Deferred)?; self.wal .append( - &Entry::RecordPut { + &Entry::RecordStored { did: did.to_owned(), collection: collection.to_owned(), rkey: rkey.clone(), - record: record.clone(), + cid: cid.to_string(), + offset: slot.offset, + len: slot.len, }, Durability::Deferred, ) .map_err(record_backend)?; - self.inner.insert_at(did, collection, rkey.clone(), record); + self.inner.insert_slot( + did, + collection, + rkey.clone(), + Held { + slot, + cid: cid.clone(), + }, + ); Ok(Written { rkey, cid }) } @@ -958,24 +1049,29 @@ impl RecordStore for FileRecordStore { record: Value, ) -> Result { validate_shape(collection, &record)?; + let (body, cid) = encode_record(&record)?; let resolved = resolve_key(collection, None)?; let _write = self.write(); let rkey = match resolved { ResolvedKey::Mint => self.inner.mint_key(), ResolvedKey::Fixed(fixed) => fixed, }; + let slot = self.store_body(&body, Durability::Sync)?; self.wal .append( - &Entry::RecordPut { + &Entry::RecordStored { did: did.to_owned(), collection: collection.to_owned(), rkey: rkey.clone(), - record: record.clone(), + cid: cid.to_string(), + offset: slot.offset, + len: slot.len, }, Durability::Sync, ) .map_err(record_backend)?; - self.inner.insert_at(did, collection, rkey.clone(), record); + self.inner + .insert_slot(did, collection, rkey.clone(), Held { slot, cid }); Ok(rkey) } @@ -1004,8 +1100,8 @@ impl RecordStore for FileRecordStore { // deletion and must not queue behind one or leave a trace of one. validate_removal(collection)?; let _write = self.write(); - let held = self.inner.get(did, collection, rkey); - expect.check(held.as_ref().map(content_id).transpose()?.as_ref())?; + let held = self.inner.cid_at(did, collection, rkey); + expect.check(held.as_ref())?; if held.is_none() { return Ok(false); } @@ -1067,7 +1163,7 @@ impl RecordStore for FileRecordStore { self.inner.stats() } - /// Plans the whole batch first, exactly as [`MemoryRecordStore`] does, + /// Plans the whole batch first, exactly as [`HeapRecordStore`] does, /// then logs and applies each planned operation under the write lock. /// /// Each operation is still its own log append rather than one entry for @@ -1091,22 +1187,34 @@ impl RecordStore for FileRecordStore { for (collection, rkey, value) in planned { match value { Some((record, cid)) => { + let (body, _) = encode_record(&record)?; + let slot = self.store_body(&body, Durability::Deferred)?; self.wal .append( - &Entry::RecordPut { + &Entry::RecordStored { did: did.to_owned(), collection: collection.clone(), rkey: rkey.clone(), - record: record.clone(), + cid: cid.to_string(), + offset: slot.offset, + len: slot.len, }, Durability::Deferred, ) .map_err(record_backend)?; - self.inner.insert_at(did, &collection, rkey.clone(), record); + self.inner.insert_slot( + did, + &collection, + rkey.clone(), + Held { + slot, + cid: cid.clone(), + }, + ); results.push(BatchOutcome::Written(Written { rkey, cid })); } None => { - if self.inner.get(did, &collection, &rkey).is_some() { + if self.inner.cid_at(did, &collection, &rkey).is_some() { if let Err(error) = self.wal.append( &Entry::RecordRemoved { did: did.to_owned(), @@ -2060,7 +2168,7 @@ mod tests { struct Replayed { dir: PathBuf, accounts: MemoryAccountStore, - records: MemoryRecordStore, + records: HeapRecordStore, blobs: FileBlobStore, ledger: MemoryLedger, commits: MemoryCommitStore, @@ -2080,10 +2188,14 @@ mod tests { let (wal, _) = Wal::open(&dir).expect("a log to hang the blob store off"); let blobs = FileBlobStore::open(&dir, Arc::new(wal), BlobLimits::default()) .expect("a blob store"); + // The bodies live outside the journal, so a replay of records + // needs the same heap the real store opens. + let heap = + Arc::new(Heap::open(&dir.join(crate::wal::RECORD_DIR)).expect("a record heap")); Self { dir, accounts: MemoryAccountStore::new(), - records: MemoryRecordStore::new(), + records: HeapRecordStore::new(heap), blobs, ledger: MemoryLedger::new(), commits: MemoryCommitStore::new(), diff --git a/crates/didbot-pds/src/heap.rs b/crates/didbot-pds/src/heap.rs new file mode 100644 index 00000000..92188061 --- /dev/null +++ b/crates/didbot-pds/src/heap.rs @@ -0,0 +1,760 @@ +//! An append-only file of bodies, and the slots a journal names them by. +//! +//! A [`Heap`] is the medium behind a [`Bodied`](crate::journal::Bodied) store +//! whose values are appended rather than filed one per name. It holds +//! `crate::wal::frame` frames — a length and a CRC-32 over both the length +//! and the payload — which is the framing `pds.wal` writes and the framing +//! [`crate::evaluation_log`] already reuses for a second file. A value is +//! appended here and the journal carries a [`Slot`] and the value's content +//! identifier in place of the bytes. +//! +//! # The ordering, and what a crash between the two costs +//! +//! One direction, and it is [`crate::FileBlobStore`]'s: **bytes reach the +//! medium before the fact that names them does.** The two crashes are not +//! symmetrical. +//! +//! - Heap bytes no journal entry names are dead space. They are reclaimed at +//! the next boot by [`Heap::reclaim`], which is the same posture +//! `.incoming/` has in the blob store: a file nothing references is a file +//! the next startup deletes. +//! - A journal entry naming bytes the heap does not hold is a read that +//! cannot be served, and it is the one worth designing against. The +//! journal's own rule decides it, in the terms `pds.wal`'s replay already +//! states: a [`Slot`] lying **beyond** the heap's length is a torn tail and +//! the value is dropped with a line naming it, and every later slot must +//! also lie beyond — a slot back inside the heap after one past its end is +//! a journal and a heap that disagree about their own order, which is +//! refused rather than guessed at. +//! +//! [`Heap::bounds`] is that decision and it costs no I/O: a slot is either +//! wholly inside the heap's length or it is not. What a boot does read is one +//! frame — [`Heap::verify`] over the highest-offset slot the journal named — +//! because an interrupted append can only have torn the last frame written. +//! Damage anywhere else is caught by the CRC on the read that wants those +//! bytes, so a damaged frame is refused to its caller and never served as a +//! value. +//! +//! # The buffer +//! +//! Every write to a repository rebuilds the whole repository — +//! `provision::commit_write` signs a commit over `RecordStore::snapshot`, so +//! a write reads every record the repository holds. With the bodies off the +//! resident heap that is a disk read per record, which is why a buffer sits +//! in front of the file: bytes keyed by offset, least-recently-used, held to +//! a budget counted in the bytes it holds ([`DEFAULT_BUFFER_BYTES`]). A +//! deployment whose repositories fit in the budget pays the disk once per +//! record and answers every rebuild after it out of memory. +//! +//! The budget is the honest form of the capacity claim: memory stops being +//! the ceiling on how many record bytes a deployment can hold and becomes a +//! number an operator sets. + +use std::collections::{BTreeMap, HashMap}; +use std::fs::{File, OpenOptions}; +use std::io::Write as _; +use std::os::unix::fs::FileExt as _; +use std::path::{Path, PathBuf}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use crate::wal::frame; +use crate::wal::{Durability, HEAP_FILE}; + +/// How many bytes of bodies the buffer in front of a heap holds by default. +/// +/// Sixty-four mebibytes, which is chosen against the read the write path +/// makes rather than against a deployment's size: a commit is signed over +/// every record in one repository, so what the buffer has to hold to keep a +/// write off the disk is one repository, not the deployment. A repository of +/// a hundred thousand records at the wire sizes this server measures fits +/// inside it. +pub const DEFAULT_BUFFER_BYTES: usize = 64 * 1024 * 1024; + +/// How long an unsynced heap waits before an append syncs it. +/// +/// Half [`crate::wal::DEFAULT_SYNC_INTERVAL`], so that in a steady stream of +/// writes the heap's own timer has already fired by the time the journal's +/// does. It is not what makes the ordering hold — [`Heap::sync`] is called +/// before every journal append that will sync, and the recovery rule above +/// covers what a power cut can still catch — it is what keeps the window +/// small enough that the recovery rule is rarely the thing doing the work. +pub const SYNC_INTERVAL: Duration = Duration::from_millis(125); + +/// Where one value's bytes are in a heap. +/// +/// An offset to the head of the value's *frame* and the length of that +/// frame's payload, which is the value itself. The frame's eight bytes of +/// header sit between them, so the bytes a slot occupies are +/// `offset .. offset + HEADER + len`. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub struct Slot { + /// Where the value's frame starts, in bytes from the head of the heap. + pub offset: u64, + /// How many bytes of payload that frame carries. + pub len: u64, +} + +impl Slot { + /// The first byte past this slot's frame. + /// + /// Saturating, so a length read out of a journal that was written by + /// something else cannot wrap into an offset that looks reachable. + #[must_use] + pub fn end(&self) -> u64 { + self.offset + .saturating_add(frame::HEADER as u64) + .saturating_add(self.len) + } +} + +/// Whether a slot is inside the heap. +/// +/// The whole of the boot-time decision, and it reads nothing: see the +/// module's ordering argument. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Bounds { + /// The slot's frame is wholly inside the bytes the heap holds. + Inside, + /// The slot's frame runs past the end of the heap. + Beyond, +} + +/// Why a heap could not answer. +#[derive(Debug, thiserror::Error)] +pub enum HeapError { + /// The file would not cooperate. + #[error("record heap at {path}: {source}")] + Io { + /// The heap file. + path: PathBuf, + /// The underlying failure. + #[source] + source: std::io::Error, + }, + /// A slot inside the heap does not name a frame. + /// + /// The bytes are there and they do not check out, which is the case + /// `pds.wal`'s replay refuses to guess about: a frame that fails its CRC + /// is not a value, and answering with what the bytes happen to say would + /// serve a record nothing wrote. + #[error( + "record heap at {path}: the frame at {offset} does not check out ({reason}), so the \ + value {len} bytes long that a journal entry names there cannot be served" + )] + Damaged { + /// The heap file. + path: PathBuf, + /// Where the frame was supposed to start. + offset: u64, + /// How long the journal said its payload was. + len: u64, + /// Which way it failed, in a phrase. + reason: &'static str, + }, + /// There is no room to append the value. + /// + /// Nothing was written and the heap is intact: this is the refusal, not + /// the damage. It carries [`crate::wal::WalError::Full`]'s meaning and + /// reaches a caller as `RecordError::StorageFull` for the same reason. + #[error("record heap at {path}: {reason}")] + Full { + /// The heap file. + path: PathBuf, + /// Which of the three it was, in a sentence. + reason: String, + }, + /// The value is larger than a frame this heap could read back. + #[error( + "record heap at {path}: a value of {size} bytes is larger than the {limit} bytes a \ + frame can carry" + )] + TooLarge { + /// The heap file. + path: PathBuf, + /// How big the value encoded to. + size: usize, + /// The longest payload a frame carries. + limit: usize, + }, +} + +/// An append-only file of framed bodies, with a buffer in front of it. +#[derive(Debug)] +pub struct Heap { + path: PathBuf, + inner: Mutex, + buffer: Mutex, +} + +/// The parts an append mutates. +#[derive(Debug)] +struct Inner { + file: File, + /// How many bytes the file holds, as of the last append that succeeded. + len: u64, + /// When the file was last synced, or `None` if nothing is outstanding. + dirty_since: Option, + /// Why the heap stopped accepting appends, if it has. + /// + /// Set only when a failed append could not be truncated back, which + /// leaves a fragment at the end. A slot after it would name bytes on the + /// far side of a frame no reader can walk past. + sealed: Option, +} + +impl Heap { + /// Opens (creating if absent) the heap under `dir`. + /// + /// `dir` is the store's own directory — `/records` for records — + /// and is created if it is not there. Nothing is read: the length of the + /// file is the whole of what a boot needs, because a slot is judged + /// against it. + pub fn open(dir: &Path) -> Result { + crate::wal::create_dir(dir).map_err(|error| HeapError::Io { + path: dir.to_path_buf(), + source: std::io::Error::other(error.to_string()), + })?; + let path = dir.join(HEAP_FILE); + let file = OpenOptions::new() + .create(true) + .append(true) + .read(true) + .open(&path) + .map_err(|source| HeapError::Io { + path: path.clone(), + source, + })?; + let len = file + .metadata() + .map_err(|source| HeapError::Io { + path: path.clone(), + source, + })? + .len(); + Ok(Self { + path, + inner: Mutex::new(Inner { + file, + len, + dirty_since: None, + sealed: None, + }), + buffer: Mutex::new(Buffer::new(DEFAULT_BUFFER_BYTES)), + }) + } + + /// Where the heap lives. + pub fn path(&self) -> &Path { + &self.path + } + + /// How many bytes the heap holds. + pub fn len(&self) -> u64 { + self.lock().len + } + + /// Whether the heap holds nothing, which it does only before the first + /// append. + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// How many bytes of bodies the buffer in front of the file may hold. + /// + /// A knob and not a ceiling: a heap with a budget of zero answers every + /// read off the disk and holds nothing resident beyond the slots. + pub fn set_buffer_bytes(&self, bytes: usize) { + self.buffer().set_budget(bytes); + } + + /// How many bytes of bodies the buffer is holding right now. + pub fn buffered_bytes(&self) -> usize { + self.buffer().held + } + + /// Whether `slot` is inside the bytes this heap holds. + /// + /// No I/O: see the module's recovery argument, which is decided on this + /// answer alone. + pub fn bounds(&self, slot: Slot) -> Bounds { + if slot.len > frame::MAX_PAYLOAD as u64 || slot.end() > self.len() { + Bounds::Beyond + } else { + Bounds::Inside + } + } + + /// Reads the value at `slot`. + /// + /// Answered out of the buffer when it is there, and off the disk when it + /// is not — and a frame that does not check out is + /// [`HeapError::Damaged`] rather than whatever the bytes happened to say. + pub fn read(&self, slot: Slot) -> Result>, HeapError> { + if let Some(body) = self.buffer().get(slot.offset) { + return Ok(body); + } + let body = Arc::new(self.read_uncached(slot)?); + self.buffer().put(slot.offset, body.clone()); + Ok(body) + } + + /// Reads the frame at `slot` and throws the bytes away. + /// + /// What a boot does to the one frame an interrupted append can have + /// torn. Nothing is buffered, because a boot's check is not a read + /// anything is about to want. + pub fn verify(&self, slot: Slot) -> Result<(), HeapError> { + self.read_uncached(slot).map(drop) + } + + /// Reads the frame at `slot` straight off the file. + fn read_uncached(&self, slot: Slot) -> Result, HeapError> { + let width = usize::try_from(slot.end().saturating_sub(slot.offset)).map_err(|_| { + HeapError::Damaged { + path: self.path.clone(), + offset: slot.offset, + len: slot.len, + reason: "the slot is longer than this machine can address", + } + })?; + let mut bytes = vec![0u8; width]; + { + let inner = self.lock(); + inner + .file + .read_exact_at(&mut bytes, slot.offset) + .map_err(|source| HeapError::Io { + path: self.path.clone(), + source, + })?; + } + match frame::decode(&bytes) { + frame::Read::Frame { payload, .. } if payload.len() as u64 == slot.len => { + Ok(payload.to_vec()) + } + frame::Read::Frame { .. } => Err(HeapError::Damaged { + path: self.path.clone(), + offset: slot.offset, + len: slot.len, + reason: "the frame's payload is not the length the slot claims", + }), + frame::Read::Torn(reason) => Err(HeapError::Damaged { + path: self.path.clone(), + offset: slot.offset, + len: slot.len, + reason: reason.reason(), + }), + } + } + + /// Appends `body` and answers where it landed. + /// + /// Refuses rather than damages the heap, exactly as an append to the + /// journal does: a value too large for a frame is never written, and an + /// append that failed part way is truncated back to the length it started + /// at — or, if even that fails, the heap is sealed rather than appended + /// after. + pub fn append(&self, body: &[u8], durability: Durability) -> Result { + if body.len() > frame::MAX_PAYLOAD { + return Err(HeapError::TooLarge { + path: self.path.clone(), + size: body.len(), + limit: frame::MAX_PAYLOAD, + }); + } + let framed = frame::encode(body); + let mut inner = self.lock(); + if let Some(reason) = &inner.sealed { + return Err(HeapError::Full { + path: self.path.clone(), + reason: reason.clone(), + }); + } + let offset = inner.len; + if let Err(error) = inner.file.write_all(&framed) { + return Err(self.rewind(&mut inner, error)); + } + inner.len += framed.len() as u64; + + let now = Instant::now(); + let overdue = inner + .dirty_since + .is_some_and(|since| now.duration_since(since) >= SYNC_INTERVAL); + if durability == Durability::Sync || overdue { + inner.file.sync_data().map_err(|source| HeapError::Io { + path: self.path.clone(), + source, + })?; + inner.dirty_since = None; + } else if inner.dirty_since.is_none() { + inner.dirty_since = Some(now); + } + drop(inner); + + let slot = Slot { + offset, + len: body.len() as u64, + }; + // Buffered on the way in, which is the whole write path's read: the + // commit this value is about to be part of reads every record in the + // repository back. + self.buffer().put(offset, Arc::new(body.to_vec())); + Ok(slot) + } + + /// Waits for everything appended so far to reach the medium. + /// + /// Called before a journal append that will sync, which is what puts the + /// bytes on the medium in front of the fact that names them. + pub fn sync(&self) -> Result<(), HeapError> { + let mut inner = self.lock(); + if inner.dirty_since.is_none() { + return Ok(()); + } + inner.file.sync_data().map_err(|source| HeapError::Io { + path: self.path.clone(), + source, + })?; + inner.dirty_since = None; + Ok(()) + } + + /// Drops every byte past `keep`, answering how many went. + /// + /// The reclaim: `keep` is the first byte past the highest-offset frame + /// any journal entry named, so everything above it is a body written by + /// an append whose journal entry never landed. It is the heap's + /// `.incoming/`, and a boot is where it is swept. + pub fn reclaim(&self, keep: u64) -> Result { + let mut inner = self.lock(); + if keep >= inner.len { + return Ok(0); + } + let reclaimed = inner.len - keep; + inner.file.set_len(keep).map_err(|source| HeapError::Io { + path: self.path.clone(), + source, + })?; + inner.file.sync_all().map_err(|source| HeapError::Io { + path: self.path.clone(), + source, + })?; + inner.len = keep; + drop(inner); + self.buffer().forget_from(keep); + Ok(reclaimed) + } + + /// Undoes a failed append, and says why it failed. + /// + /// [`crate::wal::Wal`]'s `rewind`, for the same reason and with the same + /// consequence: `write_all` says what went wrong and not how far it got, + /// so the length before the attempt is the only thing that covers every + /// case. A truncation that also fails leaves a fragment that cannot be + /// removed, and a fragment with appends after it is every later frame + /// unreadable — so the heap is sealed instead. + fn rewind(&self, inner: &mut Inner, cause: std::io::Error) -> HeapError { + // ENOSPC by number rather than by `ErrorKind`, which did not name it + // for most of this crate's supported compiler range. + let out_of_room = cause.raw_os_error() == Some(28); + if let Err(error) = inner.file.set_len(inner.len) { + let reason = format!( + "an append failed ({cause}) and the partial frame could not be removed ({error}); \ + the heap is sealed and this deployment is serving reads only until it restarts" + ); + tracing::error!(path = %self.path.display(), reason, "sealed the record heap"); + inner.sealed = Some(reason.clone()); + return HeapError::Full { + path: self.path.clone(), + reason, + }; + } + tracing::warn!( + path = %self.path.display(), + len = inner.len, + error = %cause, + "an append to the heap failed; it was truncated back and the write refused" + ); + if out_of_room { + return HeapError::Full { + path: self.path.clone(), + reason: format!("there is no room left on the medium ({cause})"), + }; + } + HeapError::Io { + path: self.path.clone(), + source: cause, + } + } + + /// Takes the file lock, recovering from a poisoned mutex. + fn lock(&self) -> std::sync::MutexGuard<'_, Inner> { + self.inner + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + } + + /// Takes the buffer lock, recovering from a poisoned mutex. + fn buffer(&self) -> std::sync::MutexGuard<'_, Buffer> { + self.buffer + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + } +} + +impl Drop for Heap { + fn drop(&mut self) { + // Best effort and unreported, the same as the journal's: a destructor + // has nobody to tell, and the alternative to trying is a clean + // shutdown that leaves the window open for no reason. + let _ = self.sync(); + } +} + +/// Bodies held in front of the file, keyed by offset and bounded in bytes. +/// +/// Least-recently-used, because the access pattern this exists for is a +/// repository rebuilt on every write: the records of the repository being +/// written are exactly the ones wanted again, and a budget that evicts them +/// in favour of another repository's would answer none of them. +#[derive(Debug)] +struct Buffer { + budget: usize, + held: usize, + /// Offset to the body and the tick it was last touched at. + bodies: HashMap>, u64)>, + /// Tick to offset, which is the eviction order. + order: BTreeMap, + tick: u64, +} + +impl Buffer { + fn new(budget: usize) -> Self { + Self { + budget, + held: 0, + bodies: HashMap::new(), + order: BTreeMap::new(), + tick: 0, + } + } + + fn set_budget(&mut self, budget: usize) { + self.budget = budget; + self.evict(); + } + + fn get(&mut self, offset: u64) -> Option>> { + let (body, was) = self.bodies.get(&offset)?; + let (body, was) = (body.clone(), *was); + self.order.remove(&was); + self.tick += 1; + let now = self.tick; + self.order.insert(now, offset); + if let Some(slot) = self.bodies.get_mut(&offset) { + slot.1 = now; + } + Some(body) + } + + fn put(&mut self, offset: u64, body: Arc>) { + // A body bigger than the whole budget is not held at all, rather than + // held by evicting everything else to make room for one value nothing + // is likely to ask for twice. + if body.len() > self.budget { + return; + } + if let Some((old, was)) = self.bodies.remove(&offset) { + self.held -= old.len(); + self.order.remove(&was); + } + self.tick += 1; + let now = self.tick; + self.held += body.len(); + self.bodies.insert(offset, (body, now)); + self.order.insert(now, offset); + self.evict(); + } + + /// Drops everything at or past `from`, for a reclaim. + fn forget_from(&mut self, from: u64) { + let gone: Vec = self + .bodies + .keys() + .copied() + .filter(|offset| *offset >= from) + .collect(); + for offset in gone { + if let Some((body, was)) = self.bodies.remove(&offset) { + self.held -= body.len(); + self.order.remove(&was); + } + } + } + + fn evict(&mut self) { + while self.held > self.budget { + let Some((&oldest, &offset)) = self.order.iter().next() else { + break; + }; + self.order.remove(&oldest); + if let Some((body, _)) = self.bodies.remove(&offset) { + self.held -= body.len(); + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn scratch(name: &str) -> PathBuf { + let dir = std::env::temp_dir().join(format!( + "didbot-heap-{name}-{}-{:?}", + std::process::id(), + std::thread::current().id() + )); + let _ = std::fs::remove_dir_all(&dir); + dir + } + + #[test] + fn a_body_comes_back_out_of_the_slot_it_went_in_at() { + let dir = scratch("round-trip"); + let heap = Heap::open(&dir).expect("a heap opens"); + let first = heap + .append(b"quernstone", Durability::Sync) + .expect("append"); + let second = heap.append(b"marlpit", Durability::Sync).expect("append"); + assert_eq!(first.offset, 0); + assert_eq!(second.offset, first.end()); + assert_eq!(**heap.read(first).expect("read"), b"quernstone".to_vec()); + assert_eq!(**heap.read(second).expect("read"), b"marlpit".to_vec()); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } + + /// The boot-time decision, in both directions and with no read behind it. + #[test] + fn a_slot_past_the_end_is_beyond_and_one_inside_is_not() { + let dir = scratch("bounds"); + let heap = Heap::open(&dir).expect("a heap opens"); + let slot = heap.append(b"sillion", Durability::Sync).expect("append"); + assert_eq!(heap.bounds(slot), Bounds::Inside); + assert_eq!( + heap.bounds(Slot { + offset: slot.end(), + len: 1 + }), + Bounds::Beyond + ); + assert_eq!( + heap.bounds(Slot { + offset: 0, + len: u64::MAX + }), + Bounds::Beyond + ); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } + + /// A frame inside the heap that does not check out is refused rather + /// than served, and the buffer is not what is answering. + #[test] + fn a_damaged_frame_is_refused_rather_than_served() { + let dir = scratch("damaged"); + let heap = Heap::open(&dir).expect("a heap opens"); + let slot = heap.append(b"weftling", Durability::Sync).expect("append"); + let path = heap.path().to_path_buf(); + drop(heap); + + let mut bytes = std::fs::read(&path).expect("the heap reads"); + let last = bytes.len() - 1; + bytes[last] ^= 0b1000_0000; + std::fs::write(&path, &bytes).expect("damage it"); + + let heap = Heap::open(&dir).expect("a heap opens"); + assert_eq!(heap.bounds(slot), Bounds::Inside); + assert!(matches!( + heap.read(slot), + Err(HeapError::Damaged { offset: 0, .. }) + )); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } + + #[test] + fn a_reclaim_drops_the_bytes_no_slot_names() { + let dir = scratch("reclaim"); + let heap = Heap::open(&dir).expect("a heap opens"); + let kept = heap + .append(b"quernstone", Durability::Sync) + .expect("append"); + let loose = heap + .append(b"farthingale", Durability::Sync) + .expect("append"); + assert_eq!( + heap.reclaim(kept.end()).expect("reclaim"), + loose.end() - kept.end() + ); + assert_eq!(heap.len(), kept.end()); + assert_eq!(heap.bounds(loose), Bounds::Beyond); + assert_eq!(**heap.read(kept).expect("read"), b"quernstone".to_vec()); + // Reclaiming again finds nothing, which is what makes a boot that + // crashed halfway through one safe to repeat. + assert_eq!(heap.reclaim(kept.end()).expect("reclaim"), 0); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } + + /// The budget is a bound on what is held, and a heap with none of it + /// still answers every read. + #[test] + fn the_buffer_holds_what_it_has_room_for_and_answers_either_way() { + let dir = scratch("buffer"); + let heap = Heap::open(&dir).expect("a heap opens"); + heap.set_buffer_bytes(32); + let mut slots = Vec::new(); + for n in 0..16u32 { + slots.push( + heap.append(&n.to_le_bytes(), Durability::Sync) + .expect("append"), + ); + } + assert!( + heap.buffered_bytes() <= 32, + "the buffer held {} bytes of its 32 byte budget", + heap.buffered_bytes() + ); + heap.set_buffer_bytes(0); + assert_eq!(heap.buffered_bytes(), 0); + for (n, slot) in slots.iter().enumerate() { + assert_eq!( + **heap.read(*slot).expect("a read off the disk"), + (n as u32).to_le_bytes().to_vec() + ); + } + assert_eq!(heap.buffered_bytes(), 0, "a budget of zero holds nothing"); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } + + /// A value that will not fit a frame is refused, and the heap is + /// unchanged. + #[test] + fn a_value_too_large_for_a_frame_is_refused_and_nothing_is_written() { + let dir = scratch("too-large"); + let heap = Heap::open(&dir).expect("a heap opens"); + let body = vec![0u8; frame::MAX_PAYLOAD + 1]; + assert!(matches!( + heap.append(&body, Durability::Sync), + Err(HeapError::TooLarge { .. }) + )); + assert!(heap.is_empty()); + drop(heap); + let _ = std::fs::remove_dir_all(&dir); + } +} diff --git a/crates/didbot-pds/src/journal.rs b/crates/didbot-pds/src/journal.rs index 9521a457..f2cd9311 100644 --- a/crates/didbot-pds/src/journal.rs +++ b/crates/didbot-pds/src/journal.rs @@ -45,7 +45,8 @@ use std::path::Path; -use crate::records::MemoryRecordStore; +use crate::heap::Heap; +use crate::records::RecordStore; use crate::wal::Entry; /// A position in the journal. @@ -72,10 +73,15 @@ impl Mark { /// Handed to [`Derived::observe`] as the stores stand *before* the entry /// being observed is applied, which is what lets a derivation be taken over /// the value an entry is about to replace. -#[derive(Debug, Clone, Copy)] +#[derive(Clone, Copy)] pub struct Stores<'a> { /// Records, under the keys they were written with. - pub records: &'a MemoryRecordStore, + /// + /// The trait rather than one implementation of it: a record's value is + /// behind a heap in the deployment this is handed to and resident in the + /// one a store's own tests build, and a derivation taken over a record + /// does not care which. + pub records: &'a dyn RecordStore, } /// A store whose state is the log replayed into it. @@ -132,11 +138,22 @@ pub trait Journaled: Send + Sync { /// naming bytes nothing has, which is a read that cannot be served. pub trait Bodied: Journaled { /// Where this store's values live under the data directory. - /// - /// The path a locator is resolved against: the root of a tree for a store - /// that files one value per name, and one file for a store that appends - /// values into it. - fn bodies(&self) -> &Path; + fn bodies(&self) -> Bodies<'_>; +} + +/// The medium a [`Bodied`] store's locators are resolved against. +/// +/// Two shapes, because the two stores that have bodies keep them differently +/// and the difference is visible to anything that wants to read one. A store +/// that files one value per name hands out a directory and a locator is a +/// path inside it; a store that appends values hands out the [`Heap`] itself, +/// and a locator is a [`Slot`](crate::heap::Slot) it can be read at. +#[derive(Debug, Clone, Copy)] +pub enum Bodies<'a> { + /// One file per value, under this root. + Tree(&'a Path), + /// Every value appended into this heap. + Heap(&'a Heap), } /// A store whose state is a function of another store's. diff --git a/crates/didbot-pds/src/lib.rs b/crates/didbot-pds/src/lib.rs index 1050584d..75bce477 100644 --- a/crates/didbot-pds/src/lib.rs +++ b/crates/didbot-pds/src/lib.rs @@ -74,6 +74,7 @@ pub mod estop; pub mod evaluation_log; pub mod export; pub mod format; +pub mod heap; pub mod history; pub mod journal; pub mod kind; @@ -127,6 +128,7 @@ pub use estop::{ Cause as EstopCause, Estop, Halted, Mode as EstopMode, Refusal as EstopRefusal, Status as EstopStatus, }; +pub use heap::{Heap, HeapError, Slot}; pub use kind::{AccountKind, NameProvenance, Retention}; pub use oauth::{ Grant, GrantError, GrantRequest, Lifetimes, MemoryOAuthGrantStore, MintedGrant, @@ -167,9 +169,9 @@ pub use provision::{ pub use records::{ blob_refs, content_id, key_strategy, resolve_key, schema_is_held, validate, validate_caller, validate_collection, validate_record_key, validate_removal, validate_shape, validate_with, - validation_status, BatchOp, BatchOutcome, KeyStrategy, ListParams, MemoryRecordStore, - Precondition, RecordError, RecordKeyError, RecordStats, RecordStore, ResolvedKey, Stance, - Written, DEFAULT_LIST_LIMIT, + validation_status, BatchOp, BatchOutcome, HeapRecordStore, Held, KeyStrategy, ListParams, + MemoryRecordStore, Precondition, RecordError, RecordKeyError, RecordStats, RecordStore, + ResolvedKey, Stance, Written, DEFAULT_LIST_LIMIT, }; pub use registration::COLLECTION as ACTOR_REGISTRATION_COLLECTION; pub use session::{ diff --git a/crates/didbot-pds/src/records.rs b/crates/didbot-pds/src/records.rs index fb2c3041..5af0b485 100644 --- a/crates/didbot-pds/src/records.rs +++ b/crates/didbot-pds/src/records.rs @@ -9,11 +9,12 @@ use std::collections::{BTreeMap, BTreeSet}; use std::ops::Bound; -use std::sync::Mutex; +use std::sync::{Arc, Mutex}; use didbot_data::Cid; use serde_json::Value; +use crate::heap::{Heap, HeapError, Slot}; use crate::wal::Entry; /// How many records a listing returns when the caller does not say. @@ -1407,16 +1408,6 @@ impl MemoryRecordStore { .unwrap_or_else(|poisoned| poisoned.into_inner()) } - /// The next record key, without writing anything. - /// - /// Split out of [`RecordStore::put`] for the durable store, which has to - /// know the key before the write: the key is in the log entry, because a - /// replay that minted a new one would move every `at://` URI a client had - /// already been handed. - pub(crate) fn mint_key(&self) -> String { - self.keys.next() - } - /// Stores a record under a key that has already been decided. /// /// Validation is the caller's, because the caller has already done it: @@ -1824,3 +1815,687 @@ mod journal_tests { ); } } + +// --------------------------------------------------------------------------- +// The record store whose values are not resident +// --------------------------------------------------------------------------- + +/// One record's location and its name. +/// +/// What the three-level map holds where [`MemoryRecordStore`] holds a +/// [`Value`]. The slot says where the bytes are; the CID says which record +/// they are, and it is carried rather than recomputed so that a +/// [`Precondition`] is a comparison against this and never a read. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct Held { + /// Where the body's frame is in the heap. + pub slot: Slot, + /// The content identifier of those bytes. + pub cid: Cid, +} + +/// One collection's records, keyed by rkey, holding slots. +type HeldCollection = BTreeMap; + +/// One repository's collections, keyed by NSID, holding slots. +type HeldRepo = BTreeMap; + +/// What a replay of a journal into a heap-backed store found. +/// +/// Accumulated during [`Journaled::apply`](crate::journal::Journaled::apply), +/// which is total and has nowhere to +/// report to, and read once afterwards by [`HeapRecordStore::settle`]. +#[derive(Debug, Default)] +struct Recovery { + /// The first slot found past the end of the heap, if there was one. + torn: Option, + /// Why the directory cannot be opened, if it cannot. + refusal: Option<&'static str>, + /// The slot reaching furthest into the heap that any entry named. + highest: Option, +} + +/// Records, with their values in a [`Heap`] rather than in memory. +/// +/// The key structure is [`MemoryRecordStore`]'s, unchanged: `(did, +/// collection, rkey)` finds a value in three nested [`BTreeMap`]s, and the +/// inner map's ordering by record key is still the chronological ordering a +/// listing pages through. What changed is what the innermost map holds — a +/// [`Held`] rather than a [`Value`] — and so what a read does: the map finds +/// a slot and the heap turns it into bytes. +/// +/// **What that claim is and is not.** The keys are resident and so are the +/// slots. A record costs its three keys, their share of three `BTreeMap` +/// nodes, a slot and a content identifier, rather than the whole of its +/// decoded value. Capacity stops being a wall at total record *bytes* and +/// becomes a slope in key *count*; it is not independence from the record +/// count, and nothing here claims it is. +#[derive(Debug)] +pub struct HeapRecordStore { + repos: Mutex>, + keys: Minter, + heap: Arc, + recovery: Mutex, +} + +impl HeapRecordStore { + /// A store over `heap`, holding no keys yet. + pub fn new(heap: Arc) -> Self { + Self { + repos: Mutex::new(BTreeMap::new()), + keys: Minter::default(), + heap, + recovery: Mutex::new(Recovery::default()), + } + } + + /// The heap the values live in. + pub fn heap(&self) -> &Arc { + &self.heap + } + + /// Takes the lock, recovering from a poisoned mutex. + fn repos(&self) -> std::sync::MutexGuard<'_, BTreeMap> { + self.repos + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + } + + /// Takes the recovery lock, recovering from a poisoned mutex. + fn recovery(&self) -> std::sync::MutexGuard<'_, Recovery> { + self.recovery + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) + } + + /// The next record key, without writing anything. + /// + /// Split out of the write for the durable store, which has to know the + /// key before the entry carrying it is appended: a replay that minted a + /// new one would move every `at://` URI a client has already been + /// handed. + pub(crate) fn mint_key(&self) -> String { + self.keys.next() + } + + /// What is at a key, without reading the body. + /// + /// The precondition check, and the reason [`Entry::RecordStored`] carries + /// a content identifier: the comparison a `swapRecord` makes is against + /// this value, so a compare-and-swap costs a map lookup rather than a + /// disk read and a rehash. + pub(crate) fn cid_at(&self, did: &str, collection: &str, rkey: &str) -> Option { + Some( + self.repos() + .get(did)? + .get(collection)? + .get(rkey)? + .cid + .clone(), + ) + } + + /// Files a slot under a key that has already been decided. + /// + /// [`MemoryRecordStore::insert_at`]'s counterpart: validation is the + /// caller's, because this is the replay and durable-write path and a + /// record in the log was validated before it was written. + pub(crate) fn insert_slot( + &self, + did: &str, + collection: &str, + rkey: String, + held: Held, + ) -> Option { + self.keys.observe(&rkey); + self.note_highest(held.slot); + self.repos() + .entry(did.to_owned()) + .or_default() + .entry(collection.to_owned()) + .or_default() + .insert(rkey, held) + } + + /// Drops a key without checking a precondition, the counterpart to + /// [`HeapRecordStore::insert_slot`]. + /// + /// The heap keeps the body. Bytes no key names are dead space rather than + /// a hole: they are reachable only through a slot, and nothing holds this + /// one any more. + pub(crate) fn remove_at(&self, did: &str, collection: &str, rkey: &str) -> Option { + let mut repos = self.repos(); + let repo = repos.get_mut(did)?; + let records = repo.get_mut(collection)?; + let gone = records.remove(rkey); + if records.is_empty() { + repo.remove(collection); + } + gone + } + + /// Notes a slot for the reclaim, which keeps everything below the highest + /// frame any entry named. + fn note_highest(&self, slot: Slot) { + let mut recovery = self.recovery(); + if recovery.highest.is_none_or(|held| slot.end() > held.end()) { + recovery.highest = Some(slot); + } + } + + /// The value at a slot, or nothing with a line saying why. + /// + /// A frame that does not check out is refused to the read that wanted it + /// rather than answered with whatever the bytes happened to say — see + /// [`crate::heap`]. A reader cannot act differently on "the key is not + /// there" and "the bytes under the key will not check out", so both are + /// `None`; the difference is in the log, and the second one is loud. + fn body(&self, held: &Held) -> Option { + let bytes = match self.heap.read(held.slot) { + Ok(bytes) => bytes, + Err(error) => { + tracing::error!( + offset = held.slot.offset, + len = held.slot.len, + cid = %held.cid, + %error, + "a record's body could not be read out of the heap" + ); + return None; + } + }; + match didbot_data::dag_cbor::decode(&bytes) { + Ok(value) => Some(value.to_json()), + Err(error) => { + tracing::error!( + offset = held.slot.offset, + len = held.slot.len, + cid = %held.cid, + error = ?error, + "a record's body checked out and is not in the data model" + ); + None + } + } + } + + /// Appends a body and files it, for a store used without a journal. + fn store( + &self, + did: &str, + collection: &str, + rkey: String, + record: &Value, + ) -> Result { + let (body, cid) = encode_record(record)?; + let slot = self + .heap + .append(&body, crate::wal::Durability::Sync) + .map_err(heap_backend)?; + self.insert_slot( + did, + collection, + rkey, + Held { + slot, + cid: cid.clone(), + }, + ); + Ok(cid) + } + + /// Closes the recovery a replay opened, and reclaims the heap's tail. + /// + /// Called once, after the whole journal has been applied and before + /// anything is served. Three things happen here and they are in this + /// order for a reason. + /// + /// A journal that named a slot inside the heap *after* one past its end + /// is a journal and a heap that disagree about their own write order, and + /// there is nothing to reconstruct from that: it is refused. + /// + /// Otherwise the highest-offset frame any entry named is read, because it + /// is the only frame an interrupted append can have torn — every earlier + /// one had a later append land on top of it. A frame that fails there is + /// dropped exactly as a slot past the end is. + /// + /// Then everything above that frame goes. Those are bodies whose journal + /// entry never landed, which is the harmless half of the ordering: bytes + /// nothing names, swept at boot the way `.incoming/` is. + pub fn settle(&self) -> Result { + let (refusal, torn, highest) = { + let mut recovery = self.recovery(); + ( + recovery.refusal.take(), + recovery.torn.take(), + recovery.highest, + ) + }; + if let Some(reason) = refusal { + return Err(HeapError::Damaged { + path: self.heap.path().to_path_buf(), + offset: torn.map_or(0, |slot| slot.offset), + len: torn.map_or(0, |slot| slot.len), + reason, + }); + } + let keep = match highest { + None => 0, + Some(slot) => match self.heap.verify(slot) { + Ok(()) => slot.end(), + Err(error) => { + tracing::warn!( + offset = slot.offset, + len = slot.len, + %error, + "the last body the journal names is an interrupted append; dropping the \ + record that names it" + ); + self.forget_slot(slot); + slot.offset + } + }, + }; + let reclaimed = self.heap.reclaim(keep)?; + if reclaimed > 0 { + tracing::info!( + reclaimed, + keep, + "reclaimed record bodies no journal entry names" + ); + } + Ok(reclaimed) + } + + /// Drops every key holding `slot`, for a body that will not read back. + fn forget_slot(&self, slot: Slot) { + let mut repos = self.repos(); + for repo in repos.values_mut() { + for records in repo.values_mut() { + records.retain(|_, held| held.slot != slot); + } + repo.retain(|_, records| !records.is_empty()); + } + } + + /// Every key it holds and where its body is, in repository, collection + /// and key order. + pub(crate) fn drain_slots(&self) -> Vec<(String, String, String, Held)> { + let repos = self.repos(); + let mut out = Vec::new(); + for (did, repo) in repos.iter() { + for (collection, records) in repo { + for (rkey, held) in records { + out.push((did.clone(), collection.clone(), rkey.clone(), held.clone())); + } + } + } + out + } +} + +/// A record's canonical DAG-CBOR and the identifier over those exact bytes. +/// +/// One function, because the two are one fact: the bytes the heap holds and +/// the name the journal carries have to be the same encoding or a read cannot +/// prove it got what it asked for. It is [`content_id`]'s encoding, and +/// `tests/durability.rs`'s +/// `a_heap_frames_payload_round_trips_as_canonical_dag_cbor` is what pins the +/// two together. +pub fn encode_record(record: &Value) -> Result<(Vec, Cid), RecordError> { + let value = didbot_data::Value::object_from_json(record) + .map_err(|source| RecordError::UnrepresentableRecord { source })?; + let body = didbot_data::dag_cbor::encode(&value); + let cid = didbot_data::Cid::of_dag_cbor(&body); + Ok((body, cid)) +} + +/// A heap failure, as the refusal a caller can act on. +/// +/// [`record_backend`](crate::durable) for the journal, and the same split for +/// the same reason: a full disk and a record too large are things a client +/// does something about, and a backend error is not. +pub(crate) fn heap_backend(error: HeapError) -> RecordError { + match error { + HeapError::TooLarge { .. } => RecordError::TooLarge { + reason: error.to_string(), + }, + HeapError::Full { .. } => RecordError::StorageFull { + reason: error.to_string(), + }, + other => RecordError::Backend(other.to_string()), + } +} + +impl RecordStore for HeapRecordStore { + fn put( + &self, + did: &str, + collection: &str, + rkey: Option<&str>, + record: Value, + expect: &Precondition, + ) -> Result { + validate(collection, &record)?; + let rkey = match resolve_key(collection, rkey)? { + ResolvedKey::Mint => self.keys.next(), + ResolvedKey::Fixed(fixed) => fixed, + }; + // The check is against the slot's stored identifier, so it costs a + // map lookup rather than a read of the body and a rehash of it. + expect.check(self.cid_at(did, collection, &rkey).as_ref())?; + let cid = self.store(did, collection, rkey.clone(), &record)?; + Ok(Written { rkey, cid }) + } + + fn put_authored( + &self, + did: &str, + collection: &str, + record: Value, + ) -> Result { + validate_shape(collection, &record)?; + let rkey = match resolve_key(collection, None)? { + ResolvedKey::Mint => self.keys.next(), + ResolvedKey::Fixed(fixed) => fixed, + }; + self.store(did, collection, rkey.clone(), &record)?; + Ok(rkey) + } + + fn list(&self, did: &str, collection: &str, params: &ListParams) -> Vec<(String, Value)> { + // The keys are taken under the lock and the bodies are read outside + // it: a page is at most `limit` records and a heap read may reach the + // disk, which is not something to hold every other writer behind. + let page: Vec<(String, Held)> = { + let repos = self.repos(); + let Some(records) = repos.get(did).and_then(|repo| repo.get(collection)) else { + return Vec::new(); + }; + let page: Box> = + match (¶ms.cursor, params.reverse) { + (Some(cursor), false) => Box::new(records.range(..cursor.clone()).rev()), + (Some(cursor), true) => { + Box::new(records.range((Bound::Excluded(cursor.clone()), Bound::Unbounded))) + } + (None, false) => Box::new(records.iter().rev()), + (None, true) => Box::new(records.iter()), + }; + page.take(params.limit) + .map(|(rkey, held)| (rkey.clone(), held.clone())) + .collect() + }; + page.into_iter() + .filter_map(|(rkey, held)| self.body(&held).map(|record| (rkey, record))) + .collect() + } + + fn get(&self, did: &str, collection: &str, rkey: &str) -> Option { + let held = self.repos().get(did)?.get(collection)?.get(rkey)?.clone(); + self.body(&held) + } + + fn remove( + &self, + did: &str, + collection: &str, + rkey: &str, + expect: &Precondition, + ) -> Result { + validate_removal(collection)?; + let mut repos = self.repos(); + let Some(repo) = repos.get_mut(did) else { + expect.check(None)?; + return Ok(false); + }; + let Some(records) = repo.get_mut(collection) else { + expect.check(None)?; + return Ok(false); + }; + expect.check(records.get(rkey).map(|held| &held.cid))?; + let removed = records.remove(rkey).is_some(); + if records.is_empty() { + repo.remove(collection); + } + Ok(removed) + } + + fn collections(&self, did: &str) -> Vec { + self.repos() + .get(did) + .map(|repo| repo.keys().cloned().collect()) + .unwrap_or_default() + } + + fn snapshot(&self, did: &str) -> Vec<(String, String, Value)> { + // The same two phases `list` has, for the same reason and with more + // riding on it: this is what a commit is signed over, so it is every + // record in the repository and it runs on every write. + let held: Vec<(String, String, Held)> = { + let repos = self.repos(); + let Some(repo) = repos.get(did) else { + return Vec::new(); + }; + repo.iter() + .flat_map(|(collection, records)| { + records + .iter() + .map(|(rkey, held)| (collection.clone(), rkey.clone(), held.clone())) + }) + .collect() + }; + held.into_iter() + .filter_map(|(collection, rkey, held)| { + self.body(&held).map(|record| (collection, rkey, record)) + }) + .collect() + } + + fn remove_repo(&self, did: &str) { + self.repos().remove(did); + } + + /// What is held, counted — with the byte total counted as stored bytes. + /// + /// [`RecordStats::bytes`] is the size of the bodies in the heap, which is + /// their canonical DAG-CBOR, rather than the compact JSON a resident + /// store measures. It is the number an operator watching this deployment + /// can act on: it is what the heap file grows by. + fn stats(&self) -> RecordStats { + let mut stats = RecordStats::default(); + for repo in self.repos().values() { + for (collection, records) in repo { + stats.records += records.len(); + stats.bytes += records + .values() + .map(|held| held.slot.len as usize) + .sum::(); + *stats.collections.entry(collection.clone()).or_default() += records.len(); + } + } + stats + } + + fn apply_batch( + &self, + did: &str, + ops: &[BatchOp], + stance: Stance, + ) -> Result, RecordError> { + let planned = plan_batch(ops, stance, || self.keys.next())?; + let mut results = Vec::with_capacity(planned.len()); + for (collection, rkey, value) in planned { + match value { + Some((record, cid)) => { + self.store(did, &collection, rkey.clone(), &record)?; + results.push(BatchOutcome::Written(Written { rkey, cid })); + } + None => { + self.remove_at(did, &collection, &rkey); + results.push(BatchOutcome::Deleted); + } + } + } + Ok(results) + } +} + +impl crate::journal::Journaled for HeapRecordStore { + const STORE: &'static str = "records"; + + fn owns(entry: &Entry) -> bool { + matches!( + entry, + Entry::RecordStored { .. } | Entry::RecordPut { .. } | Entry::RecordRemoved { .. } + ) + } + + fn apply(&self, entry: Entry) { + match entry { + Entry::RecordStored { + did, + collection, + rkey, + cid, + offset, + len, + } => self.replay_slot(&did, &collection, rkey, &cid, Slot { offset, len }), + // A log written before the bodies moved. The value is in the + // entry, so it is appended to the heap here and the key gets a + // slot like any other: a directory holding a mixture of the two + // is what every build from the format's own pull onward reads. + Entry::RecordPut { + did, + collection, + rkey, + record, + } => { + if let Err(error) = self.store(&did, &collection, rkey.clone(), &record) { + tracing::error!( + did, collection, rkey, %error, + "a record the log carries in full could not be moved into the heap" + ); + } + } + Entry::RecordRemoved { + did, + collection, + rkey, + } => { + // Replay is not conditional on anything: the log is a record + // of what already happened, and re-running the guard against + // it would let a tightened rule resurrect a record a previous + // version of this server deleted. + self.remove_at(&did, &collection, &rkey); + } + // Only what `owns` claims reaches here. + _ => {} + } + } + + /// One entry per record, naming where its body already is. + /// + /// A checkpoint of this store rewrites the journal and never the heap, so + /// the slots it writes are the slots that are live now and the bodies + /// they name do not move. + fn checkpoint(&self) -> Vec { + self.drain_slots() + .into_iter() + .map(|(did, collection, rkey, held)| Entry::RecordStored { + did, + collection, + rkey, + cid: held.cid.to_string(), + offset: held.slot.offset, + len: held.slot.len, + }) + .collect() + } + + fn durable_through(&self) -> crate::journal::Mark { + crate::journal::Mark::ORIGIN + } +} + +impl crate::journal::Bodied for HeapRecordStore { + fn bodies(&self) -> crate::journal::Bodies<'_> { + crate::journal::Bodies::Heap(&self.heap) + } +} + +impl HeapRecordStore { + /// Files one replayed slot, deciding first whether the heap holds it. + /// + /// The whole of the boot-time recovery rule, and it reads nothing: a slot + /// is either wholly inside the heap's length or it is not. See + /// [`HeapRecordStore::settle`], which acts on what this accumulates. + fn replay_slot(&self, did: &str, collection: &str, rkey: String, cid: &str, slot: Slot) { + let Ok(cid) = Cid::parse(cid) else { + tracing::error!( + did, + collection, + rkey, + "a journal entry names a record by something that is not a content identifier" + ); + return; + }; + match self.heap.bounds(slot) { + crate::heap::Bounds::Beyond => { + let mut recovery = self.recovery(); + if recovery.torn.is_none() { + tracing::warn!( + did, + collection, + rkey, + offset = slot.offset, + len = slot.len, + heap = self.heap.len(), + "the journal names a record body past the end of the heap; the write \ + never reached the medium and the record is dropped" + ); + recovery.torn = Some(slot); + } + } + crate::heap::Bounds::Inside => { + if self.recovery().torn.is_some() { + // A slot back inside the heap after one past its end. + // The two files disagree about their own order, and + // there is no state to reconstruct from that. + self.recovery().refusal = + Some("the journal names a body inside the heap after one past its end"); + return; + } + self.insert_slot(did, collection, rkey, Held { slot, cid }); + } + } + } +} + +impl HeapRecordStore { + /// The body a replayed [`Entry::RecordStored`] names, when the heap holds + /// it. + /// + /// For the one reader that has to see a replayed record as a value rather + /// than as a slot: the blob index counts the blobs a record references, + /// and it counts them by reading the record. `None` when the heap does + /// not hold the body, which is the torn tail the recovery rule drops. + pub(crate) fn replayed_body(&self, entry: &Entry) -> Option { + let Entry::RecordStored { + cid, offset, len, .. + } = entry + else { + return None; + }; + let slot = Slot { + offset: *offset, + len: *len, + }; + if self.heap.bounds(slot) != crate::heap::Bounds::Inside { + return None; + } + self.body(&Held { + slot, + cid: Cid::parse(cid).ok()?, + }) + } +} diff --git a/crates/didbot-pds/src/wal/mod.rs b/crates/didbot-pds/src/wal/mod.rs index 97cb50ab..d3161d76 100644 --- a/crates/didbot-pds/src/wal/mod.rs +++ b/crates/didbot-pds/src/wal/mod.rs @@ -740,6 +740,14 @@ pub enum WalError { /// [`MAX_ENTRY`]. limit: usize, }, + /// A store's bodies would not open, or do not agree with the log. + /// + /// A [`Bodied`](crate::journal::Bodied) store keeps its values in a + /// second file under the same directory, and a directory whose bodies + /// this binary cannot read is a directory it cannot serve — the same + /// answer this module gives for its own file. See [`crate::heap`]. + #[error(transparent)] + Bodies(#[from] crate::heap::HeapError), /// There is no room to write the entry down. /// /// A configured capacity was reached, the medium filled, or a previous @@ -1008,6 +1016,20 @@ impl Wal { Ok(()) } + /// Whether the next append will sync whatever durability it asks for. + /// + /// The deferred window, answered rather than guessed at. A + /// [`Bodied`](crate::journal::Bodied) store has to put its bytes on the + /// medium in front of the fact that names them, and the moment that + /// matters is the append that closes this window: asking here is what + /// makes the ordering exact instead of a race between two timers. + pub fn sync_due(&self) -> bool { + let inner = self.lock(); + inner + .dirty_since + .is_some_and(|since| Instant::now().duration_since(since) >= DEFAULT_SYNC_INTERVAL) + } + /// Waits for everything appended so far to reach the medium. /// /// Called on a clean shutdown, so the deferred window closes on the way diff --git a/crates/didbot-pds/tests/durability.rs b/crates/didbot-pds/tests/durability.rs index f499428a..ad546264 100644 --- a/crates/didbot-pds/tests/durability.rs +++ b/crates/didbot-pds/tests/durability.rs @@ -850,7 +850,7 @@ fn a_log_torn_mid_record_loses_only_the_torn_record() { .rposition(|payload| { serde_json::from_slice::(payload) .ok() - .and_then(|entry| entry["op"].as_str().map(|op| op == "recordPut")) + .and_then(|entry| entry["op"].as_str().map(|op| op == "recordStored")) .unwrap_or(false) }) .expect("the log holds record writes"); diff --git a/crates/didbot-pds/tests/upgrade.rs b/crates/didbot-pds/tests/upgrade.rs index d0d213ab..94bb121f 100644 --- a/crates/didbot-pds/tests/upgrade.rs +++ b/crates/didbot-pds/tests/upgrade.rs @@ -391,6 +391,12 @@ fn the_data_directory_layout_is_pinned() { "pds.layout", "pds.lock", "pds.wal", + // The record store's own directory and the heap its bodies are + // appended into. Both are `didbot_pds::layout::NAME_SHAPE` rows + // already, so the stamp names them; this is where they appear on a + // disk a server has written. + "records/", + "records/pds.heap", ] .into_iter() .map(str::to_owned) -- 2.51.2