From 322ef21090bb59ead6f9400a8eed74247a841579 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 30 Sep 2026 08:47:37 +0300 Subject: [PATCH] [stream] expose the newest stored event id at /stream/head --- docs/api/README.md | 1 + docs/api/stream.md | 12 ++++++ docs/getting-started.md | 1 + src/api/stream.rs | 72 +++++++++++++++++++++++++++++++++-- tests/stream_cursor_replay.nu | 13 +++++++ 5 files changed, 96 insertions(+), 3 deletions(-) diff --git a/docs/api/README.md b/docs/api/README.md index 3d25c0b..5bed8f3 100644 --- a/docs/api/README.md +++ b/docs/api/README.md @@ -7,6 +7,7 @@ hydrant's REST API is split into public endpoints (safe to expose) and managemen ## public - `GET /stream`: [stream API](stream.md), subscribe to the ordered websocket event stream. +- `GET /stream/head`: [stream API](stream.md#get-streamhead), get the newest durable event ID to detect when a replay has caught up. - `GET /subscribe`: [jetstream API](jetstream.md), subscribe to a filtered, jetstream-compatible websocket event stream. - `GET /stats`: get stats about the database (counts of repos, records, events; sizes of keyspaces on disk). - `GET /health` / `GET /_health` / `GET /version` / `GET /_version`: health check; returns the running version as JSON (`{"name","version","mode"}`). diff --git a/docs/api/stream.md b/docs/api/stream.md index a61ef99..cc3a47d 100644 --- a/docs/api/stream.md +++ b/docs/api/stream.md @@ -137,3 +137,15 @@ account events are live-only and are not available through cursor replay. "message": "stream socket send blocked for at least 30 seconds" } ``` + +## GET /stream/head + +get the ID of the newest durable record event. + +```json +{ "id": 289455 } +``` + +`id` is `null` while no events are stored. identity and account events consume IDs without being persisted, so this is the newest ID a cursor replay can deliver, not the ID counter. + +a consumer can snapshot the head before subscribing with `cursor=0` and treat the replay as caught up once it has seen a frame with `id` greater than or equal to the snapshot. keep an idle-timeout fallback: if the head event itself cannot be inflated, it is skipped and never delivered. diff --git a/docs/getting-started.md b/docs/getting-started.md index 412b8fd..59dc4a7 100644 --- a/docs/getting-started.md +++ b/docs/getting-started.md @@ -41,6 +41,7 @@ it is **highly recommended** to run hydrant behind a reverse proxy (like nginx o - `/xrpc/*`: XRPC endpoints. - `/stream`: hydrant's ordered event stream. +- `/stream/head`: the newest durable event ID, for detecting when a replay has caught up. - `/stats`: general database statistics. - `/health` / `/_health` / `/version` / `/_version`: health check; returns the running version as JSON. diff --git a/src/api/stream.rs b/src/api/stream.rs index b4cf76a..a643a53 100644 --- a/src/api/stream.rs +++ b/src/api/stream.rs @@ -2,17 +2,35 @@ use crate::control::Hydrant; use axum::Router; use axum::routing::get; use axum::{ + Json, extract::{Query, State}, - response::IntoResponse, + http::StatusCode, + response::{IntoResponse, Response}, }; use axum_tws::{Message, WebSocket, WebSocketUpgrade}; -use serde::Deserialize; +use serde::{Deserialize, Serialize}; use tracing::error; use super::ws::{WsAction, control_frame_only_limits, run_socket}; pub fn router() -> Router { - Router::new().route("/", get(handle_stream)) + Router::new() + .route("/", get(handle_stream)) + .route("/head", get(handle_head)) +} + +/// the newest stored event id. identity and account events take ids too but +/// are never replayed, so this is the furthest a cursor replay can reach. +#[derive(Serialize)] +struct StreamHead { + id: Option, +} + +pub async fn handle_head(State(hydrant): State) -> Response { + match hydrant.stream_head() { + Ok(id) => Json(StreamHead { id }).into_response(), + Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(), + } } #[derive(Deserialize)] @@ -67,3 +85,51 @@ async fn handle_socket(socket: WebSocket, hydrant: Hydrant, query: StreamQuery) ) .await; } + +#[cfg(test)] +mod tests { + use super::*; + use jacquard_common::types::did::Did; + use miette::IntoDiagnostic; + + async fn head(hydrant: &Hydrant) -> miette::Result { + let response = handle_head(State(hydrant.clone())).await; + assert_eq!(response.status(), StatusCode::OK); + let body = axum::body::to_bytes(response.into_body(), 1024) + .await + .into_diagnostic()?; + serde_json::from_slice(&body).into_diagnostic() + } + + #[tokio::test] + async fn head_tracks_durable_events_not_the_id_counter() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let hydrant = Hydrant::new(crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }) + .await?; + assert_eq!(head(&hydrant).await?, serde_json::json!({ "id": null })); + + let ids = &hydrant.state.db.stream.ids; + hydrant.seed_events_for_bench(3); + let last_seeded = ids.head().expect("seeded events took ids"); + assert_eq!( + head(&hydrant).await?, + serde_json::json!({ "id": last_seeded }) + ); + + let did = Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa").into_diagnostic()?; + crate::ops::announce_identity(&hydrant.state.db, &did, None).await; + assert_eq!( + ids.head(), + Some(last_seeded + 1), + "the identity event must have taken an id" + ); + assert_eq!( + head(&hydrant).await?, + serde_json::json!({ "id": last_seeded }) + ); + Ok(()) + } +} diff --git a/tests/stream_cursor_replay.nu b/tests/stream_cursor_replay.nu index 9fc2ac0..83434cd 100644 --- a/tests/stream_cursor_replay.nu +++ b/tests/stream_cursor_replay.nu @@ -18,6 +18,11 @@ def main [] { fail "api failed to start" $instance.pid } + let empty_head = (http get $"($url)/stream/head").id + if $empty_head != null { + fail $"expected no stream head before backfill, got ($empty_head)" $instance.pid + } + print $"adding repo ($did) to tracking..." http put -t application/json $"($url)/repos" [{ did: $did }] @@ -33,6 +38,9 @@ def main [] { fail "expected at least one stored event after backfill" $instance.pid } + let head = (http get $"($url)/stream/head").id + print $"stream head: ($head)" + let history_output = $"($db_path)/stream_history.txt" print "starting historical stream listener..." let history_pid = (bash -c $"websocat --no-async-stdio -n '($ws_url)?cursor=0' > '($history_output)' 2>&1 & echo $!" | str trim | into int) @@ -71,6 +79,11 @@ def main [] { print $"first event id=($first.id), type=($first.type)" print $"last event id=($last.id), type=($last.type)" + let replayed_ids = ($history_messages | each { |m| ($m | from json).id }) + if $head not-in $replayed_ids { + fail $"replay never delivered the /stream/head id ($head)" $instance.pid + } + if ($history_messages | length) < $events_count { print $"warning: replay returned ($history_messages | length)/($events_count) stored events" } -- 2.51.2