diff --git a/Cargo.lock b/Cargo.lock index eea2958..5578468 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2057,6 +2057,7 @@ version = "0.12.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d429f34c8092b2d42c7c93cec323bb4adeb7c67698f70839adec842ec10c7ceb" dependencies = [ + "async-compression", "base64", "bytes", "encoding_rs", diff --git a/Cargo.toml b/Cargo.toml index f412dba..f9b1b26 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -18,7 +18,7 @@ log = "0.4.28" native-tls = "0.2.14" poem = { version = "3.1.12", features = ["acme", "compression"] } postgres-native-tls = "0.5.1" -reqwest = { version = "0.12.23", features = ["stream", "json"] } +reqwest = { version = "0.12.23", features = ["stream", "json", "gzip"] } reqwest-middleware = "0.4.2" reqwest-retry = "0.7.0" rustls = "0.23.32" diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index ec4da15..a51148b 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -1,6 +1,6 @@ use allegedly::{Dt, bin::GlobalArgs, bin_init, pages_to_stdout, pages_to_weeks, poll_upstream}; use clap::{CommandFactory, Parser, Subcommand}; -use std::{path::PathBuf, time::Instant}; +use std::{path::PathBuf, time::Duration, time::Instant}; use tokio::fs::create_dir_all; use tokio::sync::mpsc; @@ -76,9 +76,10 @@ async fn main() -> anyhow::Result<()> { } => { let mut url = globals.upstream; url.set_path("/export"); + let throttle = Duration::from_millis(globals.upstream_throttle_ms); 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) + poll_upstream(Some(after), url, throttle, tx) .await .expect("to poll upstream") }); @@ -95,9 +96,10 @@ async fn main() -> anyhow::Result<()> { let mut url = globals.upstream; url.set_path("/export"); let start_at = after.or_else(|| Some(chrono::Utc::now())); + let throttle = Duration::from_millis(globals.upstream_throttle_ms); let (tx, rx) = mpsc::channel(1); tokio::task::spawn(async move { - poll_upstream(start_at, url, tx) + poll_upstream(start_at, url, throttle, tx) .await .expect("to poll upstream") }); diff --git a/src/bin/backfill.rs b/src/bin/backfill.rs index 0311877..36256cb 100644 --- a/src/bin/backfill.rs +++ b/src/bin/backfill.rs @@ -4,7 +4,7 @@ use allegedly::{ }; use clap::Parser; use reqwest::Url; -use std::path::PathBuf; +use std::{path::PathBuf, time::Duration}; use tokio::{ sync::{mpsc, oneshot}, task::JoinSet, @@ -53,7 +53,10 @@ pub struct Args { } pub async fn run( - GlobalArgs { upstream }: GlobalArgs, + GlobalArgs { + upstream, + upstream_throttle_ms, + }: GlobalArgs, Args { http, dir, @@ -98,7 +101,8 @@ pub async fn run( } let mut upstream = upstream; upstream.set_path("/export"); - tasks.spawn(poll_upstream(None, upstream, poll_tx)); + let throttle = Duration::from_millis(upstream_throttle_ms); + tasks.spawn(poll_upstream(None, upstream, throttle, poll_tx)); tasks.spawn(full_pages(poll_out, full_tx)); tasks.spawn(pages_to_stdout(full_out, None)); } else { @@ -128,10 +132,12 @@ pub async fn run( // and the catch-up source... if let Some(last) = found_last_out { + let throttle = Duration::from_millis(upstream_throttle_ms); tasks.spawn(async move { let mut upstream = upstream; upstream.set_path("/export"); - poll_upstream(last.await?, upstream, poll_tx).await + + poll_upstream(last.await?, upstream, throttle, poll_tx).await }); } diff --git a/src/bin/mirror.rs b/src/bin/mirror.rs index 554cc7a..c612efa 100644 --- a/src/bin/mirror.rs +++ b/src/bin/mirror.rs @@ -3,7 +3,7 @@ use allegedly::{ }; use clap::Parser; use reqwest::Url; -use std::{net::SocketAddr, path::PathBuf}; +use std::{net::SocketAddr, path::PathBuf, time::Duration}; use tokio::{fs::create_dir_all, sync::mpsc, task::JoinSet}; #[derive(Debug, clap::Args)] @@ -60,7 +60,10 @@ pub struct Args { } pub async fn run( - GlobalArgs { upstream }: GlobalArgs, + GlobalArgs { + upstream, + upstream_throttle_ms, + }: GlobalArgs, Args { wrap, wrap_pg, @@ -113,8 +116,9 @@ pub async fn run( let mut poll_url = upstream.clone(); poll_url.set_path("/export"); + let throttle = Duration::from_millis(upstream_throttle_ms); - tasks.spawn(poll_upstream(Some(latest), poll_url, send_page)); + tasks.spawn(poll_upstream(Some(latest), poll_url, throttle, send_page)); tasks.spawn(pages_to_pg(db.clone(), recv_page)); tasks.spawn(serve( upstream, diff --git a/src/bin/mod.rs b/src/bin/mod.rs index 845edd9..1ad965d 100644 --- a/src/bin/mod.rs +++ b/src/bin/mod.rs @@ -6,6 +6,12 @@ pub struct GlobalArgs { #[arg(short, long, global = true, env = "ALLEGEDLY_UPSTREAM")] #[clap(default_value = "https://plc.directory")] pub upstream: Url, + /// Self-rate-limit upstream request interval + /// + /// plc.directory's rate limiting is 500 requests per 5 mins (600ms) + #[arg(long, global = true, env = "ALLEGEDLY_UPSTREAM_THROTTLE_MS")] + #[clap(default_value = "600")] + pub upstream_throttle_ms: u64, } #[allow(dead_code)] diff --git a/src/client.rs b/src/client.rs index d154ca6..26be25f 100644 --- a/src/client.rs +++ b/src/client.rs @@ -12,6 +12,7 @@ pub const UA: &str = concat!( pub static CLIENT: LazyLock = LazyLock::new(|| { let inner = Client::builder() .user_agent(UA) + .gzip(true) .build() .expect("reqwest client to build"); diff --git a/src/poll.rs b/src/poll.rs index c997af9..ee0f65e 100644 --- a/src/poll.rs +++ b/src/poll.rs @@ -4,9 +4,6 @@ use std::time::Duration; use thiserror::Error; use tokio::sync::mpsc; -// plc.directory ratelimit on /export is 500 per 5 mins -const UPSTREAM_REQUEST_INTERVAL: Duration = Duration::from_millis(600); - #[derive(Debug, Error)] pub enum GetPageError { #[error(transparent)] @@ -141,7 +138,11 @@ pub async fn get_page(url: Url) -> Result<(ExportPage, Option), GetPageE .split('\n') .filter_map(|s| { serde_json::from_str::(s) - .inspect_err(|e| log::warn!("failed to parse op: {e} ({s})")) + .inspect_err(|e| { + if !s.is_empty() { + log::warn!("failed to parse op: {e} ({s})") + } + }) .ok() }) .collect(); @@ -154,10 +155,11 @@ pub async fn get_page(url: Url) -> Result<(ExportPage, Option), GetPageE pub async fn poll_upstream( after: Option
, base: Url, + throttle: Duration, dest: mpsc::Sender, ) -> anyhow::Result<&'static str> { - log::info!("starting upstream poller after {after:?}"); - let mut tick = tokio::time::interval(UPSTREAM_REQUEST_INTERVAL); + log::info!("starting upstream poller at {base} after {after:?}"); + let mut tick = tokio::time::interval(throttle); let mut prev_last: Option = after.map(Into::into); let mut boundary_state: Option = None; loop {