From 26a92ddb94d2e18baa4982d8ec77f3a2a2dfd697 Mon Sep 17 00:00:00 2001 From: phil Date: Mon, 15 Sep 2025 20:57:27 +0000 Subject: [PATCH] weekly exporter in rust --- .gitignore | 1 + Cargo.lock | 1 + Cargo.toml | 2 +- src/lib.rs | 4 +++- src/poll.rs | 11 +++++++++-- src/weekly.rs | 76 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/bin/bundle-weekly.rs | 53 +++++++++++++++++++++++++++++++++++++++++++++++++++++ 7 file(s) changed, 144 insertion(s)(+), 4 deletion(s)(-) diff --git a/.gitignore b/.gitignore --- a/.gitignore +++ b/.gitignore @@ -1,1 +1,2 @@ /target +weekly/ diff --git a/Cargo.lock b/Cargo.lock --- a/Cargo.lock +++ b/Cargo.lock @@ -123,6 +123,7 @@ "futures-core", "futures-io", "pin-project-lite", + "tokio", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ [dependencies] anyhow = "1.0.99" -async-compression = { version = "0.4.30", features = ["futures-io", "gzip"] } +async-compression = { version = "0.4.30", features = ["futures-io", "tokio", "gzip"] } chrono = { version = "0.4.42", features = ["serde"] } clap = { version = "4.5.47", features = ["derive", "env"] } env_logger = "0.11.8" diff --git a/src/lib.rs b/src/lib.rs --- a/src/lib.rs +++ b/src/lib.rs @@ -4,11 +4,13 @@ mod client; mod plc_pg; mod poll; +mod weekly; pub use backfill::week_to_pages; pub use client::CLIENT; pub use plc_pg::Db; -pub use poll::poll_upstream; +pub use poll::{get_page, poll_upstream}; +pub use weekly::{Week, pages_to_weeks}; pub type Dt = chrono::DateTime; diff --git a/src/poll.rs b/src/poll.rs --- a/src/poll.rs +++ b/src/poll.rs @@ -18,7 +18,7 @@ /// we assume that the order will at least be deterministic: this may be unsound #[derive(Debug, PartialEq)] pub struct LastOp { - created_at: Dt, // any op greater is definitely not duplicated + pub created_at: Dt, // any op greater is definitely not duplicated pk: (String, String), // did, cid } @@ -117,7 +117,14 @@ page.only_after_last(pl); } if !page.is_empty() { - dest.send_async(page).await?; + match dest.try_send(page) { + Ok(()) => {} + Err(flume::TrySendError::Full(page)) => { + log::warn!("export: destination channel full, awaiting..."); + dest.send_async(page).await?; + } + e => e?, + }; } prev_last = next_last.or(prev_last); diff --git a/src/weekly.rs b/src/weekly.rs new file mode 100644 --- /dev/null +++ b/src/weekly.rs @@ -0,0 +1,76 @@ +use crate::{Dt, ExportPage, Op}; +use async_compression::tokio::write::GzipEncoder; +use std::path::PathBuf; +use tokio::{fs::File, io::AsyncWriteExt}; + +const WEEK_IN_SECONDS: i64 = 7 * 86400; + +#[derive(Debug, Clone, Copy, PartialEq)] +pub struct Week(i64); + +impl From
for Week { + fn from(dt: Dt) -> Self { + let ts = dt.timestamp(); + let truncated = (ts / WEEK_IN_SECONDS) * WEEK_IN_SECONDS; + Week(truncated) + } +} + +impl From for Dt { + fn from(week: Week) -> Dt { + let Week(ts) = week; + Dt::from_timestamp(ts, 0).expect("the week to be in valid range") + } +} + +pub async fn pages_to_weeks(rx: flume::Receiver, dir: PathBuf) -> anyhow::Result<()> { + pub use std::time::Instant; + + // ...there is certainly a nicer way to write this + let mut current_week: Option = None; + let dummy_file = File::create(dir.join("_dummy")).await?; + let mut encoder = GzipEncoder::new(dummy_file); + + let mut total_ops = 0; + let total_t0 = Instant::now(); + let mut week_ops = 0; + let mut week_t0 = total_t0; + let mut week = 0; + + while let Ok(page) = rx.recv_async().await { + for mut s in page.ops { + let Ok(op) = serde_json::from_str::(&s) + .inspect_err(|e| log::error!("failed to parse plc op, ignoring: {e}")) + else { + continue; + }; + let op_week = op.created_at.into(); + if current_week.map(|w| w != op_week).unwrap_or(true) { + encoder.shutdown().await?; + let now = Instant::now(); + + log::info!( + "done week {week:3 } ({:10 }): {week_ops:7 } ({:5.0 }/s) ops, {:5 }k total ({:5.0 }/s)", + current_week.unwrap_or(Week(0)).0, + (week_ops as f64) / (now - week_t0).as_secs_f64(), + total_ops / 1000, + (total_ops as f64) / (now - total_t0).as_secs_f64(), + ); + + let file = File::create(dir.join(format!("{}.jsonl.gz", op_week.0))).await?; + encoder = GzipEncoder::with_quality(file, async_compression::Level::Best); + current_week = Some(op_week); + week_ops = 0; + week_t0 = now; + week += 1; + } + s.push('\n'); // hack + log::trace!("writing: {s}"); + encoder.write_all(s.as_bytes()).await?; + total_ops += 1; + week_ops += 1; + } + } + + Ok(()) +} diff --git a/src/bin/bundle-weekly.rs b/src/bin/bundle-weekly.rs new file mode 100644 --- /dev/null +++ b/src/bin/bundle-weekly.rs @@ -0,0 +1,53 @@ +use allegedly::{bin_init, pages_to_weeks, poll_upstream}; +use clap::Parser; +use std::path::PathBuf; +use url::Url; + +const PAGE_QUEUE_SIZE: usize = 128; + +#[derive(Parser)] +struct Args { + /// Upstream PLC server to poll + /// + /// default: https://plc.directory + #[arg(long, env)] + #[clap(default_value = "https://plc.directory")] + upstream: Url, + /// Directory to save gzipped weekly bundles + /// + /// default: ./weekly/ + #[arg(long, env)] + #[clap(default_value = "./weekly/")] + dir: PathBuf, + /// The week to start from + /// + /// Must be a week-truncated unix timestamp + #[arg(long, env)] + start_at: Option, // TODO!! +} + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + bin_init("weekly"); + let args = Args::parse(); + + let mut url = args.upstream; + url.set_path("/export"); + + log::trace!("ensure weekly output directory exists"); + std::fs::create_dir_all(&args.dir)?; + + let (tx, rx) = flume::bounded(PAGE_QUEUE_SIZE); + + tokio::task::spawn(async move { + if let Err(e) = poll_upstream(None /*todo*/, url, tx).await { + log::error!("polling failed: {e}"); + } else { + log::warn!("poller finished ok (weird?)"); + } + }); + + pages_to_weeks(rx, args.dir).await?; + + Ok(()) +} -- tangled.sh