diff --git a/consumer/src/backfill/db.rs b/consumer/src/backfill/db.rs index 71d444f0..3e299ef5 100644 --- a/consumer/src/backfill/db.rs +++ b/consumer/src/backfill/db.rs @@ -44,7 +44,10 @@ pub async fn update_actor_status( sync_state: types::ActorSyncState, ) -> QueryResult { diesel::update(schema::actors::table) - .set((schema::actors::status.eq(status), schema::actors::sync_state.eq(sync_state))) + .set(( + schema::actors::status.eq(status), + schema::actors::sync_state.eq(sync_state), + )) .filter(schema::actors::did.eq(did)) .execute(conn) .await diff --git a/consumer/src/backfill/mod.rs b/consumer/src/backfill/mod.rs index e7cf523e..0dea750f 100644 --- a/consumer/src/backfill/mod.rs +++ b/consumer/src/backfill/mod.rs @@ -6,11 +6,11 @@ use diesel_async::pooled_connection::deadpool::Pool; use diesel_async::{AsyncConnection, AsyncPgConnection}; use flume::{unbounded, Receiver, Sender}; use ipld_core::cid::Cid; +use metrics::counter; use parakeet_db::types::{ActorStatus, ActorSyncState}; use reqwest::{Client, StatusCode}; use std::str::FromStr; use std::sync::Arc; -use metrics::counter; use tracing::{instrument, Instrument}; mod db; diff --git a/consumer/src/backfill/types.rs b/consumer/src/backfill/types.rs index 51885a5e..05515dc8 100644 --- a/consumer/src/backfill/types.rs +++ b/consumer/src/backfill/types.rs @@ -1,7 +1,7 @@ +use crate::indexer::types::RecordTypes; use ipld_core::cid::Cid; use serde::Deserialize; use serde_bytes::ByteBuf; -use crate::indexer::types::RecordTypes; #[derive(Debug, Deserialize)] #[serde(untagged)] @@ -41,4 +41,4 @@ pub struct GetRepoStatusRes { pub active: bool, pub status: Option, pub rev: Option, -} \ No newline at end of file +} diff --git a/consumer/src/config.rs b/consumer/src/config.rs index 3c618e48..0f128daf 100644 --- a/consumer/src/config.rs +++ b/consumer/src/config.rs @@ -37,6 +37,10 @@ pub enum HistoryMode { Realtime, } -fn default_backfill_workers() -> u8 { 4 } +fn default_backfill_workers() -> u8 { + 4 +} -fn default_indexer_workers() -> u8 { 4 } +fn default_indexer_workers() -> u8 { + 4 +} diff --git a/consumer/src/firehose/mod.rs b/consumer/src/firehose/mod.rs index c365d0f5..28fe1443 100644 --- a/consumer/src/firehose/mod.rs +++ b/consumer/src/firehose/mod.rs @@ -21,9 +21,7 @@ impl FirehoseConsumer { let cursor = start_cursor.unwrap_or(0); let mut request = format!("{url}/xrpc/com.atproto.sync.subscribeRepos?cursor={cursor}") .into_client_request()?; - request - .headers_mut() - .insert(USER_AGENT, ua.parse()?); + request.headers_mut().insert(USER_AGENT, ua.parse()?); let (wss, _) = tokio_tungstenite::connect_async(request).await?; let (_, stream) = wss.split(); @@ -38,9 +36,7 @@ impl FirehoseConsumer { let cursor = start_cursor.unwrap_or(0); let mut request = format!("{url}/xrpc/com.atproto.label.subscribeLabels?cursor={cursor}") .into_client_request()?; - request - .headers_mut() - .insert(USER_AGENT, ua.parse()?); + request.headers_mut().insert(USER_AGENT, ua.parse()?); let (wss, _) = tokio_tungstenite::connect_async(request).await?; let (_, stream) = wss.split(); diff --git a/consumer/src/indexer/db.rs b/consumer/src/indexer/db.rs index eca5ea2f..32bd8d87 100644 --- a/consumer/src/indexer/db.rs +++ b/consumer/src/indexer/db.rs @@ -636,7 +636,10 @@ pub async fn insert_like( created_at: rec.created_at.naive_utc(), }; - diesel::insert_into(schema::likes::table).values(&data).execute(conn).await + diesel::insert_into(schema::likes::table) + .values(&data) + .execute(conn) + .await } pub async fn delete_like(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { @@ -660,7 +663,10 @@ pub async fn insert_repost( created_at: rec.created_at.naive_utc(), }; - diesel::insert_into(schema::reposts::table).values(&data).execute(conn).await + diesel::insert_into(schema::reposts::table) + .values(&data) + .execute(conn) + .await } pub async fn delete_repost(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { @@ -900,11 +906,13 @@ pub async fn upsert_verification( .on_conflict(schema::verification::at_uri) .do_update() .set(&data) - .execute(conn).await + .execute(conn) + .await } pub async fn delete_verification(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { diesel::delete(schema::verification::table) .filter(schema::verification::at_uri.eq(at_uri)) - .execute(conn).await + .execute(conn) + .await } diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs index 8babb0a3..993607bf 100644 --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -60,7 +60,10 @@ pub enum AppBskyEmbed { impl AppBskyEmbed { pub fn record_with_media_allowed(&self) -> bool { - matches!(self, AppBskyEmbed::Images(_) | AppBskyEmbed::Video(_) | AppBskyEmbed::External(_)) + matches!( + self, + AppBskyEmbed::Images(_) | AppBskyEmbed::Video(_) | AppBskyEmbed::External(_) + ) } pub fn as_str(&self) -> &'static str { @@ -196,7 +199,7 @@ pub struct AppBskyFeedPostgate { #[serde(tag = "$type")] pub enum PostgateEmbeddingRules { #[serde(rename = "app.bsky.feed.postgate#disableRule")] - Disable + Disable, } impl PostgateEmbeddingRules { @@ -237,7 +240,7 @@ pub enum ThreadgateRule { #[serde(rename = "app.bsky.feed.threadgate#followingRule")] Following, #[serde(rename = "app.bsky.feed.threadgate#listRule")] - List { list: String } + List { list: String }, } impl ThreadgateRule { diff --git a/consumer/src/indexer/types.rs b/consumer/src/indexer/types.rs index 4a98723f..835175b8 100644 --- a/consumer/src/indexer/types.rs +++ b/consumer/src/indexer/types.rs @@ -1,5 +1,5 @@ -use ipld_core::cid::Cid; use super::records; +use ipld_core::cid::Cid; use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] @@ -120,4 +120,4 @@ pub enum BackfillItemInner { Create(RecordTypes), Update(RecordTypes), Delete, -} \ No newline at end of file +} diff --git a/consumer/src/main.rs b/consumer/src/main.rs index f3da8a5d..f43902e3 100644 --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -31,7 +31,12 @@ async fn main() -> eyre::Result<()> { ..Default::default() })?); - let (label_mgr, label_svc_tx) = label_indexer::LabelServiceManager::new(&conf.database_url, resolver.clone(), user_agent.clone()).await?; + let (label_mgr, label_svc_tx) = label_indexer::LabelServiceManager::new( + &conf.database_url, + resolver.clone(), + user_agent.clone(), + ) + .await?; let relay_firehose = firehose::FirehoseConsumer::new_relay(&conf.relay_source, None, &user_agent).await?; let (backfiller, backfill_tx) = @@ -52,7 +57,10 @@ async fn main() -> eyre::Result<()> { tokio::spawn(label_mgr.run(conf.initial_label_services)), }?; - firehose_res.and(indexer_res).and(backfill_res).and(label_res) + firehose_res + .and(indexer_res) + .and(backfill_res) + .and(label_res) } async fn relay_consumer( diff --git a/consumer/src/utils.rs b/consumer/src/utils.rs index fc86a1d6..90f07e3a 100644 --- a/consumer/src/utils.rs +++ b/consumer/src/utils.rs @@ -43,8 +43,8 @@ pub fn strongref_to_parts( } pub fn empty_str_as_none(input: String) -> Option { - match input.is_empty() { + match input.is_empty() { true => None, - false => Some(input) + false => Some(input), } } diff --git a/lexica/src/app_bsky/actor.rs b/lexica/src/app_bsky/actor.rs index 798c9e10..0116d1d2 100644 --- a/lexica/src/app_bsky/actor.rs +++ b/lexica/src/app_bsky/actor.rs @@ -1,6 +1,7 @@ use crate::com_atproto::label::Label; use chrono::prelude::*; use serde::{Deserialize, Serialize}; +use std::fmt::Display; use std::str::FromStr; #[derive(Clone, Default, Debug, Serialize)] @@ -28,12 +29,12 @@ pub enum ChatAllowIncoming { Following, } -impl ToString for ChatAllowIncoming { - fn to_string(&self) -> String { +impl Display for ChatAllowIncoming { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { match self { - ChatAllowIncoming::All => "all".into(), - ChatAllowIncoming::None => "none".into(), - ChatAllowIncoming::Following => "following".into(), + ChatAllowIncoming::All => write!(f, "all"), + ChatAllowIncoming::None => write!(f, "none"), + ChatAllowIncoming::Following => write!(f, "following"), } } } @@ -46,7 +47,7 @@ impl FromStr for ChatAllowIncoming { "all" => Ok(ChatAllowIncoming::All), "none" => Ok(ChatAllowIncoming::None), "following" => Ok(ChatAllowIncoming::Following), - x => Err(format!("Unrecognized variant {}", x).into()), + x => Err(format!("Unrecognized variant {}", x)), } } } diff --git a/lexica/src/app_bsky/mod.rs b/lexica/src/app_bsky/mod.rs index 32671ba7..1ae1b370 100644 --- a/lexica/src/app_bsky/mod.rs +++ b/lexica/src/app_bsky/mod.rs @@ -14,4 +14,4 @@ pub struct RecordStats { pub repost_count: i64, pub like_count: i64, pub quote_count: i64, -} \ No newline at end of file +} diff --git a/parakeet/src/hydration/profile.rs b/parakeet/src/hydration/profile.rs index 7fa04ab3..e2885361 100644 --- a/parakeet/src/hydration/profile.rs +++ b/parakeet/src/hydration/profile.rs @@ -68,7 +68,7 @@ fn build_verification( let (verifications, verified_status) = match verification { Some(verif) => { let verifications = - get_verifications(&accept_verifiers, verif, handle, &profile.display_name); + get_verifications(accept_verifiers, verif, handle, &profile.display_name); let status = verifications .iter() diff --git a/parakeet/src/main.rs b/parakeet/src/main.rs index 043c04ff..666cd307 100644 --- a/parakeet/src/main.rs +++ b/parakeet/src/main.rs @@ -27,15 +27,21 @@ async fn main() -> eyre::Result<()> { let dataloaders = Arc::new(loaders::Dataloaders::new(pool.clone())); + #[allow(unused)] hydration::TRUSTED_VERIFIERS.set(conf.trusted_verifiers); - let cors = CorsLayer::new().allow_headers(AllowHeaders::any()).allow_origin(AllowOrigin::any()); + let cors = CorsLayer::new() + .allow_headers(AllowHeaders::any()) + .allow_origin(AllowOrigin::any()); let did_doc = did_web_doc(&conf.service); let app = axum::Router::new() .nest("/xrpc", xrpc::xrpc_routes()) - .route("/.well-known/did.json", axum::routing::get(async || axum::Json(did_doc))) + .route( + "/.well-known/did.json", + axum::routing::get(async || axum::Json(did_doc)), + ) .layer(TraceLayer::new_for_http()) .layer(cors) .with_state(GlobalState { pool, dataloaders }); diff --git a/parakeet/src/xrpc/com_atproto/identity.rs b/parakeet/src/xrpc/com_atproto/identity.rs index d61c4522..89f1c62d 100644 --- a/parakeet/src/xrpc/com_atproto/identity.rs +++ b/parakeet/src/xrpc/com_atproto/identity.rs @@ -1,8 +1,8 @@ +use crate::xrpc::error::XrpcResult; +use crate::xrpc::get_actor_did; +use crate::GlobalState; use axum::extract::{Query, State}; use axum::Json; -use crate::GlobalState; -use crate::xrpc::error::XrpcResult; -use crate::xrpc::{get_actor_did}; use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize)] @@ -15,8 +15,11 @@ pub struct ResolveHandleRes { pub did: String, } -pub async fn resolve_handle(State(state): State, Query(query): Query) -> XrpcResult> { +pub async fn resolve_handle( + State(state): State, + Query(query): Query, +) -> XrpcResult> { let did = get_actor_did(&state.dataloaders, query.handle).await?; Ok(Json(ResolveHandleRes { did })) -} \ No newline at end of file +} diff --git a/parakeet/src/xrpc/com_atproto/mod.rs b/parakeet/src/xrpc/com_atproto/mod.rs index bd3032b7..2819af25 100644 --- a/parakeet/src/xrpc/com_atproto/mod.rs +++ b/parakeet/src/xrpc/com_atproto/mod.rs @@ -1,9 +1,9 @@ -use axum::Router; use axum::routing::get; +use axum::Router; mod identity; pub fn routes() -> Router { Router::new() .route("/com.atproto.identity.resolveHandle", get(identity::resolve_handle)) -} \ No newline at end of file +}