diff --git a/Cargo.lock b/Cargo.lock index a1bab7d..345bf94 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -35,7 +35,6 @@ dependencies = [ "chrono", "clap", "env_logger", - "flume", "futures", "log", "reqwest", @@ -448,18 +447,6 @@ dependencies = [ "miniz_oxide", ] -[[package]] -name = "flume" -version = "0.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "da0e4dd2a88388a1f4ccc7c9ce104604dab68d9f408dc34cd45823d5a9069095" -dependencies = [ - "futures-core", - "futures-sink", - "nanorand", - "spin", -] - [[package]] name = "fnv" version = "1.0.7" @@ -1093,15 +1080,6 @@ dependencies = [ "windows-sys 0.59.0", ] -[[package]] -name = "nanorand" -version = "0.7.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6a51313c5820b0b02bd422f4b44776fbf47961755c74ce64afc73bfad10226c3" -dependencies = [ - "getrandom 0.2.16", -] - [[package]] name = "native-tls" version = "0.2.14" @@ -1791,15 +1769,6 @@ dependencies = [ "windows-sys 0.59.0", ] -[[package]] -name = "spin" -version = "0.9.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6980e8d7511241f8acf4aebddbb1ff938df5eebe98691418c4468d0b72a96a67" -dependencies = [ - "lock_api", -] - [[package]] name = "stable_deref_trait" version = "1.2.0" diff --git a/Cargo.toml b/Cargo.toml index ef015ff..3c90c0f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,6 @@ async-compression = { version = "0.4.30", features = ["futures-io", "tokio", "gz chrono = { version = "0.4.42", features = ["serde"] } clap = { version = "4.5.47", features = ["derive", "env"] } env_logger = "0.11.8" -flume = "0.11.1" futures = "0.3.31" log = "0.4.28" reqwest = { version = "0.12.23", features = ["stream"] } diff --git a/src/backfill.rs b/src/backfill.rs index 1bb2a0c..e3ed6dd 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -1,13 +1,16 @@ use crate::{BundleSource, Dt, ExportPage, Week, week_to_pages}; use std::sync::Arc; use std::time::Instant; -use tokio::{sync::Mutex, task::JoinSet}; +use tokio::{ + sync::{Mutex, mpsc}, + task::JoinSet, +}; const FIRST_WEEK: Week = Week::from_n(1668643200); pub async fn backfill( source: impl BundleSource + Send + 'static, - dest: flume::Sender, + dest: mpsc::Sender, source_workers: usize, until: Option
, ) -> anyhow::Result<()> { diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index fffa1be..9caee9b 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -3,8 +3,8 @@ use allegedly::{ bin_init, pages_to_pg, pages_to_weeks, poll_upstream, }; use clap::{Parser, Subcommand}; -use std::path::PathBuf; -use tokio::sync::oneshot; +use std::{path::PathBuf, time::Instant}; +use tokio::sync::{mpsc, oneshot}; use url::Url; #[derive(Debug, Parser)] @@ -80,11 +80,11 @@ enum Commands { } async fn pages_to_stdout( - rx: flume::Receiver, + mut rx: mpsc::Receiver, notify_last_at: Option>>, -) -> Result<(), flume::RecvError> { +) -> anyhow::Result<()> { let mut last_at = None; - while let Ok(page) = rx.recv_async().await { + while let Some(page) = rx.recv().await { for op in &page.ops { println!("{op}"); } @@ -107,13 +107,13 @@ async fn pages_to_stdout( /// /// PLC will return up to 1000 ops on a page, and returns full pages until it /// has caught up, so this is a (hacky?) way to stop polling once we're up. -fn full_pages(rx: flume::Receiver) -> flume::Receiver { - let (tx, fwd) = flume::bounded(0); +fn full_pages(mut rx: mpsc::Receiver) -> mpsc::Receiver { + let (tx, fwd) = mpsc::channel(1); tokio::task::spawn(async move { - while let Ok(page) = rx.recv_async().await + while let Some(page) = rx.recv().await && page.ops.len() > 900 { - tx.send_async(page).await.unwrap(); + tx.send(page).await.unwrap(); } }); fwd @@ -125,6 +125,7 @@ async fn main() { let args = Cli::parse(); + let t0 = Instant::now(); match args.command { Commands::Backfill { http, @@ -135,7 +136,7 @@ async fn main() { until, catch_up, } => { - let (tx, rx) = flume::bounded(32); // these are big pages + let (tx, rx) = mpsc::channel(32); // these are big pages tokio::task::spawn(async move { if let Some(dir) = dir { log::info!("Reading weekly bundles from local folder {dir:?}"); @@ -177,7 +178,7 @@ async fn main() { // wait until the time for `after` is known let last_at = rx_last.await.unwrap(); log::info!("beginning catch-up from {last_at:?} while the writer finalizes stuff"); - let (tx, rx) = flume::bounded(256); + let (tx, rx) = mpsc::channel(256); // these are small pages tokio::task::spawn( async move { poll_upstream(last_at, upstream, tx).await.unwrap() }, ); @@ -199,7 +200,7 @@ async fn main() { } => { let mut url = args.upstream; url.set_path("/export"); - let (tx, rx) = flume::bounded(32); // read ahead if gzip stalls for some reason + let (tx, rx) = mpsc::channel(32); // read ahead if gzip stalls for some reason tokio::task::spawn(async move { poll_upstream(Some(after), url, tx).await.unwrap() }); log::trace!("ensuring output directory exists"); std::fs::create_dir_all(&dest).unwrap(); @@ -209,9 +210,10 @@ async fn main() { let mut url = args.upstream; url.set_path("/export"); let start_at = after.or_else(|| Some(chrono::Utc::now())); - let (tx, rx) = flume::bounded(1); + let (tx, rx) = mpsc::channel(1); tokio::task::spawn(async move { poll_upstream(start_at, url, tx).await.unwrap() }); pages_to_stdout(rx, None).await.unwrap(); } } + log::info!("whew, {:?}. goodbye!", t0.elapsed()); } diff --git a/src/plc_pg.rs b/src/plc_pg.rs index d0f4697..5e9a877 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -1,7 +1,7 @@ use crate::{Dt, ExportPage, Op, PageBoundaryState}; use std::pin::pin; use std::time::Instant; -use tokio::sync::oneshot; +use tokio::sync::{mpsc, oneshot}; use tokio_postgres::{ Client, Error as PgError, NoTls, binary_copy::BinaryCopyInWriter, @@ -72,7 +72,7 @@ impl Db { } } -pub async fn pages_to_pg(db: Db, pages: flume::Receiver) -> Result<(), PgError> { +pub async fn pages_to_pg(db: Db, mut pages: mpsc::Receiver) -> Result<(), PgError> { let mut client = db.connect().await?; let ops_stmt = client @@ -89,7 +89,7 @@ pub async fn pages_to_pg(db: Db, pages: flume::Receiver) -> Result<( let mut ops_inserted = 0; let mut dids_inserted = 0; - while let Ok(page) = pages.recv_async().await { + while let Some(page) = pages.recv().await { log::trace!("writing page with {} ops", page.ops.len()); let tx = client.transaction().await?; for s in page.ops { @@ -137,7 +137,7 @@ pub async fn pages_to_pg(db: Db, pages: flume::Receiver) -> Result<( pub async fn backfill_to_pg( db: Db, reset: bool, - pages: flume::Receiver, + mut pages: mpsc::Receiver, notify_last_at: Option>>, ) -> Result<(), PgError> { let mut client = db.connect().await?; @@ -195,7 +195,7 @@ pub async fn backfill_to_pg( .await?; let mut writer = pin!(BinaryCopyInWriter::new(sync, types)); let mut last_at = None; - while let Ok(page) = pages.recv_async().await { + while let Some(page) = pages.recv().await { for s in &page.ops { let Ok(op) = serde_json::from_str::(s) else { log::warn!("ignoring unparseable op: {s:?}"); diff --git a/src/poll.rs b/src/poll.rs index c53831e..e2df0f4 100644 --- a/src/poll.rs +++ b/src/poll.rs @@ -1,6 +1,7 @@ use crate::{CLIENT, Dt, ExportPage, Op, OpKey}; use std::time::Duration; use thiserror::Error; +use tokio::sync::mpsc; use url::Url; // plc.directory ratelimit on /export is 500 per 5 mins @@ -209,7 +210,7 @@ pub async fn get_page(url: Url) -> Result<(ExportPage, Option), GetPageE pub async fn poll_upstream( after: Option
, base: Url, - dest: flume::Sender, + dest: mpsc::Sender, ) -> anyhow::Result<()> { let mut tick = tokio::time::interval(UPSTREAM_REQUEST_INTERVAL); let mut prev_last: Option = after.map(Into::into); @@ -232,9 +233,9 @@ pub async fn poll_upstream( if !page.is_empty() { match dest.try_send(page) { Ok(()) => {} - Err(flume::TrySendError::Full(page)) => { + Err(mpsc::error::TrySendError::Full(page)) => { log::warn!("export: destination channel full, awaiting..."); - dest.send_async(page).await?; + dest.send(page).await?; } e => e?, }; diff --git a/src/weekly.rs b/src/weekly.rs index 22abf61..07a2aea 100644 --- a/src/weekly.rs +++ b/src/weekly.rs @@ -8,6 +8,7 @@ use std::path::PathBuf; use tokio::{ fs::File, io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader}, + sync::mpsc, }; use tokio_stream::wrappers::LinesStream; use tokio_util::compat::FuturesAsyncReadCompatExt; @@ -120,7 +121,7 @@ impl BundleSource for HttpSource { } pub async fn pages_to_weeks( - rx: flume::Receiver, + mut rx: mpsc::Receiver, dir: PathBuf, clobber: bool, ) -> anyhow::Result<()> { @@ -136,7 +137,7 @@ pub async fn pages_to_weeks( let mut week_ops = 0; let mut week_t0 = total_t0; - while let Ok(page) = rx.recv_async().await { + while let Some(page) = rx.recv().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}")) @@ -193,7 +194,7 @@ pub async fn pages_to_weeks( pub async fn week_to_pages( source: impl BundleSource, week: Week, - dest: flume::Sender, + dest: mpsc::Sender, ) -> anyhow::Result<()> { use futures::TryStreamExt; let decoder = GzipDecoder::new(BufReader::new(source.reader_for(week).await?)); @@ -202,7 +203,7 @@ pub async fn week_to_pages( while let Some(chunk) = chunks.try_next().await? { let ops: Vec = chunk.into_iter().collect(); let page = ExportPage { ops }; - dest.send_async(page).await?; + dest.send(page).await?; } Ok(()) }