diff --git a/consumer/src/backfill/mod.rs b/consumer/src/backfill/mod.rs index 150f9235..5563ce3a 100644 --- a/consumer/src/backfill/mod.rs +++ b/consumer/src/backfill/mod.rs @@ -170,14 +170,7 @@ async fn backfill_actor( let status = repo_status .status .unwrap_or(crate::firehose::AtpAccountStatus::Deleted); - db::actor_upsert( - conn, - did, - status.into(), - ActorSyncState::Dirty, - Utc::now(), - ) - .await?; + db::actor_upsert(conn, did, status.into(), ActorSyncState::Dirty, Utc::now()).await?; return Ok(()); } @@ -336,9 +329,19 @@ async fn check_pds_repo_status( #[derive(Debug, Default)] struct CopyStore { - likes: Vec<(String, records::StrongRef, Option, DateTime)>, + likes: Vec<( + String, + records::StrongRef, + Option, + DateTime, + )>, posts: Vec<(String, Cid, records::AppBskyFeedPost)>, - reposts: Vec<(String, records::StrongRef, Option, DateTime)>, + reposts: Vec<( + String, + records::StrongRef, + Option, + DateTime, + )>, blocks: Vec<(String, String, DateTime)>, follows: Vec<(String, String, DateTime)>, list_items: Vec<(String, records::AppBskyGraphListItem)>, diff --git a/consumer/src/backfill/repo.rs b/consumer/src/backfill/repo.rs index a887fcbe..4833062c 100644 --- a/consumer/src/backfill/repo.rs +++ b/consumer/src/backfill/repo.rs @@ -131,7 +131,9 @@ async fn record_index( deltas.incr(&rec.subject.uri, AggregateType::Like).await; copies.push_record(&at_uri, cid); - copies.likes.push((at_uri, rec.subject, rec.via, rec.created_at)); + copies + .likes + .push((at_uri, rec.subject, rec.via, rec.created_at)); } RecordTypes::AppBskyFeedPost(rec) => { let maybe_reply = rec.reply.as_ref().map(|v| v.parent.uri.clone()); @@ -167,7 +169,9 @@ async fn record_index( deltas.incr(&rec.subject.uri, AggregateType::Repost).await; copies.push_record(&at_uri, cid); - copies.reposts.push((at_uri, rec.subject, rec.via, rec.created_at)); + copies + .reposts + .push((at_uri, rec.subject, rec.via, rec.created_at)); } RecordTypes::AppBskyGraphBlock(rec) => { copies.push_record(&at_uri, cid); diff --git a/consumer/src/db/backfill.rs b/consumer/src/db/backfill.rs index d7f84e56..9c62dcab 100644 --- a/consumer/src/db/backfill.rs +++ b/consumer/src/db/backfill.rs @@ -13,7 +13,11 @@ pub struct BackfillRow { pub indexed_at: NaiveDateTime, } -pub async fn backfill_job_write(conn: &mut C, did: &str, status: &str) -> PgExecResult { +pub async fn backfill_job_write( + conn: &mut C, + did: &str, + status: &str, +) -> PgExecResult { conn.execute( "INSERT INTO backfill_jobs (did, status) VALUES ($1, $2)", &[&did, &status], diff --git a/consumer/src/db/copy.rs b/consumer/src/db/copy.rs index 82ae0827..1bf09950 100644 --- a/consumer/src/db/copy.rs +++ b/consumer/src/db/copy.rs @@ -18,7 +18,12 @@ const STRONGREF_TYPES: &[Type] = &[ Type::TEXT, Type::TIMESTAMP, ]; -type StrongRefRow = (String, records::StrongRef, Option, DateTime); +type StrongRefRow = ( + String, + records::StrongRef, + Option, + DateTime, +); // SubjectRefs are used in both blocks and follows const SUBJECT_TYPES: &[Type] = &[Type::TEXT, Type::TEXT, Type::TEXT, Type::TIMESTAMP]; @@ -50,7 +55,7 @@ pub async fn copy_likes( for row in data { let (via_uri, via_cid) = strongref_to_parts(row.2.as_ref()); - + let writer = writer.as_mut(); writer .write(&[ @@ -97,7 +102,7 @@ pub async fn copy_reposts( for row in data { let (via_uri, via_cid) = strongref_to_parts(row.2.as_ref()); - + let writer = writer.as_mut(); writer .write(&[ diff --git a/consumer/src/db/record.rs b/consumer/src/db/record.rs index 0fde4652..ca73bfa2 100644 --- a/consumer/src/db/record.rs +++ b/consumer/src/db/record.rs @@ -158,7 +158,7 @@ pub async fn like_insert( rec: AppBskyFeedLike, ) -> PgExecResult { let (via_uri, via_cid) = strongref_to_parts(rec.via.as_ref()); - + conn.execute( "INSERT INTO likes (at_uri, did, subject, subject_cid, via_uri, via_cid, created_at) VALUES ($1, $2, $3, $4, $5, $6, $7)", &[&at_uri, &repo, &rec.subject.uri, &rec.subject.cid.to_string(), &via_uri, &via_cid, &rec.created_at] @@ -204,7 +204,7 @@ pub async fn list_upsert( ], ) .await - .map(|v| v == 0) + .map(|v| v == 0) } pub async fn list_delete(conn: &mut C, at_uri: &str) -> PgExecResult { @@ -528,7 +528,7 @@ pub async fn repost_insert( rec: AppBskyFeedRepost, ) -> PgExecResult { let (via_uri, via_cid) = strongref_to_parts(rec.via.as_ref()); - + conn.execute( "INSERT INTO reposts (at_uri, did, post, post_cid, via_uri, via_cid, created_at) VALUES ($1, $2, $3, $4, $5, $6, $7)", &[ @@ -587,7 +587,7 @@ pub async fn starter_pack_upsert( ], ) .await - .map(|v| v == 0) + .map(|v| v == 0) } pub async fn starter_pack_delete(conn: &mut C, at_uri: &str) -> PgExecResult { diff --git a/consumer/src/main.rs b/consumer/src/main.rs index c4f516b5..2dc660b5 100644 --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -59,12 +59,9 @@ async fn main() -> eyre::Result<()> { if cli.labels { let resume = resume.clone().unwrap(); - let label_mgr = label_indexer::LabelServiceManager::new( - pool.clone(), - resume, - user_agent.clone(), - ) - .await?; + let label_mgr = + label_indexer::LabelServiceManager::new(pool.clone(), resume, user_agent.clone()) + .await?; if let Some(label_source) = conf.label_source { join_set.spawn(label_mgr.run(label_source)); diff --git a/lexica/src/app_bsky/actor.rs b/lexica/src/app_bsky/actor.rs index 01977e0f..8c6914e4 100644 --- a/lexica/src/app_bsky/actor.rs +++ b/lexica/src/app_bsky/actor.rs @@ -215,5 +215,5 @@ pub struct StatusView { #[serde(tag = "$type")] #[serde(rename = "app.bsky.embed.external#view")] pub struct StatusViewEmbed { - pub external: External + pub external: External, } diff --git a/parakeet-db/src/models.rs b/parakeet-db/src/models.rs index c781a746..e0586721 100644 --- a/parakeet-db/src/models.rs +++ b/parakeet-db/src/models.rs @@ -727,7 +727,7 @@ pub struct Status { pub did: String, pub status: String, pub duration: Option, - + pub record: serde_json::Value, pub embed_uri: Option, @@ -739,4 +739,3 @@ pub struct Status { pub created_at: NaiveDateTime, pub indexed_at: NaiveDateTime, } - diff --git a/parakeet/src/xrpc/cdn.rs b/parakeet/src/xrpc/cdn.rs index e87c9615..4218b6f7 100644 --- a/parakeet/src/xrpc/cdn.rs +++ b/parakeet/src/xrpc/cdn.rs @@ -6,7 +6,10 @@ pub struct BskyCdn { impl BskyCdn { pub fn new(cdn_base: String, video_base: String) -> Self { - BskyCdn {cdn_base, video_base} + BskyCdn { + cdn_base, + video_base, + } } pub fn avatar(&self, did: &str, cid: &str) -> String { @@ -18,7 +21,10 @@ impl BskyCdn { } pub fn embed_thumb(&self, did: &str, cid: &str) -> String { - format!("{}/img/feed_thumbnail/plain/{did}/{cid}@jpeg", self.cdn_base) + format!( + "{}/img/feed_thumbnail/plain/{did}/{cid}@jpeg", + self.cdn_base + ) } pub fn embed_fullsize(&self, did: &str, cid: &str) -> String {