From e32ad0dfc29d6711a1e4ec2dac2efc0238509c9c Mon Sep 17 00:00:00 2001 From: phil Date: Mon, 15 Sep 2025 13:18:14 -0400 Subject: [PATCH] make tailing plc work --- src/backfill.rs | 10 ++----- src/bin/get_backfill_chunk_adsf.rs | 12 ++------ src/bin/main.rs | 30 ++++++-------------- src/bin/tail.rs | 38 ++++++++++++++++++++++++++ src/bin/tail_export.rs | 44 ------------------------------ src/client.rs | 13 +++++++++ src/lib.rs | 22 ++++++++++++++- src/poll.rs | 23 ++++++++-------- 8 files changed, 96 insertions(+), 96 deletions(-) create mode 100644 src/bin/tail.rs delete mode 100644 src/bin/tail_export.rs create mode 100644 src/client.rs diff --git a/src/backfill.rs b/src/backfill.rs index 85751ab..68bbf0d 100644 --- a/src/backfill.rs +++ b/src/backfill.rs @@ -1,15 +1,11 @@ -use crate::ExportPage; +use crate::{CLIENT, ExportPage}; use url::Url; use async_compression::futures::bufread::GzipDecoder; use futures::{AsyncBufReadExt, StreamExt, TryStreamExt, io}; -pub async fn week_to_pages( - client: &reqwest::Client, - url: Url, - dest: flume::Sender, -) -> anyhow::Result<()> { - let reader = client +pub async fn week_to_pages(url: Url, dest: flume::Sender) -> anyhow::Result<()> { + let reader = CLIENT .get(url) .send() .await? diff --git a/src/bin/get_backfill_chunk_adsf.rs b/src/bin/get_backfill_chunk_adsf.rs index be33391..4231b66 100644 --- a/src/bin/get_backfill_chunk_adsf.rs +++ b/src/bin/get_backfill_chunk_adsf.rs @@ -1,18 +1,10 @@ +use allegedly::CLIENT; use async_compression::futures::bufread::GzipDecoder; use futures::{AsyncBufReadExt, StreamExt, TryStreamExt, io}; #[tokio::main] async fn main() { - let client = reqwest::Client::builder() - .user_agent(concat!( - "allegedly (blah) v", - env!("CARGO_PKG_VERSION"), - " (from @microcosm.blue; contact @bad-example.com)" - )) - .build() - .unwrap(); - - let reader = client + 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() diff --git a/src/bin/main.rs b/src/bin/main.rs index 480fc39..50f6ebe 100644 --- a/src/bin/main.rs +++ b/src/bin/main.rs @@ -4,7 +4,7 @@ use std::time::Duration; use tokio_postgres::NoTls; use url::Url; -use allegedly::{ExportPage, poll_upstream, week_to_pages}; +use allegedly::{Dt, ExportPage, bin_init, poll_upstream, week_to_pages}; const EXPORT_PAGE_QUEUE_SIZE: usize = 0; // rendezvous for now const WEEK_IN_SECONDS: u64 = 7 * 86400; @@ -47,17 +47,13 @@ struct Args { struct Op<'a> { pub did: &'a str, pub cid: &'a str, - pub created_at: chrono::DateTime, + pub created_at: Dt, pub nullified: bool, #[serde(borrow)] pub operation: &'a serde_json::value::RawValue, } -async fn bulk_backfill( - client: reqwest::Client, - (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() @@ -68,7 +64,7 @@ async fn bulk_backfill( while week < immutable_week { log::info!("backfilling week {week_n} ({week})"); let url = upstream.join(&format!("{week}.jsonl.gz")).unwrap(); - week_to_pages(&client, url, tx.clone()).await.unwrap(); + week_to_pages(url, tx.clone()).await.unwrap(); week_n += 1; week += WEEK_IN_SECONDS; } @@ -81,21 +77,13 @@ async fn export_upstream( pg_client: tokio_postgres::Client, ) { let latest = get_latest(&pg_client).await; - let client = reqwest::Client::builder() - .user_agent(concat!( - "allegedly v", - env!("CARGO_PKG_VERSION"), - " (from @microcosm.blue; contact @bad-example.com)" - )) - .build() - .unwrap(); if latest.is_none() { - bulk_backfill(client.clone(), bulk, tx.clone()).await; + bulk_backfill(bulk, tx.clone()).await; } let mut upstream = upstream; upstream.set_path("/export"); - poll_upstream(&client, latest, upstream, tx).await.unwrap(); + poll_upstream(latest, upstream, tx).await.unwrap(); } async fn write_pages( @@ -180,7 +168,7 @@ async fn write_pages( Ok(()) } -async fn get_latest(pg_client: &tokio_postgres::Client) -> Option> { +async fn get_latest(pg_client: &tokio_postgres::Client) -> Option
{ pg_client .query_opt( r#"SELECT "createdAt" FROM operations @@ -194,9 +182,7 @@ async fn get_latest(pg_client: &tokio_postgres::Client) -> Option Vec { - client - .get(url) - .send() - .await - .unwrap() - .error_for_status() - .unwrap() - .text() - .await - .unwrap() - .trim() - .split('\n') - .map(Into::into) - .collect() -} - -#[tokio::main] -async fn main() { - let client = reqwest::Client::builder() - .user_agent(concat!( - "allegedly (export) v", - env!("CARGO_PKG_VERSION"), - " (from @microcosm.blue; contact @bad-example.com)" - )) - .build() - .unwrap(); - - let mut url = Url::parse("https://plc.directory/export").unwrap(); - let ops = get_page(&client, url.clone()).await; - - println!("first: {:?}", ops.first()); - - if let Some(last_line) = ops.last() { - let x: OpPeek = serde_json::from_str(last_line).unwrap(); - url.query_pairs_mut() - .append_pair("after", &x.created_at.to_rfc3339()); - let ops2 = get_page(&client, url).await; - println!("2nd: {:?}", ops2.first()); - } -} diff --git a/src/client.rs b/src/client.rs new file mode 100644 index 0000000..2320e50 --- /dev/null +++ b/src/client.rs @@ -0,0 +1,13 @@ +use reqwest::Client; +use std::sync::LazyLock; + +pub static CLIENT: LazyLock = LazyLock::new(|| { + Client::builder() + .user_agent(concat!( + "allegedly, v", + env!("CARGO_PKG_VERSION"), + " (from @microcosm.blue; contact @bad-example.com)" + )) + .build() + .unwrap() +}); diff --git a/src/lib.rs b/src/lib.rs index 8f8dd71..515d19b 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,11 +1,15 @@ use serde::Deserialize; mod backfill; +mod client; mod poll; pub use backfill::week_to_pages; +pub use client::CLIENT; pub use poll::poll_upstream; +pub type Dt = chrono::DateTime; + /// One page of PLC export /// /// Not limited, but expected to have up to about 1000 lines @@ -16,5 +20,21 @@ pub struct ExportPage { #[derive(Deserialize)] #[serde(rename_all = "camelCase")] pub struct OpPeek { - pub created_at: chrono::DateTime, + pub created_at: Dt, +} + +pub fn bin_init(name: &str) { + use env_logger::{Builder, Env}; + Builder::from_env(Env::new().filter_or("RUST_LOG", "info")).init(); + + log::info!( + r" + + \ | | | | + _ \ | | -_) _` | -_) _` | | | | ({name}) + _/ _\ _| _| \___| \__, | \___| \__,_| _| \_, | (v{}) + ____| __/ +", + env!("CARGO_PKG_VERSION") + ); } diff --git a/src/poll.rs b/src/poll.rs index 5bd7aa1..0c3b44d 100644 --- a/src/poll.rs +++ b/src/poll.rs @@ -1,5 +1,4 @@ -use crate::{ExportPage, OpPeek}; -use chrono::{DateTime, Utc}; +use crate::{CLIENT, Dt, ExportPage, OpPeek}; use std::time::Duration; use thiserror::Error; use url::Url; @@ -14,11 +13,10 @@ pub enum GetPageError { SerdeError(#[from] serde_json::Error), } -pub async fn get_page( - client: &reqwest::Client, - url: Url, -) -> Result<(ExportPage, Option>), GetPageError> { - let ops: Vec = client +pub async fn get_page(url: Url) -> Result<(ExportPage, Option
), GetPageError> { + log::trace!("Getting page: {url}"); + + let ops: Vec = CLIENT .get(url) .send() .await? @@ -32,16 +30,17 @@ pub async fn get_page( let last_at = ops .last() + .filter(|s| !s.is_empty()) .map(|s| serde_json::from_str::(s)) .transpose()? - .map(|o| o.created_at); + .map(|o| o.created_at) + .inspect(|at| log::trace!("new last_at: {at}")); Ok((ExportPage { ops }, last_at)) } pub async fn poll_upstream( - client: &reqwest::Client, - after: Option>, + after: Option
, base: Url, dest: flume::Sender, ) -> anyhow::Result<()> { @@ -53,8 +52,8 @@ pub async fn poll_upstream( if let Some(a) = after { url.query_pairs_mut().append_pair("after", &a.to_rfc3339()); }; - let (page, next_after) = get_page(client, url).await?; + let (page, next_after) = get_page(url).await?; dest.send_async(page).await?; - after = next_after; + after = next_after.or(after); } } -- 2.51.2