diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index a7b1787..a729b28 100644 --- a/src/bin/backfill.rs +++ b/src/bin/backfill.rs @@ -2,7 +2,7 @@ use allegedly::{ Db, Dt, ExportPage, FjallDb, FolderSource, HttpSource, backfill, backfill_to_fjall, backfill_to_pg, bin::{GlobalArgs, bin_init}, - full_pages, logo, pages_to_fjall, pages_to_pg, pages_to_stdout, poll_upstream, + fjall_to_pages, full_pages, logo, pages_to_fjall, pages_to_pg, pages_to_stdout, poll_upstream, }; use clap::Parser; use reqwest::Url; @@ -23,6 +23,9 @@ pub struct Args { /// Local folder to fetch bundles from (overrides `http`) #[arg(long)] dir: Option, + /// Local fjall database to fetch raw ops from (overrides `http` and `dir`) + #[arg(long, conflicts_with_all = ["dir"])] + from_fjall: Option, /// Don't do weekly bulk-loading at all. /// /// overrides `http` and `dir`, makes catch_up redundant @@ -72,6 +75,7 @@ pub async fn run( Args { http, dir, + from_fjall, no_bulk, source_workers, to_postgres, @@ -131,7 +135,12 @@ pub async fn run( // fun mode // set up bulk sources - if let Some(dir) = dir { + if let Some(fjall_path) = from_fjall { + log::trace!("opening source fjall db at {fjall_path:?}..."); + let db = FjallDb::open(&fjall_path)?; + log::trace!("opened source fjall db"); + tasks.spawn(fjall_to_pages(db, bulk_tx, until)); + } else if let Some(dir) = dir { if http != DEFAULT_HTTP.parse()? { anyhow::bail!( "non-default bulk http setting can't be used with bulk dir setting ({dir:?})" diff --git a/src/lib.rs b/src/lib.rs index f481635..d8edaf7 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -19,7 +19,7 @@ pub use backfill::backfill; pub use cached_value::{CachedValue, Fetcher}; pub use client::{CLIENT, UA}; pub use mirror::{ExperimentalConf, ListenConf, serve, serve_fjall}; -pub use plc_fjall::{FjallDb, backfill_to_fjall, pages_to_fjall}; +pub use plc_fjall::{FjallDb, backfill_to_fjall, fjall_to_pages, pages_to_fjall}; pub use plc_pg::{Db, backfill_to_pg, pages_to_pg}; pub use poll::{PageBoundaryState, get_page, poll_upstream}; pub use ratelimit::{CreatePlcOpLimiter, GovernorMiddleware, IpLimiters}; diff --git a/src/plc_fjall.rs b/src/plc_fjall.rs index b000b63..70ef9b0 100644 --- a/src/plc_fjall.rs +++ b/src/plc_fjall.rs @@ -1,7 +1,10 @@ use crate::{Dt, ExportPage, Op as CommonOp, PageBoundaryState}; use anyhow::Context; use data_encoding::{BASE32_NOPAD, BASE64URL_NOPAD}; -use fjall::{Database, Keyspace, KeyspaceCreateOptions, OwnedWriteBatch, PersistMode}; +use fjall::{ + Database, Keyspace, KeyspaceCreateOptions, OwnedWriteBatch, PersistMode, + config::BlockSizePolicy, +}; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; use std::fmt; @@ -694,15 +697,35 @@ struct FjallInner { impl FjallDb { pub fn open(path: impl AsRef) -> fjall::Result { + const fn kb(kb: u32) -> u32 { + kb * 1_024 + } + const fn mb(mb: u32) -> u64 { + kb(mb) as u64 * 1_024 + } + let db = Database::builder(path) - .max_journaling_size(/* 1 GiB */ 1_024 * 1_024 * 1_024) + // 32mb is too low we can afford more + // this should be configurable though! + .cache_size(mb(256)) .open()?; let opts = KeyspaceCreateOptions::default; let ops = db.keyspace("ops", || { - opts().max_memtable_size(/* 256 MiB */ 256 * 1_024 * 1_024) + opts() + // this is mainly for when backfilling + .max_memtable_size(mb(192)) + // this wont compress terribly well since its a bunch of CIDs and signatures and did:keys + // and we want to keep reads fast since we'll be reading a lot... + .data_block_size_policy(BlockSizePolicy::new([kb(4), kb(8), kb(32)])) + // this has no downsides, since the only point reads that might miss we do is on by_did + .expect_point_read_hits(true) })?; let by_did = db.keyspace("by_did", || { - opts().max_memtable_size(/* 128 MiB */ 128 * 1_024 * 1_024) + opts() + .max_memtable_size(mb(64)) + // this isn't gonna compress well anyway, since its just keys (did + timestamp + cid) + // and dids dont have many operations in the first place, so we can use small blocks + .data_block_size_policy(BlockSizePolicy::all(kb(2))) })?; Ok(Self { inner: Arc::new(FjallInner { db, ops, by_did }), @@ -948,6 +971,68 @@ pub async fn pages_to_fjall( Ok("pages_to_fjall") } +pub async fn fjall_to_pages( + db: FjallDb, + dest: mpsc::Sender, + until: Option
, +) -> anyhow::Result<&'static str> { + log::info!("starting fjall_to_pages backfill source..."); + + let t0 = Instant::now(); + + let dest_clone = dest.clone(); + let ops_sent = tokio::task::spawn_blocking(move || -> anyhow::Result { + let iter = db.export_ops(None, usize::MAX)?; + let mut current_page = Vec::with_capacity(1000); + let mut count = 0; + + for op_res in iter { + let op = op_res?; + + if let Some(u) = until { + if op.created_at >= u { + break; + } + } + + let operation_str = serde_json::to_string(&op.operation)?; + let common_op = crate::Op { + did: op.did, + cid: op.cid, + created_at: op.created_at, + nullified: op.nullified, + operation: serde_json::value::RawValue::from_string(operation_str)?, + }; + + current_page.push(common_op); + count += 1; + + if current_page.len() >= 1000 { + let page = ExportPage { + ops: std::mem::take(&mut current_page), + }; + if dest_clone.blocking_send(page).is_err() { + break; + } + } + } + + if !current_page.is_empty() { + let page = ExportPage { ops: current_page }; + let _ = dest_clone.blocking_send(page); + } + + Ok(count) + }) + .await??; + + log::info!( + "finished sending {ops_sent} ops from fjall in {:?}", + t0.elapsed() + ); + Ok("fjall_to_pages") +} + #[cfg(test)] mod tests { use super::*;