diff --git a/slingshot/src/consumer.rs b/slingshot/src/consumer.rs index ee270d2..db9f9c0 100644 --- a/slingshot/src/consumer.rs +++ b/slingshot/src/consumer.rs @@ -1,5 +1,5 @@ -use crate::CachedRecord; use crate::error::ConsumerError; +use crate::{CachedRecord, Identity, IdentityKey}; use foyer::HybridCache; use jetstream::{ DefaultJetstreamEndpoints, JetstreamCompression, JetstreamConfig, JetstreamConnector, @@ -11,6 +11,7 @@ pub async fn consume( jetstream_endpoint: String, cursor: Option, no_zstd: bool, + identity: Identity, shutdown: CancellationToken, cache: HybridCache, ) -> Result<(), ConsumerError> { @@ -46,33 +47,48 @@ pub async fn consume( break; }; - if event.kind != EventKind::Commit { - continue; - } - let Some(ref mut commit) = event.commit else { - log::warn!("consumer: commit event missing commit data, ignoring"); - continue; - }; + match event.kind { + EventKind::Commit => { + let Some(ref mut commit) = event.commit else { + log::warn!("consumer: commit event missing commit data, ignoring"); + continue; + }; - // TODO: something a bit more robust - let at_uri = format!( - "at://{}/{}/{}", - &*event.did, &*commit.collection, &*commit.rkey - ); + // TODO: something a bit more robust + let at_uri = format!( + "at://{}/{}/{}", + &*event.did, &*commit.collection, &*commit.rkey + ); - if commit.operation == CommitOp::Delete { - cache.insert(at_uri, CachedRecord::Deleted); - } else { - let Some(record) = commit.record.take() else { - log::warn!("consumer: commit insert or update missing record, ignoring"); - continue; - }; - let Some(cid) = commit.cid.take() else { - log::warn!("consumer: commit insert or update missing CID, ignoring"); - continue; - }; + if commit.operation == CommitOp::Delete { + cache.insert(at_uri, CachedRecord::Deleted); + } else { + let Some(record) = commit.record.take() else { + log::warn!("consumer: commit insert or update missing record, ignoring"); + continue; + }; + let Some(cid) = commit.cid.take() else { + log::warn!("consumer: commit insert or update missing CID, ignoring"); + continue; + }; - cache.insert(at_uri, CachedRecord::Found((cid, record).into())); + cache.insert(at_uri, CachedRecord::Found((cid, record).into())); + } + } + EventKind::Identity => { + let Some(ident) = event.identity else { + log::warn!("consumer: identity event missing identity data, ignoring"); + continue; + }; + if let Some(handle) = ident.handle { + metrics::counter!("identity_handle_refresh_queued", "reason" => "identity event").increment(1); + identity.queue_refresh(IdentityKey::Handle(handle)).await; + } + metrics::counter!("identity_did_refresh_queued", "reason" => "identity event") + .increment(1); + identity.queue_refresh(IdentityKey::Did(ident.did)).await; + } + EventKind::Account => {} // TODO: handle account events (esp hiding content on deactivate, clearing on delete) } } diff --git a/slingshot/src/identity.rs b/slingshot/src/identity.rs index 4e24b1c..5e07ce6 100644 --- a/slingshot/src/identity.rs +++ b/slingshot/src/identity.rs @@ -38,7 +38,7 @@ const MIN_TTL: Duration = Duration::from_secs(4 * 3600); // probably shoudl have const MIN_NOT_FOUND_TTL: Duration = Duration::from_secs(60); #[derive(Debug, Clone, Hash, PartialEq, Eq, Serialize, Deserialize)] -enum IdentityKey { +pub enum IdentityKey { Handle(Handle), Did(Did), } @@ -186,7 +186,7 @@ pub struct Identity { /// multi-producer *single consumer* queue refresh_queue: Arc>, /// just a lock to ensure only one refresher (queue consumer) is running (to be improved with a better refresher) - refresher: Arc>, + refresher_task: Arc>, } impl Identity { @@ -225,7 +225,7 @@ impl Identity { did_resolver: Arc::new(did_resolver), cache, refresh_queue: Default::default(), - refresher: Default::default(), + refresher_task: Default::default(), }) } @@ -293,12 +293,14 @@ impl Identity { } IdentityData::NotFound => { if (now - *last_fetch) >= MIN_NOT_FOUND_TTL { + metrics::counter!("identity_handle_refresh_queued", "reason" => "ttl", "found" => "false").increment(1); self.queue_refresh(key).await; } Ok(None) } IdentityData::Did(did) => { if (now - *last_fetch) >= MIN_TTL { + metrics::counter!("identity_handle_refresh_queued", "reason" => "ttl", "found" => "true").increment(1); self.queue_refresh(key).await; } Ok(Some(did.clone())) @@ -347,12 +349,14 @@ impl Identity { } IdentityData::NotFound => { if (now - *last_fetch) >= MIN_NOT_FOUND_TTL { + metrics::counter!("identity_did_refresh_queued", "reason" => "ttl", "found" => "false").increment(1); self.queue_refresh(key).await; } Ok(None) } IdentityData::Doc(mini_did) => { if (now - *last_fetch) >= MIN_TTL { + metrics::counter!("identity_did_refresh_queued", "reason" => "ttl", "found" => "true").increment(1); self.queue_refresh(key).await; } Ok(Some(mini_did.clone())) @@ -363,7 +367,7 @@ impl Identity { /// put a refresh task on the queue /// /// this can be safely called from multiple concurrent tasks - async fn queue_refresh(&self, key: IdentityKey) { + pub async fn queue_refresh(&self, key: IdentityKey) { // todo: max queue size let mut q = self.refresh_queue.lock().await; if !q.items.contains(&key) { @@ -440,7 +444,7 @@ impl Identity { /// run the refresh queue consumer pub async fn run_refresher(&self, shutdown: CancellationToken) -> Result<(), IdentityError> { let _guard = self - .refresher + .refresher_task .try_lock() .expect("there to only be one refresher running"); loop { @@ -462,18 +466,22 @@ impl Identity { log::trace!("refreshing handle {handle:?}"); match self.handle_resolver.resolve(handle).await { Ok(did) => { + metrics::counter!("identity_handle_refresh", "success" => "true") + .increment(1); self.cache.insert( task_key.clone(), IdentityVal(UtcDateTime::now(), IdentityData::Did(did)), ); } Err(atrium_identity::Error::NotFound) => { + metrics::counter!("identity_handle_refresh", "success" => "false", "reason" => "not found").increment(1); self.cache.insert( task_key.clone(), IdentityVal(UtcDateTime::now(), IdentityData::NotFound), ); } Err(err) => { + metrics::counter!("identity_handle_refresh", "success" => "false", "reason" => "other").increment(1); log::warn!( "failed to refresh handle: {err:?}. leaving stale (should we eventually do something?)" ); @@ -488,6 +496,7 @@ impl Identity { Ok(did_doc) => { // TODO: fix in atrium: should verify id is did if did_doc.id != did.to_string() { + metrics::counter!("identity_did_refresh", "success" => "false", "reason" => "wrong did").increment(1); log::warn!( "refreshed did doc failed: wrong did doc id. dropping refresh." ); @@ -496,24 +505,29 @@ impl Identity { let mini_doc = match did_doc.try_into() { Ok(md) => md, Err(e) => { + metrics::counter!("identity_did_refresh", "success" => "false", "reason" => "bad doc").increment(1); log::warn!( "converting mini doc failed: {e:?}. dropping refresh." ); continue; } }; + metrics::counter!("identity_did_refresh", "success" => "true") + .increment(1); self.cache.insert( task_key.clone(), IdentityVal(UtcDateTime::now(), IdentityData::Doc(mini_doc)), ); } Err(atrium_identity::Error::NotFound) => { + metrics::counter!("identity_did_refresh", "success" => "false", "reason" => "not found").increment(1); self.cache.insert( task_key.clone(), IdentityVal(UtcDateTime::now(), IdentityData::NotFound), ); } Err(err) => { + metrics::counter!("identity_did_refresh", "success" => "false", "reason" => "other").increment(1); log::warn!( "failed to refresh did doc: {err:?}. leaving stale (should we eventually do something?)" ); diff --git a/slingshot/src/lib.rs b/slingshot/src/lib.rs index 7737374..e00e566 100644 --- a/slingshot/src/lib.rs +++ b/slingshot/src/lib.rs @@ -9,6 +9,6 @@ mod server; pub use consumer::consume; pub use firehose_cache::firehose_cache; pub use healthcheck::healthcheck; -pub use identity::Identity; +pub use identity::{Identity, IdentityKey}; pub use record::{CachedRecord, ErrorResponseObject, Repo}; pub use server::serve; diff --git a/slingshot/src/main.rs b/slingshot/src/main.rs index 19b49b3..56d7500 100644 --- a/slingshot/src/main.rs +++ b/slingshot/src/main.rs @@ -154,13 +154,14 @@ async fn main() -> Result<(), String> { let repo = Repo::new(identity.clone()); + let identity_for_server = identity.clone(); let server_shutdown = shutdown.clone(); let server_cache_handle = cache.clone(); let bind = args.bind; tasks.spawn(async move { serve( server_cache_handle, - identity, + identity_for_server, repo, args.acme_domain, args.acme_contact, @@ -173,6 +174,7 @@ async fn main() -> Result<(), String> { Ok(()) }); + let identity_refreshable = identity.clone(); let consumer_shutdown = shutdown.clone(); let consumer_cache = cache.clone(); tasks.spawn(async move { @@ -180,6 +182,7 @@ async fn main() -> Result<(), String> { args.jetstream, None, args.jetstream_no_zstd, + identity_refreshable, consumer_shutdown, consumer_cache, ) diff --git a/slingshot/src/server.rs b/slingshot/src/server.rs index 1552299..26505cf 100644 --- a/slingshot/src/server.rs +++ b/slingshot/src/server.rs @@ -12,7 +12,7 @@ use std::sync::Arc; use tokio_util::sync::CancellationToken; use poem::{ - Endpoint, EndpointExt, Route, Server, + Endpoint, EndpointExt, IntoResponse, Route, Server, endpoint::{StaticFileEndpoint, make_sync}, http::Method, listener::{ @@ -772,7 +772,9 @@ where .allow_credentials(false), ) .with(CatchPanic::new()) + .around(request_counter) .with(Tracing); + Server::new(listener) .name("slingshot") .run_with_graceful_shutdown(app, shutdown.cancelled(), None) @@ -780,3 +782,17 @@ where .map_err(ServerError::ServerExited) .inspect(|()| log::info!("server ended. goodbye.")) } + +async fn request_counter(next: E, req: poem::Request) -> poem::Result { + let t0 = std::time::Instant::now(); + let method = req.method().to_string(); + let path = req.uri().path().to_string(); + let res = next.call(req).await?.into_response(); + metrics::histogram!( + "server_request", + "endpoint" => format!("{method} {path}"), + "status" => res.status().to_string(), + ) + .record(t0.elapsed()); + Ok(res) +}