Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ No version has been released yet. Pre-releases are published as `0.1.0-beta.N`;
- **Snapshots** — `snapshot_to`, `restore_from`, and incremental apply, each authenticated against the state its manifest describes. Destinations must be empty; malformed or incomplete artifacts fail closed.
- **Recovery** — open-flow GC, apply-journal replay, deep-walk `fsck`, and the `pagedb-fsck` binary.
- **Online rekey** — rekey under a new key with mixed-cipher and mixed-epoch page coexistence; no full-file migration. A rotation whose source epoch is still pinned by a reader defers retiring it and completes the retirement once the reader set drains — including when the last reader leaves during the deferral itself, so the superseded master key never stays leasable behind an `Ok(())`. A retirement that cannot be taken yet stays queued for the next attempt without holding up the others.
- **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`.
- **Handle modes** — `Standalone`, `Follower`, `ReadOnly`, and `Observer`, each holding its own sentinel. A snapshot or restored directory shares its source's identity, so it opens only as `ReadOnly` or as a `Follower` (`open_follower`, `promote_to_follower`), and a `Standalone` open reports `RestoredNotPromoted`. `rekey_into_writer` forks it into an independent `Standalone` writer under a fresh identity.
- **Open refusals name the parameter, not the store** — `KeyMismatch`, `PageSizeMismatch`, and `RealmMismatch`, each decided before anything is read or written, and none reported as corruption.
- **Failures report themselves** — an unreadable free-list chain, main file, or segment catalog fails `stats()` instead of reporting zero; compaction never skips a catalog entry whose file it cannot open; segment open distinguishes a missing file from a permission or backend error; and only genuine contention is reported as contention. Persisted named-counter rows are validated at open, and commit-history keys are rejected unless exactly eight bytes.

Expand Down
5 changes: 5 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,11 @@ defend against, stated plainly so nobody has to infer it.
capability bit this build does not understand reports
`HeaderCapabilityUnsupported`: a newer build wrote the store, and it is
untouched. Open it with a build that understands the capability.
- **A restored directory is not a writer.** A snapshot directory, and a
directory `restore_from` fills, shares the source's key and nonce space. A
Standalone open of one reports `RestoredNotPromoted`. Open it with
`open_read_only`, track the source with `open_follower`, or call
`rekey_into_writer` to fork an independent writer under a fresh identity.
- **`main.db` is not reconstructible from `seg/`.** Segment files are
identity-keyed and the mapping from embedder name to segment id lives only in
the catalog inside `main.db`. Losing `main.db` while `seg/` survives is
Expand Down
4 changes: 2 additions & 2 deletions src/compaction/full.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ use crate::catalog::codec::{Catalog, CatalogRowKind, SegmentMeta};
use crate::errors::PagedbError;
use crate::segment::reader::SegmentReader;
use crate::segment::types::SegmentPageKind;
use crate::segment::writer::SegmentWriter;
use crate::segment::writer::{STAGING_DIR, SegmentWriter};
use crate::txn::db::{Db, WriterState};
use crate::vfs::{Vfs, VfsFile};

Expand Down Expand Up @@ -149,7 +149,7 @@ async fn repack_one_segment<V: Vfs + Clone>(
mmap_limit,
)
.await?;
db.vfs.mkdir_all("seg/.staging").await?;
db.vfs.mkdir_all(STAGING_DIR).await?;
let new_segment_id = crate::crypto::random::segment_id()?;
let mut writer = SegmentWriter::create_internal(
db.pager.clone(),
Expand Down
17 changes: 8 additions & 9 deletions src/compaction/helpers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,7 @@ pub(super) async fn replace_segment_compact<V: Vfs + Clone>(
format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION,
cipher_id: db.cipher_id.as_byte(),
page_size_log2: page_size_log2(db.page_size)?,
flags: 0,
flags: db.header_flags,
file_id: db.file_id,
kek_salt: db.kek_salt,
mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst),
Expand All @@ -202,10 +202,10 @@ pub(super) async fn replace_segment_compact<V: Vfs + Clone>(
apply_journal_root_version: 0,
commit_history_root_page_id: 0,
commit_history_root_version: 0,
restore_mode: 0,
restore_mode: state.restore_mode,
next_page_id: new_next,
commit_retain_policy_tag: 0,
commit_retain_policy_value: 0,
commit_retain_policy_tag: state.commit_retain_policy_tag,
commit_retain_policy_value: state.commit_retain_policy_value,
realm_id: db.realm_id,
};

Expand Down Expand Up @@ -272,15 +272,14 @@ pub(super) fn make_header_fields<V: Vfs + Clone>(
new_next: u64,
free_list_root_page_id: u64,
) -> MainDbHeaderFields {
let _ = state;
let mut catalog_root_bytes = [0u8; 16];
catalog_root_bytes[..8].copy_from_slice(&new_cat_root.to_le_bytes());
catalog_root_bytes[8..].copy_from_slice(&new_commit_id.to_le_bytes());
MainDbHeaderFields {
format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION,
cipher_id: db.cipher_id.as_byte(),
page_size_log2: page_size_log2(db.page_size).unwrap_or(12),
flags: 0,
flags: db.header_flags,
file_id: db.file_id,
kek_salt: db.kek_salt,
mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst),
Expand All @@ -295,10 +294,10 @@ pub(super) fn make_header_fields<V: Vfs + Clone>(
apply_journal_root_version: 0,
commit_history_root_page_id: 0,
commit_history_root_version: 0,
restore_mode: 0,
restore_mode: state.restore_mode,
next_page_id: new_next,
commit_retain_policy_tag: 0,
commit_retain_policy_value: 0,
commit_retain_policy_tag: state.commit_retain_policy_tag,
commit_retain_policy_value: state.commit_retain_policy_value,
realm_id: db.realm_id,
}
}
Expand Down
58 changes: 52 additions & 6 deletions src/pager/header.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,16 +5,17 @@
//! `seq` (HK-MAC-verified). A torn write to one slot leaves the other intact.

use crate::Result;
use crate::crypto::SecretKey;
use crate::crypto::kdf::{derive_hk, derive_mk};
use crate::crypto::keys::DerivedKey;
// `decode_main_db_header` and the corruption detail it reports are needed only
// by `open_header`, which is test-only: production openers decode per slot
// because the HK has to be derived from each slot's own salt and epoch.
// `CorruptionDetail` is used only by `open_header`, which is test-only.
// Production openers authenticate per slot because each slot's HK derives from its own salt and epoch.
#[cfg(test)]
use crate::errors::CorruptionDetail;
use crate::errors::PagedbError;
#[cfg(test)]
use crate::pager::format::structural_header::decode_main_db_header;
use crate::pager::format::structural_header::{MainDbHeaderFields, encode_main_db_header};
use crate::pager::format::structural_header::{
MainDbHeaderFields, decode_main_db_header, encode_main_db_header,
};
use crate::vfs::types::OpenMode;
use crate::vfs::{Vfs, VfsFile, read_exact_at, write_all_at};

Expand Down Expand Up @@ -109,6 +110,51 @@ pub(crate) async fn read_header_slot<F: VfsFile + ?Sized>(
}
}

/// Authenticate one header slot under `hk`.
///
/// Returns `None` when the slot fails to verify. A wrong key, a torn write,
/// and damage look identical.
///
/// Returns `Err` only for `HeaderCapabilityUnsupported`. The key already
/// authenticated the slot, so no other key or slot revokes its capability.
pub(crate) fn authenticate_slot(
slot: &[u8],
hk: &DerivedKey,
page_size: usize,
) -> Result<Option<MainDbHeaderFields>> {
match decode_main_db_header(slot, hk, page_size) {
Ok(fields) => Ok(Some(fields)),
Err(error @ PagedbError::HeaderCapabilityUnsupported { .. }) => Err(error),
Err(_) => Ok(None),
}
}

/// Authenticate one header slot under `kek`, as [`authenticate_slot`] does.
///
/// The HK derives from the KEK salt and MK epoch, stored unencrypted at
/// bytes `32..48` and `48..56`. Returns the HK with the fields so the
/// caller can rewrite the slot.
pub(crate) fn authenticate_slot_with_kek(
slot: &[u8],
kek: &SecretKey,
page_size: usize,
) -> Result<Option<(MainDbHeaderFields, DerivedKey)>> {
if slot.len() < 56 {
return Ok(None);
}
let mut salt = [0u8; 16];
salt.copy_from_slice(&slot[32..48]);
let mut epoch = [0u8; 8];
epoch.copy_from_slice(&slot[48..56]);
let Ok(mk) = derive_mk(kek.as_bytes(), &salt, u64::from_le_bytes(epoch)) else {
return Ok(None);
};
let Ok(hk) = derive_hk(&mk) else {
return Ok(None);
};
Ok(authenticate_slot(slot, &hk, page_size)?.map(|fields| (fields, hk)))
}

/// Reads both slots; verifies each via HK-MAC; picks the one with the
/// greater `seq`. If only one verifies, it wins. If neither verifies, returns
/// `Corruption(HeaderUnverifiable)` — unrecoverable from inside the header
Expand Down
70 changes: 70 additions & 0 deletions src/recovery/fork.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
//! Segments a `rekey_into_writer` fork wrote before its `main.db` was live.
//!
//! The fork writes its segments in [`FORK_DIR`], which the open-time orphan
//! scan skips. A fork that never publishes leaves the restored directory
//! openable. Once the fork's `main.db` is live, its catalog names those
//! segments. Moving them into staging hands them to catalog repair, which
//! publishes each one the catalog names and sweeps the rest.

use crate::Result;
use crate::errors::PagedbError;
use crate::segment::writer::{FORK_DIR, STAGING_DIR};
use crate::vfs::Vfs;

/// Move every fork segment into staging. Returns the number moved.
///
/// Callers must hold write authority over a store whose `main.db` is the
/// fork's. Callers must run catalog repair afterwards.
pub(crate) async fn adopt_forked_segments<V: Vfs>(vfs: &V) -> Result<u64> {
let names = forked_segment_names(vfs).await?;
if names.is_empty() {
return Ok(0);
}
vfs.mkdir_all(STAGING_DIR).await?;
let mut count: u64 = 0;
for name in names {
vfs.rename(
&format!("{FORK_DIR}/{name}"),
&format!("{STAGING_DIR}/{name}"),
)
.await?;
count += 1;
}
vfs.sync_dir(STAGING_DIR).await?;
vfs.sync_dir(FORK_DIR).await?;
Ok(count)
}

/// Remove every fork segment. Returns the number removed.
///
/// A fork calls this before it starts, to drop what an earlier attempt
/// wrote without publishing.
pub(crate) async fn discard_forked_segments<V: Vfs>(vfs: &V) -> Result<u64> {
let names = forked_segment_names(vfs).await?;
if names.is_empty() {
return Ok(0);
}
let mut count: u64 = 0;
for name in names {
vfs.remove(&format!("{FORK_DIR}/{name}")).await?;
count += 1;
}
vfs.sync_dir(FORK_DIR).await?;
Ok(count)
}

/// Names in the fork directory that this crate writes: segment identities.
/// Anything else is left alone.
async fn forked_segment_names<V: Vfs>(vfs: &V) -> Result<Vec<String>> {
let entries = match vfs.list_dir(FORK_DIR).await {
Ok(entries) => entries,
Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(Vec::new());
}
Err(error) => return Err(error),
};
Ok(entries
.into_iter()
.filter(|name| crate::hex::parse_hex::<16>(name).is_some())
.collect())
}
4 changes: 2 additions & 2 deletions src/recovery/journal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -310,7 +310,7 @@ pub async fn replay_apply_journal<V: Vfs + Clone>(
#[cfg(test)]
pub async fn execute_journal_actions<V: Vfs>(vfs: &V, actions: &[JournalAction]) -> Result<()> {
vfs.mkdir_all("seg").await?;
vfs.mkdir_all("seg/.staging").await?;
vfs.mkdir_all(crate::segment::writer::STAGING_DIR).await?;
vfs.mkdir_all("seg/.tombstone").await?;
for action in actions {
match action {
Expand Down Expand Up @@ -338,7 +338,7 @@ pub async fn execute_journal_actions<V: Vfs>(vfs: &V, actions: &[JournalAction])
}
}
vfs.sync_dir("seg").await?;
vfs.sync_dir("seg/.staging").await?;
vfs.sync_dir(crate::segment::writer::STAGING_DIR).await?;
vfs.sync_dir("seg/.tombstone").await
}

Expand Down
1 change: 1 addition & 0 deletions src/recovery/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
//! tombstone GC, spill-scratch reclamation.

pub(crate) mod deep_walk;
pub(crate) mod fork;
pub(crate) mod gc;
pub(crate) mod journal;
pub(crate) mod provenance;
Expand Down
9 changes: 2 additions & 7 deletions src/recovery/quarantine/preserve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

use crate::Result;
use crate::errors::PagedbError;
use crate::segment::writer::STAGING_DIR;
use crate::vfs::Vfs;

/// Where a quarantined store went, and what moved.
Expand All @@ -43,13 +44,7 @@ pub struct QuarantineReport {
/// So the layout this crate writes is enumerated rather than discovered.
/// Anything a backend does report is still moved; this only ensures the
/// directories pagedb itself creates are never missed.
const STORE_SUBDIRS: &[&str] = &[
"seg",
"seg/.staging",
"seg/.tombstone",
"applyjournal",
"tmp",
];
const STORE_SUBDIRS: &[&str] = &["seg", STAGING_DIR, "seg/.tombstone", "applyjournal", "tmp"];

/// Move every file of the store rooted at `store_dir` into
/// `<store_dir>/quarantine/<label>/`, preserving the store's directory layout
Expand Down
12 changes: 5 additions & 7 deletions src/recovery/reconcile.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ use crate::pager::Pager;
use crate::segment::authenticated_metadata::{
ExpectedSegmentPath, authenticate_segment_metadata, validate_expected_path,
};
use crate::segment::writer::{STAGING_DIR, staging_path};
use crate::vfs::Vfs;
use crate::vfs::types::OpenMode;
use crate::{RealmId, Result};
Expand Down Expand Up @@ -181,10 +182,7 @@ async fn authenticate_row<V: Vfs + Clone>(
Ok((meta.segment_id, None))
}
Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
let staging = format!(
"seg/.staging/{}",
crate::hex::to_hex_lower(&meta.segment_id)
);
let staging = staging_path(&meta.segment_id);
let file = match vfs.open(&staging, OpenMode::Read).await {
Ok(file) => file,
Err(PagedbError::Io(error)) if error.kind() == std::io::ErrorKind::NotFound => {
Expand Down Expand Up @@ -219,7 +217,7 @@ async fn has_orphans<V: Vfs>(vfs: &V, expected: &[[u8; 16]]) -> Result<bool> {
return Ok(true);
}
}
let staging_entries = vfs.list_dir("seg/.staging").await?;
let staging_entries = vfs.list_dir(STAGING_DIR).await?;
for name in staging_entries {
let Some(id) = crate::hex::parse_hex::<16>(&name) else {
continue;
Expand Down Expand Up @@ -248,13 +246,13 @@ async fn sweep_orphans<V: Vfs>(vfs: &V, expected: &[[u8; 16]], recovery_commit:
vfs.rename(&from, &to).await?;
}
}
let staging_entries = vfs.list_dir("seg/.staging").await?;
let staging_entries = vfs.list_dir(STAGING_DIR).await?;
for name in staging_entries {
let Some(id) = crate::hex::parse_hex::<16>(&name) else {
continue;
};
if !expected_ids.contains(&id) {
vfs.remove(&format!("seg/.staging/{name}")).await?;
vfs.remove(&format!("{STAGING_DIR}/{name}")).await?;
}
}
vfs.sync_dir("seg").await?;
Expand Down
Loading
Loading