//! legacy event-body passes for v10. //! //! v10 moved permanent current heads into `records`, but legacy permanent //! stream events can still point into `blocks`. materialize only otherwise- //! unresolved pointer versions into the record-scoped `event_bodies` archive, //! then clear the legacy CAS exactly. ephemeral `Block` events remain inline: //! event TTL must delete their body bytes with the event. #[cfg(feature = "indexer_stream")] use fjall::OwnedWriteBatch; #[cfg(feature = "indexer_stream")] use miette::WrapErr; #[cfg(feature = "indexer")] use miette::{IntoDiagnostic, Result}; #[cfg(feature = "indexer_stream")] use crate::db::{Db, keys}; #[cfg(feature = "indexer_stream")] use crate::types::{StoredData, StoredEvent}; #[cfg(feature = "indexer_stream")] use super::super::{ChunkBudget, Pass}; #[cfg(feature = "indexer_stream")] pub(super) const POINTER_EVENT_BODIES: Pass = Pass { name: "pointer_event_bodies", scan: "events", visit: materialize_pointer_event, budget: ChunkBudget::DEFAULT, optional_on_absent: false, }; #[cfg(feature = "indexer_stream")] fn body_matches(value: &[u8], expected: &jacquard_common::types::cid::IpldCid) -> bool { !crate::db::indexer::is_cid_record_value(value) && jacquard_repo::mst::util::compute_cid(value).is_ok_and(|computed| computed == *expected) } #[cfg(feature = "indexer_stream")] fn stage_archive_body( db: &Db, batch: &mut OwnedWriteBatch, record_key: &[u8], cid: &jacquard_common::types::cid::IpldCid, body: &[u8], ) -> Result { let computed = jacquard_repo::mst::util::compute_cid(body) .into_diagnostic() .wrap_err("v10: cannot hash legacy event body")?; if computed != *cid { miette::bail!("v10: legacy body cid {computed} does not match event pointer {cid}"); } let key = keys::event_body_key(record_key, cid); if let Some(existing) = db.stream.event_bodies.get(&key).into_diagnostic()? { if body_matches(&existing, cid) { return Ok(0); } miette::bail!("v10: compatibility body at {cid} is corrupt"); } batch.insert(&db.stream.event_bodies, key, body); Ok(body.len()) } #[cfg(feature = "indexer_stream")] fn decode_event<'a>(value: &'a [u8]) -> Result> { rmp_serde::from_slice(value) .into_diagnostic() .wrap_err("v10: unreadable stored event") } #[cfg(feature = "indexer_stream")] fn materialize_pointer_event( db: &Db, batch: &mut OwnedWriteBatch, _scanned: &fjall::Keyspace, _key: &[u8], value: &[u8], ) -> Result { let event = decode_event(value)?; let StoredData::Ptr(cid) = &event.data else { return Ok(0); }; let record_key = keys::record_key_trimmed(&event.did, event.collection.as_str(), &event.rkey); if db .resolve_event_record_body(&record_key, &event.rev, cid)? .is_some() { return Ok(0); } // migration runs synchronously inside `Db::open`, before ingestion can add // events. every unresolved pointer here therefore predates v10 and must be // resolved from the legacy CAS. let block_key = keys::block_key(event.collection.as_str(), &cid.to_bytes()); let body = super::super::legacy_block(db, &block_key)? .ok_or_else(|| miette::miette!("v10: legacy event body {cid} is missing"))?; stage_archive_body(db, batch, &record_key, cid, &body) } /// delete the legacy CAS only after every event-materialization chunk and its /// resume cursor are durable. a crash before the migration marker merely makes /// this idempotent finalizer run again. #[cfg(feature = "indexer")] pub(super) fn retire_blocks(db: &crate::db::Db) -> Result { #[cfg(not(feature = "indexer_stream"))] if db.inner.keyspace_exists("events") { let events = db .inner .keyspace("events", Default::default) .into_diagnostic()?; if !events.is_empty().into_diagnostic()? { tracing::warn!( "v10: legacy events exist but indexer_stream is disabled; retaining blocks until a stream-enabled open" ); return Ok(false); } } // chunk commits and their resume cursor use the data WAL, while keyspace // deletion updates fjall's metadata independently. force the completed // archive durable before making the old CAS unreachable. db.persist()?; if let Some(blocks) = super::super::legacy_blocks(db)? { db.inner.delete_keyspace(blocks).into_diagnostic()?; db.persist()?; } Ok(true) } #[cfg(not(feature = "indexer"))] pub(super) fn retire_blocks(_db: &crate::db::Db) -> miette::Result { Ok(false) } #[cfg(all(test, feature = "indexer", not(feature = "indexer_stream")))] mod indexer_only_tests { use super::*; use crate::config::Config; fn config(path: &std::path::Path) -> Config { Config { database_path: path.to_path_buf(), ..Default::default() } } fn seed_dormant_events(db: &crate::db::Db, nonempty: bool) -> Result<()> { let events = db .inner .keyspace("events", Default::default) .into_diagnostic()?; if nonempty { events .insert(1_u64.to_be_bytes(), b"legacy event") .into_diagnostic()?; } let blocks = super::super::super::legacy_blocks_for_test(db)?; blocks.insert(b"sentinel", b"body").into_diagnostic()?; crate::db::migration::rewind_version_for_test(db, 9)?; db.persist() } #[test] fn empty_dormant_events_do_not_block_retirement() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); { let db = crate::db::Db::open(&cfg)?; seed_dormant_events(&db, false)?; } let db = crate::db::Db::open(&cfg)?; assert!(!db.inner.keyspace_exists("blocks")); Ok(()) } #[test] fn nonempty_dormant_events_retain_blocks_for_a_stream_build() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); { let db = crate::db::Db::open(&cfg)?; seed_dormant_events(&db, true)?; } let db = crate::db::Db::open(&cfg)?; assert!(db.inner.keyspace_exists("blocks")); Ok(()) } } #[cfg(all(test, feature = "indexer_stream"))] mod tests { use super::*; use crate::config::Config; use crate::db::types::{DbAction, DbRkey, DbTid, TrimmedDid}; #[cfg(not(feature = "backlinks"))] use crate::state::AppState; use jacquard_common::types::string::{Did, Tid}; use jacquard_common::{CowStr, IntoStatic}; use miette::IntoDiagnostic; const DID: &str = "did:plc:yk4q3id7id6p5z3bypvshc64"; const COLLECTION: &str = "app.bsky.feed.post"; const RKEY: &str = "3kzbif5moe22m"; fn config(path: &std::path::Path) -> Config { Config { database_path: path.to_path_buf(), ..Default::default() } } fn did() -> Did { Did::new_static(DID).unwrap() } fn rev() -> DbTid { DbTid::from(&Tid::new("3kzbif5moe22m").unwrap()) } fn body(text: &str) -> Vec { serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "text": text, })) .unwrap() } fn event(data: StoredData) -> StoredEvent<'static> { StoredEvent { live: false, did: TrimmedDid::from(&did()).into_static(), rev: rev(), collection: CowStr::Borrowed(COLLECTION).into_static(), rkey: DbRkey::new(RKEY), action: DbAction::Create, data, } } fn record_key() -> Vec { keys::record_key(&did(), COLLECTION, &DbRkey::new(RKEY)) } fn stage_event(db: &Db, batch: &mut OwnedWriteBatch, id: u64, event: &StoredEvent<'_>) { db.stream.stage_event( batch, keys::event_key(id), rmp_serde::to_vec(event).unwrap(), ); } fn stage_legacy_block( db: &Db, batch: &mut OwnedWriteBatch, cid: &jacquard_common::types::cid::IpldCid, body: &[u8], ) -> Result<()> { let blocks = super::super::super::legacy_blocks_for_test(db)?; batch.insert(&blocks, keys::block_key(COLLECTION, &cid.to_bytes()), body); Ok(()) } fn rewind_to_v9(db: &Db) -> Result<()> { crate::db::migration::rewind_version_for_test(db, 9)?; db.persist() } #[test] fn pointer_body_is_materialized_before_blocks_are_cleared() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body = body("legacy"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); stage_legacy_block(&db, &mut batch, &cid, &body)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; assert!(!db.inner.keyspace_exists("blocks")); assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); Ok(()) } #[cfg(not(feature = "backlinks"))] #[test] fn inline_ephemeral_body_stays_owned_by_its_event() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let mut cfg = config(tmp.path()); cfg.ephemeral = true; let body = body("inline"); { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event( &db, &mut batch, 7, &event(StoredData::Block(bytes::Bytes::copy_from_slice(&body))), ); batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; let stored = db .stream .events .get(keys::event_key(7)) .into_diagnostic()? .expect("event should remain"); let stored: StoredEvent<'_> = rmp_serde::from_slice(&stored).into_diagnostic()?; assert!(matches!(stored.data, StoredData::Block(found) if found.as_ref() == body)); assert!(db.stream.event_bodies.is_empty().into_diagnostic()?); Ok(()) } #[test] fn same_record_and_revision_can_archive_distinct_cids() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body_a = body("fork a"); let body_b = body("fork b"); let cid_a = jacquard_repo::mst::util::compute_cid(&body_a).into_diagnostic()?; let cid_b = jacquard_repo::mst::util::compute_cid(&body_b).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid_a))); stage_event(&db, &mut batch, 2, &event(StoredData::Ptr(cid_b))); stage_legacy_block(&db, &mut batch, &cid_a, &body_a)?; stage_legacy_block(&db, &mut batch, &cid_b, &body_b)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid_a)) .into_diagnostic()? .as_deref(), Some(body_a.as_slice()) ); assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid_b)) .into_diagnostic()? .as_deref(), Some(body_b.as_slice()) ); Ok(()) } #[test] fn current_head_is_not_duplicated_in_the_archive() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body = body("head"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); db.indexer.stage_record(&mut batch, record_key(), &body); stage_legacy_block(&db, &mut batch, &cid, &body)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; assert!(!db.inner.keyspace_exists("blocks")); assert!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .is_none() ); Ok(()) } #[test] fn history_before_the_event_does_not_suppress_archival() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body = body("legacy"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); db.indexer.stage_history( &mut batch, keys::history_key(&record_key(), &DbTid::new_from_bytes(0_u64.to_be_bytes())), &body, ); stage_legacy_block(&db, &mut batch, &cid, &body)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); Ok(()) } #[test] fn preexisting_archive_makes_pointer_migration_idempotent() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body = body("already archived"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); batch.insert( &db.stream.event_bodies, keys::event_body_key(&record_key(), &cid), &body, ); batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let db = Db::open(&cfg)?; assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); assert!(!db.inner.keyspace_exists("blocks")); Ok(()) } #[test] fn missing_body_aborts_without_clearing_blocks() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let missing = body("missing"); let cid = jacquard_repo::mst::util::compute_cid(&missing).into_diagnostic()?; let sentinel = body("sentinel"); let sentinel_cid = jacquard_repo::mst::util::compute_cid(&sentinel).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); stage_legacy_block(&db, &mut batch, &sentinel_cid, &sentinel)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let err = match Db::open(&cfg) { Ok(_) => panic!("missing sole body must abort migration"), Err(err) => err, }; assert!(format!("{err:?}").contains("is missing")); let raw = fjall::Database::builder(tmp.path()) .open() .into_diagnostic()?; let blocks = raw .keyspace("blocks", fjall::KeyspaceCreateOptions::default) .into_diagnostic()?; assert!(!blocks.is_empty().into_diagnostic()?); Ok(()) } #[test] fn cid_mismatch_aborts_without_clearing_blocks() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let expected_body = body("expected"); let wrong_body = body("wrong"); let cid = jacquard_repo::mst::util::compute_cid(&expected_body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); stage_legacy_block(&db, &mut batch, &cid, &wrong_body)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let err = match Db::open(&cfg) { Ok(_) => panic!("cid mismatch must abort migration"), Err(err) => err, }; assert!( format!("{err:?}").contains("legacy body cid"), "unexpected error: {err:?}" ); let raw = fjall::Database::builder(tmp.path()) .open() .into_diagnostic()?; let blocks = raw .keyspace("blocks", fjall::KeyspaceCreateOptions::default) .into_diagnostic()?; assert!(!blocks.is_empty().into_diagnostic()?); Ok(()) } #[test] fn malformed_event_aborts_without_deleting_legacy_blocks() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let sentinel = body("sentinel"); let sentinel_cid = jacquard_repo::mst::util::compute_cid(&sentinel).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); db.stream .stage_event(&mut batch, keys::event_key(1), b"not msgpack"); stage_legacy_block(&db, &mut batch, &sentinel_cid, &sentinel)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } assert!(Db::open(&cfg).is_err()); let raw = fjall::Database::builder(tmp.path()) .open() .into_diagnostic()?; let blocks = raw .keyspace("blocks", fjall::KeyspaceCreateOptions::default) .into_diagnostic()?; assert!(!blocks.is_empty().into_diagnostic()?); Ok(()) } #[test] fn migration_resumes_after_a_partial_event_chunk() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body_a = body("a"); let body_b = body("b"); let cid_a = jacquard_repo::mst::util::compute_cid(&body_a).into_diagnostic()?; let cid_b = jacquard_repo::mst::util::compute_cid(&body_b).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid_a))); stage_event(&db, &mut batch, 2, &event(StoredData::Ptr(cid_b))); stage_legacy_block(&db, &mut batch, &cid_a, &body_a)?; stage_legacy_block(&db, &mut batch, &cid_b, &body_b)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; let pass = Pass { name: "pointer_event_bodies", scan: "events", visit: materialize_pointer_event, budget: ChunkBudget { entries: 1, bytes: usize::MAX, }, optional_on_absent: false, }; let cursor = super::super::super::pass_cursor_key(9, pass.name); let events = db.keyspace_by_name("events").expect("events keyspace"); let outcome = super::super::super::run_chunk(&db, &cursor, &events, &pass)?; assert_eq!(outcome.seen, 1); assert!(!outcome.exhausted); db.persist()?; } let db = Db::open(&cfg)?; assert!(!db.inner.keyspace_exists("blocks")); for (cid, body) in [(cid_a, body_a), (cid_b, body_b)] { assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); } Ok(()) } #[test] fn migration_retries_after_keyspace_delete_before_marker() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let cfg = config(tmp.path()); let body = body("legacy"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event(&db, &mut batch, 1, &event(StoredData::Ptr(cid))); stage_legacy_block(&db, &mut batch, &cid, &body)?; batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; for pass in super::super::PASSES { assert!(super::super::super::run_pass(&db, 9, pass)?); } assert!(retire_blocks(&db)?); assert!(!db.inner.keyspace_exists("blocks")); // deliberately do not write v10's applied marker } let db = Db::open(&cfg)?; assert!(!db.inner.keyspace_exists("blocks")); assert_eq!( db.stream .event_bodies .get(keys::event_body_key(&record_key(), &cid)) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); Ok(()) } #[test] fn history_ttl_filter_never_applies_to_compatibility_bodies() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let mut cfg = config(tmp.path()); cfg.history_ttl = Some(std::time::Duration::from_secs(1)); let db = Db::open(&cfg)?; let body = body("compatibility"); let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; let key = keys::event_body_key(&record_key(), &cid); db.stream .event_bodies .insert(&key, &body) .into_diagnostic()?; db.stream .event_bodies .rotate_memtable_and_wait() .into_diagnostic()?; db.stream.event_bodies.major_compact().into_diagnostic()?; assert_eq!( db.stream .event_bodies .get(key) .into_diagnostic()? .as_deref(), Some(body.as_slice()) ); Ok(()) } #[test] #[cfg(not(feature = "backlinks"))] fn v9_ephemeral_database_upgrades_replays_and_expires() -> Result<()> { let tmp = tempfile::tempdir().into_diagnostic()?; let mut cfg = config(tmp.path()); cfg.ephemeral = true; cfg.ephemeral_ttl = std::time::Duration::from_secs(60 * 60); let legacy_body = body("legacy inline ephemeral"); let legacy_id = 7_u64; // physical v9 ephemeral shape: replay bodies live inline in events, // there are no record heads, and no storage-mode marker existed yet. { let db = Db::open(&cfg)?; let mut batch = db.inner.batch(); stage_event( &db, &mut batch, legacy_id, &event(StoredData::Block(bytes::Bytes::copy_from_slice( &legacy_body, ))), ); let blocks = super::super::super::legacy_blocks_for_test(&db)?; batch.insert(&blocks, b"legacy sentinel", b"legacy block"); batch.remove(&db.filter, crate::db::storage_mode::STORAGE_MODE_KEY); batch.commit().into_diagnostic()?; rewind_to_v9(&db)?; } let replay_body = |state: &AppState, id: u64| -> Result> { let Some(bytes) = state .db .stream .events .get(keys::event_key(id)) .into_diagnostic()? else { return Ok(None); }; let stored: StoredEvent<'_> = rmp_serde::from_slice(&bytes).into_diagnostic()?; let event = crate::control::stream::indexer::stored_to_event(state, id, stored, None) .ok_or_else(|| miette::miette!("stored event did not inflate"))?; let raw = event .record .and_then(|record| record.record) .ok_or_else(|| miette::miette!("stored event lost its record body"))?; serde_json::from_str(raw.get()).into_diagnostic().map(Some) }; let expected_legacy = serde_json::json!({ "$type": COLLECTION, "text": "legacy inline ephemeral", }); { let state = AppState::new(&cfg)?; assert!(!state.db.inner.keyspace_exists("blocks")); assert!(state.db.indexer.records.is_empty().into_diagnostic()?); assert_eq!( replay_body(&state, legacy_id)?, Some(expected_legacy.clone()) ); assert!(state.db.stream.event_bodies.is_empty().into_diagnostic()?); // an ordinary tick cannot prune an event until a watermark has // actually aged through the configured window. crate::db::ephemeral::ephemeral_ttl_tick(&state.db, &cfg.ephemeral_ttl)?; assert_eq!( replay_body(&state, legacy_id)?, Some(expected_legacy.clone()) ); state .db .stream .events .rotate_memtable_and_wait() .into_diagnostic()?; state.db.stream.events.major_compact().into_diagnostic()?; state.db.persist()?; } let state = AppState::new(&cfg)?; assert_eq!(replay_body(&state, legacy_id)?, Some(expected_legacy)); // once a v9 event's real retention window has elapsed, its inline body // disappears atomically with the event. let old_ts = (chrono::Utc::now().timestamp() as u64).saturating_sub(cfg.ephemeral_ttl.as_secs() + 1); state .db .cursors .insert( keys::event_watermark_key(old_ts), (legacy_id + 1).to_be_bytes(), ) .into_diagnostic()?; crate::db::ephemeral::ephemeral_ttl_tick(&state.db, &cfg.ephemeral_ttl)?; assert_eq!(replay_body(&state, legacy_id)?, None); assert!(state.db.stream.event_bodies.is_empty().into_diagnostic()?); // new writes stay event-only too: no head, history, redaction, or // compatibility archive can retain the body past event TTL. let new_body = bytes::Bytes::from(body("new ephemeral event")); let new_block = crate::car::CarBlock::from_body(new_body.clone()); let new_rkey = DbRkey::new("3kzbif5mof33m"); let new_key = keys::record_key(&did(), COLLECTION, &new_rkey); let repo = did(); let mut txn = crate::db::Txn::new(&state.db); let mut records = txn.records(&state, &DbTid::from(&Tid::now_0()), &repo)?; records.put_record(COLLECTION, &new_rkey, &new_block, DbAction::Create)?; records.finish()?; txn.commit()?; assert!( state .db .indexer .record(&new_key) .into_diagnostic()? .is_none() ); assert!(state.db.indexer.history.is_empty().into_diagnostic()?); assert!(state.db.indexer.redactions.is_empty().into_diagnostic()?); assert!(state.db.stream.event_bodies.is_empty().into_diagnostic()?); let new_id = legacy_id + 1; let stored = state .db .stream .events .get(keys::event_key(new_id)) .into_diagnostic()? .expect("new ephemeral event must be durable"); let stored: StoredEvent<'_> = rmp_serde::from_slice(&stored).into_diagnostic()?; assert!(matches!(stored.data, StoredData::Block(found) if found == new_body)); assert_eq!( replay_body(&state, new_id)?, Some(serde_json::json!({ "$type": COLLECTION, "text": "new ephemeral event", })) ); drop(state); let mut permanent_cfg = cfg.clone(); permanent_cfg.ephemeral = false; let err = match AppState::new(&permanent_cfg) { Ok(_) => panic!("migrated ephemeral database must reject permanent mode"), Err(err) => err, }; assert!(format!("{err}").contains("ephemeral=true")); Ok(()) } }