diff --git a/patches/components/constellation/atproto.rs.patch b/patches/components/constellation/atproto.rs.patch index 2eb8a77..5beb4bf 100644 --- a/patches/components/constellation/atproto.rs.patch +++ b/patches/components/constellation/atproto.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -0,0 +1,478 @@ +@@ -0,0 +1,458 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +use std::collections::HashSet; @@ -9,7 +9,8 @@ + +use ipc_channel::ipc::channel; +use log::{debug, error}; -+use net::async_runtime::spawn_task; ++use net::async_runtime::{spawn_blocking, spawn_task}; ++use net::atproto::db_writer; +use net::atproto::session::SessionClient; +use net::atproto::store::RepoStore; +use net::atproto::sync::sync_repo; @@ -42,6 +43,13 @@ + + debug!("ATProto session: {session:?}, config_dir: {config_dir:?}"); + ++ // Start the dedicated SQLite writer thread early so the DB is created + ++ // migrated (and its WAL kept live for read-only readers) before the first ++ // sync or query. ++ if let Some(ref dir) = config_dir { ++ db_writer::init(dir.clone()); ++ } ++ + Self { + resource_thread, + session: Arc::new(Mutex::new(session)), @@ -122,19 +130,16 @@ + return; + }; + -+ spawn_task(async move { -+ let message = match RepoStore::open(&config_dir) { -+ Ok(store) => match store.connect() { -+ Ok(conn) => { -+ let rev = RepoStore::repo_rev(&conn, &did).ok().flatten(); -+ let record_count = RepoStore::record_count(&conn, &did).unwrap_or(0); -+ AtProtoResult::RepoStatus { -+ did, -+ rev, -+ record_count, -+ } -+ }, -+ Err(_) => AtProtoResult::Error, ++ spawn_blocking(move || { ++ let message = match RepoStore::read_only(&config_dir).connect_readonly() { ++ Ok(conn) => { ++ let rev = RepoStore::repo_rev(&conn, &did).ok().flatten(); ++ let record_count = RepoStore::record_count(&conn, &did).unwrap_or(0); ++ AtProtoResult::RepoStatus { ++ did, ++ rev, ++ record_count, ++ } + }, + Err(_) => AtProtoResult::Error, + }; @@ -145,10 +150,10 @@ + } + + fn clear_repo(&self, did: String, response: GenericCallback) { -+ let Some(config_dir) = self.config_dir.clone() else { ++ if self.config_dir.is_none() { + let _ = response.send(AtProtoResult::Error); + return; -+ }; ++ } + + // Claim the DID so a clear can't race an in-flight sync of the same repo. + { @@ -162,7 +167,7 @@ + let sync_in_flight = Arc::clone(&self.sync_in_flight); + + spawn_task(async move { -+ let cleared = clear_repo_data(&config_dir, &did); ++ let cleared = db_writer::clear(did.clone()).await; + sync_in_flight.lock().remove(&did); + + let message = if cleared { @@ -188,13 +193,11 @@ + return; + }; + -+ spawn_task(async move { -+ let message = match RepoStore::open(&config_dir) { -+ Ok(store) => match store.search_text(&needle, max_results) { -+ Some(hits) => AtProtoResult::TextResults(hits), -+ None => AtProtoResult::Error, -+ }, -+ Err(_) => AtProtoResult::Error, ++ spawn_blocking(move || { ++ let message = match RepoStore::read_only(&config_dir).search_text(&needle, max_results) ++ { ++ Some(hits) => AtProtoResult::TextResults(hits), ++ None => AtProtoResult::Error, + }; + if let Err(err) = response.send(message) { + error!("Failed to send text search results: {err:?}"); @@ -208,13 +211,10 @@ + return; + }; + -+ spawn_task(async move { -+ let message = match RepoStore::open(&config_dir) { -+ Ok(store) => match store.search_url(&url, exact) { -+ Some(urls) => AtProtoResult::UrlResults(urls), -+ None => AtProtoResult::Error, -+ }, -+ Err(_) => AtProtoResult::Error, ++ spawn_blocking(move || { ++ let message = match RepoStore::read_only(&config_dir).search_url(&url, exact) { ++ Some(urls) => AtProtoResult::UrlResults(urls), ++ None => AtProtoResult::Error, + }; + if let Err(err) = response.send(message) { + error!("Failed to send url search results: {err:?}"); @@ -233,16 +233,14 @@ + return; + }; + -+ spawn_task(async move { -+ let message = match RepoStore::open(&config_dir) { -+ Ok(store) => match store.query(&sql, ¶ms) { -+ Ok(json) => AtProtoResult::QueryResult(json), -+ Err(err) => { -+ error!("queryStore failed: {err}"); -+ AtProtoResult::Error -+ }, ++ spawn_blocking(move || { ++ let store = RepoStore::read_only(&config_dir); ++ let message = match store.query(&sql, ¶ms) { ++ Ok(json) => AtProtoResult::QueryResult(json), ++ Err(err) => { ++ error!("queryStore failed: {err}"); ++ AtProtoResult::Error + }, -+ Err(_) => AtProtoResult::Error, + }; + if let Err(err) = response.send(message) { + error!("Failed to send query result: {err:?}"); @@ -461,21 +459,3 @@ + } + } +} -+ -+/// Open the repo store and wipe all data for `did` in a single transaction. -+/// Returns `true` on success. -+fn clear_repo_data(config_dir: &std::path::Path, did: &str) -> bool { -+ let Ok(store) = RepoStore::open(config_dir) else { -+ return false; -+ }; -+ let Ok(mut conn) = store.connect() else { -+ return false; -+ }; -+ let Ok(tx) = conn.transaction() else { -+ return false; -+ }; -+ if RepoStore::clear_repo(&tx, did).is_err() { -+ return false; -+ } -+ tx.commit().is_ok() -+} diff --git a/patches/components/net/async_runtime.rs.patch b/patches/components/net/async_runtime.rs.patch new file mode 100644 index 0000000..222ddce --- /dev/null +++ b/patches/components/net/async_runtime.rs.patch @@ -0,0 +1,23 @@ +--- original ++++ modified +@@ -81,6 +81,20 @@ + .spawn(task); + } + ++/// Spawn a blocking closure on the runtime's dedicated blocking-thread pool. ++/// Use this for synchronous work (e.g. SQLite) that must NOT run on the async ++/// worker threads: a blocked worker can't drive other tasks or the I/O reactor, ++/// and enough of them blocking at once stalls the whole shared `net` runtime. ++pub fn spawn_blocking(task: F) ++where ++ F: FnOnce() + Send + 'static, ++{ ++ ASYNC_RUNTIME_HANDLE ++ .get() ++ .expect("Runtime handle should be initialized on start-up") ++ .spawn_blocking(task); ++} ++ + /// Spawn a blocking task using the handle to the runtime. + pub fn spawn_blocking_task(task: F) -> F::Output + where diff --git a/patches/components/net/atproto/db_writer.rs.patch b/patches/components/net/atproto/db_writer.rs.patch new file mode 100644 index 0000000..ccbed28 --- /dev/null +++ b/patches/components/net/atproto/db_writer.rs.patch @@ -0,0 +1,126 @@ +--- original ++++ modified +@@ -0,0 +1,123 @@ ++// SPDX-License-Identifier: AGPL-3.0-or-later ++ ++//! Dedicated SQLite writer thread. ++//! ++//! All writes to the repo store run on one background thread that owns a single ++//! persistent read-write connection. Callers send a command plus a oneshot reply ++//! over a channel and await the result. ++//! ++//! Reads do not go through here: WAL allows concurrent readers, so they use ++//! read-only connections on `spawn_blocking`. Keeping this connection open also ++//! keeps the WAL live, so those read-only connections never hit WAL recovery. ++ ++use std::path::{Path, PathBuf}; ++use std::sync::OnceLock; ++use std::sync::mpsc::{Sender, channel}; ++ ++use log::error; ++use rusqlite::Connection; ++use tokio::sync::oneshot; ++ ++use super::store::RepoStore; ++use super::sync::{SyncError, SyncOutcome, index_car}; ++ ++enum WriteCommand { ++ Index { ++ did: String, ++ car_bytes: Vec, ++ reply: oneshot::Sender>, ++ }, ++ Clear { ++ did: String, ++ reply: oneshot::Sender, ++ }, ++} ++ ++static WRITER: OnceLock> = OnceLock::new(); ++ ++/// Start the writer thread for `config_dir`. Opening the connection here also creates ++/// and migrates the DB up front, before any reader needs it. ++pub fn init(config_dir: PathBuf) { ++ WRITER.get_or_init(|| { ++ let (tx, rx) = channel::(); ++ let spawned = std::thread::Builder::new() ++ .name("atproto-db-writer".to_owned()) ++ .spawn(move || { ++ let mut conn = match open_connection(&config_dir) { ++ Ok(conn) => conn, ++ Err(e) => { ++ error!("[atproto] writer thread failed to open DB: {e}"); ++ return; ++ }, ++ }; ++ // Drain commands until every Sender is dropped (never, in ++ // practice: the Sender lives in a process-lifetime OnceLock). ++ while let Ok(cmd) = rx.recv() { ++ match cmd { ++ WriteCommand::Index { ++ did, ++ car_bytes, ++ reply, ++ } => { ++ let _ = reply.send(index_car(&did, &car_bytes, &mut conn)); ++ }, ++ WriteCommand::Clear { did, reply } => { ++ let _ = reply.send(clear_repo_data(&mut conn, &did)); ++ }, ++ } ++ } ++ }); ++ if let Err(e) = spawned { ++ error!("[atproto] failed to spawn writer thread: {e}"); ++ } ++ tx ++ }); ++} ++ ++/// Decode + index a CAR on the writer thread. ++pub async fn index(did: String, car_bytes: Vec) -> Result { ++ let (reply, rx) = oneshot::channel(); ++ send(WriteCommand::Index { ++ did, ++ car_bytes, ++ reply, ++ })?; ++ rx.await ++ .map_err(|_| SyncError::Store("writer thread dropped reply".to_owned()))? ++} ++ ++/// Wipe all stored data for `did` on the writer thread. `true` on success. ++pub async fn clear(did: String) -> bool { ++ let (reply, rx) = oneshot::channel(); ++ if send(WriteCommand::Clear { did, reply }).is_err() { ++ return false; ++ } ++ rx.await.unwrap_or(false) ++} ++ ++fn send(cmd: WriteCommand) -> Result<(), SyncError> { ++ match WRITER.get() { ++ Some(tx) => tx ++ .send(cmd) ++ .map_err(|_| SyncError::Store("writer thread unavailable".to_owned())), ++ None => Err(SyncError::Store("writer thread not initialized".to_owned())), ++ } ++} ++ ++fn open_connection(config_dir: &Path) -> rusqlite::Result { ++ // RepoStore::open creates + migrates the schema; then take a fresh read-write ++ // connection for the thread to own for its lifetime. ++ let store = RepoStore::open(config_dir)?; ++ store.connect() ++} ++ ++/// Wipe a repo's rows in a single transaction on the writer's connection. ++fn clear_repo_data(conn: &mut Connection, did: &str) -> bool { ++ let Ok(tx) = conn.transaction() else { ++ return false; ++ }; ++ if RepoStore::clear_repo(&tx, did).is_err() { ++ return false; ++ } ++ tx.commit().is_ok() ++} diff --git a/patches/components/net/atproto/store.rs.patch b/patches/components/net/atproto/store.rs.patch index 20a3d6a..966190d 100644 --- a/patches/components/net/atproto/store.rs.patch +++ b/patches/components/net/atproto/store.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -0,0 +1,466 @@ +@@ -0,0 +1,559 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//! On-disk store for synced ATProto repositories. @@ -34,6 +34,9 @@ +/// the IPC payload regardless of the caller's SQL. +const MAX_QUERY_ROWS: usize = 1000; + ++/// Handle to the on-disk store. Cheap to clone (just a path), so blocking work ++/// can be moved onto the runtime's blocking pool. ++#[derive(Clone)] +pub struct RepoStore { + db_path: PathBuf, +} @@ -56,22 +59,29 @@ + pub fn connect(&self) -> rusqlite::Result { + let conn = Connection::open(&self.db_path)?; + conn.pragma_update(None, "journal_mode", "WAL")?; -+ // The sync write-phase is serialized in `sync.rs`, but other writers -+ // (e.g. clear_repo) and WAL checkpoints can still briefly hold the lock; -+ // wait rather than fail immediately. ++ conn.pragma_update(None, "synchronous", "NORMAL")?; ++ // Writes run on the single writer thread, but reads and WAL checkpoints ++ // can still briefly hold a lock; wait rather than fail immediately. + conn.busy_timeout(std::time::Duration::from_secs(30))?; + Ok(conn) + } + -+ /// Open a **read-only** connection — used by the user-facing `query` API so -+ /// arbitrary SQL can't write, DROP, or ATTACH. WAL reads work because the -+ /// indexer (a writer) keeps the `-wal`/`-shm` files present. ++ /// Open a **read-only** connection, used by the user-facing `query` API so ++ /// arbitrary SQL can't write, DROP, or ATTACH. + pub fn connect_readonly(&self) -> rusqlite::Result { + let conn = Connection::open_with_flags(&self.db_path, OpenFlags::SQLITE_OPEN_READ_ONLY)?; + conn.busy_timeout(std::time::Duration::from_secs(10))?; + Ok(conn) + } + ++ /// A store handle for read-only callers Pair with ++ /// [`Self::connect_readonly`]. ++ pub fn read_only(config_dir: &Path) -> Self { ++ Self { ++ db_path: config_dir.join(DB_FILE), ++ } ++ } ++ + fn migrate(&self, conn: &Connection) -> rusqlite::Result<()> { + conn.execute_batch( + "CREATE TABLE IF NOT EXISTS repos ( @@ -94,7 +104,11 @@ + json TEXT NOT NULL, + PRIMARY KEY (did, collection, rkey) + ); -+ CREATE INDEX IF NOT EXISTS records_collection ON records(did, collection);", ++ CREATE INDEX IF NOT EXISTS records_collection ON records(did, collection); ++ CREATE INDEX IF NOT EXISTS records_created_at ++ ON records(json_extract(json, '$.createdAt') DESC); ++ CREATE INDEX IF NOT EXISTS records_did_created_at ++ ON records(did, json_extract(json, '$.createdAt') DESC);", + ) + } + @@ -153,7 +167,7 @@ + needle: &str, + max_results: u32, + ) -> Option> { -+ let conn = self.connect().ok()?; ++ let conn = self.connect_readonly().ok()?; + super::textindex::search_text(&conn, needle, max_results).ok() + } + @@ -164,7 +178,7 @@ + url: &str, + exact: bool, + ) -> Option> { -+ let conn = self.connect().ok()?; ++ let conn = self.connect_readonly().ok()?; + super::textindex::search_url(&conn, url, exact).ok() + } + @@ -466,4 +480,83 @@ + + std::fs::remove_dir_all(&tmp).ok(); + } ++ ++ #[test] ++ fn feed_order_uses_created_at_index() { ++ let tmp = std::env::temp_dir().join(format!("beaver-feedidx-test-{}", std::process::id())); ++ std::fs::create_dir_all(&tmp).unwrap(); ++ let store = store_in(&tmp); ++ let conn = store.connect().unwrap(); ++ RepoStore::upsert_record( ++ &conn, ++ "did:plc:a", ++ "app.bsky.feed.post", ++ "r1", ++ "c1", ++ r#"{"createdAt":"2024-01-02T00:00:00Z"}"#, ++ ) ++ .unwrap(); ++ RepoStore::upsert_record( ++ &conn, ++ "did:plc:a", ++ "app.bsky.feed.post", ++ "r2", ++ "c2", ++ r#"{"createdAt":"2024-01-01T00:00:00Z"}"#, ++ ) ++ .unwrap(); ++ ++ // The feed's reverse-chron query must ride the expression index, not ++ // full-scan + sort every record. EXPLAIN QUERY PLAN's `detail` column ++ // (index 3) names the index and would say "USE TEMP B-TREE" if it sorted. ++ let detail: Vec = { ++ let mut stmt = conn ++ .prepare( ++ "EXPLAIN QUERY PLAN ++ SELECT did, collection, rkey, cid, json FROM records ++ WHERE json_extract(json, '$.createdAt') IS NOT NULL ++ ORDER BY json_extract(json, '$.createdAt') DESC ++ LIMIT 100", ++ ) ++ .unwrap(); ++ let rows = stmt.query_map([], |r| r.get::<_, String>(3)).unwrap(); ++ rows.map(Result::unwrap).collect() ++ }; ++ let plan = detail.join(" | "); ++ assert!( ++ plan.contains("records_created_at"), ++ "feed ORDER BY should use the createdAt index; plan: {plan}" ++ ); ++ assert!( ++ !plan.to_uppercase().contains("TEMP B-TREE"), ++ "feed ORDER BY should not sort; plan: {plan}" ++ ); ++ ++ // The per-DID (byline-filtered) feed must ride the composite index. ++ let detail_did: Vec = { ++ let mut stmt = conn ++ .prepare( ++ "EXPLAIN QUERY PLAN ++ SELECT did, collection, rkey, cid, json FROM records ++ WHERE did = 'did:plc:a' ++ AND json_extract(json, '$.createdAt') IS NOT NULL ++ ORDER BY json_extract(json, '$.createdAt') DESC ++ LIMIT 100", ++ ) ++ .unwrap(); ++ let rows = stmt.query_map([], |r| r.get::<_, String>(3)).unwrap(); ++ rows.map(Result::unwrap).collect() ++ }; ++ let plan_did = detail_did.join(" | "); ++ assert!( ++ plan_did.contains("records_did_created_at"), ++ "per-DID feed should use the composite index; plan: {plan_did}" ++ ); ++ assert!( ++ !plan_did.to_uppercase().contains("TEMP B-TREE"), ++ "per-DID feed should not sort; plan: {plan_did}" ++ ); ++ ++ std::fs::remove_dir_all(&tmp).ok(); ++ } +} diff --git a/patches/components/net/atproto/sync.rs.patch b/patches/components/net/atproto/sync.rs.patch index f3b0206..0aa5673 100644 --- a/patches/components/net/atproto/sync.rs.patch +++ b/patches/components/net/atproto/sync.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -0,0 +1,257 @@ +@@ -0,0 +1,259 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//! Syncing a single ATProto repository (by DID) into the local store. @@ -16,13 +16,12 @@ +//! 5. Garbage-collect the retained block set down to the nodes reachable from +//! the new root, and store the new `rev`. +//! -+//! If an incremental (`since`) walk can't be completed from the merged blocks — -+//! e.g. the PDS returned an insufficient diff — the sync transparently retries ++//! If an incremental (`since`) walk can't be completed from the merged blocks, ++//! e.g. the PDS returned an insufficient diff, the sync transparently retries +//! with a full fetch. + +use std::collections::HashMap; +use std::path::PathBuf; -+use std::sync::OnceLock; +use std::time::{SystemTime, UNIX_EPOCH}; + +use ipld_core::cid::Cid; @@ -30,9 +29,9 @@ +use log::{info, warn}; +use net_traits::CoreResourceThread; +use net_traits::fetch::utils::fetch_bytes; ++use rusqlite::Connection; +use servo_url::ServoUrl; +use sync_wrapper::SyncWrapper; -+use tokio::sync::Mutex; + +use crate::atproto::mst::MstError; +use crate::atproto::pds::get_endpoint_for_subject; @@ -40,18 +39,6 @@ +use crate::atproto::textindex::ipld_to_json; +use crate::atproto::{car, commit, mst, textindex}; + -+/// Process-wide lock serializing the decode-and-write phase of a sync. -+/// -+/// SQLite allows only a single writer at a time Syncs run concurrently for different DIDs, -+/// so without this their `index_car` write transactions race on the one writer and -+/// the losers time out with "database is locked", especially during first-time -+/// backfills, which hold the writer for the whole repo. Network fetches (the -+/// slow part) still overlap; only the fast, local write step is serialized. -+fn db_write_lock() -> &'static Mutex<()> { -+ static LOCK: OnceLock> = OnceLock::new(); -+ LOCK.get_or_init(|| Mutex::new(())) -+} -+ +#[derive(Debug)] +pub struct SyncOutcome { + pub did: String, @@ -89,19 +76,13 @@ + .await + .map_err(|e| SyncError::Pds(format!("resolve {did}: {e:?}")))?; + -+ let store = RepoStore::open(&config_dir).map_err(|e| SyncError::Store(format!("open: {e}")))?; -+ -+ let prior_rev = { -+ let conn = store -+ .connect() -+ .map_err(|e| SyncError::Store(format!("connect: {e}")))?; -+ RepoStore::repo_rev(&conn, did).map_err(|e| SyncError::Store(format!("read rev: {e}")))? -+ }; ++ // The stored rev (for an incremental `since=`) is a read, so take it on the ++ // blocking pool with a read-only connection. ++ let prior_rev = read_prior_rev(config_dir, did.to_owned()).await; + + match try_sync( + did, + &endpoint, -+ &store, + resource_thread.clone(), + prior_rev.as_deref(), + ) @@ -109,16 +90,27 @@ + { + Err(SyncError::IncompleteDiff) if prior_rev.is_some() => { + warn!("[atproto] incremental sync for {did} incomplete; retrying full fetch"); -+ try_sync(did, &endpoint, &store, resource_thread, None).await ++ try_sync(did, &endpoint, resource_thread, None).await + }, + other => other, + } +} + ++/// Read the stored rev for `did` off the async worker threads (read-only ++/// connection on the blocking pool). `None` if the repo is new or unreadable. ++async fn read_prior_rev(config_dir: PathBuf, did: String) -> Option { ++ tokio::task::spawn_blocking(move || { ++ let conn = RepoStore::read_only(&config_dir).connect_readonly().ok()?; ++ RepoStore::repo_rev(&conn, &did).ok().flatten() ++ }) ++ .await ++ .ok() ++ .flatten() ++} ++ +async fn try_sync( + did: &str, + endpoint: &ServoUrl, -+ store: &RepoStore, + resource_thread: CoreResourceThread, + since: Option<&str>, +) -> Result { @@ -135,27 +127,26 @@ + .await + .map_err(|e| SyncError::Fetch(format!("getRepo (since={since:?}): {e:?}")))?; + -+ // Serialize the decode+write phase across concurrent syncs: SQLite has a -+ // single writer, so parallel `index_car` transactions would otherwise -+ // contend and time out ("database is locked"). -+ let _write_guard = db_write_lock().lock().await; -+ index_car(did, &car_bytes, store) ++ // Hand the blocking CAR-decode + write transaction to the dedicated writer ++ // thread: it owns the single write connection (so writes serialize without a ++ // mutex) and runs off the async worker threads. ++ super::db_writer::index(did.to_owned(), car_bytes).await +} + -+/// The pure decode-and-index step (no network). Synchronous: CAR parse, MST -+/// walk, and SQLite writes. -+fn index_car(did: &str, car_bytes: &[u8], store: &RepoStore) -> Result { ++/// The pure decode-and-index step (no network), run on the writer thread with ++/// its persistent connection. Synchronous: CAR parse, MST walk, SQLite writes. ++pub(crate) fn index_car( ++ did: &str, ++ car_bytes: &[u8], ++ conn: &mut Connection, ++) -> Result { + let car = car::parse(car_bytes).map_err(|e| SyncError::Decode(format!("CAR parse: {e:?}")))?; + let root = *car.roots.first().ok_or(SyncError::NoRoot)?; + -+ let mut conn = store -+ .connect() -+ .map_err(|e| SyncError::Store(format!("connect: {e}")))?; -+ + // Merge retained MST node blocks with everything in this CAR (new nodes, + // new record blocks, the new commit). Unchanged subtrees come from the + // retained set, so the new tree is fully walkable. -+ let mut merged = RepoStore::load_node_blocks(&conn, did) ++ let mut merged = RepoStore::load_node_blocks(conn, did) + .map_err(|e| SyncError::Store(format!("load blocks: {e}")))?; + for (cid, bytes) in &car.blocks { + merged.insert(*cid, bytes.clone()); @@ -172,7 +163,7 @@ + }; + + let new_map: HashMap = walk.records.iter().cloned().collect(); -+ let old_map = RepoStore::load_record_cids(&conn, did) ++ let old_map = RepoStore::load_record_cids(conn, did) + .map_err(|e| SyncError::Store(format!("load record cids: {e}")))?; + + let mut added = 0usize; @@ -186,18 +177,24 @@ + // Added / changed records. + for (key, value_cid) in &new_map { + let cid_str = value_cid.to_string(); -+ match old_map.get(key) { ++ let is_update = match old_map.get(key) { + Some(old_cid) if old_cid == &cid_str => continue, // unchanged -+ Some(_) => updated += 1, -+ None => added += 1, -+ } ++ Some(_) => { ++ updated += 1; ++ true ++ }, ++ None => { ++ added += 1; ++ false ++ }, ++ }; + + let Some((collection, rkey)) = key.split_once('/') else { + continue; + }; + + // A changed/added record's block must be present in this CAR. If not, -+ // the diff was insufficient — bail to a full refetch. ++ // the diff was insufficient: bail to a full refetch. + let rec_bytes = merged.get(value_cid).ok_or(SyncError::IncompleteDiff)?; + let ipld: Ipld = serde_ipld_dagcbor::from_slice(rec_bytes) + .map_err(|e| SyncError::Decode(format!("record {key}: {e}")))?; @@ -206,9 +203,14 @@ + + RepoStore::upsert_record(&tx, did, collection, rkey, &cid_str, &json_str) + .map_err(|e| SyncError::Store(format!("upsert {key}: {e}")))?; -+ // Generic, lexicon-agnostic indexes (full-text + URLs). -+ textindex::index(&tx, did, collection, rkey, &json) -+ .map_err(|e| SyncError::Store(format!("index {key}: {e}")))?; ++ // Generic, lexicon-agnostic indexes (full-text + URLs). New records skip ++ // the re-index DELETE (a full FTS-table scan). ++ if is_update { ++ textindex::index(&tx, did, collection, rkey, &json) ++ } else { ++ textindex::index_new(&tx, did, collection, rkey, &json) ++ } ++ .map_err(|e| SyncError::Store(format!("index {key}: {e}")))?; + } + + // Removed records: present before, gone now. diff --git a/patches/components/net/atproto/textindex.rs.patch b/patches/components/net/atproto/textindex.rs.patch index 822a00c..ce326fd 100644 --- a/patches/components/net/atproto/textindex.rs.patch +++ b/patches/components/net/atproto/textindex.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -0,0 +1,522 @@ +@@ -0,0 +1,545 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +//! Generic, lexicon-agnostic indexing of every record. @@ -71,7 +71,7 @@ + +/// Matches an http(s) URL, greedily up to the next whitespace. Trailing +/// punctuation is trimmed by the caller. `at://` URIs are intentionally not -+/// matched — they're internal ATProto references and add noise to the index. ++/// matched since they're internal ATProto references and add noise to the index. +fn url_regex() -> &'static Regex { + static RE: OnceLock = OnceLock::new(); + RE.get_or_init(|| Regex::new(r"https?://[^\s]+").unwrap()) @@ -102,8 +102,22 @@ + ) +} + -+/// Index (or re-index) one record's text + URLs. Idempotent: clears any prior -+/// rows for the key first. ++/// Index a brand-new record's text + URLs (no prior rows to clear). Use this on ++/// the backfill / "added" path: it skips the [`remove`] that [`index`] runs, ++/// whose `records_fts` DELETE can't use an index on the UNINDEXED key columns ++/// and so full-scans the FTS table. ++pub fn index_new( ++ conn: &Connection, ++ did: &str, ++ collection: &str, ++ rkey: &str, ++ value: &Value, ++) -> rusqlite::Result<()> { ++ insert_rows(conn, did, collection, rkey, value) ++} ++ ++/// Re-index one record's text + URLs: clears any prior rows for the key first, ++/// then inserts. Use [`index_new`] when the record is known to be new. +pub fn index( + conn: &Connection, + did: &str, @@ -112,7 +126,16 @@ + value: &Value, +) -> rusqlite::Result<()> { + remove(conn, did, collection, rkey)?; ++ insert_rows(conn, did, collection, rkey, value) ++} + ++fn insert_rows( ++ conn: &Connection, ++ did: &str, ++ collection: &str, ++ rkey: &str, ++ value: &Value, ++) -> rusqlite::Result<()> { + let mut texts = Vec::new(); + collect_text(value, &mut texts); + @@ -194,7 +217,7 @@ +/// Find indexed URLs matching `url`, each paired with the record that contains +/// it. With `exact`, only URLs equal to `url`; otherwise any URL containing +/// `url` as a substring. One row per distinct `(url, did, collection, rkey)`: -+/// a URL referenced by several repos — or by several records of the same repo — ++/// a URL referenced by several repos or or by several records of the same repo +/// surfaces once per containing record. +pub fn search_url(conn: &Connection, url: &str, exact: bool) -> rusqlite::Result> { + let row = |row: &rusqlite::Row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)); diff --git a/patches/components/net/atproto/xrpc.rs.patch b/patches/components/net/atproto/xrpc.rs.patch index c165b81..d262e10 100644 --- a/patches/components/net/atproto/xrpc.rs.patch +++ b/patches/components/net/atproto/xrpc.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -0,0 +1,295 @@ +@@ -0,0 +1,296 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +use async_recursion::async_recursion; @@ -183,7 +183,8 @@ + Err(()) + } + } else { -+ if !xrpc_response.status.is_success() { ++ // Some buggy PDSes (https://atproto.brid.gy) return 0 in the success case... ++ if !xrpc_response.status.is_success() && xrpc_response.status.raw_code() != 0 { + error!( + "xrpc {} returned {}: {}", + xrpc_url, diff --git a/patches/components/net/lib.rs.patch b/patches/components/net/lib.rs.patch index df022cc..fbf026d 100644 --- a/patches/components/net/lib.rs.patch +++ b/patches/components/net/lib.rs.patch @@ -1,6 +1,6 @@ --- original +++ modified -@@ -24,7 +24,19 @@ +@@ -24,7 +24,20 @@ pub mod subresource_integrity; #[cfg(feature = "test-util")] pub mod test_util; @@ -9,6 +9,7 @@ +pub mod atproto { + pub mod car; + pub mod commit; ++ pub mod db_writer; + pub mod mst; + pub mod pds; + pub mod session; diff --git a/patches/components/net/protocols/atproto/protocol.rs.patch b/patches/components/net/protocols/atproto/protocol.rs.patch index bc15b01..b41472d 100644 --- a/patches/components/net/protocols/atproto/protocol.rs.patch +++ b/patches/components/net/protocols/atproto/protocol.rs.patch @@ -1,11 +1,11 @@ --- original +++ modified -@@ -0,0 +1,316 @@ +@@ -0,0 +1,337 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +use std::future::{self, Future}; +use std::num::NonZeroUsize; -+use std::path::Path; ++use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::sync::Arc; + @@ -57,8 +57,8 @@ + rkey: &str, + url: &ServoUrl, +) -> Option { -+ let store = RepoStore::open(config_dir?).ok()?; -+ let conn = store.connect().ok()?; ++ let store = RepoStore::read_only(config_dir?); ++ let conn = store.connect_readonly().ok()?; + let (cid, json) = RepoStore::get_record(&conn, did, collection, rkey).ok()??; + + // Stored json is the bare record value; parse it so it nests in the @@ -86,6 +86,23 @@ + Some(response) +} + ++/// Run [`try_cached_record`] on the blocking pool so the SQLite read never ++/// occupies an async runtime worker thread (a feed fires many of these at once). ++async fn cached_record( ++ config_dir: Option, ++ did: String, ++ collection: String, ++ rkey: String, ++ url: ServoUrl, ++) -> Option { ++ tokio::task::spawn_blocking(move || { ++ try_cached_record(config_dir.as_deref(), &did, &collection, &rkey, &url) ++ }) ++ .await ++ .ok() ++ .flatten() ++} ++ +/// Whether a `(collection, rkey)` getRecord can be served from the local cache. +/// Blobs aren't in the records table, and web-tile manifests must hit the +/// network so the handler's `register_webtile` still runs. @@ -173,13 +190,15 @@ + if method == Method::GET && subject.starts_with("did:") { + if let (Some(coll), Some(rkey)) = (collection, rkey) { + if is_cacheable_collection(coll) { -+ if let Some(cached) = try_cached_record( -+ context2.state.config_dir.as_deref(), -+ &subject, -+ coll, -+ rkey, -+ &url, -+ ) { ++ if let Some(cached) = cached_record( ++ context2.state.config_dir.clone(), ++ subject.to_string(), ++ coll.to_string(), ++ rkey.to_string(), ++ url.clone(), ++ ) ++ .await ++ { + return cached; + } + } @@ -223,13 +242,15 @@ + // resolved DID is now known, so try the cache once more + // before the network. + if is_cacheable_collection(coll) { -+ if let Some(cached) = try_cached_record( -+ context2.state.config_dir.as_deref(), -+ &document.id, -+ coll, -+ rkey, -+ &url, -+ ) { ++ if let Some(cached) = cached_record( ++ context2.state.config_dir.clone(), ++ document.id.to_string(), ++ coll.to_string(), ++ rkey.to_string(), ++ url.clone(), ++ ) ++ .await ++ { + return cached; + } + }