diff --git a/Cargo.lock b/Cargo.lock index 0a18077..a1bab7d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -46,6 +46,8 @@ dependencies = [ "thiserror 2.0.16", "tokio", "tokio-postgres", + "tokio-stream", + "tokio-util", "url", ] @@ -2034,6 +2036,17 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-stream" +version = "0.1.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eca58d7bba4a75707817a2c44174253f9236b2d5fbd055602e9d5c07c139a047" +dependencies = [ + "futures-core", + "pin-project-lite", + "tokio", +] + [[package]] name = "tokio-util" version = "0.7.16" @@ -2042,6 +2055,7 @@ checksum = "14307c986784f72ef81c89db7d9e28d6ac26d16213b109ea501696195e6e3ce5" dependencies = [ "bytes", "futures-core", + "futures-io", "futures-sink", "pin-project-lite", "tokio", diff --git a/Cargo.toml b/Cargo.toml index 0798857..ef015ff 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -23,4 +23,6 @@ serde_json = { version = "1.0.143", features = ["raw_value"] } thiserror = "2.0.16" tokio = { version = "1.47.1", features = ["full"] } tokio-postgres = { version = "0.7.13", features = ["with-chrono-0_4", "with-serde_json-1"] } +tokio-stream = { version = "0.1.17", features = ["io-util"] } +tokio-util = { version = "0.7.16", features = ["compat"] } url = "2.5.7" diff --git a/readme.md b/readme.md index 747afbe..e9853f8 100644 --- a/readme.md +++ b/readme.md @@ -6,6 +6,7 @@ Allegedly can - Tail PLC ops to stdout: `allegedly tail | jq` - Export PLC ops to weekly gzipped bundles: `allegdly bundle --dest ./some-folder` +- Dump bundled ops to stdout FAST: `allegedly backfill --source-workers 6 | pv -l > /ops-unordered.jsonl` (add `--help` to any command for more info about it) diff --git a/src/backfill.rs b/src/backfill.rs index 68bbf0d..744a0f3 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -1,27 +1,45 @@ -use crate::{CLIENT, ExportPage}; -use url::Url; +use crate::{BundleSource, Dt, ExportPage, Week, week_to_pages}; +use tokio::task::JoinSet; -use async_compression::futures::bufread::GzipDecoder; -use futures::{AsyncBufReadExt, StreamExt, TryStreamExt, io}; +const FIRST_WEEK: Week = Week::from_n(1668643200); -pub async fn week_to_pages(url: Url, dest: flume::Sender) -> anyhow::Result<()> { - let reader = CLIENT - .get(url) - .send() - .await? - .error_for_status()? - .bytes_stream() - .map_err(io::Error::other) - .into_async_read(); +pub async fn backfill( + source: impl BundleSource + Send + 'static, + dest: flume::Sender, + source_workers: usize, +) -> anyhow::Result<()> { + // queue up the week bundles that should be available + let (week_tx, week_rx) = flume::bounded(1024); // work queue + let mut week = FIRST_WEEK; + while week.is_immutable() { + week_tx.try_send(week)?; // if this fails, something has gone really wrong or we're farrrr in the future + week = week.next(); + } - let decoder = GzipDecoder::new(io::BufReader::new(reader)); + let mut workers: JoinSet> = JoinSet::new(); - let mut chunks = io::BufReader::new(decoder).lines().chunks(1000); + // spin up the fetchers to work in parallel + for w in 0..source_workers { + let weeks = week_rx.clone(); + let dest = dest.clone(); + let source = source.clone(); + workers.spawn(async move { + while let Ok(week) = weeks.recv_async().await { + log::info!( + "worker {w}: fetching week {} (-{})", + Into::
::into(week).to_rfc3339(), + week.n_ago(), + ); + week_to_pages(source.clone(), week, dest.clone()).await?; + } + Ok(()) + }); + } - while let Some(chunk) = chunks.next().await { - let ops = chunk.into_iter().collect::, io::Error>>()?; - let page = ExportPage { ops }; - dest.send_async(page).await?; + // wait for them to finish + while let Some(res) = workers.join_next().await { + res??; } + Ok(()) } diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index 6072d25..95be0e4 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -1,4 +1,4 @@ -use allegedly::{Dt, bin_init, pages_to_weeks, poll_upstream}; +use allegedly::{Dt, FolderSource, HttpSource, backfill, bin_init, pages_to_weeks, poll_upstream}; use clap::{Parser, Subcommand}; use std::path::PathBuf; use url::Url; @@ -15,6 +15,20 @@ struct Cli { #[derive(Debug, Subcommand)] enum Commands { + /// Use weekly bundled ops to get a complete directory mirror FAST + Backfill { + /// Remote URL prefix to fetch bundles from + #[arg(long)] + #[clap(default_value = "https://plc.t3.storage.dev/plc.directory/")] + http: Url, + /// Local folder to fetch bundles from (overrides `http`) + #[arg(long)] + dir: Option, + /// Parallel bundle fetchers + #[arg(long)] + #[clap(default_value = "4")] + source_workers: usize, + }, /// Scrape a PLC server, collecting ops into weekly bundles /// /// Bundles are gzipped files named `.jsonl.gz` where WEEK is a unix @@ -51,6 +65,31 @@ async fn main() { let args = Cli::parse(); match args.command { + Commands::Backfill { + http, + dir, + source_workers, + } => { + let (tx, rx) = flume::bounded(1024); // big pages + tokio::task::spawn(async move { + if let Some(dir) = dir { + log::info!("Reading weekly bundles from local folder {dir:?}"); + backfill(FolderSource(dir), tx, source_workers) + .await + .unwrap(); + } else { + log::info!("Fetching weekly bundles from from {http}"); + backfill(HttpSource(http), tx, source_workers) + .await + .unwrap(); + } + }); + loop { + for op in rx.recv_async().await.unwrap().ops { + println!("{op}") + } + } + } Commands::Bundle { dest, after, diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index 25fc551..87003a6 100644 --- a/src/bin/backfill.rs +++ b/src/bin/backfill.rs @@ -2,7 +2,7 @@ use clap::Parser; use std::time::Duration; use url::Url; -use allegedly::{Db, Dt, ExportPage, Op, bin_init, poll_upstream, week_to_pages}; +use allegedly::{Db, Dt, ExportPage, Op, bin_init, poll_upstream}; const EXPORT_PAGE_QUEUE_SIZE: usize = 0; // rendezvous for now const WEEK_IN_SECONDS: u64 = 7 * 86400; @@ -40,21 +40,22 @@ struct Args { postgres: String, } -async fn bulk_backfill((upstream, epoch): (Url, u64), tx: flume::Sender) { +async fn bulk_backfill((_upstream, epoch): (Url, u64), _tx: flume::Sender) { let immutable_cutoff = std::time::SystemTime::now() - Duration::from_secs((7 + 4) * 86400); let immutable_ts = (immutable_cutoff.duration_since(std::time::SystemTime::UNIX_EPOCH)) .unwrap() .as_secs(); - let immutable_week = (immutable_ts / WEEK_IN_SECONDS) * WEEK_IN_SECONDS; - let mut week = epoch; - let mut week_n = 0; - while week < immutable_week { - log::info!("backfilling week {week_n} ({week})"); - let url = upstream.join(&format!("{week}.jsonl.gz")).unwrap(); - week_to_pages(url, tx.clone()).await.unwrap(); - week_n += 1; - week += WEEK_IN_SECONDS; - } + let _immutable_week = (immutable_ts / WEEK_IN_SECONDS) * WEEK_IN_SECONDS; + let _week = epoch; + let _week_n = 0; + todo!(); + // while week < immutable_week { + // log::info!("backfilling week {week_n} ({week})"); + // let url = upstream.join(&format!("{week}.jsonl.gz")).unwrap(); + // week_to_pages(url, tx.clone()).await.unwrap(); + // week_n += 1; + // week += WEEK_IN_SECONDS; + // } } async fn export_upstream( diff --git a/src/bin/get_backfill_chunk_adsf.rs b/src/bin/get_backfill_chunk_adsf.rs index 4231b66..346d5f4 100644 --- a/src/bin/get_backfill_chunk_adsf.rs +++ b/src/bin/get_backfill_chunk_adsf.rs @@ -1,25 +1,47 @@ -use allegedly::CLIENT; -use async_compression::futures::bufread::GzipDecoder; -use futures::{AsyncBufReadExt, StreamExt, TryStreamExt, io}; +use allegedly::{HttpSource, Week, week_to_pages}; +use std::io::Write; #[tokio::main] async fn main() { - let reader = CLIENT - .get("https://plc.t3.storage.dev/plc.directory/1699488000.jsonl.gz") - // .get("https://plc.t3.storage.dev/plc.directory/1669248000.jsonl.gz") - .send() - .await - .unwrap() - .error_for_status() - .unwrap() - .bytes_stream() - .map_err(io::Error::other) - .into_async_read(); - - let decoder = GzipDecoder::new(io::BufReader::new(reader)); - let mut chunks = io::BufReader::new(decoder).lines().chunks(1000); - while let Some(ref _chunk) = chunks.next().await { + let url: url::Url = "https://plc.t3.storage.dev/plc.directory/".parse().unwrap(); + let source = HttpSource(url); + // let source = FolderSource("./weekly/".into()); + let week = Week::from_n(1699488000); + + let (tx, rx) = flume::bounded(32); + + tokio::task::spawn(async move { + week_to_pages(source, week, tx).await.unwrap(); + }); + + let mut n = 0; + + print!("receiving"); + while let Ok(page) = rx.recv_async().await { print!("."); + std::io::stdout().flush().unwrap(); + n += page.ops.len(); } println!(); + + println!("bye ({n})"); + + // let reader = CLIENT + // .get("https://plc.t3.storage.dev/plc.directory/1699488000.jsonl.gz") + // // .get("https://plc.t3.storage.dev/plc.directory/1669248000.jsonl.gz") + // .send() + // .await + // .unwrap() + // .error_for_status() + // .unwrap() + // .bytes_stream() + // .map_err(io::Error::other) + // .into_async_read(); + + // let decoder = GzipDecoder::new(io::BufReader::new(reader)); + // let mut chunks = io::BufReader::new(decoder).lines().chunks(1000); + // while let Some(ref _chunk) = chunks.next().await { + // print!("."); + // } + // println!(); } diff --git a/src/lib.rs b/src/lib.rs index c8d1b3b..78e0d5d 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -6,17 +6,17 @@ mod plc_pg; mod poll; mod weekly; -pub use backfill::week_to_pages; +pub use backfill::backfill; pub use client::CLIENT; pub use plc_pg::Db; pub use poll::{get_page, poll_upstream}; -pub use weekly::{Week, pages_to_weeks}; +pub use weekly::{BundleSource, FolderSource, HttpSource, Week, pages_to_weeks, week_to_pages}; pub type Dt = chrono::DateTime; /// One page of PLC export /// -/// Expected to have up to around 1000 lines of raw json ops +/// plc.directory caps /export at 1000 ops; backfill tasks may send more in a page. #[derive(Debug)] pub struct ExportPage { pub ops: Vec, diff --git a/src/weekly.rs b/src/weekly.rs index 4dbf356..64d80f4 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -1,15 +1,24 @@ -use crate::{Dt, ExportPage, Op}; +use crate::{CLIENT, Dt, ExportPage, Op}; +use async_compression::tokio::bufread::GzipDecoder; use async_compression::tokio::write::GzipEncoder; +use core::pin::pin; +use std::future::Future; use std::path::PathBuf; -use tokio::{fs::File, io::AsyncWriteExt}; +use tokio::{ + fs::File, + io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader}, +}; +use tokio_stream::wrappers::LinesStream; +use tokio_util::compat::FuturesAsyncReadCompatExt; +use url::Url; -const WEEK_IN_SECONDS: i64 = 7 * 86400; +const WEEK_IN_SECONDS: i64 = 7 * 86_400; #[derive(Debug, Clone, Copy, PartialEq)] pub struct Week(i64); impl Week { - pub fn from_n(n: i64) -> Self { + pub const fn from_n(n: i64) -> Self { Self(n) } pub fn n_ago(&self) -> i64 { @@ -17,6 +26,18 @@ impl Week { let Self(cur) = chrono::Utc::now().into(); (cur - us) / WEEK_IN_SECONDS } + pub fn next(&self) -> Week { + Self(self.0 + WEEK_IN_SECONDS) + } + /// is the plc log for this week entirely outside the 72h nullification window + /// + /// plus one hour for safety (week must have ended > 73 hours ago) + pub fn is_immutable(&self) -> bool { + const HOUR_IN_SECONDS: i64 = 3600; + let now = chrono::Utc::now().timestamp(); + let nullification_cutoff = now - (73 * HOUR_IN_SECONDS); + self.next().0 <= nullification_cutoff + } } impl From
for Week { @@ -34,6 +55,42 @@ impl From for Dt { } } +pub trait BundleSource: Clone { + fn reader_for( + &self, + week: Week, + ) -> impl Future> + Send; +} + +#[derive(Debug, Clone)] +pub struct FolderSource(pub PathBuf); +impl BundleSource for FolderSource { + async fn reader_for(&self, week: Week) -> anyhow::Result { + let FolderSource(dir) = self; + let path = dir.join(format!("{}.jsonl.gz", week.0)); + Ok(File::open(path).await?) + } +} + +#[derive(Debug, Clone)] +pub struct HttpSource(pub Url); +impl BundleSource for HttpSource { + async fn reader_for(&self, week: Week) -> anyhow::Result { + use futures::TryStreamExt; + let HttpSource(base) = self; + let url = base.join(&format!("{}.jsonl.gz", week.0))?; + Ok(CLIENT + .get(url) + .send() + .await? + .error_for_status()? + .bytes_stream() + .map_err(futures::io::Error::other) + .into_async_read() + .compat()) + } +} + pub async fn pages_to_weeks( rx: flume::Receiver, dir: PathBuf, @@ -104,3 +161,20 @@ pub async fn pages_to_weeks( Ok(()) } + +pub async fn week_to_pages( + source: impl BundleSource, + week: Week, + dest: flume::Sender, +) -> anyhow::Result<()> { + use futures::TryStreamExt; + let decoder = GzipDecoder::new(BufReader::new(source.reader_for(week).await?)); + let mut chunks = pin!(LinesStream::new(BufReader::new(decoder).lines()).try_chunks(10000)); + + while let Some(chunk) = chunks.try_next().await? { + let ops: Vec = chunk.into_iter().collect(); + let page = ExportPage { ops }; + dest.send_async(page).await?; + } + Ok(()) +}