//! storage-mode validation at open time. //! //! the storage layout has two immutable axes: `only_index_links` (record //! values are bodies vs cids) and `ephemeral` (bodies live only inside the //! TTL-bounded event log). neither is safely convertible in place: //! //! - flipping links-only rewrites what `records` values mean, so replay of //! events written before the flip would parse cids as bodies or vice //! versa. opening with a different value than the database was written //! with is rejected with an actionable error. //! - ephemeral databases have no record heads or history at all. interpreting //! one layout as the other would either discard permanent heads or retain //! body bytes past the promised event TTL. //! //! mismatches are detected via a marker in the `filter` keyspace written on //! first open. databases that predate the marker get it written for the //! currently configured layout: a flip that already happened silently before //! the marker existed stays silent, and future flips are protected. #[cfg(feature = "indexer_stream")] use miette::Context; use miette::{IntoDiagnostic, Result}; use serde::{Deserialize, Serialize}; use crate::config::Config; use crate::db::Db; /// persisted marker of the storage layout the db was written with, kept in /// the filter keyspace next to the other instance-level config pub(super) const STORAGE_MODE_KEY: &[u8] = b"storage_mode"; #[derive(Serialize, Deserialize, PartialEq, Eq, Debug)] struct StorageModeMarker { links_only: bool, ephemeral: bool, } /// validate an existing marker before migrations can write anything. #[cfg(feature = "indexer")] pub(crate) fn preflight(db: &Db, cfg: &Config) -> Result { let marker = StorageModeMarker { links_only: cfg.only_index_links, ephemeral: cfg.ephemeral, }; let Some(stored) = db.filter.get(STORAGE_MODE_KEY).into_diagnostic()? else { if let Some(inferred_ephemeral) = infer_existing_ephemeral(db)? && inferred_ephemeral != marker.ephemeral { return Err(miette::miette!( "legacy event bodies imply ephemeral={inferred_ephemeral} but config requests {}: \ refusing to migrate a markerless database under the wrong storage mode. keep the \ previous HYDRANT_EPHEMERAL setting or use an explicit offline migration", marker.ephemeral, )); } return Ok(false); }; let prev: StorageModeMarker = rmp_serde::from_slice(&stored).into_diagnostic()?; if prev.ephemeral && db .indexer .records .iter() .next() .map(|guard| guard.into_inner()) .transpose() .into_diagnostic()? .is_some() { miette::bail!( "ephemeral storage marker conflicts with persisted record heads; rebuild this unsupported pre-release database from v9" ); } if prev.links_only != marker.links_only { return Err(miette::miette!( "database was written with only_index_links={} but config requests {}: \ links-only and body storage are not convertible in place. keep the previous \ setting or start from a fresh database path", prev.links_only, marker.links_only, )); } if prev.ephemeral != marker.ephemeral { return Err(miette::miette!( "database was written with ephemeral={} but config requests {}: \ permanent and ephemeral storage are not convertible at startup. keep the \ previous HYDRANT_EPHEMERAL setting, use an explicit offline migration, or \ start from a fresh database path", prev.ephemeral, marker.ephemeral, )); } Ok(true) } /// infer shipped v9 mode from body-bearing stream events. ephemeral writes /// stored `Block`; permanent writes stored `Ptr`. body-less delete events do /// not identify a mode, and seeing both encodings means the database was /// already switched unsafely before markers existed. fn infer_existing_ephemeral(db: &Db) -> Result> { let inferred = db .indexer .records .iter() .next() .map(|guard| guard.into_inner()) .transpose() .into_diagnostic()? .map(|_| false); #[cfg(feature = "indexer_stream")] { let mut inferred = inferred; for guard in db.stream.events.iter() { let value = guard.value().into_diagnostic()?; let event: crate::types::StoredEvent<'_> = rmp_serde::from_slice(&value) .into_diagnostic() .wrap_err("cannot infer storage mode from a legacy event")?; let ephemeral = match event.data { crate::types::StoredData::Block(_) => true, crate::types::StoredData::Ptr(_) => false, crate::types::StoredData::Nothing => continue, }; if inferred.is_some_and(|previous| previous != ephemeral) { miette::bail!( "database contains both inline ephemeral and pointer-based permanent events; use an explicit offline migration" ); } inferred = Some(ephemeral); } Ok(inferred) } #[cfg(not(feature = "indexer_stream"))] { Ok(inferred) } } /// adopt the configured layout for a database that predates the marker. this /// runs after migrations because legacy record values alone cannot distinguish /// a body-bearing database from links-only storage. #[cfg(feature = "indexer")] pub(crate) fn adopt(db: &Db, cfg: &Config) -> Result<()> { let marker = StorageModeMarker { links_only: cfg.only_index_links, ephemeral: cfg.ephemeral, }; if let Some(inferred_links_only) = infer_existing_links_only(db)? && inferred_links_only != marker.links_only { return Err(miette::miette!( "database record values imply only_index_links={inferred_links_only} but config requests {}: \ refusing to bless an incompatible pre-marker database. keep the existing setting or \ migrate into a fresh database path", marker.links_only, )); } db.filter .insert( STORAGE_MODE_KEY, rmp_serde::to_vec_named(&marker).into_diagnostic()?, ) .into_diagnostic()?; tracing::info!( links_only = marker.links_only, ephemeral = marker.ephemeral, "storage mode marker initialized" ); Ok(()) } /// infer the layout of a database created before `STORAGE_MODE_KEY`. scanning /// every head is intentional: accepting a mixed keyspace would make whichever /// value happens to sort first define the interpretation of all other values. fn infer_existing_links_only(db: &Db) -> Result> { let mut inferred = None; for guard in db.indexer.records.iter() { let value = guard.value().into_diagnostic()?; let links_only = crate::db::is_cid_record_value(&value); if inferred.is_some_and(|previous| previous != links_only) { miette::bail!( "database contains mixed record body and cid-only values; refusing to infer only_index_links" ); } inferred = Some(links_only); } Ok(inferred) } #[cfg(all(test, feature = "indexer"))] mod tests { use super::*; use crate::db::types::{DbAction, DbRkey, DbTid}; use crate::state::AppState; use jacquard_common::types::string::{Did, Tid}; fn did() -> Did<'static> { Did::new("did:plc:ewvi7nxzyoun6zhxrhs64oiz").unwrap() } fn config_at(path: &std::path::Path, ephemeral: bool, links_only: bool) -> Config { let mut config = Config::default(); config.database_path = path.to_path_buf(); config.ephemeral = ephemeral; config.only_index_links = links_only; config } fn write_record(state: &AppState, rkey: &str, body: &[u8]) { let mut txn = crate::db::Txn::new(&state.db); let did = did(); let rev = DbTid::from(&Tid::now_0()); let mut record_txn = txn.records(state, &rev, &did).unwrap(); let block = bytes::Bytes::copy_from_slice(body); let cid = jacquard_repo::mst::util::compute_cid(block.as_ref()).unwrap(); record_txn .put_record( "app.bsky.feed.post", &DbRkey::new(rkey), &cid, &block, DbAction::Create, ) .unwrap(); record_txn.finish().unwrap(); txn.commit().unwrap(); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn ephemeral_mode_is_immutable() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), false, false)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); } let err = match AppState::new(&config_at(tmp.path(), true, false)) { Ok(_) => panic!("permanent to ephemeral flip must be rejected"), Err(err) => err, }; let msg = format!("{err}"); assert!(msg.contains("ephemeral=false"), "{msg}"); assert!(msg.contains("HYDRANT_EPHEMERAL"), "{msg}"); assert!(msg.contains("offline migration"), "{msg}"); AppState::new(&config_at(tmp.path(), false, false)).unwrap(); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn ephemeral_to_permanent_flip_is_rejected() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), true, false)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); } let err = match AppState::new(&config_at(tmp.path(), false, false)) { Ok(_) => panic!("ephemeral to permanent flip must be rejected"), Err(err) => err, }; assert!(format!("{err}").contains("ephemeral=true")); AppState::new(&config_at(tmp.path(), true, false)).unwrap(); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn empty_database_still_records_its_ephemeral_identity() { let tmp = tempfile::tempdir().unwrap(); AppState::new(&config_at(tmp.path(), true, false)).unwrap(); assert!(AppState::new(&config_at(tmp.path(), false, false)).is_err()); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn markerless_inline_events_prevent_wrong_permanent_adoption() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), true, false)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); assert!(state.db.indexer.records.is_empty().unwrap()); state.db.filter.remove(STORAGE_MODE_KEY).unwrap(); state.db.persist().unwrap(); } let err = match AppState::new(&config_at(tmp.path(), false, false)) { Ok(_) => panic!("inline v9 events must identify ephemeral mode"), Err(err) => err, }; assert!(format!("{err}").contains("imply ephemeral=true")); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn markerless_pointer_events_prevent_wrong_ephemeral_adoption() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), false, false)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); state.db.filter.remove(STORAGE_MODE_KEY).unwrap(); state.db.persist().unwrap(); } let err = match AppState::new(&config_at(tmp.path(), true, false)) { Ok(_) => panic!("pointer v9 events must identify permanent mode"), Err(err) => err, }; assert!(format!("{err}").contains("imply ephemeral=false")); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn marked_ephemeral_database_rejects_unshipped_record_heads() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), true, false)).unwrap(); let key = crate::db::keys::record_key( &did(), "app.bsky.feed.post", &DbRkey::new("3kzbif5moe22m"), ); state.db.indexer.records.insert(key, b"stale head").unwrap(); state.db.persist().unwrap(); } let err = match AppState::new(&config_at(tmp.path(), true, false)) { Ok(_) => panic!("event-only ephemeral mode must reject persisted heads"), Err(err) => err, }; assert!(format!("{err}").contains("unsupported pre-release")); } #[test] #[cfg(all(feature = "indexer_stream", not(feature = "backlinks")))] fn rejected_layout_open_cannot_arm_history_retention() { let tmp = tempfile::tempdir().unwrap(); let record_key = crate::db::keys::record_key( &did(), "app.bsky.feed.post", &DbRkey::new("3kzbif5moe22m"), ); let history_key = crate::db::keys::history_key(&record_key, &DbTid::new_from_bytes(1_u64.to_be_bytes())); { let state = AppState::new(&config_at(tmp.path(), false, false)).unwrap(); let mut batch = state.db.inner.batch(); state .db .indexer .stage_history(&mut batch, &history_key, b"must survive".as_slice()); batch.commit().unwrap(); state.db.indexer.history.rotate_memtable_and_wait().unwrap(); state.db.persist().unwrap(); } let mut wrong = config_at(tmp.path(), false, true); wrong.history_ttl = Some(std::time::Duration::from_secs(1)); assert!(AppState::new(&wrong).is_err()); let state = AppState::new(&config_at(tmp.path(), false, false)).unwrap(); state.db.indexer.history.major_compact().unwrap(); assert_eq!( state .db .indexer .history .get(history_key) .unwrap() .as_deref(), Some(b"must survive".as_slice()) ); } #[test] fn links_only_flip_is_rejected_actionably() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), false, true)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); } let err = match AppState::new(&config_at(tmp.path(), false, false)) { Ok(_) => panic!("links-only flip must be rejected"), Err(err) => err, }; let msg = format!("{err}"); assert!( msg.contains("not convertible in place"), "unexpected error: {msg}" ); assert!(msg.contains("only_index_links"), "unexpected error: {msg}"); } #[test] fn missing_marker_does_not_bless_a_links_only_database_as_body_storage() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), false, true)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); state.db.filter.remove(STORAGE_MODE_KEY).unwrap(); state.db.persist().unwrap(); } let err = match AppState::new(&config_at(tmp.path(), false, false)) { Ok(_) => panic!("pre-marker links-only mismatch must be rejected"), Err(err) => err, }; let msg = format!("{err}"); assert!(msg.contains("imply only_index_links=true"), "{msg}"); } #[test] fn missing_marker_does_not_bless_body_storage_as_links_only() { let tmp = tempfile::tempdir().unwrap(); { let state = AppState::new(&config_at(tmp.path(), false, false)).unwrap(); write_record(&state, "3kzbif5moe22m", b"v1"); state.db.filter.remove(STORAGE_MODE_KEY).unwrap(); state.db.persist().unwrap(); } let err = match AppState::new(&config_at(tmp.path(), false, true)) { Ok(_) => panic!("pre-marker body mismatch must be rejected"), Err(err) => err, }; let msg = format!("{err}"); assert!(msg.contains("imply only_index_links=false"), "{msg}"); } }