diff --git a/Cargo.lock b/Cargo.lock --- a/Cargo.lock +++ b/Cargo.lock @@ -39,6 +39,21 @@ ] [[package]] +name = "alloc-no-stdlib" +version = "2.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc7bb162ec39d46ab1ca8c77bf72e890535becd1751bb45f64c597edb4c8c6b3" + +[[package]] +name = "alloc-stdlib" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94fb8275041c72129eb51b7d0322c29b8387a0386127718b096429201a5d6ece" +dependencies = [ + "alloc-no-stdlib", +] + +[[package]] name = "android-tzdata" version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -130,6 +145,19 @@ "event-listener-strategy", "futures-core", "pin-project-lite", +] + +[[package]] +name = "async-compression" +version = "0.4.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "59a194f9d963d8099596278594b3107448656ba73831c9d8c783e613ce86da64" +dependencies = [ + "brotli", + "futures-core", + "memchr", + "pin-project-lite", + "tokio", ] [[package]] @@ -472,6 +500,27 @@ ] [[package]] +name = "brotli" +version = "7.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc97b8f16f944bba54f0433f07e30be199b6dc2bd25937444bbad560bcea29bd" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", + "brotli-decompressor", +] + +[[package]] +name = "brotli-decompressor" +version = "4.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a334ef7c9e23abf0ce748e8cd309037da93e606ad52eb372e4ce327a0dcfbdfd" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", +] + +[[package]] name = "bumpalo" version = "3.16.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -688,12 +737,10 @@ "chrono", "ciborium", "clap", + "deadpool-postgres", "did-resolver", - "diesel", - "diesel-async", "eyre", "figment", - "flume", "foldhash", "futures", "ipld-core", @@ -715,6 +762,7 @@ "tokio-postgres", "tokio-stream", "tokio-tungstenite", + "tokio-util", "tracing", "tracing-subscriber", ] @@ -800,7 +848,7 @@ checksum = "0dc92fb57ca44df6db8059111ab3af99a63d5d0f8375d9972e319a379c6bab76" dependencies = [ "generic-array", - "rand_core", + "rand_core 0.6.4", "subtle", "zeroize", ] @@ -920,7 +968,23 @@ dependencies = [ "deadpool-runtime", "num_cpus", + "serde", "tokio", +] + +[[package]] +name = "deadpool-postgres" +version = "0.14.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3d697d376cbfa018c23eb4caab1fd1883dd9c906a8c034e8d9a3cb06a7e0bef9" +dependencies = [ + "async-trait", + "deadpool", + "getrandom 0.2.15", + "serde", + "tokio", + "tokio-postgres", + "tracing", ] [[package]] @@ -928,6 +992,9 @@ version = "0.1.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "092966b41edc516079bdf31ec78a2e0588d1d0c08f78b91d8307215928642b2b" +dependencies = [ + "tokio", +] [[package]] name = "der" @@ -1114,7 +1181,7 @@ "hkdf", "pem-rfc7468", "pkcs8", - "rand_core", + "rand_core 0.6.4", "sec1", "subtle", "zeroize", @@ -1212,7 +1279,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c0b50bfb653653f9ca9095b427bed08ab8d75a137839d9ad64eb11810d5b6393" dependencies = [ - "rand_core", + "rand_core 0.6.4", "subtle", ] @@ -1241,18 +1308,6 @@ version = "0.5.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1d674e81391d1e1ab681a28d99df07927c6d4aa5b027d7da16ba32d1d21ecd99" - -[[package]] -name = "flume" -version = "0.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" -dependencies = [ - "futures-core", - "futures-sink", - "nanorand", - "spin", -] [[package]] name = "fnv" @@ -1437,8 +1492,20 @@ "cfg-if", "js-sys", "libc", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", "wasm-bindgen", +] + +[[package]] +name = "getrandom" +version = "0.3.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "26145e563e54f2cadc477553f1ec5ee650b00862f0a58bcd12cbdc5f0ea2d2f4" +dependencies = [ + "cfg-if", + "libc", + "r-efi", + "wasi 0.14.2+wasi-0.2.4", ] [[package]] @@ -1472,7 +1539,7 @@ checksum = "f0f9ef7462f7c099f518d754361858f86d8a07af53ba9af0fe635bbccb151a63" dependencies = [ "ff", - "rand_core", + "rand_core 0.6.4", "subtle", ] @@ -1572,7 +1639,7 @@ "idna", "ipnet", "once_cell", - "rand", + "rand 0.8.5", "thiserror 1.0.69", "tinyvec", "tokio", @@ -1593,7 +1660,7 @@ "lru-cache", "once_cell", "parking_lot 0.12.3", - "rand", + "rand 0.8.5", "resolv-conf", "smallvec", "thiserror 1.0.69", @@ -2072,15 +2139,15 @@ dependencies = [ "base64 0.22.1", "ed25519-dalek", - "getrandom", + "getrandom 0.2.15", "hmac", "js-sys", "k256", "p256", "p384", "pem", - "rand", - "rand_core", + "rand 0.8.5", + "rand_core 0.6.4", "rsa", "serde", "serde_json", @@ -2273,7 +2340,7 @@ "hashbrown", "metrics", "quanta", - "rand", + "rand 0.8.5", "rand_xoshiro", "sketches-ddsketch", ] @@ -2306,7 +2373,7 @@ checksum = "2886843bf800fba2e3377cff24abf6379b4c4d5c6681eaf9ea5b0d15090450bd" dependencies = [ "libc", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", "windows-sys 0.52.0", ] @@ -2337,15 +2404,6 @@ version = "0.10.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "defc4c55412d89136f966bbb339008b474350e5e6e78d2714439c386b3137a03" - -[[package]] -name = "nanorand" -version = "0.7.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a51313c5820b0b02bd422f4b44776fbf47961755c74ce64afc73bfad10226c3" -dependencies = [ - "getrandom", -] [[package]] name = "native-tls" @@ -2406,7 +2464,7 @@ "num-integer", "num-iter", "num-traits", - "rand", + "rand 0.8.5", "smallvec", "zeroize", ] @@ -2581,6 +2639,7 @@ dependencies = [ "chrono", "diesel", + "postgres-types", "serde_json", ] @@ -2843,9 +2902,9 @@ [[package]] name = "postgres-protocol" -version = "0.6.7" +version = "0.6.8" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "acda0ebdebc28befa84bee35e651e4c5f09073d668c7aed4cf7e23c3cda84b23" +checksum = "76ff0abab4a9b844b93ef7b81f1efc0a366062aaef2cd702c76256b5dc075c54" dependencies = [ "base64 0.22.1", "byteorder", @@ -2854,21 +2913,23 @@ "hmac", "md-5", "memchr", - "rand", + "rand 0.9.1", "sha2", "stringprep", ] [[package]] name = "postgres-types" -version = "0.2.8" +version = "0.2.9" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f66ea23a2d0e5734297357705193335e0a957696f34bed2f2faefacb2fec336f" +checksum = "613283563cd90e1dfc3518d548caee47e0e725455ed619881f5cf21f36de4b48" dependencies = [ "bytes", "chrono", "fallible-iterator", "postgres-protocol", + "serde", + "serde_json", ] [[package]] @@ -2989,7 +3050,7 @@ "libc", "once_cell", "raw-cpuid", - "wasi", + "wasi 0.11.0+wasi-snapshot-preview1", "web-sys", "winapi", ] @@ -3010,14 +3071,30 @@ ] [[package]] +name = "r-efi" +version = "5.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74765f6d916ee2faa39bc8e68e4f3ed8949b48cccdac59983d287a7cb71ce9c5" + +[[package]] name = "rand" version = "0.8.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34af8d1a0e25924bc5b7c43c079c942339d8f0a8b57c39049bef581b46327404" dependencies = [ "libc", - "rand_chacha", - "rand_core", + "rand_chacha 0.3.1", + "rand_core 0.6.4", +] + +[[package]] +name = "rand" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9fbfd9d094a40bf3ae768db9361049ace4c0e04a4fd6b359518bd7b73a73dd97" +dependencies = [ + "rand_chacha 0.9.0", + "rand_core 0.9.3", ] [[package]] @@ -3027,7 +3104,17 @@ checksum = "e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88" dependencies = [ "ppv-lite86", - "rand_core", + "rand_core 0.6.4", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core 0.9.3", ] [[package]] @@ -3036,7 +3123,16 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ec0be4795e2f6a28069bec0b5ff3e2ac9bafc99e6a9a7dc3547996c5c816922c" dependencies = [ - "getrandom", + "getrandom 0.2.15", +] + +[[package]] +name = "rand_core" +version = "0.9.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "99d9a13982dcf210057a8a78572b2217b667c3beacbf3a0d8b454f6f82837d38" +dependencies = [ + "getrandom 0.3.3", ] [[package]] @@ -3045,7 +3141,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6f97cdb2a36ed4183de61b2f824cc45c9f1037f28afe0a322e9fff4c108b5aaa" dependencies = [ - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -3134,6 +3230,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "43e734407157c3c2034e0258f5e4473ddb361b1e85f95a66690d67264d7cd1da" dependencies = [ + "async-compression", "base64 0.22.1", "bytes", "encoding_rs", @@ -3163,11 +3260,13 @@ "system-configuration", "tokio", "tokio-native-tls", + "tokio-util", "tower", "tower-service", "url", "wasm-bindgen", "wasm-bindgen-futures", + "wasm-streams", "web-sys", "windows-registry", ] @@ -3200,7 +3299,7 @@ dependencies = [ "cc", "cfg-if", - "getrandom", + "getrandom 0.2.15", "libc", "spin", "untrusted", @@ -3220,7 +3319,7 @@ "num-traits", "pkcs1", "pkcs8", - "rand_core", + "rand_core 0.6.4", "signature", "spki", "subtle", @@ -3571,7 +3670,7 @@ checksum = "77549399552de45a898a580c1b41d445bf730df867cc44e6c0233bbc4b8329de" dependencies = [ "digest", - "rand_core", + "rand_core 0.6.4", ] [[package]] @@ -3644,9 +3743,6 @@ version = "0.9.8" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" -dependencies = [ - "lock_api", -] [[package]] name = "spki" @@ -3747,7 +3843,7 @@ dependencies = [ "cfg-if", "fastrand", - "getrandom", + "getrandom 0.2.15", "once_cell", "rustix", "windows-sys 0.59.0", @@ -3917,7 +4013,7 @@ "pin-project-lite", "postgres-protocol", "postgres-types", - "rand", + "rand 0.8.5", "socket2", "tokio", "tokio-util", @@ -3962,9 +4058,9 @@ [[package]] name = "tokio-util" -version = "0.7.13" +version = "0.7.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d7fcaa8d55a2bdd6b83ace262b016eca0d79ee02818c5c1bcdf0305114081078" +checksum = "66a539a9ad6d5d281510d5bd368c973d636c02dbf8a67300bfb6b950696ad7df" dependencies = [ "bytes", "futures-core", @@ -4186,7 +4282,7 @@ "httparse", "log", "native-tls", - "rand", + "rand 0.8.5", "sha1", "thiserror 2.0.12", "utf-8", @@ -4337,6 +4433,15 @@ checksum = "9c8d87e72b64a3b4db28d11ce29237c246188f4f51057d65a7eab63b7987e423" [[package]] +name = "wasi" +version = "0.14.2+wasi-0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9683f9a5a998d873c0d21fcbe3c083009670149a8fab228644b8bd36b2c48cb3" +dependencies = [ + "wit-bindgen-rt", +] + +[[package]] name = "wasite" version = "0.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" @@ -4411,6 +4516,19 @@ checksum = "1a05d73b933a847d6cccdda8f838a22ff101ad9bf93e33684f39c1f5f0eece3d" dependencies = [ "unicode-ident", +] + +[[package]] +name = "wasm-streams" +version = "0.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "15053d8d85c7eccdbefef60f06769760a563c7f0a9d6902a13d35c7800b0ad65" +dependencies = [ + "futures-util", + "js-sys", + "wasm-bindgen", + "wasm-bindgen-futures", + "web-sys", ] [[package]] @@ -4687,6 +4805,15 @@ dependencies = [ "cfg-if", "windows-sys 0.48.0", +] + +[[package]] +name = "wit-bindgen-rt" +version = "0.39.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6f42320e61fe2cfd34354ecb597f86f413484a798ba44a8ca1165c58d42da6c1" +dependencies = [ + "bitflags 2.8.0", ] [[package]] diff --git a/consumer/Cargo.toml b/consumer/Cargo.toml --- a/consumer/Cargo.toml +++ b/consumer/Cargo.toml @@ -7,12 +7,10 @@ chrono = { version = "0.4.39", features = ["serde"] } ciborium = "0.2.2" clap = { version = "4.5.34", features = ["derive"] } +deadpool-postgres = { version = "0.14.1", features = ["serde"] } did-resolver = { path = "../did-resolver" } -diesel = { version = "2.2.6", features = ["chrono", "serde_json"] } -diesel-async = { version = "0.5.2", features = ["deadpool", "postgres"] } eyre = "0.6.12" figment = { version = "0.10.19", features = ["env", "toml"] } -flume = { version = "0.11.1", features = ["async"] } foldhash = "0.1.4" futures = "0.3.31" ipld-core = "0.4.1" @@ -20,10 +18,10 @@ lexica = { path = "../lexica" } metrics = "0.24.1" metrics-exporter-prometheus = "0.16.2" -parakeet-db = { path = "../parakeet-db" } +parakeet-db = { path = "../parakeet-db", default-features = false, features = ["postgres"] } parakeet-index = { path = "../parakeet-index" } redis = { version = "0.31", features = ["tokio-native-tls-comp"] } -reqwest = { version = "0.12.12", features = ["native-tls"] } +reqwest = { version = "0.12.12", features = ["native-tls", "brotli", "stream"] } serde = { version = "1.0.217", features = ["derive"] } serde_bytes = "0.11" serde_ipld_dagcbor = "0.6.1" @@ -31,8 +29,9 @@ sled = "0.34.7" thiserror = "2" tokio = { version = "1.42.0", features = ["full"] } -tokio-postgres = { version = "0.7.12", features = ["with-chrono-0_4"] } +tokio-postgres = { version = "0.7.12", features = ["with-chrono-0_4", "with-serde_json-1"] } tokio-stream = "0.1.17" tokio-tungstenite = { version = "0.26.1", features = ["native-tls"] } +tokio-util = { version = "0.7.14", features = ["io"] } tracing = "0.1.40" tracing-subscriber = "0.3.18" diff --git a/parakeet-db/Cargo.toml b/parakeet-db/Cargo.toml --- a/parakeet-db/Cargo.toml +++ b/parakeet-db/Cargo.toml @@ -5,5 +5,11 @@ [dependencies] chrono = { version = "0.4.39", features = ["serde"] } -diesel = { version = "2.2.6", features = ["chrono", "serde_json"] } -serde_json = "1.0.134" \ No newline at end of file +diesel = { version = "2.2.6", features = ["chrono", "serde_json"], optional = true } +postgres-types = { version = "0.2.9", optional = true } +serde_json = "1.0.134" + +[features] +default = ["diesel"] +diesel = ["dep:diesel"] +postgres = ["dep:postgres-types"] \ No newline at end of file diff --git a/consumer/src/config.rs b/consumer/src/config.rs --- a/consumer/src/config.rs +++ b/consumer/src/config.rs @@ -14,13 +14,11 @@ #[derive(Debug, Deserialize)] pub struct Config { pub index_uri: String, - pub database_url: String, + pub database: deadpool_postgres::Config, pub redis_uri: String, pub plc_directory: Option, /// Adds contact details (email / bluesky handle / website) to the UA header. pub ua_contact: Option, - #[serde(default = "default_backfill_workers")] - pub backfill_workers: u8, /// DIDs of label services to force subscription to. #[serde(default)] pub initial_label_services: Vec, @@ -29,6 +27,8 @@ /// Configuration items specific to indexer pub indexer: Option, + /// Configuration items specific to backfill + pub backfill: Option, } #[derive(Debug, Deserialize)] @@ -51,6 +51,16 @@ BackfillHistory, /// Discover new accounts as they come and do not import history Realtime, +} + +#[derive(Clone, Debug, Deserialize)] +pub struct BackfillConfig { + #[serde(default = "default_backfill_workers")] + pub backfill_workers: u8, + #[serde(default)] + pub skip_aggregation: bool, + #[serde(default)] + pub skip_handle_validation: bool, } fn default_backfill_workers() -> u8 { diff --git a/consumer/src/main.rs b/consumer/src/main.rs --- a/consumer/src/main.rs +++ b/consumer/src/main.rs @@ -1,14 +1,14 @@ +use deadpool_postgres::Runtime; use did_resolver::{Resolver, ResolverOpts}; -use diesel_async::pooled_connection::deadpool::Pool; -use diesel_async::pooled_connection::AsyncDieselConnectionManager; -use diesel_async::AsyncPgConnection; use eyre::OptionExt; use metrics_exporter_prometheus::PrometheusBuilder; use std::sync::Arc; +use tokio_postgres::NoTls; mod backfill; mod cmd; mod config; +mod db; mod firehose; mod indexer; mod label_indexer; @@ -24,8 +24,7 @@ let user_agent = build_ua(&conf.ua_contact); - let db_mgr = AsyncDieselConnectionManager::::new(&conf.database_url); - let pool = Pool::builder(db_mgr).build()?; + let pool = conf.database.create_pool(Some(Runtime::Tokio1), NoTls)?; let (redis_conn, redis_fut) = redis::Client::open(conf.redis_uri)? .create_multiplexed_tokio_connection() @@ -61,7 +60,7 @@ let resume = resume.clone().unwrap(); let (label_mgr, _label_svc_tx) = label_indexer::LabelServiceManager::new( - &conf.database_url, + pool.clone(), resolver.clone(), resume, user_agent.clone(), @@ -72,12 +71,16 @@ } if cli.backfill { + let bf_cfg = conf + .backfill + .ok_or_eyre("Config item [backfill] must be specified when using --backfill")?; + let backfiller = backfill::BackfillManager::new( pool.clone(), redis_conn.clone(), resolver.clone(), - index_client.clone(), - conf.backfill_workers, + (!bf_cfg.skip_aggregation).then_some(index_client.clone()), + bf_cfg, ) .await?; diff --git a/migrations/2025-01-29-213341_follows_and_blocks/up.sql b/migrations/2025-01-29-213341_follows_and_blocks/up.sql --- a/migrations/2025-01-29-213341_follows_and_blocks/up.sql +++ b/migrations/2025-01-29-213341_follows_and_blocks/up.sql @@ -6,8 +6,8 @@ created_at timestamptz not null ); -create index blocks_did_index on blocks using hash (did); -create index blocks_subject_index on blocks using hash (subject); +create index blocks_did_index on blocks (did); +create index blocks_subject_index on blocks (subject); create table follows ( @@ -17,5 +17,5 @@ created_at timestamptz not null ); -create index follow_did_index on follows using hash (did); -create index follow_subject_index on follows using hash (subject); \ No newline at end of file +create index follow_did_index on follows (did); +create index follow_subject_index on follows (subject); \ No newline at end of file diff --git a/migrations/2025-02-07-203450_lists/up.sql b/migrations/2025-02-07-203450_lists/up.sql --- a/migrations/2025-02-07-203450_lists/up.sql +++ b/migrations/2025-02-07-203450_lists/up.sql @@ -25,8 +25,8 @@ indexed_at timestamp not null default now() ); -create index listitems_list_index on list_items using hash (list_uri); -create index listitems_subject_index on list_items using hash (subject); +create index listitems_list_index on list_items (list_uri); +create index listitems_subject_index on list_items (subject); create table list_blocks ( diff --git a/migrations/2025-02-16-142357_posts/up.sql b/migrations/2025-02-16-142357_posts/up.sql --- a/migrations/2025-02-16-142357_posts/up.sql +++ b/migrations/2025-02-16-142357_posts/up.sql @@ -22,15 +22,15 @@ indexed_at timestamp not null default now() ); -create index posts_did_index on posts using hash (did); -create index posts_parent_index on posts using hash (parent_uri); -create index posts_root_index on posts using hash (root_uri); +create index posts_did_index on posts (did); +create index posts_parent_index on posts (parent_uri); +create index posts_root_index on posts (root_uri); create index posts_lang_index on posts using gin (languages); create index posts_tags_index on posts using gin (tags); create table post_embed_images ( - post_uri text not null references posts (at_uri) on delete cascade, + post_uri text not null references posts (at_uri) on delete cascade deferrable, seq smallint not null, mime_type text not null, @@ -47,7 +47,7 @@ create table post_embed_video ( - post_uri text primary key references posts (at_uri) on delete cascade, + post_uri text primary key references posts (at_uri) on delete cascade deferrable, mime_type text not null, cid text not null, @@ -61,7 +61,7 @@ create table post_embed_video_captions ( - post_uri text not null references posts (at_uri) on delete cascade, + post_uri text not null references posts (at_uri) on delete cascade deferrable, language text not null, mime_type text not null, @@ -74,7 +74,7 @@ create table post_embed_ext ( - post_uri text primary key references posts (at_uri) on delete cascade, + post_uri text primary key references posts (at_uri) on delete cascade deferrable, uri text not null, title text not null, @@ -87,7 +87,7 @@ create table post_embed_record ( - post_uri text primary key references posts (at_uri) on delete cascade, + post_uri text primary key references posts (at_uri) on delete cascade deferrable, record_type text not null, uri text not null, @@ -101,7 +101,7 @@ ( at_uri text primary key, cid text not null, - post_uri text not null references posts (at_uri) on delete cascade, + post_uri text not null references posts (at_uri) on delete cascade deferrable, detached text[] not null, rules text[] not null, @@ -118,7 +118,7 @@ ( at_uri text primary key, cid text not null, - post_uri text not null references posts (at_uri) on delete cascade, + post_uri text not null references posts (at_uri) on delete cascade deferrable, hidden_replies text[] not null, allow text[] not null, diff --git a/migrations/2025-04-05-114428_likes_and_reposts/up.sql b/migrations/2025-04-05-114428_likes_and_reposts/up.sql --- a/migrations/2025-04-05-114428_likes_and_reposts/up.sql +++ b/migrations/2025-04-05-114428_likes_and_reposts/up.sql @@ -8,8 +8,8 @@ indexed_at timestamp not null default now() ); -create index likes_did_index on likes using hash (did); -create index likes_subject_index on likes using hash (subject); +create index likes_did_index on likes (did); +create index likes_subject_index on likes (subject); create table reposts ( @@ -21,5 +21,5 @@ indexed_at timestamp not null default now() ); -create index reposts_did_index on reposts using hash (did); -create index reposts_post_index on reposts using hash (post); \ No newline at end of file +create index reposts_did_index on reposts (did); +create index reposts_post_index on reposts (post); \ No newline at end of file diff --git a/migrations/2025-04-18-185717_verification/up.sql b/migrations/2025-04-18-185717_verification/up.sql --- a/migrations/2025-04-18-185717_verification/up.sql +++ b/migrations/2025-04-18-185717_verification/up.sql @@ -12,5 +12,5 @@ indexed_at timestamp not null default now() ); -create index verification_verifier_index on verification using hash (verifier); -create index verification_subject_index on verification using hash (subject); \ No newline at end of file +create index verification_verifier_index on verification (verifier); +create index verification_subject_index on verification (subject); \ No newline at end of file diff --git a/parakeet-db/src/lib.rs b/parakeet-db/src/lib.rs --- a/parakeet-db/src/lib.rs +++ b/parakeet-db/src/lib.rs @@ -1,3 +1,5 @@ +#[cfg(feature = "diesel")] pub mod models; +#[cfg(feature = "diesel")] pub mod schema; pub mod types; diff --git a/parakeet-db/src/types.rs b/parakeet-db/src/types.rs --- a/parakeet-db/src/types.rs +++ b/parakeet-db/src/types.rs @@ -1,81 +1,112 @@ -use diesel::backend::Backend; -use diesel::deserialize::FromSql; -use diesel::pg::Pg; -use diesel::serialize::{Output, ToSql}; -use diesel::{AsExpression, FromSqlRow}; - -#[derive(Debug, PartialOrd, PartialEq, AsExpression, FromSqlRow)] -#[diesel(sql_type = diesel::sql_types::Text)] -pub enum ActorStatus { - Active, - Takendown, - Suspended, - Deleted, - Deactivated, -} - -impl FromSql for ActorStatus -where - DB: Backend, - String: FromSql, -{ - fn from_sql(bytes: DB::RawValue<'_>) -> diesel::deserialize::Result { - match String::from_sql(bytes)?.as_str() { - "active" => Ok(ActorStatus::Active), - "takendown" => Ok(ActorStatus::Takendown), - "suspended" => Ok(ActorStatus::Suspended), - "deleted" => Ok(ActorStatus::Deleted), - "deactivated" => Ok(ActorStatus::Deactivated), - x => Err(format!("Unrecognized variant {}", x).into()), +macro_rules! text_enum { + (enum $name:ident {$($variant:ident = $value:expr,)*}) => { + #[derive(Debug, PartialOrd, PartialEq)] + #[cfg_attr(feature = "diesel", derive(diesel::AsExpression, diesel::FromSqlRow))] + #[cfg_attr(feature = "diesel", diesel(sql_type = diesel::sql_types::Text))] + pub enum $name { + $($variant,)* } - } -} -impl ToSql for ActorStatus { - fn to_sql<'b>(&'b self, out: &mut Output<'b, '_, Pg>) -> diesel::serialize::Result { - let val = match self { - ActorStatus::Active => "active", - ActorStatus::Takendown => "takendown", - ActorStatus::Suspended => "suspended", - ActorStatus::Deleted => "deleted", - ActorStatus::Deactivated => "deactivated", - }; - - >::to_sql(val, out) - } -} - -#[derive(Debug, PartialOrd, PartialEq, AsExpression, FromSqlRow)] -#[diesel(sql_type = diesel::sql_types::Text)] -pub enum ActorSyncState { - Synced, - Dirty, - Processing, -} - -impl FromSql for ActorSyncState -where - DB: Backend, - String: FromSql, -{ - fn from_sql(bytes: DB::RawValue<'_>) -> diesel::deserialize::Result { - match String::from_sql(bytes)?.as_str() { - "synced" => Ok(ActorSyncState::Synced), - "dirty" => Ok(ActorSyncState::Dirty), - "processing" => Ok(ActorSyncState::Processing), - x => Err(format!("Unrecognized variant {}", x).into()), + impl std::fmt::Display for $name { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + $(Self::$variant => write!(f, $value),)* + } + } } - } + + impl std::str::FromStr for $name { + type Err = String; + + fn from_str(s: &str) -> Result { + match s { + $($value => Ok(Self::$variant),)* + x => Err(format!("Unrecognized variant {}", x).into()), + } + } + } + + #[cfg(feature = "postgres")] + impl postgres_types::FromSql<'_> for $name { + fn from_sql( + ty: &postgres_types::Type, + raw: &[u8] + ) -> Result> { + Ok(String::from_sql(ty, raw)?.parse()?) + } + + fn accepts(ty: &postgres_types::Type) -> bool { + ty == &postgres_types::Type::TEXT + } + } + + #[cfg(feature = "postgres")] + impl postgres_types::ToSql for $name { + fn to_sql( + &self, + ty: &postgres_types::Type, + out: &mut postgres_types::private::BytesMut + ) -> Result> + where + Self: Sized, + { + self.to_string().to_sql(ty, out) + } + + fn accepts(ty: &postgres_types::Type) -> bool + where + Self: Sized, + { + ty == &postgres_types::Type::TEXT + } + + postgres_types::to_sql_checked!(); + } + + #[cfg(feature = "diesel")] + impl diesel::deserialize::FromSql for $name + where + DB: diesel::backend::Backend, + String: diesel::deserialize::FromSql, + { + fn from_sql(bytes: DB::RawValue<'_>) -> diesel::deserialize::Result { + use std::str::FromStr; + + let st = >::from_sql(bytes)?; + + let out = Self::from_str(&st)?; + Ok(out) + } + } + + #[cfg(feature = "diesel")] + impl diesel::serialize::ToSql for $name { + fn to_sql<'b>(&'b self, out: &mut diesel::serialize::Output<'b, '_, diesel::pg::Pg>) -> diesel::serialize::Result { + use std::io::Write; + let val = self.to_string(); + + out.write(val.as_bytes())?; + Ok(diesel::serialize::IsNull::No) + } + } + + }; } -impl ToSql for ActorSyncState { - fn to_sql<'b>(&'b self, out: &mut Output<'b, '_, Pg>) -> diesel::serialize::Result { - let val = match self { - ActorSyncState::Synced => "synced", - ActorSyncState::Dirty => "dirty", - ActorSyncState::Processing => "processing", - }; - - >::to_sql(val, out) +text_enum!( + enum ActorStatus { + Active = "active", + Takendown = "takendown", + Suspended = "suspended", + Deleted = "deleted", + Deactivated = "deactivated", } -} +); + +text_enum!( + enum ActorSyncState { + Synced = "synced", + Dirty = "dirty", + Processing = "processing", + } +); diff --git a/consumer/src/backfill/db.rs b/consumer/src/backfill/db.rs deleted file mode 100644 --- a/consumer/src/backfill/db.rs +++ /dev/null @@ -1,99 +0,0 @@ -use diesel::prelude::*; -use diesel_async::{AsyncPgConnection, RunQueryDsl}; -use parakeet_db::{models, schema, types}; - -pub async fn write_backfill_job( - conn: &mut AsyncPgConnection, - repo: &str, - status: &str, -) -> QueryResult { - diesel::insert_into(schema::backfill_jobs::table) - .values(( - schema::backfill_jobs::did.eq(repo), - schema::backfill_jobs::status.eq(status), - )) - .on_conflict_do_nothing() - .execute(conn) - .await -} - -pub async fn get_actor_status( - conn: &mut AsyncPgConnection, - did: &str, -) -> QueryResult<(types::ActorStatus, types::ActorSyncState)> { - schema::actors::table - .select((schema::actors::status, schema::actors::sync_state)) - .find(&did) - .get_result(conn) - .await -} - -pub async fn update_repo_sync_state( - conn: &mut AsyncPgConnection, - did: &str, - sync_state: types::ActorSyncState, -) -> QueryResult { - diesel::update(schema::actors::table) - .set(schema::actors::sync_state.eq(sync_state)) - .filter(schema::actors::did.eq(did)) - .execute(conn) - .await -} - -pub async fn update_handle( - conn: &mut AsyncPgConnection, - did: &str, - handle: Option, -) -> QueryResult { - diesel::update(schema::actors::table) - .set(schema::actors::handle.eq(handle)) - .filter(schema::actors::did.eq(did)) - .execute(conn) - .await -} - -pub async fn update_actor_status( - conn: &mut AsyncPgConnection, - did: &str, - status: types::ActorStatus, - sync_state: types::ActorSyncState, -) -> QueryResult { - diesel::update(schema::actors::table) - .set(( - schema::actors::status.eq(status), - schema::actors::sync_state.eq(sync_state), - )) - .filter(schema::actors::did.eq(did)) - .execute(conn) - .await -} - -pub async fn defer(conn: &mut AsyncPgConnection) -> QueryResult { - diesel::sql_query("SET CONSTRAINTS ALL DEFERRED") - .execute(conn) - .await -} - -pub async fn pull_backfill_rows( - conn: &mut AsyncPgConnection, - repo: &str, - rev: &str, -) -> QueryResult> { - schema::backfill::table - .select(models::BackfillRow::as_select()) - .filter( - schema::backfill::repo - .eq(repo) - .and(schema::backfill::repo_ver.gt(rev)), - ) - .order(schema::backfill::repo_ver) - .load(conn) - .await -} - -pub async fn clear_backfill_rows(conn: &mut AsyncPgConnection, repo: &str) -> QueryResult { - diesel::delete(schema::backfill::table) - .filter(schema::backfill::repo.eq(repo)) - .execute(conn) - .await -} diff --git a/consumer/src/backfill/mod.rs b/consumer/src/backfill/mod.rs --- a/consumer/src/backfill/mod.rs +++ b/consumer/src/backfill/mod.rs @@ -1,21 +1,21 @@ +use crate::config::BackfillConfig; +use crate::db; use crate::indexer::types::{AggregateDeltaStore, BackfillItem, BackfillItemInner}; -use crate::indexer::{self, db as indexer_db}; +use crate::indexer::{self, records}; +use chrono::prelude::*; +use deadpool_postgres::{Object, Pool, Transaction}; use did_resolver::Resolver; -use diesel_async::pooled_connection::deadpool::Pool; -use diesel_async::AsyncPgConnection; use ipld_core::cid::Cid; use metrics::counter; use parakeet_db::types::{ActorStatus, ActorSyncState}; use redis::aio::MultiplexedConnection; use redis::{AsyncCommands, Direction}; use reqwest::{Client, StatusCode}; -use std::collections::HashMap; use std::str::FromStr; use std::sync::Arc; use tokio::sync::Semaphore; use tracing::instrument; -mod db; mod repo; mod types; @@ -28,11 +28,12 @@ pub struct BackfillManagerInner { resolver: Arc, client: Client, - index_client: parakeet_index::Client, + index_client: Option, + opts: BackfillConfig, } pub struct BackfillManager { - pool: Pool, + pool: Pool, redis: MultiplexedConnection, semaphore: Arc, inner: BackfillManagerInner, @@ -40,14 +41,14 @@ impl BackfillManager { pub async fn new( - pool: Pool, + pool: Pool, redis: MultiplexedConnection, resolver: Arc, - index_client: parakeet_index::Client, - threads: u8, + index_client: Option, + opts: BackfillConfig, ) -> eyre::Result { - let client = Client::new(); - let semaphore = Arc::new(Semaphore::new(threads as usize)); + let client = Client::builder().brotli(true).build()?; + let semaphore = Arc::new(Semaphore::new(opts.backfill_workers as usize)); Ok(BackfillManager { pool, @@ -57,6 +58,7 @@ resolver, client, index_client, + opts, }, }) } @@ -93,13 +95,13 @@ tracing::error!(did = &job, "backfill failed: {e}"); counter!("backfill_failure").increment(1); - db::write_backfill_job(&mut conn, &job, "failed") + db::backfill_job_write(&mut conn, &job, "failed") .await .unwrap(); } else { counter!("backfill_success").increment(1); - db::write_backfill_job(&mut conn, &job, "successful") + db::backfill_job_write(&mut conn, &job, "successful") .await .unwrap(); } @@ -119,11 +121,14 @@ #[instrument(skip(conn, inner))] async fn backfill_actor( - conn: &mut AsyncPgConnection, + conn: &mut Object, inner: &mut BackfillManagerInner, did: &str, ) -> eyre::Result<()> { - let (status, sync_state) = db::get_actor_status(conn, did).await?; + let Some((status, sync_state)) = db::actor_get_statuses(conn, did).await? else { + tracing::error!("skipping backfill on unknown repo"); + return Ok(()); + }; if sync_state != ActorSyncState::Dirty || status != ActorStatus::Active { tracing::debug!("skipping non-dirty or inactive repo"); @@ -135,13 +140,6 @@ eyre::bail!("missing did doc"); }; - let Some(handle) = did_doc - .also_known_as - .and_then(|aka| aka.first().cloned()) - .and_then(|handle| handle.strip_prefix("at://").map(String::from)) - else { - eyre::bail!("DID doc contained no handle"); - }; let Some(service) = did_doc .service .and_then(|services| services.into_iter().find(|svc| svc.id == PDS_SERVICE_ID)) @@ -156,7 +154,14 @@ let Some(repo_status) = check_pds_repo_status(&inner.client, &pds_url, did).await? else { // this repo can't be found - set dirty and assume deleted. tracing::debug!("repo was deleted"); - db::update_actor_status(conn, did, ActorStatus::Deleted, ActorSyncState::Dirty).await?; + db::actor_upsert( + conn, + did, + ActorStatus::Deleted, + ActorSyncState::Dirty, + Utc::now(), + ) + .await?; return Ok(()); }; @@ -165,103 +170,120 @@ let status = repo_status .status .unwrap_or(crate::firehose::AtpAccountStatus::Deleted); - db::update_actor_status(conn, did, status.into(), ActorSyncState::Dirty).await?; + db::actor_upsert( + conn, + did, + status.into(), + ActorSyncState::Dirty, + Utc::now(), + ) + .await?; return Ok(()); } - // at this point, the account will be active and we can attempt to resolve the handle. + if !inner.opts.skip_handle_validation { + // at this point, the account will be active and we can attempt to resolve the handle. + let Some(handle) = did_doc + .also_known_as + .and_then(|aka| aka.first().cloned()) + .and_then(|handle| handle.strip_prefix("at://").map(String::from)) + else { + eyre::bail!("DID doc contained no handle"); + }; - // in theory, we can use com.atproto.identity.resolveHandle against a PDS, but that seems - // like a way to end up with really sus handles. - let Some(handle_did) = inner.resolver.resolve_handle(&handle).await? else { - eyre::bail!("Failed to resolve did for handle {handle}"); - }; + // in theory, we can use com.atproto.identity.resolveHandle against a PDS, but that seems + // like a way to end up with really sus handles. + let Some(handle_did) = inner.resolver.resolve_handle(&handle).await? else { + eyre::bail!("Failed to resolve did for handle {handle}"); + }; - if handle_did != did { - eyre::bail!("requested DID doesn't match handle"); + if handle_did != did { + eyre::bail!("requested DID doesn't match handle"); + } + + // set the handle from above + db::actor_upsert_handle( + conn, + did, + ActorSyncState::Processing, + Some(handle), + Utc::now(), + ) + .await?; } - // set the handle from above - db::update_handle(conn, did, Some(handle)).await?; - // now we can start actually backfilling - db::update_repo_sync_state(conn, did, ActorSyncState::Processing).await?; + db::actor_set_sync_status(conn, did, ActorSyncState::Processing, Utc::now()).await?; + + let mut t = conn.transaction().await?; + t.execute("SET CONSTRAINTS ALL DEFERRED", &[]).await?; tracing::trace!("pulling repo"); - let (rev, cid, records) = repo::pull_repo(&inner.client, did, &pds_url).await?; + let (commit, mut deltas, copies) = + repo::stream_and_insert_repo(&mut t, &inner.client, did, &pds_url).await?; - tracing::trace!("repo pulled - inserting"); + db::actor_set_repo_state(&mut t, did, &commit.rev, commit.data).await?; - let mut delta_store = HashMap::new(); + copies.submit(&mut t, did).await?; - db::defer(conn).await?; + t.execute( + "UPDATE actors SET sync_state=$2, last_indexed=$3 WHERE did=$1", + &[&did, &ActorSyncState::Synced, &Utc::now().naive_utc()], + ) + .await?; - indexer_db::update_repo_version(conn, did, &rev, cid).await?; - - for (path, (cid, record)) in records { - let Some((collection, rkey)) = path.split_once("/") else { - tracing::warn!("record contained invalid path {}", path); - return Err(diesel::result::Error::RollbackTransaction.into()); - }; - - counter!("backfilled_commits", "collection" => collection.to_string()).increment(1); - - let full_path = format!("at://{did}/{path}"); - - indexer::index_op(conn, &mut delta_store, did, cid, record, &full_path, rkey).await? - } - - db::update_repo_sync_state(conn, did, ActorSyncState::Synced).await?; - - handle_backfill_rows(conn, &mut delta_store, did, &rev).await?; + handle_backfill_rows(&mut t, &mut deltas, did, &commit.rev).await?; tracing::trace!("insertion finished"); - // submit the deltas - let delta_store = delta_store - .into_iter() - .map(|((uri, typ), delta)| parakeet_index::AggregateDeltaReq { - typ, - uri: uri.to_string(), - delta, - }) - .collect::>(); + if let Some(index_client) = &mut inner.index_client { + // submit the deltas + let delta_store = deltas + .into_iter() + .map(|((uri, typ), delta)| parakeet_index::AggregateDeltaReq { + typ, + uri: uri.to_string(), + delta, + }) + .collect::>(); - let mut read = 0; + let mut read = 0; - while read < delta_store.len() { - let rem = delta_store.len() - read; - let take = DELTA_BATCH_SIZE.min(rem); + while read < delta_store.len() { + let rem = delta_store.len() - read; + let take = DELTA_BATCH_SIZE.min(rem); - tracing::debug!("reading & submitting {take} deltas"); + tracing::debug!("reading & submitting {take} deltas"); - let deltas = delta_store[read..read + take].to_vec(); - inner - .index_client - .submit_aggregate_delta_batch(parakeet_index::AggregateDeltaBatchReq { deltas }) - .await?; + let deltas = delta_store[read..read + take].to_vec(); + index_client + .submit_aggregate_delta_batch(parakeet_index::AggregateDeltaBatchReq { deltas }) + .await?; - read += take; - tracing::debug!("read {read} of {} deltas", delta_store.len()); + read += take; + tracing::debug!("read {read} of {} deltas", delta_store.len()); + } } + + t.commit().await?; Ok(()) } async fn handle_backfill_rows( - conn: &mut AsyncPgConnection, + conn: &mut Transaction<'_>, deltas: &mut impl AggregateDeltaStore, repo: &str, rev: &str, -) -> diesel::QueryResult<()> { +) -> Result<(), tokio_postgres::Error> { // `pull_backfill_rows` filters out anything before the last commit we pulled - let backfill_rows = db::pull_backfill_rows(conn, repo, rev).await?; + let backfill_rows = db::backfill_rows_get(conn, repo, rev).await?; for row in backfill_rows { // blindly unwrap-ing this CID as we've already parsed it and re-serialized it let repo_cid = Cid::from_str(&row.cid).unwrap(); - indexer_db::update_repo_version(conn, repo, &row.repo_ver, repo_cid).await?; + db::actor_set_repo_state(conn, repo, &row.repo_ver, repo_cid).await?; // again, we've serialized this. let items: Vec = serde_json::from_value(row.data).unwrap(); @@ -288,7 +310,7 @@ } // finally, clear the backfill table entries for this actor - db::clear_backfill_rows(conn, repo).await?; + db::backfill_delete_rows(conn, repo).await?; Ok(()) } @@ -310,4 +332,35 @@ } Ok(res.json().await?) +} + +#[derive(Debug, Default)] +struct CopyStore { + likes: Vec<(String, records::StrongRef, DateTime)>, + posts: Vec<(String, Cid, records::AppBskyFeedPost)>, + reposts: Vec<(String, records::StrongRef, DateTime)>, + blocks: Vec<(String, String, DateTime)>, + follows: Vec<(String, String, DateTime)>, + list_items: Vec<(String, records::AppBskyGraphListItem)>, + verifications: Vec<(String, Cid, records::AppBskyGraphVerification)>, + records: Vec<(String, Cid)>, +} + +impl CopyStore { + async fn submit(self, t: &mut Transaction<'_>, did: &str) -> Result<(), tokio_postgres::Error> { + db::copy::copy_likes(t, did, self.likes).await?; + db::copy::copy_posts(t, did, self.posts).await?; + db::copy::copy_reposts(t, did, self.reposts).await?; + db::copy::copy_blocks(t, did, self.blocks).await?; + db::copy::copy_follows(t, did, self.follows).await?; + db::copy::copy_list_items(t, self.list_items).await?; + db::copy::copy_verification(t, did, self.verifications).await?; + db::copy::copy_records(t, did, self.records).await?; + + Ok(()) + } + + fn push_record(&mut self, at_uri: &str, cid: Cid) { + self.records.push((at_uri.to_string(), cid)) + } } diff --git a/consumer/src/backfill/repo.rs b/consumer/src/backfill/repo.rs --- a/consumer/src/backfill/repo.rs +++ b/consumer/src/backfill/repo.rs @@ -1,55 +1,76 @@ -use super::types::{CarCommitEntry, CarEntry}; -use crate::indexer::types::RecordTypes; -use futures::{StreamExt, TryStreamExt}; +use super::{ + types::{CarCommitEntry, CarEntry}, + CopyStore, +}; +use crate::indexer::records; +use crate::indexer::types::{AggregateDeltaStore, RecordTypes}; +use crate::{db, indexer}; +use deadpool_postgres::Transaction; +use futures::TryStreamExt; use ipld_core::cid::Cid; use iroh_car::CarReader; +use metrics::counter; +use parakeet_index::AggregateType; use reqwest::Client; use std::collections::HashMap; +use std::io::ErrorKind; +use tokio::io::BufReader; +use tokio_util::io::StreamReader; -pub async fn pull_repo<'a>( +type BackfillDeltaStore = HashMap<(String, i32), i32>; + +pub async fn stream_and_insert_repo( + t: &mut Transaction<'_>, client: &Client, repo: &str, pds: &str, -) -> eyre::Result<(String, Cid, HashMap)> { +) -> eyre::Result<(CarCommitEntry, BackfillDeltaStore, CopyStore)> { let res = client .get(format!("{pds}/xrpc/com.atproto.sync.getRepo?did={repo}")) .send() .await? .error_for_status()?; - let body = res.bytes().await?; + let strm = res + .bytes_stream() + .map_err(|err| std::io::Error::new(ErrorKind::Other, err)); + let reader = StreamReader::new(strm); + let mut car_stream = CarReader::new(BufReader::new(reader)).await?; - let (commit, records) = read_car(&body).await?; - - Ok((commit.rev, commit.data, records)) -} - -// beware: this is probably: 1. insecure, 2. slow, 3. a/n other crimes -async fn read_car( - data: &[u8], -) -> eyre::Result<(CarCommitEntry, HashMap)> { - let car = CarReader::new(data).await?; - - let entries = car - .stream() - .map_ok( - |(cid, block)| match serde_ipld_dagcbor::from_slice::(&block) { - Ok(decoded) => Some((cid, decoded)), - Err(_) => None, - }, - ) - .filter_map(|v| async move { v.ok().flatten() }) - .collect::>() - .await; + // the root should be the commit block + let root = car_stream.header().roots().first().cloned().unwrap(); let mut commit = None; - let mut mst_nodes = Vec::new(); - let mut records = HashMap::new(); + let mut mst_nodes: HashMap = HashMap::new(); + let mut records: HashMap = HashMap::new(); + let mut deltas = HashMap::new(); + let mut copies = CopyStore::default(); - for (cid, entry) in entries { - match entry { - CarEntry::Record(rec) => { - records.insert(cid, rec); + while let Some((cid, block)) = car_stream.next_block().await? { + let Ok(block) = serde_ipld_dagcbor::from_slice::(&block) else { + tracing::warn!("failed to parse block {cid}"); + continue; + }; + + if root == cid { + if let CarEntry::Commit(commit_entry) = block { + commit = Some(commit_entry); + } else { + tracing::warn!("root did not point to a commit entry"); + } + continue; + } + + match block { + CarEntry::Commit(_) => { + tracing::warn!("got commit entry that was not in root") + } + CarEntry::Record(record) => { + if let Some(path) = mst_nodes.remove(&cid) { + record_index(t, &mut copies, &mut deltas, repo, &path, cid, record).await?; + } else { + records.insert(cid, record); + } } CarEntry::Mst(mst) => { let mut out = Vec::with_capacity(mst.e.len()); @@ -60,29 +81,122 @@ let key = if node.p == 0 { ks.to_string() } else { - let (prev, _): &(String, Cid) = out.last().unwrap(); + let (_, prev): &(Cid, String) = out.last().unwrap(); let prefix = &prev[..node.p as usize]; format!("{prefix}{ks}") }; - out.push((key, node.v)); + out.push((node.v, key.to_string())); } mst_nodes.extend(out); } - CarEntry::Commit(car_commit) => { - commit = Some(car_commit); - } } } - let records_out = mst_nodes - .into_iter() - .filter_map(|(key, cid)| records.remove(&cid).map(|v| (key, (cid, v)))) - .collect::>(); + for (cid, record) in records { + if let Some(path) = mst_nodes.remove(&cid) { + record_index(t, &mut copies, &mut deltas, repo, &path, cid, record).await?; + } else { + tracing::warn!("couldn't find MST node for record {cid}") + } + } - let commit = commit.ok_or(eyre::eyre!("no commit found"))?; + let commit = commit.unwrap(); - Ok((commit, records_out)) + Ok((commit, deltas, copies)) +} + +async fn record_index( + t: &mut Transaction<'_>, + copies: &mut CopyStore, + deltas: &mut BackfillDeltaStore, + did: &str, + path: &str, + cid: Cid, + record: RecordTypes, +) -> eyre::Result<()> { + let Some((collection_raw, rkey)) = path.split_once("/") else { + tracing::warn!("op contained invalid path {path}"); + return Ok(()); + }; + + counter!("backfilled_commits", "collection" => collection_raw.to_string()).increment(1); + + let at_uri = format!("at://{did}/{path}"); + + match record { + RecordTypes::AppBskyFeedLike(rec) => { + deltas.incr(&rec.subject.uri, AggregateType::Like).await; + + copies.push_record(&at_uri, cid); + copies.likes.push((at_uri, rec.subject, rec.created_at)); + } + RecordTypes::AppBskyFeedPost(rec) => { + let maybe_reply = rec.reply.as_ref().map(|v| v.parent.uri.clone()); + let maybe_embed = rec + .embed + .as_ref() + .and_then(|v| v.as_bsky()) + .and_then(|v| match v { + records::AppBskyEmbed::Record(r) => Some(r.record.uri.clone()), + records::AppBskyEmbed::RecordWithMedia(r) => Some(r.record.record.uri.clone()), + _ => None, + }); + + if let Some(labels) = rec.labels.clone() { + db::maintain_self_labels(t, did, Some(cid), &at_uri, labels).await?; + } + if let Some(embed) = rec.embed.clone().and_then(|embed| embed.into_bsky()) { + db::post_embed_insert(t, &at_uri, embed, rec.created_at).await?; + } + + deltas.incr(did, AggregateType::ProfilePost).await; + if let Some(reply) = maybe_reply { + deltas.incr(&reply, AggregateType::Reply).await; + } + if let Some(embed) = maybe_embed { + deltas.incr(&embed, AggregateType::Embed).await; + } + + copies.push_record(&at_uri, cid); + copies.posts.push((at_uri, cid, rec)); + } + RecordTypes::AppBskyFeedRepost(rec) => { + deltas.incr(&rec.subject.uri, AggregateType::Repost).await; + + copies.push_record(&at_uri, cid); + copies.reposts.push((at_uri, rec.subject, rec.created_at)); + } + RecordTypes::AppBskyGraphBlock(rec) => { + copies.push_record(&at_uri, cid); + copies.blocks.push((at_uri, rec.subject, rec.created_at)); + } + RecordTypes::AppBskyGraphFollow(rec) => { + deltas.incr(did, AggregateType::Follow).await; + deltas.incr(&rec.subject, AggregateType::Follower).await; + + copies.push_record(&at_uri, cid); + copies.follows.push((at_uri, rec.subject, rec.created_at)); + } + RecordTypes::AppBskyGraphListItem(rec) => { + let split_aturi = rec.list.rsplitn(4, '/').collect::>(); + if did != split_aturi[2] { + // it's also probably a bad idea to log *all* the attempts to do this... + tracing::warn!("tried to create a listitem on a list we don't control!"); + return Ok(()); + } + + copies.push_record(&at_uri, cid); + copies.list_items.push((at_uri, rec)); + } + RecordTypes::AppBskyGraphVerification(rec) => { + copies.push_record(&at_uri, cid); + copies.verifications.push((at_uri, cid, rec)); + } + _ => indexer::index_op(t, deltas, did, cid, record, &at_uri, rkey).await?, + } + + Ok(()) } diff --git a/consumer/src/db/actor.rs b/consumer/src/db/actor.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/actor.rs @@ -0,0 +1,101 @@ +use super::{PgExecResult, PgOptResult}; +use chrono::{DateTime, Utc}; +use deadpool_postgres::GenericClient; +use ipld_core::cid::Cid; +use parakeet_db::types::{ActorStatus, ActorSyncState}; + +pub async fn actor_upsert( + conn: &mut C, + did: &str, + status: ActorStatus, + sync_state: ActorSyncState, + time: DateTime, +) -> PgExecResult { + conn.execute( + r#"INSERT INTO actors (did, status, sync_state, last_indexed) VALUES ($1, $2, $3, $4) + ON CONFLICT (did) DO UPDATE SET status=EXCLUDED.status, last_indexed=EXCLUDED.last_indexed"#, + &[&did, &status, &sync_state, &time.naive_utc()], + ).await +} + +pub async fn actor_upsert_handle( + conn: &mut C, + did: &str, + sync_state: ActorSyncState, + handle: Option, + time: DateTime, +) -> PgExecResult { + conn.execute( + r#"INSERT INTO actors (did, handle, sync_state, last_indexed) VALUES ($1, $2, $3, $4) + ON CONFLICT (did) DO UPDATE SET handle=EXCLUDED.handle, last_indexed=EXCLUDED.last_indexed"#, + &[&did, &handle, &sync_state, &time.naive_utc()] + ).await +} + +pub async fn actor_set_sync_status( + conn: &mut C, + did: &str, + sync_state: ActorSyncState, + time: DateTime, +) -> PgExecResult { + conn.execute( + "UPDATE actors SET sync_state=$2, last_indexed=$3 WHERE did=$1", + &[&did, &sync_state, &time.naive_utc()], + ) + .await +} + +pub async fn actor_set_repo_state( + conn: &mut C, + did: &str, + rev: &str, + cid: Cid, +) -> PgExecResult { + conn.execute( + "UPDATE actors SET repo_rev=$2, repo_cid=$3 WHERE did=$1", + &[&did, &rev, &cid.to_string()], + ) + .await +} + +pub async fn actor_get_status_and_rev( + conn: &mut C, + did: &str, +) -> PgOptResult<(ActorStatus, Option)> { + let res = conn + .query_opt( + "SELECT status, repo_rev FROM actors WHERE did=$1 LIMIT 1", + &[&did], + ) + .await?; + + Ok(res.map(|v| (v.get(0), v.get(1)))) +} + +pub async fn actor_get_repo_status( + conn: &mut C, + did: &str, +) -> PgOptResult<(ActorSyncState, Option)> { + let res = conn + .query_opt( + "SELECT sync_state, repo_rev FROM actors WHERE did=$1 LIMIT 1", + &[&did], + ) + .await?; + + Ok(res.map(|v| (v.get(0), v.get(1)))) +} + +pub async fn actor_get_statuses( + conn: &mut C, + did: &str, +) -> PgOptResult<(ActorStatus, ActorSyncState)> { + let res = conn + .query_opt( + "SELECT status, sync_state FROM actors WHERE did=$1 LIMIT 1", + &[&did], + ) + .await?; + + Ok(res.map(|v| (v.get(0), v.get(1)))) +} diff --git a/consumer/src/db/backfill.rs b/consumer/src/db/backfill.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/backfill.rs @@ -0,0 +1,65 @@ +use super::{PgExecResult, PgResult}; +use chrono::NaiveDateTime; +use deadpool_postgres::GenericClient; +use ipld_core::cid::Cid; + +pub struct BackfillRow { + pub repo: String, + pub repo_ver: String, + pub cid: String, + + pub data: serde_json::Value, + + pub indexed_at: NaiveDateTime, +} + +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], + ) + .await +} + +pub async fn backfill_write_row( + conn: &mut C, + repo: &str, + rev: &str, + cid: Cid, + data: serde_json::Value, +) -> PgExecResult { + conn.execute( + "INSERT INTO backfill (repo, repo_ver, cid, data) VALUES ($1, $2, $3, $4)", + &[&repo, &rev, &cid.to_string(), &data], + ) + .await +} + +pub async fn backfill_rows_get( + conn: &mut C, + repo: &str, + rev: &str, +) -> PgResult> { + let res = conn + .query( + "SELECT * FROM backfill WHERE repo=$1 AND repo_ver > $2 ORDER BY repo_ver", + &[&repo, &rev], + ) + .await?; + + Ok(res + .into_iter() + .map(|row| BackfillRow { + repo: row.get(0), + repo_ver: row.get(1), + cid: row.get(2), + data: row.get(3), + indexed_at: row.get(4), + }) + .collect()) +} + +pub async fn backfill_delete_rows(conn: &mut C, repo: &str) -> PgExecResult { + conn.execute("DELETE FROM backfill WHERE repo=$1", &[&repo]) + .await +} diff --git a/consumer/src/db/copy.rs b/consumer/src/db/copy.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/copy.rs @@ -0,0 +1,385 @@ +use super::PgExecResult; +use crate::indexer::records; +use crate::utils::strongref_to_parts; +use chrono::prelude::*; +use deadpool_postgres::Transaction; +use futures::pin_mut; +use ipld_core::cid::Cid; +use tokio_postgres::binary_copy::BinaryCopyInWriter; +use tokio_postgres::types::Type; + +// StrongRefs are used in both likes and reposts +const STRONGREF_TYPES: &[Type] = &[ + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TIMESTAMP, +]; +type StrongRefRow = (String, records::StrongRef, DateTime); + +// SubjectRefs are used in both blocks and follows +const SUBJECT_TYPES: &[Type] = &[Type::TEXT, Type::TEXT, Type::TEXT, Type::TIMESTAMP]; +type SubjectRefRow = (String, String, DateTime); + +pub async fn copy_likes( + conn: &mut Transaction<'_>, + did: &str, + data: Vec, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE likes_tmp (LIKE likes INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in( + "COPY likes_tmp (at_uri, did, subject, subject_cid, created_at) FROM STDIN (FORMAT binary)", + ) + .await?; + let writer = BinaryCopyInWriter::new(writer, STRONGREF_TYPES); + + pin_mut!(writer); + + for row in data { + let writer = writer.as_mut(); + writer + .write(&[ + &row.0, + &did, + &row.1.uri, + &row.1.cid.to_string(), + &row.2.naive_utc(), + ]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO likes (SELECT * FROM likes_tmp)", &[]) + .await +} + +pub async fn copy_reposts( + conn: &mut Transaction<'_>, + did: &str, + data: Vec, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE reposts_tmp (LIKE reposts INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in( + "COPY reposts_tmp (at_uri, did, post, post_cid, created_at) FROM STDIN (FORMAT binary)", + ) + .await?; + let writer = BinaryCopyInWriter::new(writer, STRONGREF_TYPES); + + pin_mut!(writer); + + for row in data { + let writer = writer.as_mut(); + writer + .write(&[ + &row.0, + &did, + &row.1.uri, + &row.1.cid.to_string(), + &row.2.naive_utc(), + ]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO reposts (SELECT * FROM reposts_tmp)", &[]) + .await +} + +const POST_STMT: &str = "COPY posts_tmp (at_uri, cid, did, record, content, facets, languages, tags, parent_uri, parent_cid, root_uri, root_cid, embed, embed_subtype, created_at) FROM STDIN (FORMAT binary)"; +const POST_TYPES: &[Type] = &[ + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::JSONB, + Type::TEXT, + Type::JSONB, + Type::TEXT_ARRAY, + Type::TEXT_ARRAY, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TIMESTAMP, +]; +pub async fn copy_posts( + conn: &mut Transaction<'_>, + did: &str, + data: Vec<(String, Cid, records::AppBskyFeedPost)>, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE posts_tmp (LIKE posts INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn.copy_in(POST_STMT).await?; + let writer = BinaryCopyInWriter::new(writer, POST_TYPES); + + pin_mut!(writer); + + for (at_uri, cid, post) in data { + let record = serde_json::to_value(&post).unwrap(); + let facets = post.facets.and_then(|v| serde_json::to_value(v).ok()); + let embed = post.embed.as_ref().map(|v| v.as_str()); + let embed_subtype = post.embed.as_ref().and_then(|v| v.subtype()); + let (parent_uri, parent_cid) = strongref_to_parts(post.reply.as_ref().map(|v| &v.parent)); + let (root_uri, root_cid) = strongref_to_parts(post.reply.as_ref().map(|v| &v.root)); + + let writer = writer.as_mut(); + writer + .write(&[ + &at_uri, + &cid.to_string(), + &did, + &record, + &post.text, + &facets, + &post.langs.unwrap_or_default(), + &post.tags.unwrap_or_default(), + &parent_uri, + &parent_cid, + &root_uri, + &root_cid, + &embed, + &embed_subtype, + &post.created_at.naive_utc(), + ]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO posts (SELECT * FROM posts_tmp)", &[]) + .await +} + +pub async fn copy_blocks( + conn: &mut Transaction<'_>, + did: &str, + data: Vec, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE blocks_tmp (LIKE blocks INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in("COPY blocks_tmp (at_uri, did, subject, created_at) FROM STDIN (FORMAT binary)") + .await?; + let writer = BinaryCopyInWriter::new(writer, SUBJECT_TYPES); + + pin_mut!(writer); + + for row in data { + let writer = writer.as_mut(); + writer + .write(&[&row.0, &did, &row.1, &row.2.naive_utc()]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO blocks (SELECT * FROM blocks_tmp)", &[]) + .await +} + +pub async fn copy_list_items( + conn: &mut Transaction<'_>, + data: Vec<(String, records::AppBskyGraphListItem)>, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE list_items_tmp (LIKE list_items INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in( + "COPY list_items_tmp (at_uri, list_uri, subject, created_at) FROM STDIN (FORMAT binary)", + ) + .await?; + let writer = BinaryCopyInWriter::new( + writer, + &[Type::TEXT, Type::TEXT, Type::TEXT, Type::TIMESTAMP], + ); + + pin_mut!(writer); + + for (at_uri, record) in data { + let writer = writer.as_mut(); + writer + .write(&[ + &at_uri, + &record.list, + &record.subject, + &record.created_at.naive_utc(), + ]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO list_items (SELECT * FROM list_items_tmp)", &[]) + .await +} + +pub async fn copy_follows( + conn: &mut Transaction<'_>, + did: &str, + data: Vec, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE follows_tmp (LIKE follows INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in("COPY follows_tmp (at_uri, did, subject, created_at) FROM STDIN (FORMAT binary)") + .await?; + let writer = BinaryCopyInWriter::new(writer, SUBJECT_TYPES); + + pin_mut!(writer); + + for row in data { + let writer = writer.as_mut(); + writer + .write(&[&row.0, &did, &row.1, &row.2.naive_utc()]) + .await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO follows (SELECT * FROM follows_tmp)", &[]) + .await +} + +const VERIFICATION_TYPES: &[Type] = &[ + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TEXT, + Type::TIMESTAMP, +]; +pub async fn copy_verification( + conn: &mut Transaction<'_>, + did: &str, + data: Vec<(String, Cid, records::AppBskyGraphVerification)>, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE verification_tmp (LIKE verification INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in("COPY verification_tmp (at_uri, cid, verifier, subject, handle, display_name, created_at) FROM STDIN (FORMAT binary)") + .await?; + let writer = BinaryCopyInWriter::new(writer, VERIFICATION_TYPES); + + pin_mut!(writer); + + for (at_uri, cid, record) in data { + let writer = writer.as_mut(); + writer + .write(&[ + &at_uri, + &cid.to_string(), + &did, + &record.subject, + &record.handle, + &record.display_name, + &record.created_at.naive_utc(), + ]) + .await?; + } + + writer.finish().await?; + + conn.execute( + "INSERT INTO verification (SELECT * FROM verification_tmp)", + &[], + ) + .await +} + +pub async fn copy_records( + conn: &mut Transaction<'_>, + did: &str, + data: Vec<(String, Cid)>, +) -> PgExecResult { + if data.is_empty() { + return Ok(0); + } + + conn.execute( + "CREATE TEMP TABLE records_tmp (LIKE records INCLUDING DEFAULTS) ON COMMIT DROP", + &[], + ) + .await?; + + let writer = conn + .copy_in("COPY records_tmp (at_uri, cid, did) FROM STDIN (FORMAT binary)") + .await?; + let writer = BinaryCopyInWriter::new(writer, &[Type::TEXT, Type::TEXT, Type::TEXT]); + + pin_mut!(writer); + + for (at_uri, cid) in data { + let writer = writer.as_mut(); + writer.write(&[&at_uri, &did, &cid.to_string()]).await?; + } + + writer.finish().await?; + + conn.execute("INSERT INTO records (SELECT * FROM records_tmp)", &[]) + .await +} diff --git a/consumer/src/db/labels.rs b/consumer/src/db/labels.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/labels.rs @@ -0,0 +1,79 @@ +use super::PgExecResult; +use crate::indexer::records::AppBskyLabelerService; +use deadpool_postgres::GenericClient; +use ipld_core::cid::Cid; +use lexica::com_atproto::label::{LabelValueDefinition, SelfLabels}; +use std::collections::HashMap; + +pub async fn maintain_label_defs( + conn: &mut C, + repo: &str, + rec: &AppBskyLabelerService, +) -> PgExecResult { + // drop any label defs not currently in the list + conn.execute( + "DELETE FROM labeler_defs WHERE labeler=$1 AND NOT label_identifier = any($2)", + &[&repo, &rec.policies.label_values], + ) + .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 severity = definition.map(|v| v.severity.to_string()); + let blurs = definition.map(|v| v.blurs.to_string()); + let default_setting = definition + .and_then(|v| v.default_setting) + .map(|v| v.to_string()); + let adult_only = definition.and_then(|v| v.adult_only); + let locales = definition.and_then(|v| serde_json::to_value(&v.locales).ok()); + + conn.execute( + include_str!("sql/label_defs_upsert.sql"), + &[ + &repo, + &label, + &severity, + &blurs, + &default_setting, + &adult_only, + &locales, + ], + ) + .await?; + } + + Ok(0) +} + +pub async fn maintain_self_labels( + conn: &mut C, + repo: &str, + cid: Option, + at_uri: &str, + self_labels: SelfLabels, +) -> PgExecResult { + conn.execute( + "DELETE FROM labels WHERE self_label=TRUE AND uri=$1", + &[&at_uri], + ) + .await?; + + let cid = cid.map(|cid| cid.to_string()); + + let stmt = conn.prepare_cached("INSERT INTO labels (labeler, label, uri, self_label, cid, created_at) VALUES ($1, $2, $3, TRUE, $4, NOW())").await?; + + for label in self_labels.values { + conn.execute(&stmt, &[&repo, &label.val, &at_uri, &cid.clone()]) + .await?; + } + + Ok(0) +} diff --git a/consumer/src/db/mod.rs b/consumer/src/db/mod.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/mod.rs @@ -0,0 +1,16 @@ +use tokio_postgres::Error as PgError; + +type PgResult = Result; +type PgExecResult = PgResult; +type PgOptResult = PgResult>; + +mod actor; +mod backfill; +pub mod copy; +mod labels; +mod record; + +pub use actor::*; +pub use backfill::*; +pub use labels::*; +pub use record::*; diff --git a/consumer/src/db/record.rs b/consumer/src/db/record.rs new file mode 100644 --- /dev/null +++ b/consumer/src/db/record.rs @@ -0,0 +1,651 @@ +use super::{PgExecResult, PgOptResult}; +use crate::indexer::records::*; +use crate::utils::{blob_ref, strongref_to_parts}; +use chrono::prelude::*; +use deadpool_postgres::GenericClient; +use ipld_core::cid::Cid; + +pub async fn record_upsert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, +) -> PgExecResult { + conn.execute( + "INSERT INTO records (at_uri, did, cid) VALUES ($1, $2, $3) ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid", + &[&repo, &at_uri, &cid.to_string()], + ).await +} + +pub async fn record_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM records WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn block_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + rec: AppBskyGraphBlock, +) -> PgExecResult { + conn.execute( + "INSERT INTO blocks (at_uri, did, subject, created_at) VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING", + &[&at_uri, &repo, &rec.subject, &rec.created_at], + ).await +} + +pub async fn block_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM blocks WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn chat_decl_upsert( + conn: &mut C, + repo: &str, + rec: ChatBskyActorDeclaration, +) -> PgExecResult { + conn.execute( + "INSERT INTO chat_decls (did, allow_incoming) VALUES ($1, $2) ON CONFLICT (did) DO UPDATE SET allow_incoming=EXCLUDED.allow_incoming", + &[&repo, &rec.allow_incoming.to_string()] + ).await +} + +pub async fn chat_decl_delete(conn: &mut C, repo: &str) -> PgExecResult { + conn.execute("DELETE FROM chat_decls WHERE did=$1", &[&repo]) + .await +} + +pub async fn feedgen_upsert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, + rec: AppBskyFeedGenerator, +) -> PgExecResult { + let cid = cid.to_string(); + let description_facets = rec + .description_facets + .and_then(|v| serde_json::to_value(v).ok()); + let avatar = blob_ref(rec.avatar); + + conn.execute( + include_str!("sql/feedgen_upsert.sql"), + &[ + &at_uri, + &repo, + &cid, + &rec.did, + &rec.content_mode, + &rec.display_name, + &rec.description, + &description_facets, + &avatar, + &rec.created_at, + ], + ) + .await +} + +pub async fn feedgen_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM feedgens WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn follow_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + rec: AppBskyGraphFollow, +) -> PgExecResult { + conn.execute( + "INSERT INTO follows (at_uri, did, subject, created_at) VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING", + &[&at_uri, &repo, &rec.subject, &rec.created_at], + ).await +} + +pub async fn follow_delete(conn: &mut C, at_uri: &str) -> PgOptResult { + let res = conn + .query_opt( + "DELETE FROM follows WHERE at_uri=$1 RETURNING subject", + &[&at_uri], + ) + .await?; + + Ok(res.map(|v| v.get(0))) +} + +pub async fn labeler_upsert( + conn: &mut C, + repo: &str, + cid: Cid, + rec: AppBskyLabelerService, +) -> PgExecResult { + let cid = cid.to_string(); + let reasons = rec + .reason_types + .as_ref() + .map(|v| v.iter().map(|v| v.to_string()).collect::>()); + let subject_types = rec + .subject_types + .as_ref() + .map(|v| v.iter().map(|v| v.to_string()).collect::>()); + + conn.execute( + include_str!("sql/label_service_upsert.sql"), + &[ + &repo, + &cid, + &reasons, + &subject_types, + &rec.subject_collections, + ], + ) + .await?; + + super::maintain_label_defs(conn, repo, &rec).await +} + +pub async fn labeler_delete(conn: &mut C, repo: &str) -> PgExecResult { + conn.execute("DELETE FROM labelers WHERE did=$1", &[&repo]) + .await +} + +pub async fn like_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + rec: AppBskyFeedLike, +) -> PgExecResult { + conn.execute( + "INSERT INTO likes (at_uri, did, subject, subject_cid, created_at) VALUES ($1, $2, $3, $4, $5)", + &[&at_uri, &repo, &rec.subject.uri, &rec.subject.cid.to_string(), &rec.created_at] + ).await +} + +pub async fn like_delete(conn: &mut C, at_uri: &str) -> PgOptResult { + let res = conn + .query_opt( + "DELETE FROM likes WHERE at_uri=$1 RETURNING subject", + &[&at_uri], + ) + .await?; + + Ok(res.map(|v| v.get(0))) +} + +pub async fn list_upsert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, + rec: AppBskyGraphList, +) -> PgExecResult { + let cid = cid.to_string(); + let description_facets = rec + .description_facets + .and_then(|v| serde_json::to_value(v).ok()); + let avatar = blob_ref(rec.avatar); + + conn.execute( + include_str!("sql/list_upsert.sql"), + &[ + &at_uri, + &repo, + &cid, + &rec.purpose, + &rec.name, + &rec.description, + &description_facets, + &avatar, + &rec.created_at, + ], + ) + .await +} + +pub async fn list_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM lists WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn list_block_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + rec: AppBskyGraphListBlock, +) -> PgExecResult { + conn.execute( + "INSERT INTO list_blocks (at_uri, did, list_uri, created_at) VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING", + &[&at_uri, &repo, &rec.subject, &rec.created_at], + ).await +} + +pub async fn list_block_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM list_blocks WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn list_item_insert( + conn: &mut C, + at_uri: &str, + rec: AppBskyGraphListItem, +) -> PgExecResult { + conn.execute( + "INSERT INTO list_items (at_uri, list_uri, subject, created_at) VALUES ($1, $2, $3, $4) ON CONFLICT DO NOTHING", + &[&at_uri, &rec.list, &rec.subject, &rec.created_at], + ).await +} + +pub async fn list_item_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM list_items WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn post_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, + rec: AppBskyFeedPost, +) -> PgExecResult { + let cid = cid.to_string(); + let record = serde_json::to_value(&rec).unwrap(); + let facets = rec.facets.and_then(|v| serde_json::to_value(v).ok()); + let (parent_uri, parent_cid) = strongref_to_parts(rec.reply.as_ref().map(|v| &v.parent)); + let (root_uri, root_cid) = strongref_to_parts(rec.reply.as_ref().map(|v| &v.root)); + let embed = rec.embed.as_ref().map(|v| v.as_str()); + let embed_subtype = rec.embed.as_ref().and_then(|v| v.subtype()); + + let count = conn + .execute( + include_str!("sql/post_insert.sql"), + &[ + &at_uri, + &repo, + &cid, + &record, + &rec.text, + &facets, + &rec.langs.unwrap_or_default(), + &rec.tags.unwrap_or_default(), + &parent_uri, + &parent_cid, + &root_uri, + &root_cid, + &embed, + &embed_subtype, + &rec.created_at, + ], + ) + .await?; + + if let Some(embed) = rec.embed.and_then(|embed| embed.into_bsky()) { + post_embed_insert(conn, at_uri, embed, rec.created_at).await?; + } + + Ok(count) +} + +pub async fn post_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM posts WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn post_get_info_for_delete( + conn: &mut C, + at_uri: &str, +) -> PgOptResult<(Option, Option)> { + let res = conn + .query_opt( + "SELECT parent_uri, per.uri FROM posts LEFT JOIN post_embed_record per on at_uri = per.post_uri WHERE at_uri = $1", + &[&at_uri], + ) + .await?; + + Ok(res.map(|row| (row.get(0), row.get(1)))) +} + +pub async fn post_embed_insert( + conn: &mut C, + post: &str, + embed: AppBskyEmbed, + created_at: DateTime, +) -> PgExecResult { + match embed { + AppBskyEmbed::Images(embed) => post_embed_image_insert(conn, post, embed).await, + AppBskyEmbed::Video(embed) => post_embed_video_insert(conn, post, embed).await, + AppBskyEmbed::External(embed) => post_embed_external_insert(conn, post, embed).await, + AppBskyEmbed::Record(embed) => { + post_embed_record_insert(conn, post, embed, created_at).await + } + AppBskyEmbed::RecordWithMedia(embed) => { + post_embed_record_insert(conn, post, embed.record, created_at).await?; + match *embed.media { + AppBskyEmbed::Images(embed) => post_embed_image_insert(conn, post, embed).await, + AppBskyEmbed::Video(embed) => post_embed_video_insert(conn, post, embed).await, + AppBskyEmbed::External(embed) => { + post_embed_external_insert(conn, post, embed).await + } + _ => unreachable!(), + } + } + } +} + +async fn post_embed_image_insert( + conn: &mut C, + post: &str, + embed: AppBskyEmbedImages, +) -> PgExecResult { + let stmt = conn.prepare("INSERT INTO post_embed_images (post_uri, seq, cid, mime_type, alt, width, height) VALUES ($1, $2, $3, $4, $5, $6, $7)").await?; + + for (idx, image) in embed.images.iter().enumerate() { + let cid = image.image.r#ref.to_string(); + let width = image.aspect_ratio.as_ref().map(|v| v.width); + let height = image.aspect_ratio.as_ref().map(|v| v.height); + + conn.execute( + &stmt, + &[ + &post, + &(idx as i16), + &cid, + &image.image.mime_type, + &image.alt, + &width, + &height, + ], + ) + .await?; + } + + Ok(0) +} + +async fn post_embed_video_insert( + conn: &mut C, + post: &str, + embed: AppBskyEmbedVideo, +) -> PgExecResult { + let cid = embed.video.r#ref.to_string(); + let width = embed.aspect_ratio.as_ref().map(|v| v.width); + let height = embed.aspect_ratio.as_ref().map(|v| v.height); + + let count = conn.execute( + "INSERT INTO post_embed_video (post_uri, cid, mime_type, alt, width, height) VALUES ($1, $2, $3, $4, $5, $6)", + &[&post, &cid, &embed.video.mime_type, &embed.alt, &width, &height], + ).await?; + + if let Some(captions) = embed.captions { + let stmt = conn.prepare_cached("INSERT INTO post_embed_video_captions (post_uri, cid, mime_type, language) VALUES ($1, $2, $3, $4)").await?; + + for caption in captions { + let cid = caption.file.r#ref.to_string(); + conn.execute( + &stmt, + &[&post, &cid, &caption.file.mime_type, &caption.lang], + ) + .await?; + } + } + + Ok(count) +} + +async fn post_embed_external_insert( + conn: &mut C, + post: &str, + embed: AppBskyEmbedExternal, +) -> PgExecResult { + let thumb_mime = embed.external.thumb.as_ref().map(|v| v.mime_type.clone()); + let thumb_cid = embed.external.thumb.as_ref().map(|v| v.r#ref.to_string()); + + conn.execute( + "INSERT INTO post_embed_ext (post_uri, uri, title, description, thumb_mime_type, thumb_cid) VALUES ($1, $2, $3, $4, $5, $6)", + &[&post, &embed.external.uri, &embed.external.title, &embed.external.description, &thumb_mime, &thumb_cid], + ).await +} + +async fn post_embed_record_insert( + conn: &mut C, + post: &str, + embed: AppBskyEmbedRecord, + post_created_at: DateTime, +) -> PgExecResult { + // strip "at://" then break into parts by '/' + let parts = embed.record.uri[5..].split('/').collect::>(); + + let detached = if parts[1] == "app.bsky.feed.post" { + let postgate_effective: Option> = conn + .query_opt( + "SELECT created_at FROM postgates WHERE post_uri=$1", + &[&post], + ) + .await? + .map(|v| v.get(0)); + + postgate_effective + .map(|v| Utc::now().min(post_created_at) > v) + .unwrap_or_default() + } else { + false + }; + + conn.execute( + "INSERT INTO post_embed_record (post_uri, record_type, uri, cid, detached) VALUES ($1, $2, $3, $4, $5)", + &[&post, &parts[1], &embed.record.uri, &embed.record.cid.to_string(), &detached], + ).await +} + +pub async fn postgate_upsert( + conn: &mut C, + at_uri: &str, + cid: Cid, + rec: &AppBskyFeedPostgate, +) -> PgExecResult { + let rules = rec + .embedding_rules + .iter() + .map(|v| v.as_str().to_string()) + .collect::>(); + + conn.execute( + include_str!("sql/postgate_upsert.sql"), + &[ + &at_uri, + &cid.to_string(), + &rec.post, + &rec.detached_embedding_uris, + &rules, + &rec.created_at, + ], + ) + .await +} + +pub async fn postgate_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM postgates WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn postgate_maintain_detaches( + conn: &mut C, + post: &str, + detached: &[String], + disable_effective: Option, +) -> PgExecResult { + conn.execute( + "SELECT maintain_postgates($1, $2, $3)", + &[&post, &detached, &disable_effective], + ) + .await +} + +pub async fn profile_upsert( + conn: &mut C, + repo: &str, + cid: Cid, + rec: AppBskyActorProfile, +) -> PgExecResult { + let cid = cid.to_string(); + let avatar = blob_ref(rec.avatar); + let banner = blob_ref(rec.banner); + let (pinned_uri, pinned_cid) = strongref_to_parts(rec.pinned_post.as_ref()); + let (joined_sp_uri, joined_sp_cid) = strongref_to_parts(rec.joined_via_starter_pack.as_ref()); + + conn.execute( + include_str!("sql/profile_upsert.sql"), + &[ + &repo, + &cid, + &avatar, + &banner, + &rec.display_name, + &rec.description, + &pinned_uri, + &pinned_cid, + &joined_sp_uri, + &joined_sp_cid, + &rec.created_at.unwrap_or(Utc::now()).naive_utc(), + ], + ) + .await +} + +pub async fn profile_delete(conn: &mut C, repo: &str) -> PgExecResult { + conn.execute("DELETE FROM profiles WHERE did=$1", &[&repo]) + .await +} + +pub async fn repost_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + rec: AppBskyFeedRepost, +) -> PgExecResult { + conn.execute( + "INSERT INTO reposts (at_uri, did, post, post_cid, created_at) VALUES ($1, $2, $3, $4, $5)", + &[ + &at_uri, + &repo, + &rec.subject.uri, + &rec.subject.cid.to_string(), + &rec.created_at, + ], + ) + .await +} + +pub async fn repost_delete(conn: &mut C, at_uri: &str) -> PgOptResult { + let res = conn + .query_opt( + "DELETE FROM reposts WHERE at_uri=$1 RETURNING post", + &[&at_uri], + ) + .await?; + + Ok(res.map(|v| v.get(0))) +} + +pub async fn starter_pack_upsert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, + rec: AppBskyGraphStarterPack, +) -> PgExecResult { + let cid = cid.to_string(); + let record = serde_json::to_value(&rec).unwrap(); + let description_facets = rec + .description_facets + .and_then(|v| serde_json::to_value(v).ok()); + let feeds = rec + .feeds + .map(|v| v.into_iter().map(|item| item.uri).collect::>()); + + conn.execute( + include_str!("sql/starterpack_upsert.sql"), + &[ + &at_uri, + &repo, + &cid, + &record, + &rec.name, + &rec.description, + &description_facets, + &rec.list, + &feeds, + &rec.created_at, + ], + ) + .await +} + +pub async fn starter_pack_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM starterpacks WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn threadgate_upsert( + conn: &mut C, + at_uri: &str, + cid: Cid, + rec: AppBskyFeedThreadgate, +) -> PgExecResult { + let record = serde_json::to_value(&rec).unwrap(); + + let allowed_lists = rec + .allow + .iter() + .filter_map(|rule| match rule { + ThreadgateRule::List { list } => Some(list.clone()), + _ => None, + }) + .collect::>(); + + let allow = rec + .allow + .into_iter() + .map(|v| v.as_str().to_string()) + .collect::>(); + + conn.execute( + include_str!("sql/threadgate_upsert.sql"), + &[ + &at_uri, + &cid.to_string(), + &rec.post, + &rec.hidden_replies, + &allow, + &allowed_lists, + &record, + &rec.created_at, + ], + ) + .await +} + +pub async fn threadgate_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM threadgates WHERE at_uri=$1", &[&at_uri]) + .await +} + +pub async fn verification_insert( + conn: &mut C, + at_uri: &str, + repo: &str, + cid: Cid, + rec: AppBskyGraphVerification, +) -> PgExecResult { + let cid = cid.to_string(); + + conn.execute( + "INSERT INTO verification (at_uri, verifier, cid, subject, handle, display_name, created_at) VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT DO NOTHING", + &[&at_uri, &repo, &cid, &rec.subject, &rec.handle, &rec.display_name, &rec.created_at], + ).await +} + +pub async fn verification_delete(conn: &mut C, at_uri: &str) -> PgExecResult { + conn.execute("DELETE FROM verification WHERE at_uri=$1", &[&at_uri]) + .await +} diff --git a/consumer/src/firehose/error.rs b/consumer/src/firehose/error.rs --- a/consumer/src/firehose/error.rs +++ b/consumer/src/firehose/error.rs @@ -1,5 +1,5 @@ -use thiserror::Error; use std::io::Error as IoError; +use thiserror::Error; #[derive(Debug, Error)] pub enum FirehoseError { @@ -9,4 +9,4 @@ IpldCbor(#[from] serde_ipld_dagcbor::error::DecodeError), #[error("{0}")] Websocket(#[from] tokio_tungstenite::tungstenite::error::Error), -} \ No newline at end of file +} diff --git a/consumer/src/firehose/mod.rs b/consumer/src/firehose/mod.rs --- a/consumer/src/firehose/mod.rs +++ b/consumer/src/firehose/mod.rs @@ -140,7 +140,10 @@ match err { WsError::Protocol(ProtocolError::ResetWithoutClosingHandshake) | WsError::ConnectionClosed => true, - WsError::Io(ioerr) => matches!(ioerr.kind(), ErrorKind::BrokenPipe | ErrorKind::ConnectionReset), + WsError::Io(ioerr) => matches!( + ioerr.kind(), + ErrorKind::BrokenPipe | ErrorKind::ConnectionReset + ), _ => false, } } diff --git a/consumer/src/indexer/db.rs b/consumer/src/indexer/db.rs deleted file mode 100644 --- a/consumer/src/indexer/db.rs +++ /dev/null @@ -1,971 +0,0 @@ -use super::records::{self, AppBskyEmbed}; -use crate::utils::{blob_ref, empty_str_as_none, strongref_to_parts}; -use chrono::prelude::*; -use diesel::prelude::*; -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_record( - conn: &mut AsyncPgConnection, - at_uri: &str, - repo: &str, - cid: Cid, -) -> QueryResult { - let cid = cid.to_string(); - - diesel::insert_into(schema::records::table) - .values(( - schema::records::at_uri.eq(at_uri), - schema::records::did.eq(repo), - schema::records::cid.eq(&cid), - )) - .on_conflict(schema::records::at_uri) - .do_update() - .set(schema::records::cid.eq(&cid)) - .execute(conn) - .await -} - -pub async fn delete_record(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::records::table) - .filter(schema::records::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn write_backfill_row( - conn: &mut AsyncPgConnection, - repo: &str, - rev: &str, - cid: Cid, - data: serde_json::Value, -) -> QueryResult { - diesel::insert_into(schema::backfill::table) - .values(models::NewBackfillRow { - repo, - repo_ver: rev, - cid: cid.to_string(), - data, - }) - .execute(conn) - .await -} - -pub async fn get_repo_info( - conn: &mut AsyncPgConnection, - repo: &str, -) -> QueryResult, types::ActorSyncState)>> { - schema::actors::table - .select((schema::actors::repo_rev, schema::actors::sync_state)) - .find(repo) - .get_result(conn) - .await - .optional() -} - -pub async fn upsert_actor( - conn: &mut AsyncPgConnection, - did: &str, - handle: Option>, - status: Option, - sync_state: Option, - time: DateTime, -) -> QueryResult { - let data = models::NewActor { - did, - handle, - status, - sync_state, - last_indexed: Some(time.naive_utc()), - }; - - diesel::insert_into(schema::actors::table) - .values(&data) - .on_conflict(schema::actors::did) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn account_status_and_rev( - conn: &mut AsyncPgConnection, - did: &str, -) -> QueryResult)>> { - schema::actors::table - .select((schema::actors::status, schema::actors::repo_rev)) - .for_update() - .find(did) - .get_result(conn) - .await - .optional() -} - -/// Attempts to update a repo to the given version. -/// returns false if the repo doesn't exist or is too new, or true if the update succeeded. -pub async fn update_repo_version( - conn: &mut AsyncPgConnection, - repo: &str, - rev: &str, - cid: Cid, -) -> QueryResult { - diesel::update(schema::actors::table) - .set(( - schema::actors::repo_rev.eq(rev), - schema::actors::repo_cid.eq(cid.to_string()), - )) - .filter(schema::actors::did.eq(repo)) - .execute(conn) - .await -} - -pub async fn insert_block( - conn: &mut AsyncPgConnection, - repo: &str, - at_uri: &str, - rec: records::AppBskyGraphBlock, -) -> QueryResult { - diesel::insert_into(schema::blocks::table) - .values(&models::NewBlock { - at_uri, - did: repo, - subject: &rec.subject, - created_at: rec.created_at.naive_utc(), - }) - .on_conflict_do_nothing() - .execute(conn) - .await -} - -pub async fn delete_block(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::blocks::table) - .filter(schema::blocks::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn insert_follow( - conn: &mut AsyncPgConnection, - repo: &str, - at_uri: &str, - rec: records::AppBskyGraphFollow, -) -> QueryResult { - diesel::insert_into(schema::follows::table) - .values(&models::NewFollow { - at_uri, - did: repo, - subject: &rec.subject, - created_at: rec.created_at.naive_utc(), - }) - .on_conflict_do_nothing() - .execute(conn) - .await -} - -pub async fn delete_follow( - conn: &mut AsyncPgConnection, - at_uri: &str, -) -> QueryResult> { - diesel::delete(schema::follows::table) - .filter(schema::follows::at_uri.eq(at_uri)) - .returning(schema::follows::subject) - .get_result(conn) - .await - .optional() -} - -pub async fn upsert_profile( - conn: &mut AsyncPgConnection, - repo: &str, - cid: Cid, - rec: records::AppBskyActorProfile, -) -> QueryResult { - let (pinned_uri, pinned_cid) = strongref_to_parts(rec.pinned_post.as_ref()); - let (joined_sp_uri, joined_sp_cid) = strongref_to_parts(rec.joined_via_starter_pack.as_ref()); - - let data = models::UpsertProfile { - did: repo, - cid: cid.to_string(), - avatar_cid: blob_ref(rec.avatar), - banner_cid: blob_ref(rec.banner), - display_name: rec.display_name, - description: rec.description, - pinned_uri, - pinned_cid, - joined_sp_uri, - joined_sp_cid, - created_at: rec.created_at.map(|val| val.naive_utc()), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::profiles::table) - .values(&data) - .on_conflict(schema::profiles::did) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_profile(conn: &mut AsyncPgConnection, repo: &str) -> QueryResult { - diesel::delete(schema::profiles::table) - .filter(schema::profiles::did.eq(repo)) - .execute(conn) - .await -} - -pub async fn upsert_list( - conn: &mut AsyncPgConnection, - repo: &str, - at_uri: &str, - cid: Cid, - rec: records::AppBskyGraphList, -) -> QueryResult { - let description_facets = rec - .description_facets - .and_then(|v| serde_json::to_value(v).ok()); - - let data = models::UpsertList { - at_uri, - owner: repo, - cid: cid.to_string(), - list_type: &rec.purpose, - name: &rec.name, - description: rec.description, - description_facets, - avatar_cid: blob_ref(rec.avatar), - created_at: rec.created_at.naive_utc(), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::lists::table) - .values(&data) - .on_conflict(schema::lists::at_uri) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_list(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::lists::table) - .filter(schema::lists::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn insert_list_block( - conn: &mut AsyncPgConnection, - repo: &str, - at_uri: &str, - rec: records::AppBskyGraphListBlock, -) -> QueryResult { - let data = models::NewListBlock { - at_uri, - did: repo, - list_uri: &rec.subject, - created_at: rec.created_at.naive_utc(), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::list_blocks::table) - .values(&data) - .execute(conn) - .await -} - -pub async fn delete_list_block(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::list_blocks::table) - .filter(schema::list_blocks::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn insert_list_item( - conn: &mut AsyncPgConnection, - at_uri: &str, - rec: records::AppBskyGraphListItem, -) -> QueryResult { - let data = models::NewListItem { - at_uri, - list_uri: &rec.list, - subject: &rec.subject, - created_at: rec.created_at.naive_utc(), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::list_items::table) - .values(&data) - .execute(conn) - .await -} - -pub async fn delete_list_item(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::list_items::table) - .filter(schema::list_items::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn upsert_feedgen( - conn: &mut AsyncPgConnection, - repo: &str, - cid: Cid, - at_uri: &str, - rec: records::AppBskyFeedGenerator, -) -> QueryResult { - let description_facets = rec - .description_facets - .and_then(|v| serde_json::to_value(v).ok()); - - let data = models::UpsertFeedGen { - at_uri, - cid: &cid.to_string(), - owner: repo, - service_did: &rec.did, - content_mode: rec.content_mode, - name: &rec.display_name, - description: rec.description, - description_facets, - avatar_cid: blob_ref(rec.avatar), - accepts_interactions: rec.accepts_interactions, - created_at: rec.created_at.naive_utc(), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::feedgens::table) - .values(&data) - .on_conflict(schema::feedgens::at_uri) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_feedgen(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::feedgens::table) - .filter(schema::feedgens::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn insert_post( - conn: &mut AsyncPgConnection, - did: &str, - cid: Cid, - at_uri: &str, - rec: records::AppBskyFeedPost, -) -> QueryResult { - let record = serde_json::to_value(&rec).unwrap(); - let facets = rec.facets.and_then(|v| serde_json::to_value(v).ok()); - - let embed = rec.embed.as_ref().map(|v| v.as_str()); - let embed_subtype = rec.embed.as_ref().and_then(|v| v.subtype()); - - let (parent_uri, parent_cid) = strongref_to_parts(rec.reply.as_ref().map(|v| &v.parent)); - let (root_uri, root_cid) = strongref_to_parts(rec.reply.as_ref().map(|v| &v.root)); - - let res = diesel::insert_into(schema::posts::table) - .values(models::NewPost { - at_uri, - cid: cid.to_string(), - did, - record, - content: &rec.text, - facets, - languages: rec.langs.unwrap_or_default(), - tags: rec.tags.unwrap_or_default(), - parent_uri, - parent_cid, - root_uri, - root_cid, - embed, - embed_subtype, - created_at: rec.created_at.naive_utc(), - }) - .execute(conn) - .await?; - - match rec.embed.and_then(|v| v.into_bsky()) { - Some(AppBskyEmbed::Images(embed)) => insert_post_embed_images(conn, at_uri, embed).await, - Some(AppBskyEmbed::Video(embed)) => insert_post_embed_video(conn, at_uri, embed).await, - Some(AppBskyEmbed::External(embed)) => insert_post_embed_ext(conn, at_uri, embed).await, - Some(AppBskyEmbed::Record(embed)) => { - insert_post_embed_record(conn, at_uri, embed, rec.created_at).await - } - Some(AppBskyEmbed::RecordWithMedia(embed)) => { - insert_post_embed_record(conn, at_uri, embed.record, rec.created_at).await?; - match *embed.media { - AppBskyEmbed::Images(embed) => insert_post_embed_images(conn, at_uri, embed).await, - AppBskyEmbed::Video(embed) => insert_post_embed_video(conn, at_uri, embed).await, - AppBskyEmbed::External(embed) => insert_post_embed_ext(conn, at_uri, embed).await, - _ => unreachable!(), - } - } - _ => Ok(res), - } -} - -async fn insert_post_embed_images( - conn: &mut AsyncPgConnection, - at_uri: &str, - rec: records::AppBskyEmbedImages, -) -> QueryResult { - let images = rec - .images - .into_iter() - .enumerate() - .map(|(idx, img)| models::NewPostEmbedImage { - post_uri: at_uri, - seq: idx as i16, - mime_type: img.image.mime_type, - cid: img.image.r#ref.to_string(), - alt: empty_str_as_none(img.alt), - width: img.aspect_ratio.as_ref().map(|v| v.width), - height: img.aspect_ratio.map(|v| v.height), - }) - .collect::>(); - - diesel::insert_into(schema::post_embed_images::table) - .values(images) - .execute(conn) - .await -} - -async fn insert_post_embed_video( - conn: &mut AsyncPgConnection, - at_uri: &str, - rec: records::AppBskyEmbedVideo, -) -> QueryResult { - let res = diesel::insert_into(schema::post_embed_video::table) - .values(models::NewPostEmbedVideo { - post_uri: at_uri, - mime_type: &rec.video.mime_type, - cid: rec.video.r#ref.to_string(), - alt: rec.alt, - width: rec.aspect_ratio.as_ref().map(|v| v.width), - height: rec.aspect_ratio.map(|v| v.height), - }) - .execute(conn) - .await?; - - match rec.captions { - Some(captions) => insert_post_embed_video_captions(conn, at_uri, &captions).await, - None => Ok(res), - } -} - -async fn insert_post_embed_video_captions( - conn: &mut AsyncPgConnection, - at_uri: &str, - captions: &[records::EmbedVideoCaptions], -) -> QueryResult { - let captions = captions - .iter() - .map(|caption| models::NewPostEmbedVideoCaption { - post_uri: at_uri, - language: caption.lang.clone(), - mime_type: caption.file.mime_type.clone(), - cid: caption.file.r#ref.to_string(), - }) - .collect::>(); - - diesel::insert_into(schema::post_embed_video_captions::table) - .values(captions) - .execute(conn) - .await -} - -async fn insert_post_embed_ext( - conn: &mut AsyncPgConnection, - at_uri: &str, - rec: records::AppBskyEmbedExternal, -) -> QueryResult { - diesel::insert_into(schema::post_embed_ext::table) - .values(models::NewPostEmbedExt { - post_uri: at_uri, - uri: &rec.external.uri, - title: &rec.external.title, - description: &rec.external.description, - thumb_mime_type: rec.external.thumb.as_ref().map(|v| v.mime_type.clone()), - thumb_cid: rec.external.thumb.as_ref().map(|v| v.r#ref.to_string()), - }) - .execute(conn) - .await -} - -async fn insert_post_embed_record( - conn: &mut AsyncPgConnection, - at_uri: &str, - rec: records::AppBskyEmbedRecord, - post_created_at: DateTime, -) -> QueryResult { - // strip "at://" then break into parts by '/' - let parts = rec.record.uri[5..].split('/').collect::>(); - - let detached = if parts[1] == "app.bsky.feed.post" { - // do a lookup on if we have a postgate for this record - let postgate_effective = schema::postgates::table - .select(schema::postgates::created_at) - .filter(schema::postgates::post_uri.eq(at_uri)) - .get_result::>(conn) - .await - .optional()?; - - postgate_effective.map(|v| Utc::now().min(post_created_at) > v) - } else { - None - }; - - diesel::insert_into(schema::post_embed_record::table) - .values(models::NewPostEmbedRecord { - post_uri: at_uri, - record_type: parts[1], - uri: &rec.record.uri, - cid: rec.record.cid.to_string(), - detached, - }) - .execute(conn) - .await -} - -pub async fn delete_post(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::posts::table) - .filter(schema::posts::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn get_post_info_for_delete( - conn: &mut AsyncPgConnection, - at_uri: &str, -) -> QueryResult, Option)>> { - schema::posts::table - .left_join( - schema::post_embed_record::table - .on(schema::posts::at_uri.eq(schema::post_embed_record::post_uri)), - ) - .select(( - schema::posts::parent_uri, - schema::post_embed_record::uri.nullable(), - )) - .filter(schema::posts::at_uri.eq(at_uri)) - .get_result(conn) - .await - .optional() -} - -pub async fn upsert_postgate( - conn: &mut AsyncPgConnection, - at_uri: &str, - cid: Cid, - rec: &records::AppBskyFeedPostgate, -) -> QueryResult { - let rules = rec - .embedding_rules - .iter() - .map(|v| v.as_str().to_string()) - .collect(); - - let data = models::UpsertPostgate { - at_uri, - cid: cid.to_string(), - post_uri: &rec.post, - detached: &rec.detached_embedding_uris, - rules, - created_at: rec.created_at.naive_utc(), - }; - - diesel::insert_into(schema::postgates::table) - .values(&data) - .on_conflict(schema::postgates::at_uri) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_postgate(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::postgates::table) - .filter(schema::postgates::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -define_sql_function! {fn maintain_postgates(post: Text, detached: Array, effective: Nullable)} - -pub async fn postgate_maintain_detaches( - conn: &mut AsyncPgConnection, - post: &str, - detached: &[String], - disable_effective: Option, -) -> QueryResult { - diesel::select(maintain_postgates(post, detached, disable_effective)) - .execute(conn) - .await -} - -pub async fn upsert_threadgate( - conn: &mut AsyncPgConnection, - at_uri: &str, - cid: Cid, - rec: records::AppBskyFeedThreadgate, -) -> QueryResult { - let record = serde_json::to_value(&rec).unwrap(); - - let allowed_lists = rec - .allow - .iter() - .filter_map(|rule| match rule { - records::ThreadgateRule::List { list } => Some(list.clone()), - _ => None, - }) - .collect(); - - let allow = rec - .allow - .into_iter() - .map(|v| v.as_str().to_string()) - .collect(); - - let data = models::UpsertThreadgate { - at_uri, - cid: cid.to_string(), - post_uri: &rec.post, - hidden_replies: rec.hidden_replies, - allow, - allowed_lists, - record, - created_at: rec.created_at.naive_utc(), - }; - - diesel::insert_into(schema::threadgates::table) - .values(&data) - .on_conflict(schema::threadgates::at_uri) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_threadgate(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::threadgates::table) - .filter(schema::threadgates::at_uri.eq(at_uri)) - .execute(conn) - .await -} - -pub async fn insert_like( - conn: &mut AsyncPgConnection, - did: &str, - at_uri: &str, - rec: records::AppBskyFeedLike, -) -> QueryResult { - let data = models::NewLike { - at_uri, - did, - subject: &rec.subject.uri, - subject_cid: rec.subject.cid.to_string(), - created_at: rec.created_at.naive_utc(), - }; - - diesel::insert_into(schema::likes::table) - .values(&data) - .execute(conn) - .await -} - -pub async fn delete_like( - conn: &mut AsyncPgConnection, - at_uri: &str, -) -> QueryResult> { - diesel::delete(schema::likes::table) - .filter(schema::likes::at_uri.eq(at_uri)) - .returning(schema::likes::subject) - .get_result(conn) - .await - .optional() -} - -pub async fn insert_repost( - conn: &mut AsyncPgConnection, - did: &str, - at_uri: &str, - rec: records::AppBskyFeedRepost, -) -> QueryResult { - let data = models::NewRepost { - at_uri, - did, - post: &rec.subject.uri, - post_cid: rec.subject.cid.to_string(), - created_at: rec.created_at.naive_utc(), - }; - - diesel::insert_into(schema::reposts::table) - .values(&data) - .execute(conn) - .await -} - -pub async fn delete_repost( - conn: &mut AsyncPgConnection, - at_uri: &str, -) -> QueryResult> { - diesel::delete(schema::reposts::table) - .filter(schema::reposts::at_uri.eq(at_uri)) - .returning(schema::reposts::post) - .get_result(conn) - .await - .optional() -} - -pub async fn upsert_chat_decl( - conn: &mut AsyncPgConnection, - did: &str, - rec: records::ChatBskyActorDeclaration, -) -> QueryResult { - let data = models::NewChatDecl { - did, - allow_incoming: rec.allow_incoming.to_string(), - }; - - diesel::insert_into(schema::chat_decls::table) - .values(&data) - .on_conflict(schema::chat_decls::did) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_chat_decl(conn: &mut AsyncPgConnection, did: &str) -> QueryResult { - diesel::delete(schema::chat_decls::table) - .filter(schema::chat_decls::did.eq(did)) - .execute(conn) - .await -} - -pub async fn upsert_starterpack( - conn: &mut AsyncPgConnection, - did: &str, - cid: Cid, - at_uri: &str, - rec: records::AppBskyGraphStarterPack, -) -> QueryResult { - let record = serde_json::to_value(&rec).unwrap(); - - let feeds = rec - .feeds - .map(|v| v.into_iter().map(|item| item.uri).collect()); - - let description_facets = rec - .description_facets - .and_then(|v| serde_json::to_value(v).ok()); - - let data = models::NewStarterPack { - at_uri, - cid: cid.to_string(), - owner: did, - record, - name: &rec.name, - description: rec.description, - description_facets, - list: &rec.list, - feeds, - created_at: rec.created_at.naive_utc(), - indexed_at: Utc::now().naive_utc(), - }; - - diesel::insert_into(schema::starterpacks::table) - .values(&data) - .on_conflict(schema::starterpacks::at_uri) - .do_update() - .set(&data) - .execute(conn) - .await -} - -pub async fn delete_starterpack(conn: &mut AsyncPgConnection, at_uri: &str) -> QueryResult { - diesel::delete(schema::starterpacks::table) - .filter(schema::starterpacks::at_uri.eq(at_uri)) - .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 subject_types = rec - .subject_types - .as_ref() - .map(|v| v.iter().map(|v| v.to_string()).collect()); - - let data = models::UpsertLabelerService { - did: repo, - cid: cid.to_string(), - reasons, - subject_types, - subject_collections: rec.subject_collections.as_ref(), - 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 -} - -pub async fn upsert_verification( - conn: &mut AsyncPgConnection, - did: &str, - cid: Cid, - at_uri: &str, - rec: records::AppBskyGraphVerification, -) -> QueryResult { - let data = models::NewVerificationEntry { - at_uri, - cid: cid.to_string(), - verifier: did, - subject: &rec.subject, - handle: &rec.handle, - display_name: &rec.display_name, - created_at: rec.created_at.naive_utc(), - indexed_at: None, - }; - - diesel::insert_into(schema::verification::table) - .values(&data) - .on_conflict(schema::verification::at_uri) - .do_update() - .set(&data) - .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 -} diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -1,4 +1,5 @@ use crate::config::HistoryMode; +use crate::db; use crate::firehose::{ AtpAccountEvent, AtpCommitEvent, AtpIdentityEvent, CommitOp, FirehoseConsumer, FirehoseEvent, FirehoseOutput, @@ -6,9 +7,8 @@ use crate::indexer::types::{ AggregateDeltaStore, BackfillItem, BackfillItemInner, CollectionType, RecordTypes, }; +use deadpool_postgres::{Object, Pool, Transaction}; use did_resolver::Resolver; -use diesel_async::pooled_connection::deadpool::Pool; -use diesel_async::{AsyncConnection, AsyncPgConnection}; use foldhash::quality::RandomState; use futures::StreamExt; use ipld_core::cid::Cid; @@ -23,7 +23,6 @@ use tokio::sync::mpsc::{channel, Sender}; use tracing::instrument; -pub mod db; pub mod records; pub mod types; @@ -41,7 +40,7 @@ } pub struct RelayIndexer { - pool: Pool, + pool: Pool, redis: MultiplexedConnection, state: RelayIndexerState, firehose: FirehoseConsumer, @@ -51,7 +50,7 @@ impl RelayIndexer { pub async fn new( - pool: Pool, + pool: Pool, redis: MultiplexedConnection, idxc_tx: Sender, resolver: Arc, @@ -197,25 +196,20 @@ #[instrument(skip_all, fields(seq = identity.seq, repo = identity.did))] async fn index_identity( state: &RelayIndexerState, - conn: &mut AsyncPgConnection, + conn: &mut Object, identity: AtpIdentityEvent, ) -> eyre::Result<()> { let new_handle = match state.do_handle_res { true => resolve_handle(state, &identity.did, identity.handle).await?, - false => Some(identity.handle), + false => identity.handle, }; - let sync_state = (!state.do_backfill).then_some(ActorSyncState::Synced); + let sync_state = match state.do_backfill { + true => ActorSyncState::Dirty, + false => ActorSyncState::Synced, + }; - db::upsert_actor( - conn, - &identity.did, - new_handle, - None, - sync_state, - identity.time, - ) - .await?; + db::actor_upsert_handle(conn, &identity.did, sync_state, new_handle, identity.time).await?; Ok(()) } @@ -224,7 +218,7 @@ state: &RelayIndexerState, did: &str, expected_handle: Option, -) -> eyre::Result>> { +) -> eyre::Result> { // Resolve the did doc let Some(did_doc) = state.resolver.resolve_did(did).await? else { eyre::bail!("missing did doc"); @@ -232,7 +226,7 @@ // if there's no handles in aka or the expected is none, set to none in DB. if did_doc.also_known_as.as_ref().is_none_or(|v| v.is_empty()) || expected_handle.is_none() { - return Ok(Some(None)); + return Ok(None); } let expected = expected_handle.unwrap(); @@ -241,33 +235,33 @@ let expected_in_doc = did_doc.also_known_as.is_some_and(|v| { v.iter() .filter_map(|v| v.strip_prefix("at://")) - .any(|v| v == &expected) + .any(|v| v == expected) }); // if it isn't, set to invalid. if !expected_in_doc { tracing::warn!("Handle not in DID doc"); - return Ok(Some(None)); + return Ok(None); } // in theory, we can use com.atproto.identity.resolveHandle against a PDS, but that seems // like a way to end up with really sus handles. let Some(handle_did) = state.resolver.resolve_handle(&expected).await? else { - return Ok(Some(None)); + return Ok(None); }; // finally, check if the event did matches the handle, if not, set invalid, otherwise set the handle. if handle_did != did { - Ok(Some(None)) + Ok(None) } else { - Ok(Some(Some(expected))) + Ok(Some(expected)) } } #[instrument(skip_all, fields(seq = account.seq, repo = account.did))] async fn index_account( state: &RelayIndexerState, - conn: &mut AsyncPgConnection, + conn: &mut Object, rc: &mut MultiplexedConnection, account: AtpAccountEvent, ) -> eyre::Result<()> { @@ -279,7 +273,7 @@ let trigger_bf = if state.do_backfill && status == ActorStatus::Active { // check old status - if they exist (Some(*)), AND were previously != Active but not Deleted, // AND have a rev == null, then trigger backfill. - db::account_status_and_rev(conn, &account.did) + db::actor_get_status_and_rev(conn, &account.did) .await? .is_some_and(|(old_status, old_rev)| { old_rev.is_none() @@ -290,17 +284,12 @@ false }; - let sync_state = (!state.do_backfill).then_some(ActorSyncState::Synced); + let sync_state = match state.do_backfill { + true => ActorSyncState::Dirty, + false => ActorSyncState::Synced, + }; - db::upsert_actor( - conn, - &account.did, - None, - Some(status), - sync_state, - account.time, - ) - .await?; + db::actor_upsert(conn, &account.did, status, sync_state, account.time).await?; if trigger_bf { tracing::debug!("triggering backfill due to account coming out of inactive state"); @@ -313,11 +302,11 @@ #[instrument(skip_all, fields(seq = commit.seq, repo = commit.repo, rev = commit.rev))] async fn index_commit( state: &mut RelayIndexerState, - conn: &mut AsyncPgConnection, + conn: &mut Object, rc: &mut MultiplexedConnection, commit: AtpCommitEvent, ) -> eyre::Result<()> { - let (current_rev, sync_status) = db::get_repo_info(conn, &commit.repo).await?.unzip(); + let (sync_status, current_rev) = db::actor_get_repo_status(conn, &commit.repo).await?.unzip(); // what's the backfill status of this account? this respects locks held by the backfiller. // we should drop events for 'dirty' and queue 'processing' @@ -338,15 +327,8 @@ } // this is the first commit in an actor's repo - set them to Synced. - db::upsert_actor( - conn, - &commit.repo, - None, - None, - Some(ActorSyncState::Synced), - commit.time, - ) - .await?; + db::actor_set_sync_status(conn, &commit.repo, ActorSyncState::Synced, commit.time) + .await?; true } @@ -354,9 +336,19 @@ tracing::debug!("found new repo from commit"); let trigger_backfill = state.do_backfill && commit.since.is_some(); - let sync_state = (!trigger_backfill).then_some(ActorSyncState::Synced); + let sync_state = match trigger_backfill { + true => ActorSyncState::Dirty, + false => ActorSyncState::Synced, + }; - db::upsert_actor(conn, &commit.repo, None, None, sync_state, commit.time).await?; + db::actor_upsert( + conn, + &commit.repo, + ActorStatus::Active, + sync_state, + commit.time, + ) + .await?; if trigger_backfill { rc.rpush::<_, _, i32>("backfill_queue", commit.repo).await?; @@ -383,17 +375,12 @@ .await; if is_active { - conn.transaction::<_, diesel::result::Error, _>(|t| { - Box::pin(async move { - db::update_repo_version(t, &commit.repo, &commit.rev, commit.commit).await?; + let mut t = conn.transaction().await?; + db::actor_set_repo_state(&mut t, &commit.repo, &commit.rev, commit.commit).await?; - for op in &commit.ops { - process_op(t, &mut state.idxc_tx, &commit.repo, op, &blocks).await?; - } - Ok(true) - }) - }) - .await?; + for op in &commit.ops { + process_op(&mut t, &mut state.idxc_tx, &commit.repo, op, &blocks).await?; + } } else { let items = commit .ops @@ -402,7 +389,7 @@ .collect::>(); let items = serde_json::to_value(items).unwrap_or_default(); - db::write_backfill_row(conn, &commit.repo, &commit.rev, commit.commit, items).await?; + db::backfill_write_row(conn, &commit.repo, &commit.rev, commit.commit, items).await?; } Ok(()) @@ -456,12 +443,12 @@ #[inline(always)] async fn process_op( - conn: &mut AsyncPgConnection, + conn: &mut Transaction<'_>, deltas: &mut impl AggregateDeltaStore, repo: &str, op: &CommitOp, blocks: &HashMap>, -) -> diesel::QueryResult<()> { +) -> Result<(), tokio_postgres::Error> { let Some((collection_raw, rkey)) = op.path.split_once("/") else { tracing::warn!("op contained invalid path {}", op.path); return Ok(()); @@ -512,19 +499,19 @@ } pub async fn index_op( - conn: &mut AsyncPgConnection, + conn: &mut Transaction<'_>, deltas: &mut impl AggregateDeltaStore, repo: &str, cid: Cid, record: RecordTypes, at_uri: &str, rkey: &str, -) -> diesel::QueryResult<()> { +) -> Result<(), tokio_postgres::Error> { match record { RecordTypes::AppBskyActorProfile(record) => { if rkey == "self" { let labels = record.labels.clone(); - db::upsert_profile(conn, repo, cid, record).await?; + db::profile_upsert(conn, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; @@ -533,7 +520,7 @@ } RecordTypes::AppBskyFeedGenerator(record) => { let labels = record.labels.clone(); - let count = db::upsert_feedgen(conn, repo, cid, at_uri, record).await?; + let count = db::feedgen_upsert(conn, at_uri, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; @@ -545,7 +532,7 @@ } RecordTypes::AppBskyFeedLike(record) => { let subject = record.subject.uri.clone(); - let count = db::insert_like(conn, repo, at_uri, record).await?; + let count = db::like_insert(conn, at_uri, repo, record).await?; deltas .add_delta(&subject, AggregateType::Like, count as i32) @@ -575,7 +562,7 @@ }); let labels = record.labels.clone(); - db::insert_post(conn, repo, cid, at_uri, record).await?; + db::post_insert(conn, at_uri, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; } @@ -600,7 +587,7 @@ .contains(&records::PostgateEmbeddingRules::Disable); let disable_effective = has_disable_rule.then_some(record.created_at.naive_utc()); - db::upsert_postgate(conn, at_uri, cid, &record).await?; + db::postgate_upsert(conn, at_uri, cid, &record).await?; db::postgate_maintain_detaches( conn, @@ -614,7 +601,7 @@ deltas .incr(&record.subject.uri, AggregateType::Repost) .await; - db::insert_repost(conn, repo, at_uri, record).await?; + db::repost_insert(conn, at_uri, repo, record).await?; } RecordTypes::AppBskyFeedThreadgate(record) => { let split_aturi = record.post.rsplitn(4, '/').collect::>(); @@ -623,14 +610,14 @@ return Ok(()); } - db::upsert_threadgate(conn, at_uri, cid, record).await?; + db::threadgate_upsert(conn, at_uri, cid, record).await?; } RecordTypes::AppBskyGraphBlock(record) => { - db::insert_block(conn, repo, at_uri, record).await?; + db::block_insert(conn, at_uri, repo, record).await?; } RecordTypes::AppBskyGraphFollow(record) => { let subject = record.subject.clone(); - let count = db::insert_follow(conn, repo, at_uri, record).await?; + let count = db::follow_insert(conn, at_uri, repo, record).await?; deltas .add_delta(repo, AggregateType::Follow, count as i32) @@ -641,7 +628,7 @@ } RecordTypes::AppBskyGraphList(record) => { let labels = record.labels.clone(); - let count = db::upsert_list(conn, repo, at_uri, cid, record).await?; + let count = db::list_upsert(conn, at_uri, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; @@ -652,7 +639,7 @@ .await; } RecordTypes::AppBskyGraphListBlock(record) => { - db::insert_list_block(conn, repo, at_uri, record).await?; + db::list_block_insert(conn, at_uri, repo, record).await?; } RecordTypes::AppBskyGraphListItem(record) => { let split_aturi = record.list.rsplitn(4, '/').collect::>(); @@ -662,21 +649,21 @@ return Ok(()); } - db::insert_list_item(conn, at_uri, record).await?; + db::list_item_insert(conn, at_uri, record).await?; } RecordTypes::AppBskyGraphStarterPack(record) => { - let count = db::upsert_starterpack(conn, repo, cid, at_uri, record).await?; + let count = db::starter_pack_upsert(conn, at_uri, repo, cid, record).await?; deltas .add_delta(repo, AggregateType::ProfileStarterpack, count as i32) .await; } RecordTypes::AppBskyGraphVerification(record) => { - db::upsert_verification(conn, repo, cid, at_uri, record).await?; + db::verification_insert(conn, at_uri, repo, cid, record).await?; } RecordTypes::AppBskyLabelerService(record) => { if rkey == "self" { let labels = record.labels.clone(); - db::upsert_label_service(conn, repo, cid, record).await?; + db::labeler_upsert(conn, repo, cid, record).await?; if let Some(labels) = labels { db::maintain_self_labels(conn, repo, Some(cid), at_uri, labels).await?; @@ -685,43 +672,43 @@ } RecordTypes::ChatBskyActorDeclaration(record) => { if rkey == "self" { - db::upsert_chat_decl(conn, repo, record).await?; + db::chat_decl_upsert(conn, repo, record).await?; } } } - db::write_record(conn, at_uri, repo, cid).await?; + db::record_upsert(conn, at_uri, repo, cid).await?; Ok(()) } pub async fn index_op_delete( - conn: &mut AsyncPgConnection, + conn: &mut Transaction<'_>, deltas: &mut impl AggregateDeltaStore, repo: &str, collection: CollectionType, at_uri: &str, -) -> diesel::QueryResult<()> { +) -> Result<(), tokio_postgres::Error> { match collection { - CollectionType::BskyProfile => db::delete_profile(conn, repo).await?, - CollectionType::BskyBlock => db::delete_block(conn, at_uri).await?, + CollectionType::BskyProfile => db::profile_delete(conn, repo).await?, + CollectionType::BskyBlock => db::block_delete(conn, at_uri).await?, CollectionType::BskyFeedGen => { - let count = db::delete_feedgen(conn, at_uri).await?; + let count = db::feedgen_delete(conn, at_uri).await?; deltas .add_delta(repo, AggregateType::ProfileFeed, -(count as i32)) .await; count } CollectionType::BskyFeedLike => { - if let Some(subject) = db::delete_like(conn, at_uri).await? { + if let Some(subject) = db::like_delete(conn, at_uri).await? { deltas.decr(&subject, AggregateType::Like).await; } 0 } CollectionType::BskyFeedPost => { - let post_info = db::get_post_info_for_delete(conn, at_uri).await?; + let post_info = db::post_get_info_for_delete(conn, at_uri).await?; - db::delete_post(conn, at_uri).await?; + db::post_delete(conn, at_uri).await?; if let Some((reply_to, embed)) = post_info { deltas.decr(repo, AggregateType::ProfilePost).await; @@ -735,44 +722,44 @@ 0 } - CollectionType::BskyFeedPostgate => db::delete_postgate(conn, at_uri).await?, + CollectionType::BskyFeedPostgate => db::postgate_delete(conn, at_uri).await?, CollectionType::BskyFeedRepost => { - if let Some(subject) = db::delete_repost(conn, at_uri).await? { + if let Some(subject) = db::repost_delete(conn, at_uri).await? { deltas.decr(&subject, AggregateType::Repost).await; } 0 } - CollectionType::BskyFeedThreadgate => db::delete_threadgate(conn, at_uri).await?, + CollectionType::BskyFeedThreadgate => db::threadgate_delete(conn, at_uri).await?, CollectionType::BskyFollow => { - if let Some(followee) = db::delete_follow(conn, at_uri).await? { + if let Some(followee) = db::follow_delete(conn, at_uri).await? { deltas.decr(&followee, AggregateType::Follower).await; deltas.decr(repo, AggregateType::Follow).await; } 0 } CollectionType::BskyList => { - let count = db::delete_list(conn, at_uri).await?; + let count = db::list_delete(conn, at_uri).await?; deltas .add_delta(repo, AggregateType::ProfileList, -(count as i32)) .await; count } - CollectionType::BskyListBlock => db::delete_list_block(conn, at_uri).await?, - CollectionType::BskyListItem => db::delete_list_item(conn, at_uri).await?, + CollectionType::BskyListBlock => db::list_block_delete(conn, at_uri).await?, + CollectionType::BskyListItem => db::list_item_delete(conn, at_uri).await?, CollectionType::BskyStarterPack => { - let count = db::delete_starterpack(conn, at_uri).await?; + let count = db::starter_pack_delete(conn, at_uri).await?; deltas .add_delta(repo, AggregateType::ProfileStarterpack, -(count as i32)) .await; count } - CollectionType::BskyVerification => db::delete_verification(conn, at_uri).await?, - CollectionType::BskyLabelerService => db::delete_label_service(conn, at_uri).await?, - CollectionType::ChatActorDecl => db::delete_chat_decl(conn, at_uri).await?, + CollectionType::BskyVerification => db::verification_delete(conn, at_uri).await?, + CollectionType::BskyLabelerService => db::labeler_delete(conn, at_uri).await?, + CollectionType::ChatActorDecl => db::chat_decl_delete(conn, at_uri).await?, _ => unreachable!(), }; - db::delete_record(conn, at_uri).await?; + db::record_delete(conn, at_uri).await?; Ok(()) } diff --git a/consumer/src/indexer/records.rs b/consumer/src/indexer/records.rs --- a/consumer/src/indexer/records.rs +++ b/consumer/src/indexer/records.rs @@ -9,7 +9,7 @@ use lexica::com_atproto::moderation::{ReasonType, SubjectType}; use serde::{Deserialize, Serialize}; -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct StrongRef { #[serde( deserialize_with = "utils::cid_from_string", @@ -19,7 +19,7 @@ pub uri: String, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(tag = "$type")] #[serde(rename = "blob")] #[serde(rename_all = "camelCase")] @@ -43,7 +43,7 @@ pub created_at: Option>, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(untagged)] pub enum EmbedOuter { Bsky(AppBskyEmbed), @@ -83,7 +83,7 @@ } } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(tag = "$type")] pub enum AppBskyEmbed { #[serde(rename = "app.bsky.embed.images")] @@ -117,13 +117,13 @@ } } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct AppBskyEmbedImages { pub images: Vec, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct EmbedImage { pub image: Blob, @@ -132,7 +132,7 @@ pub aspect_ratio: Option, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] #[serde(rename_all = "camelCase")] pub struct AppBskyEmbedVideo { pub video: Blob, @@ -144,18 +144,18 @@ pub aspect_ratio: Option, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct EmbedVideoCaptions { pub lang: String, pub file: Blob, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct AppBskyEmbedExternal { pub external: EmbedExternal, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct EmbedExternal { pub uri: String, pub title: String, @@ -164,12 +164,12 @@ pub thumb: Option, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct AppBskyEmbedRecord { pub record: StrongRef, } -#[derive(Debug, Deserialize, Serialize)] +#[derive(Clone, Debug, Deserialize, Serialize)] pub struct AppBskyEmbedRecordWithMedia { pub record: AppBskyEmbedRecord, pub media: Box, diff --git a/consumer/src/label_indexer/mod.rs b/consumer/src/label_indexer/mod.rs --- a/consumer/src/label_indexer/mod.rs +++ b/consumer/src/label_indexer/mod.rs @@ -6,18 +6,17 @@ use std::sync::Arc; use std::time::Duration; use tokio::sync::mpsc::{channel, Receiver, Sender}; +use tokio::sync::watch::Receiver as WatchReceiver; use tokio::task::JoinHandle; use tokio::time::Instant; use tokio_postgres::binary_copy::BinaryCopyInWriter; use tokio_postgres::types::Type; -use tokio_postgres::NoTls; use tracing::instrument; -use tokio::sync::watch::Receiver as WatchReceiver; const LABELER_SERVICE_ID: &str = "#atproto_labeler"; pub struct LabelServiceManager { - client: tokio_postgres::Client, + conn: deadpool_postgres::Object, rx: Receiver, resolver: Arc, services: HashMap>, @@ -27,23 +26,16 @@ impl LabelServiceManager { pub async fn new( - pg_url: &str, + pool: deadpool_postgres::Pool, resolver: Arc, resume: sled::Db, 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 conn = pool.get().await?; let (tx, rx) = channel(8); let lsm = LabelServiceManager { - client, + conn, rx, resolver, resume, @@ -112,7 +104,7 @@ continue; } tracing::debug!("got {} labels", buf.len()); - store_labels(&mut self.client, &buf).await + store_labels(&mut self.conn, &buf).await } }; @@ -169,7 +161,7 @@ let count = binary_writer.finish().await?; - t.execute(include_str!("../sql/label_copy_upsert.sql"), &[]) + t.execute(include_str!("../db/sql/label_copy_upsert.sql"), &[]) .await?; t.commit().await?; @@ -186,7 +178,11 @@ user_agent: String, db_tx: Sender, ) { - let start_seq = resume.get(&service_did).ok().flatten().and_then(crate::utils::u64_from_ivec); + let start_seq = resume + .get(&service_did) + .ok() + .flatten() + .and_then(crate::utils::u64_from_ivec); if let Some(start_seq) = start_seq { tracing::info!("starting {service_did} label consumer from {start_seq}"); diff --git a/consumer/src/sql/label_copy_upsert.sql b/consumer/src/sql/label_copy_upsert.sql deleted file mode 100644 --- a/consumer/src/sql/label_copy_upsert.sql +++ /dev/null @@ -1,8 +0,0 @@ -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/consumer/src/db/sql/feedgen_upsert.sql b/consumer/src/db/sql/feedgen_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/feedgen_upsert.sql @@ -0,0 +1,11 @@ +INSERT INTO feedgens (at_uri, owner, cid, service_did, content_mode, name, description, description_facets, avatar_cid, + created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) +ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, + service_did=EXCLUDED.service_did, + content_mode=EXCLUDED.content_mode, + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + avatar_cid=EXCLUDED.avatar_cid, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/label_copy_upsert.sql b/consumer/src/db/sql/label_copy_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/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/consumer/src/db/sql/label_defs_upsert.sql b/consumer/src/db/sql/label_defs_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/label_defs_upsert.sql @@ -0,0 +1,9 @@ +INSERT INTO labeler_defs (labeler, label_identifier, severity, blurs, default_setting, adult_only, locales) +VALUES ($1, $2, $3, $4, $5, $6, $7) +ON CONFLICT (labeler, label_identifier) DO UPDATE + SET severity=EXCLUDED.severity, + blurs=EXCLUDED.blurs, + default_setting=EXCLUDED.default_setting, + adult_only=EXCLUDED.adult_only, + locales=EXCLUDED.locales, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/label_service_upsert.sql b/consumer/src/db/sql/label_service_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/label_service_upsert.sql @@ -0,0 +1,7 @@ +INSERT INTO labelers (did, cid, reasons, subject_types, subject_collections) +VALUES ($1, $2, $3, $4, $5) +ON CONFLICT (did) DO UPDATE SET cid=EXCLUDED.cid, + reasons=EXCLUDED.reasons, + subject_types=EXCLUDED.subject_types, + subject_collections=EXCLUDED.subject_collections, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/list_upsert.sql b/consumer/src/db/sql/list_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/list_upsert.sql @@ -0,0 +1,9 @@ +INSERT INTO lists (at_uri, owner, cid, list_type, name, description, description_facets, avatar_cid, created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) +ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, + list_type=EXCLUDED.list_type, + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + avatar_cid=EXCLUDED.avatar_cid, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/post_insert.sql b/consumer/src/db/sql/post_insert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/post_insert.sql @@ -0,0 +1,4 @@ +INSERT INTO posts (at_uri, did, cid, record, content, facets, languages, tags, parent_uri, parent_cid, root_uri, + root_cid, embed, embed_subtype, created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15) +ON CONFLICT DO NOTHING \ No newline at end of file diff --git a/consumer/src/db/sql/postgate_upsert.sql b/consumer/src/db/sql/postgate_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/postgate_upsert.sql @@ -0,0 +1,7 @@ +INSERT INTO postgates (at_uri, cid, post_uri, detached, rules, created_at) +VALUES ($1, $2, $3, $4, $5, $6) +ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, + post_uri=EXCLUDED.post_uri, + detached=EXCLUDED.detached, + rules=EXCLUDED.rules, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/profile_upsert.sql b/consumer/src/db/sql/profile_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/profile_upsert.sql @@ -0,0 +1,13 @@ +INSERT INTO profiles (did, cid, avatar_cid, banner_cid, display_name, description, pinned_uri, pinned_cid, + joined_sp_uri, joined_sp_cid, created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) +ON CONFLICT (did) DO UPDATE SET cid=EXCLUDED.cid, + avatar_cid=EXCLUDED.avatar_cid, + banner_cid=EXCLUDED.banner_cid, + display_name=EXCLUDED.display_name, + description=EXCLUDED.description, + pinned_uri=EXCLUDED.pinned_uri, + pinned_cid=EXCLUDED.pinned_cid, + joined_sp_uri=EXCLUDED.joined_sp_uri, + joined_sp_cid=EXCLUDED.joined_sp_cid, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/starterpack_upsert.sql b/consumer/src/db/sql/starterpack_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/starterpack_upsert.sql @@ -0,0 +1,10 @@ +INSERT INTO starterpacks (at_uri, owner, cid, record, name, description, description_facets, list, feeds, created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10) +ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, + record=EXCLUDED.record, + name=EXCLUDED.name, + description=EXCLUDED.description, + description_facets=EXCLUDED.description_facets, + list=EXCLUDED.list, + feeds=EXCLUDED.feeds, + indexed_at=NOW() \ No newline at end of file diff --git a/consumer/src/db/sql/threadgate_upsert.sql b/consumer/src/db/sql/threadgate_upsert.sql new file mode 100644 --- /dev/null +++ b/consumer/src/db/sql/threadgate_upsert.sql @@ -0,0 +1,8 @@ +INSERT INTO threadgates (at_uri, cid, post_uri, hidden_replies, allow, allowed_lists, record, created_at) +VALUES ($1, $2, $3, $4, $5, $6, $7, $8) +ON CONFLICT (at_uri) DO UPDATE SET cid=EXCLUDED.cid, + hidden_replies=EXCLUDED.hidden_replies, + allow=EXCLUDED.allow, + allowed_lists=EXCLUDED.allowed_lists, + record=EXCLUDED.record, + indexed_at=NOW() \ No newline at end of file