From ec8b7fcfd25a5c0eea3c99066cf4cfc712d14bcb Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 7 Oct 2026 00:41:58 +0300 Subject: [PATCH] [db] document and test ephemeral ttl changes across restarts HYDRANT_EPHEMERAL_TTL is applied on every tick against watermarks persisted in cursors, so changing it between restarts changes the cutoff right away, but nothing said what shrinking or growing it does or covered either direction. docs/concepts/retention.md now explains the watermark mechanism, what a shrink and a grow do across a restart, why the cut is only as fine as the watermarks (one a minute, none while down), and the boundary table lag from drop_range. no behavior changes. the new event_ttl_tests in src/db/ephemeral.rs write real ephemeral events into separate tables with the watermarks a running worker would have left, then reopen the database with a shorter or longer ttl and tick. they check that a shrink prunes on the first tick and a grow keeps what is left and brings nothing back, that inline event bodies replay while their events are kept and go with them, and that a table straddling the cutoff outlives the ttl. picking the oldest watermark instead of the newest one past the cutoff fails the shrink and grow tests (event 3 should be pruned), and storing ephemeral events as pointers instead of inline blocks fails all three (event lost its body). covers hydrant-4r4. --- docs/concepts/retention.md | 18 +++- docs/configuration.md | 2 +- src/db/ephemeral.rs | 203 +++++++++++++++++++++++++++++++++++++ 3 files changed, 221 insertions(+), 2 deletions(-) diff --git a/docs/concepts/retention.md b/docs/concepts/retention.md index 0325917..46a7647 100644 --- a/docs/concepts/retention.md +++ b/docs/concepts/retention.md @@ -2,7 +2,7 @@ title: retention --- -hydrant keeps two things on a clock: stream events, bounded by `EPHEMERAL_TTL`, and in permanent databases the replaced record versions in `history`, bounded by `HISTORY_TTL`. +hydrant keeps two things on a clock: stream events, bounded by `EPHEMERAL_TTL` in ephemeral indexers and in relay mode, and in permanent databases the replaced record versions in `history`, bounded by `HISTORY_TTL`. ## record history @@ -11,3 +11,19 @@ when a record is updated or deleted, its previous body moves to the `history` ke without `HISTORY_TTL` history is kept forever. with it, entries older than the ttl are dropped by a compaction filter, so they go when compaction next rewrites their part of the keyspace, not right at the deadline. replaying an event whose version was trimmed gives the event without its record body, and `/stats` counts those as `counts.history_trim_misses`. `counts.history` in `/stats` is approximate. it's read from the storage engine's table metadata rather than kept as a counter, so it does go down when retention or a repo erase removes entries, but only once compaction has rewritten the tables holding them. until then removed and overwritten entries still count. read it as about how many versions are on disk, not an exact number. + +## stream events + +an indexer only expires its events when `EPHEMERAL` is on, while a relay always expires its `subscribeRepos` events. both work the same way. events are numbered in commit order, and once a minute the ttl worker writes a watermark into the `cursors` keyspace saying "at this second, the next event id was N", then looks for the newest watermark at least `EPHEMERAL_TTL` old and drops every event below its id, deleting the watermarks it has used up. the watermarks only record time and position, not the ttl, so the ttl in effect is whatever hydrant was started with, applied fresh on every tick against whatever watermarks are on disk. + +in ephemeral databases an event carries its record body inline, so that body lives exactly as long as its event: it replays with the event while the event is kept and goes in the same drop. nothing else (no record head, history entry or body archive) holds a copy. + +### changing `EPHEMERAL_TTL` across a restart + +- **shrinking** takes effect on the first tick, about a minute after startup. the newest watermark past the new, shorter cutoff picks where to cut, and everything older goes in one go. +- **growing** never brings anything back, since already pruned events are gone. what's still stored is kept until it's older than the new ttl, so pruning pauses until the oldest remaining watermark ages past it, and from then on the window is the new length. +- either way the cut is only as fine as the watermarks: one a minute while hydrant runs and none while it's down. if no watermark falls exactly at the cutoff the newest one before it is used, so hydrant keeps a little more, never less. after downtime that can be up to the length of the gap. + +### boundary lag + +events are dropped a whole table at a time (`drop_range`), so a storage table holding both events older than the cutoff and newer ones stays until a later cutoff passes its newest event. expect a bit more than `EPHEMERAL_TTL` worth of events on disk, by about one table, and expired events in that table still replay with their bodies until it goes. counts in `/stats` are approximate for the same reason. jetstream replay metadata is pruned on the same clock, by timestamp rather than watermark, with the same table granularity. diff --git a/docs/configuration.md b/docs/configuration.md index 2a4b29e..25a11da 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -22,7 +22,7 @@ hydrant is configured via environment variables, all prefixed with `HYDRANT_` (e | :--- | :--- | :--- | | `FULL_NETWORK` | `false` (indexer), `true` (relay) | if `true`, discover and index all repos in the network | | `EPHEMERAL` | `false` (indexer), `true` (relay) | if enabled, no records are stored (in indexer mode). events are deleted after a certain duration (`EPHEMERAL_TTL`). immutable per database: the setting is persisted on first open and a later flip is rejected at startup — keep the previous setting or use a fresh database path | -| `EPHEMERAL_TTL` | `60min`, `3d` (relay) | how long to keep events before deletion. when built with `jetstream`, retained Jetstream replay metadata is pruned on the same schedule | +| `EPHEMERAL_TTL` | `60min`, `3d` (relay) | how long to keep events before deletion. when built with `jetstream`, retained Jetstream replay metadata is pruned on the same schedule. can change between restarts, see [retention](concepts/retention.md) for what shrinking or growing it does | | `HISTORY_TTL` | unset | indexer only, permanent databases only. how long replaced record versions stay in `history` so replays of older events can still show them. humantime duration like `2h`, unset keeps them forever. see [retention](concepts/retention.md) | | `ONLY_INDEX_LINKS` | `false` | indexer only. if enabled, record blocks are not stored, only the index (records, counts, events) is kept. `getRecord`, `listRecords`, and `getRepo` will return errors. the event stream and Jetstream stream still work, but create/update events will not include record values. immutable per database: the setting is persisted on first open and a later flip is rejected at startup — keep the previous setting or use a fresh database path | diff --git a/src/db/ephemeral.rs b/src/db/ephemeral.rs index c210897..62fc415 100644 --- a/src/db/ephemeral.rs +++ b/src/db/ephemeral.rs @@ -544,3 +544,206 @@ mod tests { Ok(()) } } + +#[cfg(all(test, feature = "indexer_stream", not(feature = "backlinks")))] +mod event_ttl_tests { + use super::*; + use crate::config::Config; + use crate::db::types::{DbAction, DbRkey, DbTid}; + use crate::state::AppState; + use crate::types::StoredEvent; + use jacquard_common::types::string::{Did, Tid}; + use std::ops::Range; + + const COLLECTION: &str = "app.bsky.feed.post"; + const MINUTE: u64 = 60; + const HOUR: u64 = 60 * MINUTE; + + fn open(dir: &std::path::Path, ttl: u64) -> miette::Result { + AppState::new(&Config { + database_path: dir.to_path_buf(), + ephemeral: true, + ephemeral_ttl: Duration::from_secs(ttl), + ..Config::default() + }) + } + + fn body(text: &str) -> Vec { + serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "text": text })) + .unwrap() + } + + /// commits one create per rkey and returns the event ids they took + fn write_events(state: &AppState, rkeys: &[&str]) -> miette::Result> { + let did = Did::new_static("did:plc:ewvi7nxzyoun6zhxrhs64oiz").into_diagnostic()?; + let first = state.db.stream.ids.next(); + for rkey in rkeys { + let mut txn = crate::db::Txn::new(&state.db); + let mut records = txn.records(state, &DbTid::from(&Tid::now_0()), &did)?; + let block = crate::car::CarBlock::from_body(body(rkey)); + records.put_record(COLLECTION, &DbRkey::new(rkey), &block, DbAction::Create)?; + records.finish()?; + txn.commit()?; + } + Ok(first..state.db.stream.ids.next()) + } + + /// seals the memtable so the next events land in a table of their own + fn seal_table(state: &AppState) -> miette::Result<()> { + state + .db + .stream + .events + .rotate_memtable_and_wait() + .into_diagnostic() + } + + /// the watermark a tick `age` ago would have written: events from + /// `next_seq` on did not exist yet + fn plant_watermark(state: &AppState, age: u64, next_seq: u64) -> miette::Result<()> { + let ts = chrono::Utc::now().timestamp() as u64 - age; + state + .db + .cursors + .insert(keys::event_watermark_key(ts), next_seq.to_be_bytes()) + .into_diagnostic() + } + + fn watermark_seqs(state: &AppState) -> miette::Result> { + state + .db + .cursors + .prefix(keys::EVENT_WATERMARK_PREFIX) + .map(|guard| { + let value = guard.value().into_diagnostic()?; + Ok(u64::from_be_bytes( + value.as_ref().try_into().into_diagnostic()?, + )) + }) + .collect() + } + + /// what replay hands out for `id`: its record text, read from the body + /// the event row carries + fn replayed_text(state: &AppState, id: u64) -> miette::Result> { + let Some(row) = state + .db + .stream + .events + .get(keys::event_key(id)) + .into_diagnostic()? + else { + return Ok(None); + }; + let stored: StoredEvent<'_> = rmp_serde::from_slice(&row).into_diagnostic()?; + let event = crate::control::stream::indexer::stored_to_event(state, id, stored, None) + .ok_or_else(|| miette::miette!("event {id} did not inflate"))?; + let raw = event + .record + .and_then(|record| record.record) + .ok_or_else(|| miette::miette!("event {id} lost its body"))?; + let json: serde_json::Value = serde_json::from_str(raw.get()).into_diagnostic()?; + Ok(json["text"].as_str().map(str::to_owned)) + } + + fn assert_retained(state: &AppState, ids: Range) -> miette::Result<()> { + for id in ids { + assert!( + replayed_text(state, id)?.is_some(), + "event {id} should replay with its body" + ); + } + Ok(()) + } + + fn assert_pruned(state: &AppState, ids: Range) -> miette::Result<()> { + for id in ids { + assert_eq!( + replayed_text(state, id)?, + None, + "event {id} should be pruned" + ); + } + Ok(()) + } + + /// three tables written 150, 90 and 30 minutes ago and one written now, + /// with the watermarks the ttl worker left after each. returns their ids + fn seed_history(state: &AppState) -> miette::Result<[Range; 4]> { + let oldest = write_events(state, &["3kzbif5moa22m", "3kzbif5moa23m"])?; + seal_table(state)?; + plant_watermark(state, 150 * MINUTE, oldest.end)?; + let older = write_events(state, &["3kzbif5mob22m", "3kzbif5mob23m"])?; + seal_table(state)?; + plant_watermark(state, 90 * MINUTE, older.end)?; + let recent = write_events(state, &["3kzbif5moc22m", "3kzbif5moc23m"])?; + seal_table(state)?; + plant_watermark(state, 30 * MINUTE, recent.end)?; + let fresh = write_events(state, &["3kzbif5mod22m"])?; + Ok([oldest, older, recent, fresh]) + } + + #[test] + fn shrinking_the_ttl_across_a_restart_prunes_on_the_first_tick() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let [oldest, older, recent, fresh] = { + let state = open(tmp.path(), 3 * HOUR)?; + let seeded = seed_history(&state)?; + // nothing is three hours old yet + ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; + assert_retained(&state, seeded[0].start..seeded[3].end)?; + state.db.persist()?; + seeded + }; + + let state = open(tmp.path(), HOUR)?; + ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; + // the 90 minute watermark is the newest one past the new cutoff, so + // everything written before it goes, inline bodies with it + assert_pruned(&state, oldest.start..older.end)?; + assert_retained(&state, recent.start..fresh.end)?; + assert!(state.db.stream.event_bodies.is_empty().into_diagnostic()?); + // the consumed watermarks are gone, the 30 minute one waits its turn + assert_eq!(watermark_seqs(&state)?.first(), Some(&recent.end)); + Ok(()) + } + + #[test] + fn growing_the_ttl_across_a_restart_keeps_what_is_left() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let [oldest, older, recent, fresh] = { + let state = open(tmp.path(), HOUR)?; + let seeded = seed_history(&state)?; + ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; + assert_pruned(&state, seeded[0].start..seeded[1].end)?; + state.db.persist()?; + seeded + }; + + let state = open(tmp.path(), 3 * HOUR)?; + ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; + // pruned events stay gone, and no watermark is three hours old, so + // nothing else goes until the 30 minute one ages past the new ttl + assert_pruned(&state, oldest.start..older.end)?; + assert_retained(&state, recent.start..fresh.end)?; + assert_eq!(watermark_seqs(&state)?.first(), Some(&recent.end)); + Ok(()) + } + + #[test] + fn a_table_straddling_the_cutoff_outlives_the_ttl() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let state = open(tmp.path(), HOUR)?; + let expired = write_events(&state, &["3kzbif5moa22m", "3kzbif5moa23m"])?; + plant_watermark(&state, 90 * MINUTE, expired.end)?; + let live = write_events(&state, &["3kzbif5mob22m"])?; + seal_table(&state)?; + + ephemeral_ttl_tick(&state.db, &state.ephemeral_ttl)?; + // drop_range only drops whole tables, so expired events sharing a + // table with live ones stay, bodies included, until a later cutoff + // passes the whole table + assert_retained(&state, expired.start..live.end)?; + Ok(()) + } +} -- 2.51.2