From 57c0d71d2ed8e74e2354df08acb99bb1fa66707c Mon Sep 17 00:00:00 2001 From: phil Date: Wed, 25 Jun 2025 17:05:31 -0400 Subject: [PATCH] wow fmt --- Makefile | 2 +- spacedust/src/consumer.rs | 18 ++++++---- spacedust/src/delay.rs | 4 +-- spacedust/src/lib.rs | 27 +++++++++----- spacedust/src/main.rs | 17 +++++---- spacedust/src/removable_delay_queue.rs | 26 ++++++++------ spacedust/src/server.rs | 36 ++++++++++--------- spacedust/src/subscriber.rs | 49 ++++++++++++-------------- who-am-i/Cargo.toml | 2 +- who-am-i/src/lib.rs | 2 +- who-am-i/src/main.rs | 2 +- who-am-i/src/oauth.rs | 31 ++++++++-------- who-am-i/src/server.rs | 34 +++++++++--------- 13 files changed, 135 insertions(+), 115 deletions(-) diff --git a/Makefile b/Makefile index 0fd5e48..4803c3c 100644 --- a/Makefile +++ b/Makefile @@ -5,7 +5,7 @@ test: cargo test --all-features fmt: - cargo fmt --package links --package constellation --package ufos + cargo fmt --package links --package constellation --package ufos --package spacedust --package who-am-i cargo +nightly fmt --package jetstream clippy: diff --git a/spacedust/src/consumer.rs b/spacedust/src/consumer.rs index bf78686..fdad8ac 100644 --- a/spacedust/src/consumer.rs +++ b/spacedust/src/consumer.rs @@ -1,5 +1,3 @@ -use std::sync::Arc; -use tokio_util::sync::CancellationToken; use crate::ClientMessage; use crate::error::ConsumerError; use crate::removable_delay_queue; @@ -8,7 +6,9 @@ use jetstream::{ events::{CommitOp, Cursor, EventKind}, }; use links::collect_links; +use std::sync::Arc; use tokio::sync::broadcast; +use tokio_util::sync::CancellationToken; const MAX_LINKS_PER_EVENT: usize = 100; @@ -61,12 +61,16 @@ pub async fn consume( }; // TODO: something a bit more robust - let at_uri = format!("at://{}/{}/{}", &*event.did, &*commit.collection, &*commit.rkey); + let at_uri = format!( + "at://{}/{}/{}", + &*event.did, &*commit.collection, &*commit.rkey + ); // TODO: keep a buffer and remove quick deletes to debounce notifs // for now we just drop all deletes eek if commit.operation == CommitOp::Delete { - d.remove_range((at_uri.clone(), 0)..=(at_uri.clone(), MAX_LINKS_PER_EVENT)).await; + d.remove_range((at_uri.clone(), 0)..=(at_uri.clone(), MAX_LINKS_PER_EVENT)) + .await; continue; } let Some(ref record) = commit.record else { @@ -86,7 +90,8 @@ pub async fn consume( if i >= MAX_LINKS_PER_EVENT { // todo: indicate if the link limit was reached (-> links omitted) log::warn!("consumer: event has too many links, ignoring the rest"); - metrics::counter!("consumer_dropped_links", "reason" => "too_many_links").increment(1); + metrics::counter!("consumer_dropped_links", "reason" => "too_many_links") + .increment(1); break; } let client_message = match ClientMessage::new_link(link, &at_uri, commit) { @@ -94,7 +99,8 @@ pub async fn consume( Err(e) => { // TODO indicate to clients that a link has been dropped log::warn!("consumer: failed to serialize link to json: {e:?}"); - metrics::counter!("consumer_dropped_links", "reason" => "failed_to_serialize").increment(1); + metrics::counter!("consumer_dropped_links", "reason" => "failed_to_serialize") + .increment(1); continue; } }; diff --git a/spacedust/src/delay.rs b/spacedust/src/delay.rs index 54eb97a..94044df 100644 --- a/spacedust/src/delay.rs +++ b/spacedust/src/delay.rs @@ -1,7 +1,7 @@ +use crate::error::DelayError; use crate::removable_delay_queue; -use tokio_util::sync::CancellationToken; use tokio::sync::broadcast; -use crate::error::DelayError; +use tokio_util::sync::CancellationToken; pub async fn to_broadcast( source: removable_delay_queue::Output<(String, usize), T>, diff --git a/spacedust/src/lib.rs b/spacedust/src/lib.rs index 049aa86..915f061 100644 --- a/spacedust/src/lib.rs +++ b/spacedust/src/lib.rs @@ -1,15 +1,15 @@ pub mod consumer; pub mod delay; pub mod error; +pub mod removable_delay_queue; pub mod server; pub mod subscriber; -pub mod removable_delay_queue; -use links::CollectedLink; use jetstream::events::CommitEvent; -use tokio_tungstenite::tungstenite::Message; +use links::CollectedLink; use serde::{Deserialize, Serialize}; use server::MultiSubscribeQuery; +use tokio_tungstenite::tungstenite::Message; #[derive(Debug)] pub struct FilterableProperties { @@ -32,7 +32,11 @@ pub struct ClientMessage { } impl ClientMessage { - pub fn new_link(link: CollectedLink, at_uri: &str, commit: &CommitEvent) -> Result { + pub fn new_link( + link: CollectedLink, + at_uri: &str, + commit: &CommitEvent, + ) -> Result { let subject_did = link.target.did(); let subject = link.target.into_string(); @@ -61,16 +65,23 @@ impl ClientMessage { let message = Message::Text(client_event_json.into()); - let properties = FilterableProperties { subject, subject_did, source }; + let properties = FilterableProperties { + subject, + subject_did, + source, + }; - Ok(ClientMessage { message, properties }) + Ok(ClientMessage { + message, + properties, + }) } } #[derive(Debug, Serialize)] -#[serde(rename_all="snake_case")] +#[serde(rename_all = "snake_case")] pub struct ClientEvent { - kind: &'static str, // "link" + kind: &'static str, // "link" origin: &'static str, // "live", "replay", "backfill" link: ClientLinkEvent, } diff --git a/spacedust/src/main.rs b/spacedust/src/main.rs index 5e138e1..54db30e 100644 --- a/spacedust/src/main.rs +++ b/spacedust/src/main.rs @@ -1,14 +1,14 @@ -use spacedust::error::MainTaskError; use spacedust::consumer; -use spacedust::server; use spacedust::delay; +use spacedust::error::MainTaskError; use spacedust::removable_delay_queue::removable_delay_queue; +use spacedust::server; use clap::Parser; use metrics_exporter_prometheus::PrometheusBuilder; +use std::time::Duration; use tokio::sync::broadcast; use tokio_util::sync::CancellationToken; -use std::time::Duration; /// Aggregate links in the at-mosphere #[derive(Parser, Debug, Clone)] @@ -80,15 +80,20 @@ async fn main() -> Result<(), String> { args.jetstream, None, args.jetstream_no_zstd, - consumer_shutdown + consumer_shutdown, ) - .await?; + .await?; Ok(()) }); let delay_shutdown = shutdown.clone(); tasks.spawn(async move { - delay::to_broadcast(delay_queue_receiver, consumer_delayed_sender, delay_shutdown).await?; + delay::to_broadcast( + delay_queue_receiver, + consumer_delayed_sender, + delay_shutdown, + ) + .await?; Ok(()) }); diff --git a/spacedust/src/removable_delay_queue.rs b/spacedust/src/removable_delay_queue.rs index cb2b09a..f528677 100644 --- a/spacedust/src/removable_delay_queue.rs +++ b/spacedust/src/removable_delay_queue.rs @@ -1,9 +1,9 @@ -use std::ops::RangeBounds; use std::collections::{BTreeMap, VecDeque}; -use std::time::{Duration, Instant}; -use tokio::sync::Mutex; +use std::ops::RangeBounds; use std::sync::Arc; +use std::time::{Duration, Instant}; use thiserror::Error; +use tokio::sync::Mutex; #[derive(Debug, Error)] pub enum EnqueueError { @@ -17,7 +17,7 @@ impl Key for T {} #[derive(Debug)] struct Queue { queue: VecDeque<(Instant, K)>, - items: BTreeMap + items: BTreeMap, } pub struct Input { @@ -49,7 +49,12 @@ impl Input { pub async fn remove_range(&self, range: impl RangeBounds) { let n = { let mut q = self.q.lock().await; - let keys = q.items.range(range).map(|(k, _)| k).cloned().collect::>(); + let keys = q + .items + .range(range) + .map(|(k, _)| k) + .cloned() + .collect::>(); for k in &keys { q.items.remove(k); } @@ -94,22 +99,21 @@ impl Output { } else { let overshoot = now.saturating_duration_since(expected_release); metrics::counter!("delay_queue_emit_total", "early" => "no").increment(1); - metrics::histogram!("delay_queue_emit_overshoot").record(overshoot.as_secs_f64()); + metrics::histogram!("delay_queue_emit_overshoot") + .record(overshoot.as_secs_f64()); } - return Some(item) + return Some(item); } else if Arc::strong_count(&self.q) == 1 { return None; } // the queue is *empty*, so we need to wait at least as long as the current delay tokio::time::sleep(self.delay).await; metrics::counter!("delay_queue_entirely_empty_total").increment(1); - }; + } } } -pub fn removable_delay_queue( - delay: Duration, -) -> (Input, Output) { +pub fn removable_delay_queue(delay: Duration) -> (Input, Output) { let q: Arc>> = Arc::new(Mutex::new(Queue { queue: VecDeque::new(), items: BTreeMap::new(), diff --git a/spacedust/src/server.rs b/spacedust/src/server.rs index 37bee28..a74b207 100644 --- a/spacedust/src/server.rs +++ b/spacedust/src/server.rs @@ -1,28 +1,26 @@ +use crate::ClientMessage; use crate::error::ServerError; use crate::subscriber::Subscriber; -use metrics::{histogram, counter}; -use std::sync::Arc; -use crate::ClientMessage; +use dropshot::{ + ApiDescription, ApiEndpointBodyContentType, Body, ConfigDropshot, ConfigLogging, + ConfigLoggingLevel, ExtractorMetadata, HttpError, HttpResponse, Query, RequestContext, + ServerBuilder, ServerContext, SharedExtractor, WebsocketConnection, channel, endpoint, +}; use http::{ - header::{ORIGIN, USER_AGENT}, Response, StatusCode, + header::{ORIGIN, USER_AGENT}, }; -use dropshot::{ - Body, - ApiDescription, ConfigDropshot, ConfigLogging, ConfigLoggingLevel, Query, RequestContext, - ServerBuilder, WebsocketConnection, channel, endpoint, HttpResponse, - ApiEndpointBodyContentType, ExtractorMetadata, HttpError, ServerContext, - SharedExtractor, -}; +use metrics::{counter, histogram}; +use std::sync::Arc; +use async_trait::async_trait; use schemars::JsonSchema; use serde::{Deserialize, Serialize}; +use std::collections::HashSet; use tokio::sync::broadcast; use tokio::time::Instant; use tokio_tungstenite::tungstenite::protocol::{Role, WebSocketConfig}; use tokio_util::sync::CancellationToken; -use async_trait::async_trait; -use std::collections::HashSet; const INDEX_HTML: &str = include_str!("../static/index.html"); const FAVICON: &[u8] = include_bytes!("../static/favicon.ico"); @@ -30,7 +28,7 @@ const FAVICON: &[u8] = include_bytes!("../static/favicon.ico"); pub async fn serve( b: broadcast::Sender>, d: broadcast::Sender>, - shutdown: CancellationToken + shutdown: CancellationToken, ) -> Result<(), ServerError> { let config_logging = ConfigLogging::StderrTerminal { level: ConfigLoggingLevel::Info, @@ -65,7 +63,12 @@ pub async fn serve( ); let sub_shutdown = shutdown.clone(); - let ctx = Context { spec, b, d, shutdown: sub_shutdown }; + let ctx = Context { + spec, + b, + d, + shutdown: sub_shutdown, + }; let server = ServerBuilder::new(api, ctx, log) .config(ConfigDropshot { @@ -162,7 +165,6 @@ where // TODO: cors for HttpError - /// Serve index page as html #[endpoint { method = GET, @@ -316,7 +318,7 @@ async fn subscribe( upgraded.into_inner(), Role::Server, Some(WebSocketConfig::default().max_message_size( - Some(10 * 2_usize.pow(20)) // 10MiB, matching jetstream + Some(10 * 2_usize.pow(20)), // 10MiB, matching jetstream )), ) .await; diff --git a/spacedust/src/subscriber.rs b/spacedust/src/subscriber.rs index afa4717..2e6a963 100644 --- a/spacedust/src/subscriber.rs +++ b/spacedust/src/subscriber.rs @@ -1,16 +1,16 @@ use crate::error::SubscriberUpdateError; -use std::sync::Arc; -use tokio::time::interval; -use std::time::Duration; -use futures::StreamExt; -use crate::{ClientMessage, FilterableProperties, SubscriberSourcedMessage}; use crate::server::MultiSubscribeQuery; +use crate::{ClientMessage, FilterableProperties, SubscriberSourcedMessage}; +use dropshot::WebsocketConnectionRaw; use futures::SinkExt; +use futures::StreamExt; use std::error::Error; +use std::sync::Arc; +use std::time::Duration; use tokio::sync::broadcast::{self, error::RecvError}; +use tokio::time::interval; use tokio_tungstenite::{WebSocketStream, tungstenite::Message}; use tokio_util::sync::CancellationToken; -use dropshot::WebsocketConnectionRaw; const PING_PERIOD: Duration = Duration::from_secs(30); @@ -20,17 +20,14 @@ pub struct Subscriber { } impl Subscriber { - pub fn new( - query: MultiSubscribeQuery, - shutdown: CancellationToken, - ) -> Self { + pub fn new(query: MultiSubscribeQuery, shutdown: CancellationToken) -> Self { Self { query, shutdown } } pub async fn start( mut self, ws: WebSocketStream, - mut receiver: broadcast::Receiver> + mut receiver: broadcast::Receiver>, ) -> Result<(), Box> { let mut ping_state = None; let (mut ws_sender, mut ws_receiver) = ws.split(); @@ -83,6 +80,7 @@ impl Subscriber { // TODO: send client an explanation self.shutdown.cancel(); } + log::trace!("subscriber updated with opts: {:?}", self.query); }, Some(Ok(m)) => log::trace!("subscriber sent an unexpected message: {m:?}"), Some(Err(e)) => { @@ -122,36 +120,35 @@ impl Subscriber { Ok(()) } - fn filter( - &self, - properties: &FilterableProperties, - ) -> bool { + fn filter(&self, properties: &FilterableProperties) -> bool { let query = &self.query; // subject + subject DIDs are logical OR - if !( - query.wanted_subjects.is_empty() && query.wanted_subject_dids.is_empty() || - query.wanted_subjects.contains(&properties.subject) || - properties.subject_did.as_ref().map(|did| query.wanted_subject_dids.contains(did)).unwrap_or(false) - ) { // wowwww ^^ fix that - return false + if !(query.wanted_subjects.is_empty() && query.wanted_subject_dids.is_empty() + || query.wanted_subjects.contains(&properties.subject) + || properties + .subject_did + .as_ref() + .map(|did| query.wanted_subject_dids.contains(did)) + .unwrap_or(false)) + { + // wowwww ^^ fix that + return false; } // subjects together with sources are logical AND if !(query.wanted_sources.is_empty() || query.wanted_sources.contains(&properties.source)) { - return false + return false; } true } } - - impl MultiSubscribeQuery { pub fn update_from_raw(&mut self, s: &str) -> Result<(), SubscriberUpdateError> { - let SubscriberSourcedMessage::OptionsUpdate(opts) = serde_json::from_str(s) - .map_err(SubscriberUpdateError::FailedToParseMessage)?; + let SubscriberSourcedMessage::OptionsUpdate(opts) = + serde_json::from_str(s).map_err(SubscriberUpdateError::FailedToParseMessage)?; if opts.wanted_sources.len() > 1_000 { return Err(SubscriberUpdateError::TooManySourcesWanted); } diff --git a/who-am-i/Cargo.toml b/who-am-i/Cargo.toml index b15b2d3..47e6742 100644 --- a/who-am-i/Cargo.toml +++ b/who-am-i/Cargo.toml @@ -4,7 +4,7 @@ version = "0.1.0" edition = "2024" [dependencies] -atrium-api = { version = "0.25.4", default-features = false, features = ["tokio", "agent"] } +atrium-api = { version = "0.25.4", default-features = false } atrium-identity = "0.1.5" atrium-oauth = "0.1.3" clap = { version = "4.5.40", features = ["derive"] } diff --git a/who-am-i/src/lib.rs b/who-am-i/src/lib.rs index 597444c..b112d21 100644 --- a/who-am-i/src/lib.rs +++ b/who-am-i/src/lib.rs @@ -3,5 +3,5 @@ mod oauth; mod server; pub use dns_resolver::HickoryDnsTxtResolver; +pub use oauth::{Client, authorize, client}; pub use server::serve; -pub use oauth::{Client, client, authorize}; diff --git a/who-am-i/src/main.rs b/who-am-i/src/main.rs index d4ebeee..6258eea 100644 --- a/who-am-i/src/main.rs +++ b/who-am-i/src/main.rs @@ -1,5 +1,5 @@ -use who_am_i::serve; use tokio_util::sync::CancellationToken; +use who_am_i::serve; #[tokio::main] async fn main() { diff --git a/who-am-i/src/oauth.rs b/who-am-i/src/oauth.rs index 37e797c..5bcd821 100644 --- a/who-am-i/src/oauth.rs +++ b/who-am-i/src/oauth.rs @@ -1,15 +1,14 @@ +use crate::HickoryDnsTxtResolver; use atrium_identity::{ did::{CommonDidResolver, CommonDidResolverConfig, DEFAULT_PLC_DIRECTORY_URL}, handle::{AtprotoHandleResolver, AtprotoHandleResolverConfig}, }; use atrium_oauth::{ - AuthorizeOptions, + AtprotoLocalhostClientMetadata, AuthorizeOptions, DefaultHttpClient, KnownScope, OAuthClient, + OAuthClientConfig, OAuthResolverConfig, Scope, store::{session::MemorySessionStore, state::MemoryStateStore}, - AtprotoLocalhostClientMetadata, DefaultHttpClient, KnownScope, OAuthClient, OAuthClientConfig, - OAuthResolverConfig, Scope, }; use std::sync::Arc; -use crate::HickoryDnsTxtResolver; pub type Client = OAuthClient< MemoryStateStore, @@ -23,9 +22,7 @@ pub fn client() -> Client { let config = OAuthClientConfig { client_metadata: AtprotoLocalhostClientMetadata { redirect_uris: Some(vec![String::from("http://127.0.0.1:9997/authorized")]), - scopes: Some(vec![ - Scope::Known(KnownScope::Atproto), - ]), + scopes: Some(vec![Scope::Known(KnownScope::Atproto)]), }, keys: None, resolver: OAuthResolverConfig { @@ -52,16 +49,16 @@ pub fn client() -> Client { } pub async fn authorize(client: &Client, handle: &str) -> String { - let Ok(url) = client.authorize( - handle, - AuthorizeOptions { - scopes: vec![ - Scope::Known(KnownScope::Atproto), - ], - ..Default::default() - }, - ) - .await else { + let Ok(url) = client + .authorize( + handle, + AuthorizeOptions { + scopes: vec![Scope::Known(KnownScope::Atproto)], + ..Default::default() + }, + ) + .await + else { panic!("failed to authorize"); }; url diff --git a/who-am-i/src/server.rs b/who-am-i/src/server.rs index 5a0fc1f..bb08d3c 100644 --- a/who-am-i/src/server.rs +++ b/who-am-i/src/server.rs @@ -1,17 +1,16 @@ - use atrium_api::agent::SessionManager; -use std::error::Error; -use metrics::{histogram, counter}; -use std::sync::Arc; +use dropshot::{ + ApiDescription, Body, ConfigDropshot, ConfigLogging, ConfigLoggingLevel, HttpError, + HttpResponse, HttpResponseSeeOther, Query, RequestContext, ServerBuilder, ServerContext, + endpoint, http_response_see_other, +}; use http::{ - header::{ORIGIN, USER_AGENT}, Response, StatusCode, + header::{ORIGIN, USER_AGENT}, }; -use dropshot::{ - Body, HttpResponseSeeOther, http_response_see_other, - ApiDescription, ConfigDropshot, ConfigLogging, ConfigLoggingLevel, RequestContext, - ServerBuilder, endpoint, HttpResponse, HttpError, ServerContext, Query, -}; +use metrics::{counter, histogram}; +use std::error::Error; +use std::sync::Arc; use atrium_oauth::CallbackParams; use schemars::JsonSchema; @@ -19,20 +18,17 @@ use serde::{Deserialize, Serialize}; use tokio::time::Instant; use tokio_util::sync::CancellationToken; -use crate::{Client, client, authorize}; +use crate::{Client, authorize, client}; const INDEX_HTML: &str = include_str!("../static/index.html"); const FAVICON: &[u8] = include_bytes!("../static/favicon.ico"); -pub async fn serve( - shutdown: CancellationToken -) -> Result<(), Box> { +pub async fn serve(shutdown: CancellationToken) -> Result<(), Box> { let config_logging = ConfigLogging::StderrTerminal { level: ConfigLoggingLevel::Info, }; - let log = config_logging - .to_logger("example-basic")?; + let log = config_logging.to_logger("example-basic")?; let mut api = ApiDescription::new(); api.register(index).unwrap(); @@ -58,7 +54,10 @@ pub async fn serve( .json()?, ); - let ctx = Context { spec, client: client().into() }; + let ctx = Context { + spec, + client: client().into(), + }; let server = ServerBuilder::new(api, ctx, log) .config(ConfigDropshot { @@ -153,7 +152,6 @@ where // TODO: cors for HttpError - /// Serve index page as html #[endpoint { method = GET, -- 2.51.2