very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342//! a C api for running hydrant inside another program.//!//! hydrant keeps serving its usual http api, and the embedder talks to it over that, ideally//! on a `unix:` bind. the C side only starts and stops it and learns when it stopped, so no//! events or callbacks ever cross the boundary.
use std::any::Any;use std::collections::HashMap;use std::ffi::{CStr, CString, c_char, c_int};use std::io::IsTerminal;use std::panic::{AssertUnwindSafe, catch_unwind};use std::sync::{Arc, Condvar, Mutex, PoisonError};use std::time::Duration;
use hydrant::config::Config;use hydrant::control::{ApiBinds, Hydrant};use hydrant::deps::futures::FutureExt;
#[path = "../../src/allocator.rs"]mod allocator;
/// a hydrant started by `hydrant_start`, until `hydrant_free`.pub struct HydrantHandle { exit: Arc<Exit>, // taken by the first `hydrant_shutdown` running: Mutex<Option<Running>>,}
struct Running { runtime: tokio::runtime::Runtime, hydrant: Hydrant,}
impl Running { fn shutdown(self, timeout: Duration) -> Result<(), String> { // not a plain drop, because that waits on blocking tasks, and some of those wait on // threads only `Hydrant::shutdown` stops self.runtime.shutdown_background(); self.hydrant.shutdown(timeout).map_err(report) }}
#[derive(Default)]struct Exit { result: Mutex<Option<Result<(), String>>>, changed: Condvar,}
impl Exit { /// the first result is the one `hydrant_wait` reports, so a failure isn't /// overwritten by the shutdown that follows it fn finish(&self, result: Result<(), String>) { let mut current = self.result.lock().unwrap_or_else(PoisonError::into_inner); if current.is_none() { *current = Some(result); self.changed.notify_all(); } }}
/// starts hydrant from a json object of the `HYDRANT_*` settings the binary reads from its/// environment, plus `RUST_LOG`. the process environment and `.env` are ignored./// `HYDRANT_API_BIND` is required, since the binary's default would expose the api on every/// interface. set it to `none` to run without an api.////// returns null and sets `*err` (free it with `hydrant_free_string`) on failure.////// # Safety/// `settings_json` must be a valid nul-terminated string, and `err` null or writable.#[unsafe(no_mangle)]pub unsafe extern "C" fn hydrant_start( settings_json: *const c_char, err: *mut *mut c_char,) -> *mut HydrantHandle { let started = guard(|| { if settings_json.is_null() { return Err("settings_json is null".to_string()); } // SAFETY: the caller promises a valid nul-terminated string let settings = unsafe { CStr::from_ptr(settings_json) }; start( settings .to_str() .map_err(|e| format!("settings_json: {e}"))?, ) }); match started { Ok(handle) => Box::into_raw(Box::new(handle)), Err(e) => { // SAFETY: the caller promises `err` is null or writable unsafe { set_err(err, e) }; std::ptr::null_mut() } }}
/// blocks until hydrant stops. returns -1 with `*err` set to why when it failed, or 0 once/// `hydrant_shutdown` stopped it. safe to call from several threads, and again after it/// returned. a failed hydrant still holds its database until `hydrant_shutdown`.////// # Safety/// `handle` must come from `hydrant_start` and not be freed yet, and `err` must be null or/// writable.#[unsafe(no_mangle)]pub unsafe extern "C" fn hydrant_wait( handle: *const HydrantHandle, err: *mut *mut c_char,) -> c_int { let result = guard(|| { // SAFETY: the caller promises a live handle from `hydrant_start` let handle = unsafe { handle.as_ref() }.ok_or("handle is null")?; let exit = &handle.exit; let result = exit.result.lock().unwrap_or_else(PoisonError::into_inner); let result = exit .changed .wait_while(result, |result| result.is_none()) .unwrap_or_else(PoisonError::into_inner); result .clone() .expect("wait_while only returns once a result is set") }); match result { Ok(()) => 0, Err(e) => { // SAFETY: the caller promises `err` is null or writable unsafe { set_err(err, e) }; -1 } }}
/// stops hydrant and waits up to `timeout_ms` for its database to close, which gives its/// memory back and lets `hydrant_start` open the same database again. `hydrant_wait` returns/// once this does. whatever hydrant was in the middle of is cut off like a crash, which it/// recovers from on the next start.////// returns -1 with `*err` set if the database is still open after the timeout, or if it was/// already shut down. the handle stays valid until `hydrant_free` either way.////// # Safety/// `handle` must come from `hydrant_start` and not be freed yet, and `err` must be null or/// writable.#[unsafe(no_mangle)]pub unsafe extern "C" fn hydrant_shutdown( handle: *const HydrantHandle, timeout_ms: u64, err: *mut *mut c_char,) -> c_int { let result = guard(|| { // SAFETY: the caller promises a live handle from `hydrant_start` let handle = unsafe { handle.as_ref() }.ok_or("handle is null")?; let running = handle .running .lock() .unwrap_or_else(PoisonError::into_inner) .take() .ok_or("hydrant was already shut down")?; let result = running.shutdown(Duration::from_millis(timeout_ms)); handle.exit.finish(Ok(())); result }); match result { Ok(()) => 0, Err(e) => { // SAFETY: the caller promises `err` is null or writable unsafe { set_err(err, e) }; -1 } }}
/// frees a handle. one that wasn't shut down is told to stop without waiting for it, so its/// database closes a moment later.////// # Safety/// `handle` must be null or come from `hydrant_start` and not be freed yet, and nothing may/// still be blocked in `hydrant_wait` on it.#[unsafe(no_mangle)]pub unsafe extern "C" fn hydrant_free(handle: *mut HydrantHandle) { if handle.is_null() { return; } // SAFETY: the caller promises it came from `Box::into_raw` in `hydrant_start` let handle = unsafe { Box::from_raw(handle) }; let running = handle .running .into_inner() .unwrap_or_else(PoisonError::into_inner); if let Some(running) = running { let _ = guard(|| running.shutdown(Duration::ZERO)); }}
/// frees an error string returned by this library.////// # Safety/// `s` must be null or a string from this library that wasn't freed yet.#[unsafe(no_mangle)]pub unsafe extern "C" fn hydrant_free_string(s: *mut c_char) { if !s.is_null() { // SAFETY: the caller promises the string came from `CString::into_raw` in `set_err` drop(unsafe { CString::from_raw(s) }); }}
fn start(settings_json: &str) -> Result<HydrantHandle, String> { let settings: HashMap<String, String> = serde_json::from_str(settings_json) .map_err(|e| format!("settings must be a json object of strings: {e}"))?; let lookup = |key: &str| settings.get(key).cloned(); if !settings.contains_key("HYDRANT_API_BIND") { return Err("HYDRANT_API_BIND is required, set it to `none` for no api".to_string()); } // before reading the config, so what that warns about gets printed init_logging( settings .get("RUST_LOG") .map_or(hydrant::DEFAULT_LOG_FILTER, String::as_str), ); let binds = ApiBinds::from_lookup(lookup, None).map_err(report)?; let config = Config::from_lookup(lookup).map_err(report)?; // fails when the embedder or an earlier start already installed one, which is fine let _ = hydrant::deps::rustls::crypto::aws_lc_rs::default_provider().install_default();
let runtime = tokio::runtime::Builder::new_multi_thread() .enable_all() .thread_name("hydrant") .build() .map_err(|e| format!("failed to start the tokio runtime: {e}"))?; let hydrant = runtime.block_on(Hydrant::new(config)).map_err(report)?;
let exit = Arc::new(Exit::default()); runtime.spawn({ let exit = exit.clone(); let hydrant = hydrant.clone(); async move { let serve = binds .map(|binds| hydrant.serve(binds).boxed()) .unwrap_or_else(|| std::future::pending().boxed()); let result = AssertUnwindSafe(async { tokio::select! { r = hydrant.run()? => r, r = serve => r, } }) // tokio would swallow a panic here and hydrant_wait would block forever .catch_unwind() .await .unwrap_or_else(|panic| { Err(miette::miette!( "hydrant panicked: {}", panic_message(&*panic) )) }); exit.finish(result.map_err(report)); } });
Ok(HydrantHandle { exit, running: Mutex::new(Some(Running { runtime, hydrant })), })}
fn init_logging(filter: &str) { let filter = tracing_subscriber::EnvFilter::builder() .with_default_directive(tracing_subscriber::filter::LevelFilter::INFO.into()) .parse_lossy(filter); // a global subscriber can only be set once per process, so a second start (or an embedder // with its own) keeps the first one let _ = tracing_subscriber::fmt() .with_env_filter(filter) .with_writer(std::io::stderr) .with_ansi(std::io::stderr().is_terminal()) .try_init();}
/// plain text, because miette's debug rendering is meant for terminals and may carry colorsfn report(e: miette::Report) -> String { e.chain() .map(ToString::to_string) .collect::<Vec<_>>() .join(": ")}
/// unwinding into C is undefined behavior, so a panic becomes an error instead.fn guard<T>(f: impl FnOnce() -> Result<T, String>) -> Result<T, String> { catch_unwind(AssertUnwindSafe(f)) .unwrap_or_else(|panic| Err(format!("hydrant panicked: {}", panic_message(&*panic))))}
fn panic_message(panic: &(dyn Any + Send)) -> String { panic .downcast_ref::<&str>() .map(|s| s.to_string()) .or_else(|| panic.downcast_ref::<String>().cloned()) .unwrap_or_else(|| "unknown panic".to_string())}
/// # Safety/// `err` must be null or writable.unsafe fn set_err(err: *mut *mut c_char, msg: String) { if err.is_null() { return; } let msg = CString::new(msg.replace('\0', "\\0")).expect("nul bytes were escaped"); // SAFETY: checked non-null above, the caller promises it's writable unsafe { *err = msg.into_raw() };}
#[cfg(test)]mod tests { use super::*;
#[test] fn bad_settings_are_errors_naming_the_variable() { for (key, value) in [ ("HYDRANT_ENABLE_FIREHOSE", "no"), ("HYDRANT_FIREHOSE_WORKERS", "0"), ("HYDRANT_CURSOR_SAVE_INTERVAL", "300"), ] { // if a bad value ever gets through, hydrant starts on a throwaway database let db = std::env::temp_dir().join(format!("hydrant-ffi-test-{}", std::process::id())); let db = db.to_string_lossy(); let settings = HashMap::from([ ("HYDRANT_API_BIND", "none"), ("HYDRANT_DATABASE_PATH", &db), ("HYDRANT_RELAY_HOSTS", ""), ("HYDRANT_CRAWLER_URLS", ""), ("HYDRANT_ENABLE_CRAWLER", "false"), (key, value), ]); let Err(err) = start(&serde_json::to_string(&settings).unwrap()) else { panic!("{key}={value:?} was accepted"); }; assert!( err.contains(key) && err.contains(value), "{key}={value:?}: {err}" ); } }}