From c30958a32dbdc85244ddac563245b06ed0af8e71 Mon Sep 17 00:00:00 2001 From: phil Date: Tue, 5 Aug 2025 16:11:44 -0400 Subject: [PATCH] healthcheck --- slingshot/src/error.rs | 8 ++++++++ slingshot/src/healthcheck.rs | 32 ++++++++++++++++++++++++++++++++ slingshot/src/lib.rs | 2 ++ slingshot/src/main.rs | 15 ++++++++++++++- slingshot/src/server.rs | 3 ++- 5 files changed, 58 insertions(+), 2 deletions(-) create mode 100644 slingshot/src/healthcheck.rs diff --git a/slingshot/src/error.rs b/slingshot/src/error.rs index 17eb0ee..67b3e07 100644 --- a/slingshot/src/error.rs +++ b/slingshot/src/error.rs @@ -46,6 +46,12 @@ pub enum IdentityError { RefreshQueueKeyError(&'static str), } +#[derive(Debug, Error)] +pub enum HealthCheckError { + #[error("failed to send checkin: {0}")] + HealthCheckError(#[from] reqwest::Error), +} + #[derive(Debug, Error)] pub enum MainTaskError { #[error(transparent)] @@ -54,6 +60,8 @@ pub enum MainTaskError { ServerTaskError(#[from] ServerError), #[error(transparent)] IdentityTaskError(#[from] IdentityError), + #[error(transparent)] + HealthCheckError(#[from] HealthCheckError), #[error("firehose cache failed to close: {0}")] FirehoseCacheCloseError(foyer::Error), } diff --git a/slingshot/src/healthcheck.rs b/slingshot/src/healthcheck.rs new file mode 100644 index 0000000..fd60758 --- /dev/null +++ b/slingshot/src/healthcheck.rs @@ -0,0 +1,32 @@ +use crate::error::HealthCheckError; +use reqwest::Client; +use std::time::Duration; +use tokio::time::sleep; +use tokio_util::sync::CancellationToken; + +pub async fn healthcheck( + endpoint: String, + shutdown: CancellationToken, +) -> Result<(), HealthCheckError> { + let client = Client::builder() + .user_agent(format!( + "microcosm slingshot v{} (dev: @bad-example.com)", + env!("CARGO_PKG_VERSION") + )) + .no_proxy() + .timeout(Duration::from_secs(10)) + .build()?; + + loop { + tokio::select! { + res = client.get(&endpoint).send() => { + let _ = res + .and_then(|r| r.error_for_status()) + .inspect_err(|e| log::error!("failed to send healthcheck: {e}")); + }, + _ = shutdown.cancelled() => break, + } + sleep(Duration::from_secs(51)).await; + } + Ok(()) +} diff --git a/slingshot/src/lib.rs b/slingshot/src/lib.rs index 2e5e687..7737374 100644 --- a/slingshot/src/lib.rs +++ b/slingshot/src/lib.rs @@ -1,12 +1,14 @@ mod consumer; pub mod error; mod firehose_cache; +mod healthcheck; mod identity; mod record; mod server; pub use consumer::consume; pub use firehose_cache::firehose_cache; +pub use healthcheck::healthcheck; pub use identity::Identity; pub use record::{CachedRecord, ErrorResponseObject, Repo}; pub use server::serve; diff --git a/slingshot/src/main.rs b/slingshot/src/main.rs index a02d02e..ae5fdbb 100644 --- a/slingshot/src/main.rs +++ b/slingshot/src/main.rs @@ -1,7 +1,9 @@ // use foyer::HybridCache; // use foyer::{Engine, DirectFsDeviceOptions, HybridCacheBuilder}; use metrics_exporter_prometheus::PrometheusBuilder; -use slingshot::{Identity, Repo, consume, error::MainTaskError, firehose_cache, serve}; +use slingshot::{ + Identity, Repo, consume, error::MainTaskError, firehose_cache, healthcheck, serve, +}; use std::path::PathBuf; use clap::Parser; @@ -44,6 +46,9 @@ struct Args { /// recommended in production, but mind the file permissions. #[arg(long)] certs: Option, + /// an web address to send healtcheck pings to every ~51s or so + #[arg(long)] + healthcheck: Option, } #[tokio::main] @@ -127,6 +132,14 @@ async fn main() -> Result<(), String> { Ok(()) }); + if let Some(hc) = args.healthcheck { + let healthcheck_shutdown = shutdown.clone(); + tasks.spawn(async move { + healthcheck(hc, healthcheck_shutdown).await?; + Ok(()) + }); + } + tokio::select! { _ = shutdown.cancelled() => log::warn!("shutdown requested"), Some(r) = tasks.join_next() => { diff --git a/slingshot/src/server.rs b/slingshot/src/server.rs index db198a4..9066da1 100644 --- a/slingshot/src/server.rs +++ b/slingshot/src/server.rs @@ -19,7 +19,7 @@ use poem::{ Listener, TcpListener, acme::{AutoCert, LETS_ENCRYPT_PRODUCTION}, }, - middleware::{Cors, Tracing}, + middleware::{CatchPanic, Cors, Tracing}, }; use poem_openapi::{ ApiResponse, ContactObject, ExternalDocumentObject, Object, OpenApi, OpenApiService, Tags, @@ -758,6 +758,7 @@ where .allow_methods([Method::GET]) .allow_credentials(false), ) + .with(CatchPanic::new()) .with(Tracing); Server::new(listener) .name("slingshot") -- 2.51.2