From 09482a61a7308e2bd84a83e3b8e557321410b307 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sun, 21 Jun 2026 22:14:10 +0300 Subject: [PATCH] [all] fmt --- src/backfill/worker.rs | 23 +++++----- src/backfill/worker/process.rs | 4 +- src/backfill/worker/task.rs | 12 ++--- src/config/types.rs | 6 +-- src/control/db.rs | 8 ++-- src/control/hosts.rs | 6 +-- src/control/hydrant.rs | 10 ++--- src/control/hydrant/run.rs | 14 +++--- src/control/mod.rs | 20 ++++----- src/control/stats.rs | 4 +- src/control/stream.rs | 4 +- src/control/stream/engine.rs | 4 +- src/control/stream/indexer.rs | 2 +- src/control/stream/jetstream.rs | 13 +++--- src/control/stream/relay.rs | 2 +- src/crawler/list_repos/checker.rs | 14 +++--- src/crawler/list_repos/producer.rs | 12 ++--- src/crawler/list_repos/retry.rs | 12 ++--- src/db/counts.rs | 8 +++- src/db/mod.rs | 6 +-- src/db/open.rs | 26 +++++------ src/db/train.rs | 2 +- src/ingest/firehose_stats.rs | 12 ++--- src/ingest/firehose_stats/relay.rs | 6 +-- src/ingest/firehose_stats/source.rs | 9 ++-- src/ingest/indexer.rs | 4 +- src/ingest/indexer/message.rs | 4 +- src/ingest/indexer/shard.rs | 22 ++++----- src/ingest/indexer/worker.rs | 10 ++--- src/ingest/relay.rs | 12 ++--- src/ingest/relay/context.rs | 69 +++++++++++++++++------------ src/ingest/relay/handlers.rs | 20 ++++----- src/ingest/relay/worker.rs | 46 ++++++++++++------- src/ingest/stream.rs | 14 +++--- src/ingest/stream/codec.rs | 4 +- src/types.rs | 8 ++-- src/types/event.rs | 6 +-- src/types/v2.rs | 6 +-- src/types/v4.rs | 8 ++-- src/types/v7.rs | 8 ++-- 40 files changed, 255 insertions(+), 225 deletions(-) diff --git a/src/backfill/worker.rs b/src/backfill/worker.rs index 93d71f7..b8d6105 100644 --- a/src/backfill/worker.rs +++ b/src/backfill/worker.rs @@ -1,19 +1,19 @@ -use std::sync::Arc; -use std::time::Duration; -use tokio::sync::Semaphore; -use tracing::{debug, error, info, Instrument}; -use jacquard_common::types::did::Did; -use jacquard_common::IntoStatic; use crate::backfill::client::ThrottledHttpClient; use crate::backfill::error::BackfillError; use crate::config::BackfillStrategy; use crate::db::{self, types::TrimmedDid}; -use crate::state::AppState; use crate::ingest::indexer::IndexerTx; +use crate::state::AppState; use crate::util::WatchEnabledExt; +use jacquard_common::IntoStatic; +use jacquard_common::types::did::Did; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::Semaphore; +use tracing::{Instrument, debug, error, info}; -pub mod task; pub mod process; +pub mod task; pub struct BackfillWorker { state: Arc, @@ -134,9 +134,10 @@ impl BackfillWorker { tokio::spawn( async move { let _guard = guard; - let res = - task::did_task(&state, http, buffer_tx, &did, key, permit, verify, strategy) - .await; + let res = task::did_task( + &state, http, buffer_tx, &did, key, permit, verify, strategy, + ) + .await; if let Err(e) = res { match &e { diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index bfafa5e..8e70012 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -4,8 +4,8 @@ use std::time::Instant; use fjall::Slice; use miette::{IntoDiagnostic, Result}; -use tracing::{debug, error, trace, warn}; use smol_str::{SmolStr, ToSmolStr}; +use tracing::{debug, error, trace, warn}; use jacquard_api::com_atproto::sync::get_repo::{GetRepo, GetRepoError}; use jacquard_common::IntoStatic; @@ -287,7 +287,7 @@ pub(crate) async fn process_did( state.status = status; batch.insert(&db.resync, key, resync_bytes); Ok((true, ())) - }, + }, )?; } let lifecycle_reservation = lifecycle_counts.stage(&mut batch); diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 53475e1..1bc9a12 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -1,17 +1,17 @@ -use std::sync::Arc; use fjall::Slice; -use miette::{Result, IntoDiagnostic}; -use tracing::{debug, error, warn}; -use jacquard_common::types::did::Did; use jacquard_common::IntoStatic; +use jacquard_common::types::did::Did; +use miette::{IntoDiagnostic, Result}; +use std::sync::Arc; +use tracing::{debug, error, warn}; use crate::backfill::client::ThrottledHttpClient; use crate::backfill::error::BackfillError; use crate::config::BackfillStrategy; use crate::db::{Db, keys}; -use crate::types::{RepoState, RepoStatus, ResyncErrorKind, ResyncState, GaugeState}; -use crate::state::AppState; use crate::ingest::indexer::{IndexerMessage, IndexerTx}; +use crate::state::AppState; +use crate::types::{GaugeState, RepoState, RepoStatus, ResyncErrorKind, ResyncState}; use super::process::process_did; diff --git a/src/config/types.rs b/src/config/types.rs index ced1579..837c003 100644 --- a/src/config/types.rs +++ b/src/config/types.rs @@ -1,9 +1,9 @@ +use miette::Result; +use serde::{Deserialize, Serialize}; +use smol_str::ToSmolStr; use std::fmt; use std::str::FromStr; use url::Url; -use serde::{Deserialize, Serialize}; -use smol_str::ToSmolStr; -use miette::Result; /// rate limit parameters for a named tier of PDS connections. /// diff --git a/src/control/db.rs b/src/control/db.rs index 4e1c46e..5bef111 100644 --- a/src/control/db.rs +++ b/src/control/db.rs @@ -1,8 +1,8 @@ -use std::sync::Arc; -use miette::{IntoDiagnostic, Result}; -use futures::FutureExt; -use crate::state::AppState; use super::Hydrant; +use crate::state::AppState; +use futures::FutureExt; +use miette::{IntoDiagnostic, Result}; +use std::sync::Arc; /// control over database maintenance operations. /// diff --git a/src/control/hosts.rs b/src/control/hosts.rs index d09da5b..18822e2 100644 --- a/src/control/hosts.rs +++ b/src/control/hosts.rs @@ -1,11 +1,11 @@ -use std::collections::BTreeSet; -use smol_str::{SmolStr, ToSmolStr}; use miette::{IntoDiagnostic, Result, WrapErr}; +use smol_str::{SmolStr, ToSmolStr}; +use std::collections::BTreeSet; use url::Url; +use super::Hydrant; use crate::db::keys; use crate::pds_meta::HostStatus; -use super::Hydrant; #[derive(Debug, Clone)] pub struct ApiBinds { diff --git a/src/control/hydrant.rs b/src/control/hydrant.rs index be5bdaa..488ee17 100644 --- a/src/control/hydrant.rs +++ b/src/control/hydrant.rs @@ -1,8 +1,8 @@ +use futures::FutureExt; +use miette::{IntoDiagnostic, Result}; use std::future::Future; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; -use futures::FutureExt; -use miette::{IntoDiagnostic, Result}; use tokio::sync::watch; use tracing::{debug, error, info, warn}; use url::Url; @@ -12,17 +12,17 @@ use crate::db::{self, filter as db_filter, load_persisted_firehose_sources}; use crate::filter::FilterMode; use crate::state::AppState; -use super::{FilterControl, pds::PdsControl, ReposControl, DbControl}; use super::firehose::{FirehoseHandle, FirehoseShared}; +use super::{DbControl, FilterControl, ReposControl, pds::PdsControl}; +#[cfg(feature = "indexer")] +use super::{BackfillHandle, crawler}; #[cfg(feature = "indexer")] use crate::backfill::BackfillWorker; #[cfg(feature = "indexer")] use crate::db::load_persisted_crawler_sources; #[cfg(feature = "indexer")] use crate::ingest::indexer::FirehoseWorker; -#[cfg(feature = "indexer")] -use super::{crawler, BackfillHandle}; #[cfg(feature = "backlinks")] use crate::backlinks::BacklinksControl; diff --git a/src/control/hydrant/run.rs b/src/control/hydrant/run.rs index 6412fdd..5bc2479 100644 --- a/src/control/hydrant/run.rs +++ b/src/control/hydrant/run.rs @@ -1,27 +1,27 @@ +use futures::FutureExt; +use miette::{IntoDiagnostic, Result}; use std::future::Future; use std::sync::Arc; use std::sync::atomic::Ordering; -use futures::FutureExt; -use miette::{IntoDiagnostic, Result}; use tokio::sync::watch; use tracing::{debug, error, info, warn}; use url::Url; +use super::super::firehose::FirehoseShared; +use super::super::seed; +use super::Hydrant; use crate::config::SignatureVerification; use crate::db::{self, load_persisted_firehose_sources}; use crate::state::AppState; -use super::super::seed; -use super::super::firehose::FirehoseShared; -use super::Hydrant; +#[cfg(feature = "indexer")] +use super::super::crawler; #[cfg(feature = "indexer")] use crate::backfill::BackfillWorker; #[cfg(feature = "indexer")] use crate::db::load_persisted_crawler_sources; #[cfg(feature = "indexer")] use crate::ingest::indexer::FirehoseWorker; -#[cfg(feature = "indexer")] -use super::super::crawler; impl Hydrant { /// start all background components and return a future that resolves when any diff --git a/src/control/mod.rs b/src/control/mod.rs index f6fdda4..66f41fc 100644 --- a/src/control/mod.rs +++ b/src/control/mod.rs @@ -9,9 +9,9 @@ pub(crate) mod repos; mod seed; pub(crate) mod stream; -pub mod hydrant; -pub mod hosts; pub mod db; +pub mod hosts; +pub mod hydrant; pub mod stats; #[cfg(feature = "indexer")] @@ -35,9 +35,9 @@ pub use firehose::{FirehoseHandle, FirehoseSourceInfo}; pub use pds::{PdsControl, PdsTierAssignment, PdsTierDefinition}; pub use repos::{ListedRecord, Record, RecordList, RepoHandle, RepoInfo, ReposControl}; -pub use hydrant::Hydrant; -pub use hosts::{ApiBinds, Host}; pub use db::DbControl; +pub use hosts::{ApiBinds, Host}; +pub use hydrant::Hydrant; pub use stats::StatsResponse; #[cfg(feature = "indexer_stream")] @@ -55,16 +55,16 @@ use crate::types::MarshallableEvt; pub type Event = MarshallableEvt<'static>; // Crate-internal re-exports for submodules using `super::*` +pub(crate) use crate::state::AppState; +pub(crate) use futures::Stream; +pub(crate) use std::pin::Pin; pub(crate) use std::sync::Arc; pub(crate) use std::sync::atomic::{AtomicBool, Ordering}; -pub(crate) use std::pin::Pin; pub(crate) use std::task::{Context, Poll}; -pub(crate) use futures::Stream; -pub(crate) use tokio::sync::mpsc; -pub(crate) use crate::state::AppState; #[cfg(feature = "indexer_stream")] use stream::event_stream_thread; -#[cfg(feature = "relay")] -use stream::relay_stream_thread; #[cfg(feature = "jetstream")] use stream::jetstream_stream_thread; +#[cfg(feature = "relay")] +use stream::relay_stream_thread; +pub(crate) use tokio::sync::mpsc; diff --git a/src/control/stats.rs b/src/control/stats.rs index 7a491df..4926873 100644 --- a/src/control/stats.rs +++ b/src/control/stats.rs @@ -1,6 +1,6 @@ -use std::collections::BTreeMap; -use miette::{IntoDiagnostic, Result}; use super::Hydrant; +use miette::{IntoDiagnostic, Result}; +use std::collections::BTreeMap; /// database statistics returned by [`Hydrant::stats`]. #[derive(serde::Serialize)] diff --git a/src/control/stream.rs b/src/control/stream.rs index 622423f..51c56ff 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -1,5 +1,5 @@ -pub(crate) mod types; pub(crate) mod engine; +pub(crate) mod types; #[cfg(feature = "indexer_stream")] pub(crate) mod indexer; @@ -20,5 +20,5 @@ pub(crate) use jetstream::{ JetstreamAccount, JetstreamCommit, JetstreamEvent, JetstreamIdentity, JetstreamPayload, }; -pub(crate) use types::*; pub(crate) use engine::*; +pub(crate) use types::*; diff --git a/src/control/stream/engine.rs b/src/control/stream/engine.rs index 16e6824..994844b 100644 --- a/src/control/stream/engine.rs +++ b/src/control/stream/engine.rs @@ -8,8 +8,8 @@ use tokio::sync::{broadcast, mpsc}; use tracing::warn; use super::types::{ - PendingLiveEvents, ReplayChunk, SendOutcome, StreamBroadcast, StreamOptions, StreamTooSlow, - STREAM_SEND_RETRY_PAUSE, + PendingLiveEvents, ReplayChunk, STREAM_SEND_RETRY_PAUSE, SendOutcome, StreamBroadcast, + StreamOptions, StreamTooSlow, }; pub(crate) fn run_ordered_stream( diff --git a/src/control/stream/indexer.rs b/src/control/stream/indexer.rs index 432e88c..35c3d16 100644 --- a/src/control/stream/indexer.rs +++ b/src/control/stream/indexer.rs @@ -13,8 +13,8 @@ use jacquard_common::{CowStr, IntoStatic, RawData}; use jacquard_repo::DAG_CBOR_CID_CODEC; use sha2::{Digest, Sha256}; -use crate::control::{Event, StreamError}; use super::{ReplayChunk, StreamBroadcast, StreamOptions, run_ordered_stream, stream_seq_after}; +use crate::control::{Event, StreamError}; pub(crate) fn event_stream_thread( state: Arc, diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index 020d3fb..553d12f 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -1,24 +1,23 @@ +use bytes::Bytes; use std::sync::Arc; use std::time::{Duration, Instant}; use tokio::sync::{broadcast, mpsc}; use tracing::{error, warn}; -use bytes::Bytes; use crate::db::keys; use crate::state::AppState; use crate::types::{JetstreamBroadcast, StoredJetstreamEvent}; -#[cfg(feature = "indexer_stream")] -use crate::types::StoredEvent; #[cfg(feature = "relay")] use crate::ingest::stream::{SubscribeReposMessage, decode_frame}; +#[cfg(feature = "indexer_stream")] +use crate::types::StoredEvent; -use crate::control::{JetstreamStreamError, JetstreamFilter}; use super::{ - StreamOptions, ReplayChunk, StreamTooSlow, - STREAM_SEND_RETRY_PAUSE, replay_chunk_size_for, note_replay_blocked, clear_replay_blocked, - send_stream_error, + ReplayChunk, STREAM_SEND_RETRY_PAUSE, StreamOptions, StreamTooSlow, clear_replay_blocked, + note_replay_blocked, replay_chunk_size_for, send_stream_error, }; +use crate::control::{JetstreamFilter, JetstreamStreamError}; #[cfg(feature = "indexer_stream")] use super::indexer::stored_to_event; diff --git a/src/control/stream/relay.rs b/src/control/stream/relay.rs index c739cd0..ffce9c2 100644 --- a/src/control/stream/relay.rs +++ b/src/control/stream/relay.rs @@ -7,8 +7,8 @@ use crate::db::keys; use crate::state::AppState; use crate::types::RelayBroadcast; -use crate::control::RelayStreamError; use super::{ReplayChunk, StreamOptions, run_ordered_stream, stream_seq_after}; +use crate::control::RelayStreamError; pub(crate) fn relay_stream_thread( state: Arc, diff --git a/src/crawler/list_repos/checker.rs b/src/crawler/list_repos/checker.rs index 8ba2a5d..31f9832 100644 --- a/src/crawler/list_repos/checker.rs +++ b/src/crawler/list_repos/checker.rs @@ -1,7 +1,3 @@ -use std::collections::{HashMap, HashSet}; -use std::ops::{Add, Sub}; -use std::sync::Arc; -use std::time::Duration; use chrono::{DateTime, TimeDelta, Utc}; use fjall::OwnedWriteBatch; use futures::{Future, FutureExt}; @@ -11,17 +7,21 @@ use miette::{IntoDiagnostic, Result}; use reqwest::StatusCode; use serde::{Deserialize, Serialize}; use smol_str::{SmolStr, ToSmolStr}; +use std::collections::{HashMap, HashSet}; +use std::ops::{Add, Sub}; +use std::sync::Arc; +use std::time::Duration; use tracing::{Instrument, error, info_span, trace, warn}; use url::Url; +use super::super::InFlightGuard; use crate::db::{Db, keys}; use crate::state::AppState; use crate::util::throttle::{OrFailure, ThrottleHandle, Throttler}; use crate::util::{ - ErrorForStatus, RetryOutcome, RetryWithBackoff, is_io_error_their_fault, - is_status_their_fault, is_tls_cert_error, parse_retry_after, + ErrorForStatus, RetryOutcome, RetryWithBackoff, is_io_error_their_fault, is_status_their_fault, + is_tls_cert_error, parse_retry_after, }; -use super::super::InFlightGuard; pub(super) const MAX_RETRY_ATTEMPTS: u32 = 5; diff --git a/src/crawler/list_repos/producer.rs b/src/crawler/list_repos/producer.rs index 52bc7d7..c53c8ec 100644 --- a/src/crawler/list_repos/producer.rs +++ b/src/crawler/list_repos/producer.rs @@ -1,23 +1,23 @@ -use std::collections::HashMap; -use std::sync::Arc; -use std::time::Duration; use jacquard_api::com_atproto::sync::list_repos::ListReposOutput; use jacquard_common::{IntoStatic, types::string::Did}; use miette::{Context, IntoDiagnostic, Result}; use reqwest::StatusCode; use smol_str::{SmolStr, ToSmolStr}; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; use tokio::sync::{mpsc, watch}; use tracing::{debug, error, info, warn}; use url::Url; +use super::super::worker::{CrawlerBatch, CursorUpdate}; +use super::super::{CrawlerStats, InFlight, base_url}; +use super::checker::SignalChecker; use crate::db::keys::crawler_cursor_key; use crate::db::{Db, keys}; use crate::state::AppState; use crate::util::WatchEnabledExt; use crate::util::{ErrorForStatus, RetryOutcome, RetryWithBackoff}; -use super::super::worker::{CrawlerBatch, CursorUpdate}; -use super::super::{CrawlerStats, InFlight, base_url}; -use super::checker::SignalChecker; const BLOCKING_TASK_TIMEOUT: Duration = Duration::from_secs(30); diff --git a/src/crawler/list_repos/retry.rs b/src/crawler/list_repos/retry.rs index b9125da..2346e76 100644 --- a/src/crawler/list_repos/retry.rs +++ b/src/crawler/list_repos/retry.rs @@ -1,18 +1,18 @@ -use std::collections::HashMap; -use std::ops::Mul; -use std::time::Duration; use chrono::TimeDelta; use jacquard_common::types::string::Did; use miette::{IntoDiagnostic, Result}; use rand::RngExt; use rand::rngs::SmallRng; +use std::collections::HashMap; +use std::ops::Mul; +use std::time::Duration; use tokio::sync::mpsc; use tracing::{debug, error, info}; -use crate::db::keys; -use super::super::worker::CrawlerBatch; use super::super::InFlight; -use super::checker::{RetryState, SignalChecker, MAX_RETRY_ATTEMPTS}; +use super::super::worker::CrawlerBatch; +use super::checker::{MAX_RETRY_ATTEMPTS, RetryState, SignalChecker}; +use crate::db::keys; const MAX_RETRY_BATCH: usize = 1000; diff --git a/src/db/counts.rs b/src/db/counts.rs index 1fdd374..ca6c135 100644 --- a/src/db/counts.rs +++ b/src/db/counts.rs @@ -4,8 +4,8 @@ use miette::{Context, IntoDiagnostic, Result}; use smol_str::SmolStr; use std::collections::{BTreeMap, BTreeSet}; use std::sync::Arc; -use std::sync::atomic::Ordering; use std::sync::Mutex; +use std::sync::atomic::Ordering; use tracing::error; #[derive(Debug, Clone, Default)] @@ -49,7 +49,11 @@ impl CountDeltas { } #[cfg(feature = "indexer")] - pub(crate) fn add_gauge_diff(&mut self, old: &crate::types::GaugeState, new: &crate::types::GaugeState) { + pub(crate) fn add_gauge_diff( + &mut self, + old: &crate::types::GaugeState, + new: &crate::types::GaugeState, + ) { use crate::types::GaugeState; if old == new { diff --git a/src/db/mod.rs b/src/db/mod.rs index 5120272..3903b05 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -21,9 +21,9 @@ use url::Url; pub mod compaction; pub mod counts; -pub use counts::{CountDeltas, load_count_delta_watermark, set_ks_count}; #[cfg(feature = "indexer")] pub(crate) use counts::CountDeltaReservation; +pub use counts::{CountDeltas, load_count_delta_watermark, set_ks_count}; pub mod ephemeral; pub mod filter; pub mod keys; @@ -251,8 +251,6 @@ pub fn check_poisoned_report(e: &miette::Report) { self::check_poisoned(err); } - - /// load the persisted (day, count) pair for the daily PDS add counter, if present. /// returns `None` if no entry exists or the stored data is malformed. #[cfg(feature = "relay")] @@ -303,5 +301,3 @@ pub fn load_persisted_firehose_sources( } Ok(sources) } - - diff --git a/src/db/open.rs b/src/db/open.rs index 9ee5f2a..e7b7807 100644 --- a/src/db/open.rs +++ b/src/db/open.rs @@ -1,24 +1,24 @@ -use std::sync::{Arc, Mutex}; -use std::sync::atomic::{AtomicU64, Ordering}; -#[cfg(feature = "jetstream")] -use std::sync::atomic::AtomicI64; -use std::collections::BTreeSet; -use tokio::sync::broadcast; -use scc::HashMap; -use smol_str::SmolStr; -use miette::{Result, IntoDiagnostic, Context}; use fjall::{ - config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}, CompressionType, Database, KeyspaceCreateOptions, + config::{BlockSizePolicy, CompressionPolicy, RestartIntervalPolicy}, }; use lsm_tree::compaction::Factory; +use miette::{Context, IntoDiagnostic, Result}; +use scc::HashMap; +use smol_str::SmolStr; +use std::collections::BTreeSet; +#[cfg(feature = "jetstream")] +use std::sync::atomic::AtomicI64; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex}; +use tokio::sync::broadcast; -use crate::config::{Config, Compression}; +use crate::config::{Compression, Config}; -use super::{Db, migration}; use super::compaction::CountsGcFilterFactory; -use super::counts::{load_count_delta_watermark, replay_count_deltas, read_u64_counter}; +use super::counts::{load_count_delta_watermark, read_u64_counter, replay_count_deltas}; use super::keys; +use super::{Db, migration}; const fn kb(v: u32) -> u32 { v * 1024 diff --git a/src/db/train.rs b/src/db/train.rs index 1dcdb87..059986f 100644 --- a/src/db/train.rs +++ b/src/db/train.rs @@ -1,6 +1,6 @@ +use miette::{IntoDiagnostic, Result}; use std::cell::RefCell; use std::collections::HashSet; -use miette::{Result, IntoDiagnostic}; use super::Db; #[cfg(feature = "indexer")] diff --git a/src/ingest/firehose_stats.rs b/src/ingest/firehose_stats.rs index 5d5287d..d7cc208 100644 --- a/src/ingest/firehose_stats.rs +++ b/src/ingest/firehose_stats.rs @@ -1,15 +1,15 @@ -mod source; mod relay; +mod source; -pub use source::{FirehoseStats, FirehoseSourceStats, FirehoseStatsSnapshot, FirehoseMessageStats}; pub use relay::{ - RelayWorkerStats, RelayShardStats, RelayShardStatsSnapshot, RelayShardTimings, - RelayMessageKind, RepoStateLoadOutcome, HostAuthorityStatsOutcome, ValidationStatsOutcome, - RelayWorkerStatsSnapshot, + HostAuthorityStatsOutcome, RelayMessageKind, RelayShardStats, RelayShardStatsSnapshot, + RelayShardTimings, RelayWorkerStats, RelayWorkerStatsSnapshot, RepoStateLoadOutcome, + ValidationStatsOutcome, }; +pub use source::{FirehoseMessageStats, FirehoseSourceStats, FirehoseStats, FirehoseStatsSnapshot}; +use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::time::Duration; -use std::sync::atomic::{AtomicU64, AtomicI64, Ordering}; fn now_ts() -> i64 { chrono::Utc::now().timestamp() diff --git a/src/ingest/firehose_stats/relay.rs b/src/ingest/firehose_stats/relay.rs index 51d7c50..008f6b4 100644 --- a/src/ingest/firehose_stats/relay.rs +++ b/src/ingest/firehose_stats/relay.rs @@ -1,11 +1,11 @@ +use parking_lot::Mutex; +use serde::Serialize; use std::collections::BTreeMap; use std::sync::Arc; use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::time::Duration; -use parking_lot::Mutex; -use serde::Serialize; -use super::{now_ts, add_duration, add_duration_with_max, nonzero_i64}; +use super::{add_duration, add_duration_with_max, nonzero_i64, now_ts}; #[derive(Clone, Copy, Debug)] pub enum RelayMessageKind { diff --git a/src/ingest/firehose_stats/source.rs b/src/ingest/firehose_stats/source.rs index 22eaa80..09993a2 100644 --- a/src/ingest/firehose_stats/source.rs +++ b/src/ingest/firehose_stats/source.rs @@ -1,11 +1,14 @@ +use parking_lot::Mutex; +use serde::Serialize; use std::sync::Arc; use std::sync::atomic::{AtomicI64, AtomicU64, Ordering}; use std::time::Duration; -use parking_lot::Mutex; -use serde::Serialize; use url::Url; -use super::{now_ts, duration_micros, nonzero_i64, RelayWorkerStats, RelayWorkerStatsSnapshot, RelayShardStats}; +use super::{ + RelayShardStats, RelayWorkerStats, RelayWorkerStatsSnapshot, duration_micros, nonzero_i64, + now_ts, +}; #[derive(Default)] pub struct FirehoseStats { diff --git a/src/ingest/indexer.rs b/src/ingest/indexer.rs index 1cb3da8..67736a3 100644 --- a/src/ingest/indexer.rs +++ b/src/ingest/indexer.rs @@ -1,9 +1,9 @@ pub mod message; -pub mod worker; pub mod shard; +pub mod worker; pub use message::{ - IndexerCommitData, IndexerIdentityData, IndexerAccountData, IndexerEventData, IndexerEvent, + IndexerAccountData, IndexerCommitData, IndexerEvent, IndexerEventData, IndexerIdentityData, IndexerMessage, IndexerTx, }; pub use worker::FirehoseWorker; diff --git a/src/ingest/indexer/message.rs b/src/ingest/indexer/message.rs index 9725ca2..7d04616 100644 --- a/src/ingest/indexer/message.rs +++ b/src/ingest/indexer/message.rs @@ -1,8 +1,8 @@ -use url::Url; +use crate::ingest::mailbox::{ShardedMessage, ShardedReceiver, ShardedSender}; use crate::ingest::stream; -use crate::ingest::mailbox::{ShardedMessage, ShardedSender, ShardedReceiver}; use jacquard_common::types::did::Did; use miette::Result; +use url::Url; #[derive(Debug)] pub struct IndexerCommitData { diff --git a/src/ingest/indexer/shard.rs b/src/ingest/indexer/shard.rs index 87cb6cf..aabdf26 100644 --- a/src/ingest/indexer/shard.rs +++ b/src/ingest/indexer/shard.rs @@ -1,20 +1,20 @@ +use fjall::OwnedWriteBatch; +use jacquard_common::IntoStatic; +use jacquard_common::types::did::Did; +use miette::{IntoDiagnostic, Result}; use std::sync::Arc; use std::sync::atomic::Ordering::SeqCst; use tokio::runtime::Handle as TokioHandle; use tracing::{debug, error, warn}; -use fjall::OwnedWriteBatch; -use jacquard_common::IntoStatic; -use jacquard_common::types::did::Did; -use miette::{Result, IntoDiagnostic}; use crate::db::{self, CountDeltas, keys, ser_repo_meta}; -use crate::state::AppState; -use crate::types::{GaugeState, RepoMetadata, RepoState, RepoStatus}; -use crate::ops; -use crate::ingest::stream::{Commit, Identity, Account}; +use crate::ingest::stream::types::AccountStatus; +use crate::ingest::stream::{Account, Commit, Identity}; use crate::ingest::validation; +use crate::ops; use crate::resolver::ResolverError; -use crate::ingest::stream::types::AccountStatus; +use crate::state::AppState; +use crate::types::{GaugeState, RepoMetadata, RepoState, RepoStatus}; #[cfg(feature = "jetstream")] use crate::{ @@ -28,8 +28,8 @@ use { }; use super::message::{ - IndexerRx, IndexerMessage, IndexerEvent, IndexerEventData, IndexerCommitData, - IndexerIdentityData, IndexerAccountData, + IndexerAccountData, IndexerCommitData, IndexerEvent, IndexerEventData, IndexerIdentityData, + IndexerMessage, IndexerRx, }; use super::worker::{FirehoseWorker, IngestError, RepoProcessResult}; diff --git a/src/ingest/indexer/worker.rs b/src/ingest/indexer/worker.rs index a5ddc4d..6e512ec 100644 --- a/src/ingest/indexer/worker.rs +++ b/src/ingest/indexer/worker.rs @@ -1,16 +1,16 @@ +use miette::{Diagnostic, IntoDiagnostic, Result}; use std::sync::Arc; -use miette::{Diagnostic, Result, IntoDiagnostic}; use thiserror::Error; use tokio::runtime::Handle as TokioHandle; use tracing::info; +use super::message::{IndexerRx, IndexerTx}; +use crate::ingest::mailbox::ShardedSender; +use crate::ingest::stream::Commit; use crate::resolver::{NoSigningKeyError, ResolverError}; -use jacquard_repo::error::CommitError; use crate::state::AppState; use crate::types::RepoState; -use crate::ingest::stream::Commit; -use crate::ingest::mailbox::ShardedSender; -use super::message::{IndexerTx, IndexerRx}; +use jacquard_repo::error::CommitError; #[derive(Debug, Diagnostic, Error)] pub(crate) enum IngestError { diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index 1c4c63b..55e3928 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -6,11 +6,9 @@ use jacquard_api::com_atproto::sync::get_repo_status::{ use smol_str::SmolStr; #[cfg(feature = "firehose-diagnostics")] -use crate::ingest::stream::SubscribeReposMessage; +use crate::ingest::firehose_stats::{RelayMessageKind, ValidationStatsOutcome}; #[cfg(feature = "firehose-diagnostics")] -use crate::ingest::firehose_stats::{ - RelayMessageKind, ValidationStatsOutcome, -}; +use crate::ingest::stream::SubscribeReposMessage; #[cfg(feature = "firehose-diagnostics")] use crate::ingest::validation::{ CommitValidationError, SyncValidationError, ValidatedCommit, ValidatedSync, @@ -22,13 +20,15 @@ pub mod handlers; pub mod worker; pub(crate) use context::WorkerContext; -pub(crate) use worker::WorkerMessage; pub use worker::RelayWorker; +pub(crate) use worker::WorkerMessage; pub(crate) const WRONG_HOST_AUTHORITY_RECHECK_INTERVAL: Duration = Duration::from_secs(60); pub(crate) const WRONG_HOST_AUTHORITY_CACHE_PRUNE_AT: usize = 8192; -pub(crate) fn map_repo_status_probe(output: Option>) -> Option> { +pub(crate) fn map_repo_status_probe( + output: Option>, +) -> Option> { let output = output?; let mut repo_state = RepoState::backfilling(); diff --git a/src/ingest/relay/context.rs b/src/ingest/relay/context.rs index c1b43d4..2dcfe8f 100644 --- a/src/ingest/relay/context.rs +++ b/src/ingest/relay/context.rs @@ -4,37 +4,35 @@ use std::sync::Arc; use std::time::Instant; use fjall::OwnedWriteBatch; -use jacquard_api::com_atproto::sync::get_repo_status::{ - GetRepoStatus, GetRepoStatusError, -}; +use jacquard_api::com_atproto::sync::get_repo_status::{GetRepoStatus, GetRepoStatusError}; +use jacquard_common::IntoStatic; use jacquard_common::types::crypto::PublicKey; use jacquard_common::types::did::Did; use jacquard_common::xrpc::{XrpcError, XrpcExt}; -use jacquard_common::IntoStatic; use miette::{IntoDiagnostic, Result}; +use smol_str::{SmolStr, ToSmolStr}; use tokio::runtime::Handle; use tracing::{debug, trace, warn}; use url::Url; -use smol_str::{SmolStr, ToSmolStr}; use crate::db::{self, CountDeltas, keys}; use crate::state::AppState; use crate::types::{RepoState, RepoStatus}; use crate::util; -#[cfg(feature = "relay")] -use crate::types::RelayBroadcast; -#[cfg(feature = "relay")] -use std::sync::atomic::Ordering; +use super::{ + AuthorityOutcome, RelayWorker, WRONG_HOST_AUTHORITY_CACHE_PRUNE_AT, + WRONG_HOST_AUTHORITY_RECHECK_INTERVAL, WorkerMessage, map_repo_status_probe, +}; use crate::ingest::stream::AccountStatus; use crate::ingest::stream::SubscribeReposMessage; use crate::ingest::validation::{ CommitValidationError, SyncValidationError, ValidatedCommit, ValidatedSync, ValidationContext, }; -use super::{ - AuthorityOutcome, RelayWorker, WorkerMessage, WRONG_HOST_AUTHORITY_RECHECK_INTERVAL, - WRONG_HOST_AUTHORITY_CACHE_PRUNE_AT, map_repo_status_probe, -}; +#[cfg(feature = "relay")] +use crate::types::RelayBroadcast; +#[cfg(feature = "relay")] +use std::sync::atomic::Ordering; #[cfg(feature = "firehose-diagnostics")] use super::{commit_validation_outcome, sync_validation_outcome}; @@ -339,7 +337,10 @@ impl WorkerContext<'_> { Ok(map_repo_status_probe(Some(output))) } - pub(crate) fn load_repo_state(&mut self, msg: &WorkerMessage) -> Result>> { + pub(crate) fn load_repo_state( + &mut self, + msg: &WorkerMessage, + ) -> Result>> { let db = &self.state.db; let did = msg.msg.did().expect("we checked if valid"); let repo_key = keys::repo_key(did); @@ -355,8 +356,9 @@ impl WorkerContext<'_> { if metadata.is_some_and(|m| !m.tracked) { trace!(did = %did, "ignoring message, repo is explicitly untracked"); #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, + ); return Ok(None); } @@ -369,8 +371,9 @@ impl WorkerContext<'_> { if let Some(repo_state) = repo_state_opt { #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Hit); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Hit, + ); return Ok(Some(repo_state)); } @@ -382,8 +385,9 @@ impl WorkerContext<'_> { SubscribeReposMessage::Commit(c) => c, _ => { #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, + ); return Ok(None); } }; @@ -404,8 +408,9 @@ impl WorkerContext<'_> { if !touches_signal { trace!(did = %did, "dropping commit, no signal-matching ops"); #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, + ); return Ok(None); } } @@ -432,8 +437,9 @@ impl WorkerContext<'_> { if pds_host.as_deref() != msg.firehose.host_str() { warn!(did = %did, got = ?pds_host, expected = ?msg.firehose.host_str(), "message rejected: wrong host for new account"); #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, + ); return Ok(None); } @@ -445,8 +451,9 @@ impl WorkerContext<'_> { if self.state.is_over_account_limit(host, count) { warn!(did = %did, host, count, "account limit reached for host, dropping new account"); #[cfg(feature = "firehose-diagnostics")] - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Drop, + ); return Ok(None); } } @@ -497,8 +504,9 @@ impl WorkerContext<'_> { #[cfg(feature = "firehose-diagnostics")] { - self.stats - .record_repo_state_outcome(crate::ingest::firehose_stats::RepoStateLoadOutcome::Miss); + self.stats.record_repo_state_outcome( + crate::ingest::firehose_stats::RepoStateLoadOutcome::Miss, + ); self.stats.record_new_account(new_account_started.elapsed()); } @@ -506,7 +514,10 @@ impl WorkerContext<'_> { } #[cfg(feature = "relay")] - pub(crate) fn queue_emit(&mut self, make_frame: impl FnOnce(i64) -> Result) -> Result { + pub(crate) fn queue_emit( + &mut self, + make_frame: impl FnOnce(i64) -> Result, + ) -> Result { #[cfg(feature = "firehose-diagnostics")] let started = Instant::now(); let result = (|| { diff --git a/src/ingest/relay/handlers.rs b/src/ingest/relay/handlers.rs index f10bc11..935ea16 100644 --- a/src/ingest/relay/handlers.rs +++ b/src/ingest/relay/handlers.rs @@ -1,29 +1,29 @@ use miette::Result; -use tokio::runtime::Handle; -use tracing::{debug, warn}; -use url::Url; use smol_str::SmolStr; #[cfg(all(feature = "relay", feature = "jetstream"))] use smol_str::ToSmolStr; +use tokio::runtime::Handle; +use tracing::{debug, warn}; +use url::Url; use crate::db::{self, keys}; use crate::types::{RepoState, RepoStatus}; -#[cfg(feature = "relay")] -use crate::types::RelayBroadcast; -#[cfg(all(feature = "relay", feature = "jetstream"))] -use crate::types::StoredJetstreamEvent; #[cfg(all(feature = "relay", feature = "jetstream"))] use crate::db::types::TrimmedDid; #[cfg(feature = "relay")] use crate::ingest::stream::encode_frame; -use crate::ingest::stream::{Commit, Sync, Identity, Account, AccountStatus}; +use crate::ingest::stream::{Account, AccountStatus, Commit, Identity, Sync}; use crate::ingest::validation::ValidatedCommit; -use jacquard_common::IntoStatic; +#[cfg(feature = "relay")] +use crate::types::RelayBroadcast; +#[cfg(all(feature = "relay", feature = "jetstream"))] +use crate::types::StoredJetstreamEvent; #[cfg(all(feature = "relay", feature = "jetstream"))] use jacquard_common::CowStr; +use jacquard_common::IntoStatic; -use super::{WorkerContext, RelayWorker}; +use super::{RelayWorker, WorkerContext}; impl RelayWorker { pub(crate) fn handle_commit( diff --git a/src/ingest/relay/worker.rs b/src/ingest/relay/worker.rs index 3b44c50..a566c9d 100644 --- a/src/ingest/relay/worker.rs +++ b/src/ingest/relay/worker.rs @@ -10,12 +10,12 @@ use tracing::{debug, error, info, info_span, trace, warn}; use url::Url; use crate::db::CountDeltas; -use crate::ingest::{BufferRx, BufferTx, IngestMessage}; -use crate::ingest::stream::{SubscribeReposMessage, InfoName}; +use crate::ingest::stream::{InfoName, SubscribeReposMessage}; use crate::ingest::validation::ValidationOptions; +use crate::ingest::{BufferRx, BufferTx, IngestMessage}; use crate::state::AppState; -use super::{WorkerContext, AuthorityOutcome}; +use super::{AuthorityOutcome, WorkerContext}; #[cfg(feature = "firehose-diagnostics")] use super::relay_message_kind; @@ -156,7 +156,9 @@ impl RelayWorker { match inf.name { InfoName::OutdatedCursor => {} InfoName::Other(name) => { - let message = inf.message.unwrap_or(jacquard_common::CowStr::Borrowed("")); + let message = inf + .message + .unwrap_or(jacquard_common::CowStr::Borrowed("")); info!(name = %name, "relay sent info: {message}"); } } @@ -310,9 +312,15 @@ impl RelayWorker { ctx.stats.record_host_authority( authority_started.elapsed(), match &outcome_result { - Ok(AuthorityOutcome::Authorized) => crate::ingest::firehose_stats::HostAuthorityStatsOutcome::Authorized, - Ok(AuthorityOutcome::WasStale) => crate::ingest::firehose_stats::HostAuthorityStatsOutcome::WasStale, - Ok(AuthorityOutcome::WrongHost { .. }) => crate::ingest::firehose_stats::HostAuthorityStatsOutcome::WrongHost, + Ok(AuthorityOutcome::Authorized) => { + crate::ingest::firehose_stats::HostAuthorityStatsOutcome::Authorized + } + Ok(AuthorityOutcome::WasStale) => { + crate::ingest::firehose_stats::HostAuthorityStatsOutcome::WasStale + } + Ok(AuthorityOutcome::WrongHost { .. }) => { + crate::ingest::firehose_stats::HostAuthorityStatsOutcome::WrongHost + } Err(_) => crate::ingest::firehose_stats::HostAuthorityStatsOutcome::Error, }, ); @@ -333,8 +341,10 @@ impl RelayWorker { let started = Instant::now(); let result = Self::handle_commit(ctx, &mut repo_state, &msg.firehose, *commit); #[cfg(feature = "firehose-diagnostics")] - ctx.stats - .record_handle_message(crate::ingest::firehose_stats::RelayMessageKind::Commit, started.elapsed()); + ctx.stats.record_handle_message( + crate::ingest::firehose_stats::RelayMessageKind::Commit, + started.elapsed(), + ); result } SubscribeReposMessage::Sync(sync) => { @@ -343,8 +353,10 @@ impl RelayWorker { let started = Instant::now(); let result = Self::handle_sync(ctx, &mut repo_state, &msg.firehose, *sync); #[cfg(feature = "firehose-diagnostics")] - ctx.stats - .record_handle_message(crate::ingest::firehose_stats::RelayMessageKind::Sync, started.elapsed()); + ctx.stats.record_handle_message( + crate::ingest::firehose_stats::RelayMessageKind::Sync, + started.elapsed(), + ); result } SubscribeReposMessage::Identity(identity) => { @@ -359,8 +371,10 @@ impl RelayWorker { msg.is_pds, ); #[cfg(feature = "firehose-diagnostics")] - ctx.stats - .record_handle_message(crate::ingest::firehose_stats::RelayMessageKind::Identity, started.elapsed()); + ctx.stats.record_handle_message( + crate::ingest::firehose_stats::RelayMessageKind::Identity, + started.elapsed(), + ); result } SubscribeReposMessage::Account(account) => { @@ -370,8 +384,10 @@ impl RelayWorker { let result = Self::handle_account(ctx, &mut repo_state, &msg.firehose, *account, msg.is_pds); #[cfg(feature = "firehose-diagnostics")] - ctx.stats - .record_handle_message(crate::ingest::firehose_stats::RelayMessageKind::Account, started.elapsed()); + ctx.stats.record_handle_message( + crate::ingest::firehose_stats::RelayMessageKind::Account, + started.elapsed(), + ); result } _ => Ok(()), diff --git a/src/ingest/stream.rs b/src/ingest/stream.rs index 777f3e1..5b430ea 100644 --- a/src/ingest/stream.rs +++ b/src/ingest/stream.rs @@ -12,18 +12,18 @@ use tokio_websockets::{ClientBuilder, Message as WsMsg, WebSocketStream}; use tracing::trace; use url::Url; -pub mod types; pub mod codec; +pub mod types; +pub use codec::decode_frame; #[allow(unused_imports)] pub use types::{ - Datetime, RepoOpAction, RepoOp, Commit, Identity, AccountStatus, Account, Sync, InfoName, - Info, SubscribeReposMessage, + Account, AccountStatus, Commit, Datetime, Identity, Info, InfoName, RepoOp, RepoOpAction, + SubscribeReposMessage, Sync, }; -pub use codec::decode_frame; #[cfg(feature = "relay")] -pub use codec::{encode_frame, encode_error_frame}; +pub use codec::{encode_error_frame, encode_frame}; #[derive(Debug, Error, Diagnostic)] pub enum FirehoseError { @@ -148,14 +148,14 @@ impl FirehoseStream { mod test { #[cfg(feature = "relay")] use super::FirehoseError; + #[cfg(feature = "relay")] + use super::types::{Datetime, RepoOp}; use super::{SubscribeReposMessage, decode_frame}; #[cfg(feature = "relay")] use jacquard_common::types::{ cid::CidLink, string::{Did, Tid}, }; - #[cfg(feature = "relay")] - use super::types::{RepoOp, Datetime}; #[cfg(feature = "relay")] #[derive(serde::Serialize)] diff --git a/src/ingest/stream/codec.rs b/src/ingest/stream/codec.rs index 74062c3..be98145 100644 --- a/src/ingest/stream/codec.rs +++ b/src/ingest/stream/codec.rs @@ -1,6 +1,6 @@ -use serde::{Deserialize, Serialize}; use super::FirehoseError; -use super::types::{SubscribeReposMessage, Info, InfoName}; +use super::types::{Info, InfoName, SubscribeReposMessage}; +use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] struct EventHeader { diff --git a/src/types.rs b/src/types.rs index 0e14dbc..9b3b928 100644 --- a/src/types.rs +++ b/src/types.rs @@ -9,13 +9,13 @@ use smol_str::ToSmolStr; use crate::db::types::DbTid; use crate::resolver::MiniDoc; +pub(crate) mod event; pub(crate) mod v2; pub(crate) mod v4; pub(crate) mod v7; -pub(crate) mod event; -pub(crate) use v7::*; pub(crate) use event::*; +pub(crate) use v7::*; impl<'c> From> for Commit { fn from(value: AtpCommit<'c>) -> Self { @@ -224,9 +224,9 @@ pub(crate) enum ResyncState { #[cfg(test)] mod tests { use super::*; - use miette::IntoDiagnostic; - use jacquard_common::types::string::Handle; use crate::db::types::DidKey; + use jacquard_common::types::string::Handle; + use miette::IntoDiagnostic; #[test] fn identity_dedupe_does_not_depend_on_commit_clock() { diff --git a/src/types/event.rs b/src/types/event.rs index 1a14208..f7bdd05 100644 --- a/src/types/event.rs +++ b/src/types/event.rs @@ -1,13 +1,13 @@ -use serde::{Deserialize, Serialize, Serializer}; -use serde_json::Value; use bytes::Bytes; -use jacquard_common::{CowStr, types::string::Handle}; #[cfg(feature = "jetstream")] use jacquard_common::IntoStatic; use jacquard_common::types::cid::IpldCid; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::{Did, Rkey}; use jacquard_common::types::tid::Tid; +use jacquard_common::{CowStr, types::string::Handle}; +use serde::{Deserialize, Serialize, Serializer}; +use serde_json::Value; use std::fmt::Debug; #[cfg(any(feature = "indexer_stream", feature = "jetstream"))] diff --git a/src/types/v2.rs b/src/types/v2.rs index 316b72f..f66b6c8 100644 --- a/src/types/v2.rs +++ b/src/types/v2.rs @@ -1,10 +1,10 @@ -use serde::{Deserialize, Serialize}; +use crate::db::types::{DbTid, DidKey}; use bytes::Bytes; +use jacquard_common::CowStr; use jacquard_common::types::cid::IpldCid; use jacquard_common::types::string::Handle; -use jacquard_common::CowStr; +use serde::{Deserialize, Serialize}; use smol_str::SmolStr; -use crate::db::types::{DbTid, DidKey}; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub enum RepoStatus { diff --git a/src/types/v4.rs b/src/types/v4.rs index 46cbf0b..38b68df 100644 --- a/src/types/v4.rs +++ b/src/types/v4.rs @@ -1,9 +1,9 @@ -use serde::{Deserialize, Serialize}; -use jacquard_common::types::string::Handle; +pub(crate) use super::v2::Commit; +use crate::db::types::DidKey; use jacquard_common::CowStr; +use jacquard_common::types::string::Handle; +use serde::{Deserialize, Serialize}; use smol_str::SmolStr; -use crate::db::types::DidKey; -pub(crate) use super::v2::Commit; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub enum RepoStatus { diff --git a/src/types/v7.rs b/src/types/v7.rs index 77f2eb2..e177c01 100644 --- a/src/types/v7.rs +++ b/src/types/v7.rs @@ -1,8 +1,8 @@ -use serde::{Deserialize, Serialize}; -use jacquard_common::types::string::Handle; -use jacquard_common::CowStr; -use crate::db::types::DidKey; pub(crate) use super::v4::{Commit, RepoMetadata, RepoStatus}; +use crate::db::types::DidKey; +use jacquard_common::CowStr; +use jacquard_common::types::string::Handle; +use serde::{Deserialize, Serialize}; #[derive(Debug, Clone, Serialize, Deserialize)] #[serde(bound(deserialize = "'i: 'de"))] -- 2.51.2