From 567d97e92d458f616c8948ea9f85e9cf9cdf5e8a Mon Sep 17 00:00:00 2001 From: dawn <90008@gaze.systems> Date: Wed, 25 Feb 2026 22:26:40 +0300 Subject: [PATCH] update from-fjall backfill to use the weekly code by implementing BundleSource --- src/bin/backfill.rs | 4 +- src/lib.rs | 2 +- src/mirror/fjall.rs | 4 +- src/plc_fjall.rs | 137 ++++++++++++++++++++++---------------------- src/weekly.rs | 12 ++-- 5 files changed, 79 insertions(+), 80 deletions(-) diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index a729b28..3b90425 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}, - fjall_to_pages, full_pages, logo, pages_to_fjall, pages_to_pg, pages_to_stdout, poll_upstream, + full_pages, logo, pages_to_fjall, pages_to_pg, pages_to_stdout, poll_upstream, }; use clap::Parser; use reqwest::Url; @@ -139,7 +139,7 @@ pub async fn run( 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)); + tasks.spawn(backfill(db, bulk_tx, source_workers.unwrap_or(4), until)); } else if let Some(dir) = dir { if http != DEFAULT_HTTP.parse()? { anyhow::bail!( diff --git a/src/lib.rs b/src/lib.rs index d8edaf7..f481635 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, fjall_to_pages, pages_to_fjall}; +pub use plc_fjall::{FjallDb, backfill_to_fjall, 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/mirror/fjall.rs b/src/mirror/fjall.rs index 3ab3843..bf4e304 100644 --- a/src/mirror/fjall.rs +++ b/src/mirror/fjall.rs @@ -275,8 +275,8 @@ async fn fjall_export( let db = fjall.clone(); let ops = tokio::task::spawn_blocking(move || { - let iter = db.export_ops(after, limit)?; - iter.collect::>>() + let iter = db.export_ops(after.unwrap_or(Dt::UNIX_EPOCH)..)?; + iter.take(limit).collect::>>() }) .await .map_err(|e| Error::from_string(e.to_string(), StatusCode::INTERNAL_SERVER_ERROR))? diff --git a/src/plc_fjall.rs b/src/plc_fjall.rs index df74e77..8c2f18f 100644 --- a/src/plc_fjall.rs +++ b/src/plc_fjall.rs @@ -1,3 +1,4 @@ +use crate::{BundleSource, Week}; use crate::{Dt, ExportPage, Op as CommonOp, PageBoundaryState}; use anyhow::Context; use data_encoding::{BASE32_NOPAD, BASE64URL_NOPAD}; @@ -5,12 +6,14 @@ use fjall::{ Database, Keyspace, KeyspaceCreateOptions, OwnedWriteBatch, PersistMode, config::BlockSizePolicy, }; +use futures::Future; use serde::{Deserialize, Serialize}; use std::collections::BTreeMap; use std::fmt; use std::path::Path; use std::sync::Arc; use std::time::Instant; +use tokio::io::{AsyncRead, AsyncWriteExt}; use tokio::sync::{mpsc, oneshot}; const SEP: u8 = 0; @@ -1022,17 +1025,21 @@ impl FjallDb { pub fn export_ops( &self, - after: Option
, - limit: usize, + range: impl std::ops::RangeBounds
, ) -> anyhow::Result> + '_> { - let iter = if let Some(after) = after { - let start = (after.timestamp_micros() as u64).to_be_bytes(); - self.inner.ops.range(start..) - } else { - self.inner.ops.iter() + use std::ops::Bound; + let map_bound = |b: Bound<&Dt>| -> Bound<[u8; 8]> { + match b { + Bound::Included(dt) => Bound::Included(dt.timestamp_micros().to_be_bytes()), + Bound::Excluded(dt) => Bound::Excluded(dt.timestamp_micros().to_be_bytes()), + Bound::Unbounded => Bound::Unbounded, + } }; + let range = (map_bound(range.start_bound()), map_bound(range.end_bound())); + + let iter = self.inner.ops.range(range); - Ok(iter.take(limit).map(|item| { + Ok(iter.map(|item| { let (key, value) = item .into_inner() .map_err(|e| anyhow::anyhow!("fjall read error: {e}"))?; @@ -1060,6 +1067,58 @@ impl FjallDb { }) })) } + + pub fn export_ops_week( + &self, + week: Week, + ) -> anyhow::Result> + '_> { + let after: Dt = week.into(); + let before: Dt = week.next().into(); + + self.export_ops(after..before) + } +} + +impl BundleSource for FjallDb { + fn reader_for( + &self, + week: Week, + ) -> impl Future> + Send { + let db = self.clone(); + + async move { + let (mut tx, rx) = tokio::io::duplex(1024 * 1024 * 64); + + tokio::task::spawn_blocking(move || -> anyhow::Result<()> { + let iter = db.export_ops_week(week)?; + + let rt = tokio::runtime::Handle::current(); + + for op_res in iter { + let op = op_res?; + 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)?, + }; + + let mut json_bytes = serde_json::to_vec(&common_op)?; + json_bytes.push(b'\n'); + + if rt.block_on(tx.write_all(&json_bytes)).is_err() { + break; + } + } + + Ok(()) + }); + + Ok(rx) + } + } } pub async fn backfill_to_fjall( @@ -1151,68 +1210,6 @@ 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::*; diff --git a/src/weekly.rs b/src/weekly.rs index a7b5b3b..39dfdbd 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -101,7 +101,8 @@ impl BundleSource for FolderSource { let file = File::open(path) .await .inspect_err(|e| log::error!("failed to open file: {e}"))?; - Ok(file) + let decoder = GzipDecoder::new(BufReader::new(file)); + Ok(decoder) } } @@ -112,7 +113,7 @@ impl BundleSource for HttpSource { use futures::TryStreamExt; let HttpSource(base) = self; let url = base.join(&format!("{}.jsonl.gz", week.0))?; - Ok(CLIENT + let stream = CLIENT .get(url) .send() .await? @@ -120,7 +121,9 @@ impl BundleSource for HttpSource { .bytes_stream() .map_err(futures::io::Error::other) .into_async_read() - .compat()) + .compat(); + let decoder = GzipDecoder::new(BufReader::new(stream)); + Ok(decoder) } } @@ -213,8 +216,7 @@ pub async fn week_to_pages( } }; - let decoder = GzipDecoder::new(BufReader::new(reader)); - let mut chunks = pin!(LinesStream::new(BufReader::new(decoder).lines()).try_chunks(10000)); + let mut chunks = pin!(LinesStream::new(BufReader::new(reader).lines()).try_chunks(10000)); let mut success = true; while let Some(chunk) = match chunks.as_mut().try_next().await { -- 2.51.2