From 99cfce490e280b407c4dd9b12071fe081f179236 Mon Sep 17 00:00:00 2001 From: Mia Date: Fri, 18 Apr 2025 12:51:42 +0000 Subject: [PATCH] feat: labels --- Cargo.lock | 2 + consumer/Cargo.toml | 1 + consumer/src/config.rs | 3 + consumer/src/indexer/db.rs | 137 ++++++++++- consumer/src/indexer/mod.rs | 31 +++ consumer/src/indexer/records.rs | 22 +- consumer/src/indexer/types.rs | 5 + consumer/src/label_indexer/mod.rs | 227 ++++++++++++++++++ consumer/src/main.rs | 10 +- consumer/src/sql/label_copy_upsert.sql | 8 + lexica/src/app_bsky/actor.rs | 12 +- lexica/src/app_bsky/embed.rs | 9 +- lexica/src/app_bsky/feed.rs | 7 +- lexica/src/app_bsky/graph.rs | 13 +- lexica/src/app_bsky/labeler.rs | 51 ++++ lexica/src/app_bsky/mod.rs | 1 + lexica/src/com_atproto/label.rs | 147 ++++++++++++ lexica/src/com_atproto/mod.rs | 2 + lexica/src/com_atproto/moderation.rs | 60 +++++ lexica/src/lib.rs | 1 + migrations/2025-04-11-155301_labels/down.sql | 3 + migrations/2025-04-11-155301_labels/up.sql | 53 ++++ parakeet-db/src/models.rs | 97 ++++++++ parakeet-db/src/schema.rs | 44 ++++ parakeet/src/hydration/embed.rs | 62 +++-- parakeet/src/hydration/feedgen.rs | 26 +- parakeet/src/hydration/labeler.rs | 154 ++++++++++++ parakeet/src/hydration/list.rs | 50 +++- parakeet/src/hydration/mod.rs | 23 ++ parakeet/src/hydration/posts.rs | 60 +++-- parakeet/src/hydration/profile.rs | 91 +++++-- parakeet/src/hydration/starter_packs.rs | 45 ++-- parakeet/src/loaders.rs | 127 +++++++++- parakeet/src/xrpc/app_bsky/actor.rs | 8 +- parakeet/src/xrpc/app_bsky/feed/feedgen.rs | 11 +- parakeet/src/xrpc/app_bsky/feed/likes.rs | 4 +- parakeet/src/xrpc/app_bsky/feed/posts.rs | 25 +- parakeet/src/xrpc/app_bsky/graph/lists.rs | 10 +- parakeet/src/xrpc/app_bsky/graph/relations.rs | 13 +- .../src/xrpc/app_bsky/graph/starter_packs.rs | 18 +- parakeet/src/xrpc/app_bsky/labeler.rs | 52 ++++ parakeet/src/xrpc/app_bsky/mod.rs | 2 + parakeet/src/xrpc/extract.rs | 59 +++++ parakeet/src/xrpc/mod.rs | 1 + 44 files changed, 1647 insertions(+), 140 deletions(-) create mode 100644 consumer/src/label_indexer/mod.rs create mode 100644 consumer/src/sql/label_copy_upsert.sql create mode 100644 lexica/src/app_bsky/labeler.rs create mode 100644 lexica/src/com_atproto/label.rs create mode 100644 lexica/src/com_atproto/mod.rs create mode 100644 lexica/src/com_atproto/moderation.rs create mode 100644 migrations/2025-04-11-155301_labels/down.sql create mode 100644 migrations/2025-04-11-155301_labels/up.sql create mode 100644 parakeet/src/hydration/labeler.rs create mode 100644 parakeet/src/xrpc/app_bsky/labeler.rs create mode 100644 parakeet/src/xrpc/extract.rs diff --git a/Cargo.lock b/Cargo.lock index c4aec5d6..fe319cb5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -662,6 +662,7 @@ dependencies = [ "serde_ipld_dagcbor", "serde_json", "tokio", + "tokio-postgres", "tokio-tungstenite", "tracing", "tracing-subscriber", @@ -2330,6 +2331,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f66ea23a2d0e5734297357705193335e0a957696f34bed2f2faefacb2fec336f" dependencies = [ "bytes", + "chrono", "fallible-iterator", "postgres-protocol", ] diff --git a/consumer/Cargo.toml b/consumer/Cargo.toml index 59534be9..24c6d1f9 100644 --- a/consumer/Cargo.toml +++ b/consumer/Cargo.toml @@ -26,6 +26,7 @@ serde_bytes = "0.11" serde_ipld_dagcbor = "0.6.1" serde_json = "1.0.134" tokio = { version = "1.42.0", features = ["full"] } +tokio-postgres = { version = "0.7.12", features = ["with-chrono-0_4"] } tokio-tungstenite = { version = "0.26.1", features = ["native-tls"] } tracing = "0.1.40" tracing-subscriber = "0.3.18" diff --git a/consumer/src/config.rs b/consumer/src/config.rs index 13ff721d..3c618e48 100644 --- a/consumer/src/config.rs +++ b/consumer/src/config.rs @@ -23,6 +23,9 @@ pub struct Config { pub backfill_workers: u8, #[serde(default = "default_indexer_workers")] pub indexer_workers: u8, + /// DIDs of label services to force subscription to. + #[serde(default)] + pub initial_label_services: Vec, } #[derive(Copy, Clone, Debug, PartialEq, PartialOrd, Deserialize)] diff --git a/consumer/src/indexer/db.rs b/consumer/src/indexer/db.rs index db4640d4..e623e7bf 100644 --- a/consumer/src/indexer/db.rs +++ b/consumer/src/indexer/db.rs @@ -7,7 +7,9 @@ use diesel::serialize::ToSql; use diesel::sql_types::{Array, Nullable, Text, Timestamp}; use diesel_async::{AsyncPgConnection, RunQueryDsl}; use ipld_core::cid::Cid; +use lexica::com_atproto::label::{LabelValueDefinition, SelfLabels}; use parakeet_db::{models, schema, types}; +use std::collections::HashMap; pub async fn write_backfill_row( conn: &mut AsyncPgConnection, @@ -711,7 +713,6 @@ pub async fn upsert_starterpack( .description_facets .and_then(|v| serde_json::to_value(v).ok()); - let data = models::NewStarterPack { at_uri, cid: cid.to_string(), @@ -741,3 +742,137 @@ pub async fn delete_starterpack(conn: &mut AsyncPgConnection, at_uri: &str) -> Q .execute(conn) .await } + +pub async fn upsert_label_service( + conn: &mut AsyncPgConnection, + repo: &str, + cid: Cid, + rec: records::AppBskyLabelerService, +) -> QueryResult { + let reasons = rec + .reason_types + .as_ref() + .map(|v| v.iter().map(|v| v.to_string()).collect()); + + let data = models::UpsertLabelerService { + did: repo, + cid: cid.to_string(), + reasons, + indexed_at: Utc::now().naive_utc(), + }; + + let res = diesel::insert_into(schema::labelers::table) + .values(&data) + .on_conflict(schema::labelers::did) + .do_update() + .set(&data) + .execute(conn) + .await?; + + maintain_label_defs(conn, repo, &rec).await?; + + Ok(res) +} + +pub async fn delete_label_service(conn: &mut AsyncPgConnection, repo: &str) -> QueryResult { + diesel::delete(schema::labelers::table) + .filter(schema::labelers::did.eq(repo)) + .execute(conn) + .await +} + +pub async fn maintain_label_defs( + conn: &mut AsyncPgConnection, + repo: &str, + rec: &records::AppBskyLabelerService, +) -> QueryResult<()> { + // drop any label defs not currently in the list + diesel::delete(schema::labeler_defs::table) + .filter( + schema::labeler_defs::labeler + .eq(repo) + .and(schema::labeler_defs::label_identifier.ne_all(&rec.policies.label_values)), + ) + .execute(conn) + .await?; + + let definitions = rec + .policies + .label_value_definitions + .iter() + .map(|def| (def.identifier.clone(), def)) + .collect::>(); + + for label in &rec.policies.label_values { + let definition = definitions.get(label); + + let locales = definition.and_then(|v| serde_json::to_value(&v.locales).ok()); + + let data = models::UpsertLabelDefinition { + labeler: repo, + label_identifier: label, + severity: definition.map(|v| v.severity.to_string()), + blurs: definition.map(|v| v.blurs.to_string()), + default_setting: definition + .and_then(|v| v.default_setting) + .map(|v| v.to_string()), + adult_only: definition.and_then(|v| v.adult_only), + locales, + indexed_at: Utc::now().naive_utc(), + }; + + diesel::insert_into(schema::labeler_defs::table) + .values(&data) + .on_conflict(( + schema::labeler_defs::labeler, + schema::labeler_defs::label_identifier, + )) + .do_update() + .set(&data) + .execute(conn) + .await?; + } + + Ok(()) +} + +pub async fn maintain_self_labels( + conn: &mut AsyncPgConnection, + repo: &str, + cid: Option, + at_uri: &str, + self_labels: SelfLabels, +) -> QueryResult { + // purge any existing self-labels + diesel::delete(schema::labels::table) + .filter( + schema::labels::self_label + .eq(true) + .and(schema::labels::uri.eq(at_uri)), + ) + .execute(conn) + .await?; + + let cid = cid.map(|cid| cid.to_string()); + let now = Utc::now().naive_utc(); + + let labels = self_labels + .values + .iter() + .map(|v| models::NewLabel { + labeler: repo, + label: &v.val, + uri: at_uri, + self_label: true, + cid: cid.clone(), + expires: None, + sig: None, + created_at: now, + }) + .collect::>(); + + diesel::insert_into(schema::labels::table) + .values(&labels) + .execute(conn) + .await +} diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index 0f77a8dc..89d91dde 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -397,11 +397,21 @@ pub async fn index_op( match record { RecordTypes::AppBskyActorProfile(record) => { if rkey == "self" { + let labels = record.labels.clone(); db::upsert_profile(conn, repo, cid, record).await?; + + if let Some(labels) = labels { + db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; + } } } RecordTypes::AppBskyFeedGenerator(record) => { + let labels = record.labels.clone(); db::upsert_feedgen(conn, repo, cid, at_uri, record).await?; + + if let Some(labels) = labels { + db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; + } } RecordTypes::AppBskyFeedLike(record) => { db::insert_like(conn, repo, at_uri, record).await?; @@ -413,7 +423,11 @@ pub async fn index_op( } } + let labels = record.labels.clone(); db::insert_post(conn, repo, cid, at_uri, record).await?; + if let Some(labels) = labels { + db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; + } } RecordTypes::AppBskyFeedPostgate(record) => { let split_aturi = record.post.rsplitn(4, '/').collect::>(); @@ -456,7 +470,13 @@ pub async fn index_op( db::insert_follow(conn, repo, at_uri, record).await?; } RecordTypes::AppBskyGraphList(record) => { + let labels = record.labels.clone(); db::upsert_list(conn, repo, at_uri, cid, record).await?; + + if let Some(labels) = labels { + db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; + } + // todo: when we have profile stats, update them. } RecordTypes::AppBskyGraphListBlock(record) => { @@ -475,6 +495,16 @@ pub async fn index_op( RecordTypes::AppBskyGraphStarterPack(record) => { db::upsert_starterpack(conn, repo, cid, at_uri, record).await?; } + RecordTypes::AppBskyLabelerService(record) => { + if rkey == "self" { + let labels = record.labels.clone(); + db::upsert_label_service(conn, repo, cid, record).await?; + + if let Some(labels) = labels { + db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; + } + } + } RecordTypes::ChatBskyActorDeclaration(record) => { if rkey == "self" { db::upsert_chat_decl(conn, repo, record).await?; @@ -511,6 +541,7 @@ pub async fn index_op_delete( CollectionType::BskyListBlock => db::delete_list_block(conn, at_uri).await?, CollectionType::BskyListItem => db::delete_list_item(conn, at_uri).await?, CollectionType::BskyStarterPack => db::delete_starterpack(conn, at_uri).await?, + CollectionType::BskyLabelerService => db::delete_label_service(conn, at_uri).await?, CollectionType::ChatActorDecl => db::delete_chat_decl(conn, at_uri).await?, _ => unreachable!(), }; diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs index 89680cb6..fadbfbf6 100644 --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -3,7 +3,10 @@ use chrono::{DateTime, Utc}; use ipld_core::cid::Cid; use lexica::app_bsky::actor::ChatAllowIncoming; use lexica::app_bsky::embed::AspectRatio; +use lexica::app_bsky::labeler::LabelerPolicy; use lexica::app_bsky::richtext::FacetMain; +use lexica::com_atproto::label::SelfLabels; +use lexica::com_atproto::moderation::{ReasonType, SubjectType}; use serde::{Deserialize, Serialize}; #[derive(Debug, Deserialize, Serialize)] @@ -34,6 +37,7 @@ pub struct AppBskyActorProfile { pub description: Option, pub avatar: Option, pub banner: Option, + pub labels: Option, pub joined_via_starter_pack: Option, pub pinned_post: Option, pub created_at: Option>, @@ -144,7 +148,7 @@ pub struct AppBskyFeedGenerator { pub description_facets: Option>, pub avatar: Option, pub accepts_interactions: Option, - // pub labels: Option>, + pub labels: Option, pub content_mode: Option, pub created_at: DateTime, } @@ -170,7 +174,8 @@ pub struct AppBskyFeedPost { pub embed: Option, #[serde(skip_serializing_if = "Option::is_none")] pub langs: Option>, - // pub labels: Option>, + #[serde(skip_serializing_if = "Option::is_none")] + pub labels: Option, #[serde(skip_serializing_if = "Option::is_none")] pub tags: Option>, pub created_at: DateTime, @@ -274,7 +279,7 @@ pub struct AppBskyGraphList { pub description: Option, pub description_facets: Option>, pub avatar: Option, - // pub labels: Option>, + pub labels: Option, pub created_at: DateTime, } @@ -311,6 +316,17 @@ pub struct StarterPackFeedItem { pub uri: String, } +#[derive(Debug, Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct AppBskyLabelerService { + pub policies: LabelerPolicy, + pub labels: Option, + pub reason_types: Option>, + pub subject_types: Option>, + pub subject_collections: Option>, + pub created_at: DateTime, +} + #[derive(Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct ChatBskyActorDeclaration { diff --git a/consumer/src/indexer/types.rs b/consumer/src/indexer/types.rs index 8e810576..4005f7a7 100644 --- a/consumer/src/indexer/types.rs +++ b/consumer/src/indexer/types.rs @@ -31,6 +31,8 @@ pub enum RecordTypes { AppBskyGraphListItem(records::AppBskyGraphListItem), #[serde(rename = "app.bsky.graph.starterpack")] AppBskyGraphStarterPack(records::AppBskyGraphStarterPack), + #[serde(rename = "app.bsky.labeler.service")] + AppBskyLabelerService(records::AppBskyLabelerService), #[serde(rename = "chat.bsky.actor.declaration")] ChatBskyActorDeclaration(records::ChatBskyActorDeclaration), } @@ -50,6 +52,7 @@ pub enum CollectionType { BskyListBlock, BskyListItem, BskyStarterPack, + BskyLabelerService, ChatActorDecl, Unsupported, } @@ -70,6 +73,7 @@ impl CollectionType { "app.bsky.graph.listblock" => CollectionType::BskyListBlock, "app.bsky.graph.listitem" => CollectionType::BskyListItem, "app.bsky.graph.starterpack" => CollectionType::BskyStarterPack, + "app.bsky.labeler.service" => CollectionType::BskyLabelerService, "chat.bsky.actor.declaration" => CollectionType::ChatActorDecl, _ => CollectionType::Unsupported, } @@ -91,6 +95,7 @@ impl CollectionType { CollectionType::BskyListItem => false, CollectionType::ChatActorDecl => true, CollectionType::BskyStarterPack => true, + CollectionType::BskyLabelerService => true, CollectionType::Unsupported => false, } } diff --git a/consumer/src/label_indexer/mod.rs b/consumer/src/label_indexer/mod.rs new file mode 100644 index 00000000..210c2a4f --- /dev/null +++ b/consumer/src/label_indexer/mod.rs @@ -0,0 +1,227 @@ +use crate::firehose::{AtpLabel, FirehoseConsumer, FirehoseEvent, FirehoseOutput}; +use did_resolver::Resolver; +use futures::pin_mut; +use metrics::counter; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; +use tokio::sync::mpsc::{channel, Receiver, Sender}; +use tokio::task::JoinHandle; +use tokio_postgres::binary_copy::BinaryCopyInWriter; +use tokio_postgres::types::Type; +use tokio_postgres::NoTls; +use tracing::instrument; + +const LABELER_SERVICE_ID: &str = "#atproto_labeler"; + +pub struct LabelServiceManager { + client: tokio_postgres::Client, + rx: Receiver, + resolver: Arc, + services: HashMap>, + user_agent: String, +} + +impl LabelServiceManager { + pub async fn new( + pg_url: &str, + resolver: Arc, + user_agent: String, + ) -> eyre::Result<(Self, Sender)> { + let (client, connection) = tokio_postgres::connect(pg_url, NoTls).await?; + + tokio::spawn(async move { + if let Err(e) = connection.await { + tracing::error!("connection error: {}", e); + } + }); + + let (tx, rx) = channel(8); + + let lsm = LabelServiceManager { + client, + rx, + resolver, + services: HashMap::new(), + user_agent, + }; + + Ok((lsm, tx)) + } + + pub async fn run(mut self, inital_services: Vec) -> eyre::Result<()> { + let (db_tx, mut db_rx) = channel(8192); + + for service in inital_services { + let service_endpoint = resolve_service_endpoint(&self.resolver, &service).await?; + let service_endpoint = service_endpoint.replace("https://", "wss://"); + + let handle = tokio::spawn(label_consumer( + service.clone(), + service_endpoint, + self.user_agent.clone(), + db_tx.clone(), + )); + self.services.insert(service, handle); + } + + let mut timer = tokio::time::interval(Duration::from_millis(250)); + + let mut buf = Vec::with_capacity(8192); + + loop { + let res = tokio::select! { + Some(service) = self.rx.recv() => { + if self.services.contains_key(&service) { + continue; + } + + let service_endpoint = resolve_service_endpoint(&self.resolver, &service).await?; + let service_endpoint = service_endpoint.replace("https://", "wss://"); + + let handle = tokio::spawn(label_consumer( + service.clone(), + service_endpoint, + self.user_agent.clone(), + db_tx.clone(), + )); + + self.services.insert(service, handle); + + continue; + } + _ = timer.tick() => { + db_rx.recv_many(&mut buf, db_rx.len()).await; + if buf.is_empty() { + continue; + } + tracing::debug!("got {} labels", buf.len()); + store_labels(&mut self.client, &buf).await + } + }; + + match res { + Ok(count) => counter!("inserted_labels").increment(count), + Err(e) => tracing::error!("failed to store labels {}", e), + } + + buf.clear(); + } + + Ok(()) + } +} + +async fn store_labels(conn: &mut tokio_postgres::Client, labels: &[AtpLabel]) -> eyre::Result { + let t = conn.transaction().await?; + t.execute( + "CREATE TEMP TABLE label_tmp (LIKE labels INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let sink = t.copy_in("COPY label_tmp (labeler, label, uri, cid, negated, expires, sig, created_at) FROM STDIN (FORMAT binary)").await?; + let binary_writer = BinaryCopyInWriter::new( + sink, + &[ + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::BOOL, + Type::TIMESTAMP, + Type::BYTEA, + Type::TIMESTAMP, + ], + ); + + pin_mut!(binary_writer); + + for label in labels { + let exp = label.exp.map(|v| v.naive_utc()); + let sig = label.sig.as_ref().map(|v| v.as_slice()); + let created = label.cts.naive_utc(); + let neg = label.neg.unwrap_or_default(); + + binary_writer + .as_mut() + .write(&[ + &label.src, &label.val, &label.uri, &label.cid, &neg, &exp, &sig, &created, + ]) + .await?; + } + + let count = binary_writer.finish().await?; + + t.execute(include_str!("../sql/label_copy_upsert.sql"), &[]) + .await?; + + t.commit().await?; + + Ok(count) +} + +#[instrument(skip(service_did, user_agent, db_tx))] +async fn label_consumer( + service_did: String, + service: String, + user_agent: String, + db_tx: Sender, +) { + let mut consumer = match FirehoseConsumer::new_labeler(&service, None, &user_agent).await { + Ok(consumer) => consumer, + Err(err) => { + tracing::error!("Failed to connect to labeler: {err}"); + return; + } + }; + + while let Ok(output) = consumer.drive().await { + match output { + FirehoseOutput::Close => break, + FirehoseOutput::Continue => continue, + FirehoseOutput::Error(err) => { + tracing::error!("Firehose sent an error, exiting: {err:?}"); + break; + } + FirehoseOutput::Event(event) => { + if let FirehoseEvent::Label(label_event) = *event { + let count = label_event.labels.len() as u64; + + for label in label_event.labels { + if label.src != service_did { + tracing::warn!("labeler sent incorrect src"); + continue; + } + + if let Err(e) = db_tx.send(label).await { + tracing::error!("Failed to send label event: {e:?}"); + return; + } + } + + counter!("seen_labels", "service" => service_did.clone()).increment(count); + } else { + tracing::warn!("got non #label event from a labeler"); + } + } + } + } +} + +async fn resolve_service_endpoint(resolver: &Resolver, did: &str) -> eyre::Result { + // resolve the did to a PDS (also validates the handle) + let Some(did_doc) = resolver.resolve_did(did).await? else { + eyre::bail!("missing did doc"); + }; + + let Some(service) = did_doc.service.and_then(|services| { + services + .into_iter() + .find(|svc| svc.id == LABELER_SERVICE_ID) + }) else { + eyre::bail!("DID doc contained no service endpoint"); + }; + + Ok(service.service_endpoint) +} diff --git a/consumer/src/main.rs b/consumer/src/main.rs index be7ff58f..f3da8a5d 100644 --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -10,6 +10,7 @@ mod backfill; mod config; mod firehose; mod indexer; +mod label_indexer; mod utils; #[tokio::main] @@ -30,10 +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 relay_firehose = firehose::FirehoseConsumer::new_relay(&conf.relay_source, None, &user_agent).await?; let (backfiller, backfill_tx) = backfill::BackfillManager::new(pool.clone(), conf.history_mode, resolver.clone()).await?; + let (relay_indexer, tx) = indexer::RelayIndexer::new( pool.clone(), backfill_tx, @@ -42,13 +45,14 @@ async fn main() -> eyre::Result<()> { ) .await?; - let (firehose_res, indexer_res, backfill_res) = tokio::try_join! { + let (firehose_res, indexer_res, backfill_res, label_res) = tokio::try_join! { tokio::spawn(relay_consumer(relay_firehose, tx)), tokio::spawn(relay_indexer.run(conf.indexer_workers)), - tokio::spawn(backfiller.run(conf.backfill_workers)) + tokio::spawn(backfiller.run(conf.backfill_workers)), + tokio::spawn(label_mgr.run(conf.initial_label_services)), }?; - firehose_res.and(indexer_res).and(backfill_res) + firehose_res.and(indexer_res).and(backfill_res).and(label_res) } async fn relay_consumer( diff --git a/consumer/src/sql/label_copy_upsert.sql b/consumer/src/sql/label_copy_upsert.sql new file mode 100644 index 00000000..82458906 --- /dev/null +++ b/consumer/src/sql/label_copy_upsert.sql @@ -0,0 +1,8 @@ +INSERT INTO labels +SELECT DISTINCT on (labeler, label, uri) * +FROM label_tmp +ON CONFLICT (labeler, label, uri) DO UPDATE + SET negated=EXCLUDED.negated, + expires=EXCLUDED.expires, + sig=EXCLUDED.sig, + created_at=excluded.created_at \ No newline at end of file diff --git a/lexica/src/app_bsky/actor.rs b/lexica/src/app_bsky/actor.rs index 6856a7fa..b30ea27c 100644 --- a/lexica/src/app_bsky/actor.rs +++ b/lexica/src/app_bsky/actor.rs @@ -1,6 +1,7 @@ -use std::str::FromStr; +use crate::com_atproto::label::Label; use chrono::prelude::*; use serde::{Deserialize, Serialize}; +use std::str::FromStr; #[derive(Clone, Default, Debug, Serialize)] #[serde(rename_all = "camelCase")] @@ -64,7 +65,8 @@ pub struct ProfileViewBasic { pub associated: Option, // #[serde(skip_serializing_if = "Option::is_none")] // pub viewer: Option<()>, - // pub labels: Vec<()>, + #[serde(skip_serializing_if = "Vec::is_empty")] + pub labels: Vec