Something went wrong. Try again.
atproto Thingiverse but good
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146//! Server-side indexing infrastructure.//!//! [Hydrant](hydrant) consumes the ATProto firehose for `space.polymodel.*`//! collections. A Tokio task projects each event into the SQLite derived tables//! (see [`projection`]) and persists a durable cursor so the pipeline resumes//! after a crash without skipping or re-processing beyond the last success.//!//! Two durable stores back the server://! 1. Hydrant's fjall store (`HYDRANT_DATABASE_PATH`) for raw ATProto records.//! 2. The SQLite projection database (`DATABASE_URL`) for app-specific views.
#![cfg_attr(test, allow(dead_code))]
pub mod config;pub mod db;pub mod projection;pub mod sample_data;pub mod setup;
use std::time::Duration;
use futures::StreamExt;use hydrant::control::Hydrant;use sqlx::SqlitePool;
/// Start the indexing pipeline: initialize Hydrant, resume from the persisted/// cursor, and spawn the firehose driver plus the projection consumer.////// Both background tasks run for the lifetime of the server's Tokio runtime, so/// this must be called from within that runtime (the `dioxus::serve` closure in/// `main` runs there).pub async fn start_indexing(db: SqlitePool) -> anyhow::Result<()> { let hydrant = setup::init_hydrant().await?;
// Drive the firehose + backfill for the process lifetime. run() takes &self // and returns a future that resolves when a fatal component exits; it is // called inside the spawned task because the returned future borrows the // handle (Rust 2024 lifetime capture). Hydrant's fatal paths (crawler, // firehose worker, stats ticker) call process::abort() on unexpected exit, so // a non-aborting resolution here is rare; in any case the supervised consumer // keeps indexing alive by resubscribing from the persisted cursor. let runner = hydrant.clone(); tokio::spawn(async move { let fut = match runner.run() { Ok(fut) => fut, Err(e) => { tracing::error!(error = %e, "hydrant::run() failed to start"); return; } }; if let Err(e) = fut.await { tracing::error!(error = %e, "hydrant driver exited with error"); } });
// Supervised projection consumer. The stream can terminate (a slow consumer // gets a StreamError and the channel closes); when it does, resubscribe from // the persisted cursor after backoff instead of stopping. A process restart // also resumes correctly, but supervision keeps indexing live without one. tokio::spawn(consume_loop(hydrant.clone(), db.clone()));
Ok(())}
/// Maximum backoff for resubscribe / projection retry (exponential, capped).const MAX_BACKOFF: Duration = Duration::from_secs(60);
/// Supervised consumer loop: subscribe, project each event until the stream/// ends, then resubscribe from the persisted cursor with exponential backoff.////// The cursor never advances past an event whose projection failed: that event/// is retried in place before any later event is read. Projection is idempotent,/// so replaying across a resubscribe is safe.async fn consume_loop(hydrant: Hydrant, db: SqlitePool) { let mut reconnect = Duration::from_secs(2); loop { // Resume from the last successfully projected event. let cursor = match db::load_cursor(&db).await { Ok(c) => c.unwrap_or(0), Err(e) => { tracing::error!(error = %e, "cursor load failed; will retry"); tokio::time::sleep(reconnect).await; reconnect = (reconnect * 2).min(MAX_BACKOFF); continue; } };
let mut events = hydrant.subscribe(Some(cursor)); tracing::info!(cursor, "(re)subscribed to hydrant event stream"); reconnect = Duration::from_secs(2); // reset after a successful subscribe
// Consume until the stream terminates. while let Some(result) = events.next().await { let event = match result { Ok(event) => event, Err(e) => { // Transient stream error: keep draining. If the channel has // closed, the next poll returns None and we resubscribe. tracing::warn!(error = %e, "hydrant stream error"); continue; } };
let pe = match projection::ProjectionInput::from_hydrant(&event) { Ok(Some(pe)) => pe, Ok(None) => { // Account/non-projectable event: not projected, but still advance // the durable cursor so a reconnect/resume doesn't replay it. if let Err(e) = db::advance_cursor(&db, event.id).await { tracing::warn!(error = %e, seq = event.id, "cursor advance for skipped event failed"); } continue; } Err(e) => { tracing::warn!(error = %e, "dropping unconvertible event"); continue; } };
// Retry the same event until it projects. A persistently-failing // ("poison") event wedges here and eventually stalls the stream, // triggering a resubscribe; it never advances the cursor, so no data // is lost. Loud per-retry logging keeps it observable. let mut retry = Duration::from_secs(1); loop { match projection::project_event(&db, &pe).await { Ok(()) => break, Err(e) => { tracing::error!(error = %e, seq = pe.seq(), "projection failed; retrying"); tokio::time::sleep(retry).await; retry = (retry * 2).min(MAX_BACKOFF); } } } }
// Stream ended (next() returned None): back off and resubscribe. tracing::warn!( reconnect_secs = reconnect.as_secs(), "hydrant event stream ended; resubscribing" ); tokio::time::sleep(reconnect).await; reconnect = (reconnect * 2).min(MAX_BACKOFF); }}