From bc437466979ddd8a3d4939aec2a4eb0a57818c83 Mon Sep 17 00:00:00 2001 From: "Willow (GHOST)" Date: Sun, 26 Jul 2026 18:31:55 +0100 Subject: [PATCH] perf(uplc): only read state from db on startup --- src/routes/plc/CountsGraph.svelte | 14 +- src/routes/plc/PdsDistribution.svelte | 21 ++- src/routes/plc/PdsShare.svelte | 28 +++- src/routes/plc/plc.remote.ts | 137 ++++++++++-------- uplc/src/api.rs | 170 ++-------------------- uplc/src/main.rs | 18 ++- uplc/src/plc.rs | 11 +- uplc/src/stats.rs | 201 ++++++++++++++++++++++++++ 8 files changed, 350 insertions(+), 250 deletions(-) create mode 100644 uplc/src/stats.rs diff --git a/src/routes/plc/CountsGraph.svelte b/src/routes/plc/CountsGraph.svelte index d7a7a98..218abba 100644 --- a/src/routes/plc/CountsGraph.svelte +++ b/src/routes/plc/CountsGraph.svelte @@ -1,18 +1,18 @@

Operations & Dids over time

- {#if countsGraph.loading} + {#if stats.loading}

Loading...

- {:else if countsGraph.current} + {:else if stats.current} - {:else if countsGraph.error} -

{countsGraph.error}

+ {:else if stats.error} +

{stats.error}

{/if}
diff --git a/src/routes/plc/PdsDistribution.svelte b/src/routes/plc/PdsDistribution.svelte index 0a9876d..eca0f49 100644 --- a/src/routes/plc/PdsDistribution.svelte +++ b/src/routes/plc/PdsDistribution.svelte @@ -1,26 +1,26 @@
-

PDS Distribution

+

PDS Distribution (approx)

- {#if pdsDistribution.loading} + {#if stats.loading}

Loading...

- {:else if pdsDistribution.current} + {:else if stats.current} - {:else if pdsDistribution.error} -

{pdsDistribution.error}

+ {:else if stats.error} +

{stats.error}

{/if}
@@ -35,5 +35,10 @@ margin-top: 4px; margin-bottom: 16px; } + + small { + font-size: 0.9rem; + color: var(--text-grey); + } } diff --git a/src/routes/plc/PdsShare.svelte b/src/routes/plc/PdsShare.svelte index 96352e7..232b275 100644 --- a/src/routes/plc/PdsShare.svelte +++ b/src/routes/plc/PdsShare.svelte @@ -1,26 +1,33 @@
-

Did Allocations

+

Did Allocations (approx)

- {#if pdsShare.loading} + {#if stats.loading}

Loading...

- {:else if pdsShare.current} + {:else if stats.current} - {:else if pdsShare.error} -

{pdsShare.error}

+ {:else if stats.error} +

{stats.error}

{/if}
@@ -34,5 +41,10 @@ margin-top: 4px; margin-bottom: 16px; } + + small { + font-size: 0.9rem; + color: var(--text-grey); + } } diff --git a/src/routes/plc/plc.remote.ts b/src/routes/plc/plc.remote.ts index a87d778..facd896 100644 --- a/src/routes/plc/plc.remote.ts +++ b/src/routes/plc/plc.remote.ts @@ -1,88 +1,105 @@ import { UPLC_ENDPOINT } from '$app/env/private'; import { query } from '$app/server'; -async function grab(path: string) { - const url = new URL(UPLC_ENDPOINT!); - url.pathname = path; - - const response = await fetch(url); - if (!response.ok) throw new Error('broked'); +interface HistoryEntry { + dids: number; + operations: number; +} - const data = await response.json(); - return data as T; +interface PdsShare { + deactivated: number; + withPds: number; + withoutPds: number; } -export interface Counts { +interface RawStats { + dids: number; lastOperation: string; operations: number; - dids: number; -} - -export const getCounts = query.live(async function* () { - while (true) { - yield await grab('/v0/counts'); - await new Promise((resolve) => setTimeout(resolve, 10_000)); - } -}); - -type PatchDate = Omit & { date: Date }; - -function patchDate(data: T): PatchDate { - // @ts-expect-error shhh - data.date = new Date(data.date); - return data as unknown as PatchDate; + history: Record; + pdsDistribution: Record; + pdsShare: PdsShare; } -interface RawCountsGraphPoint { - date: string; +interface HistoryGraphPoint { + date: Date; dids: number; operations: number; } -export const getCountsGraph = query(async () => { - const raw = await grab('/v0/counts/history'); - return raw.map(patchDate); -}); - -interface PdsShare { - deactivated: number; - withPds: number; - withoutPds: number; -} +async function fetchRawStats() { + const url = new URL(UPLC_ENDPOINT!); + url.pathname = '/v0/stats'; -export const getPdsShare = query(async () => { - const raw = await grab('/v0/pds/share'); + const response = await fetch(url); + if (!response.ok) throw new Error('broked'); - return [ - { key: 'PDS', value: raw.withPds }, - { key: 'No PDS', value: raw.withoutPds }, - { key: 'Deactivated', value: raw.deactivated }, - ]; -}); + const data = await response.json(); + return data as RawStats; +} -interface PdsDistribution { - pds: string; - count: number; +interface Stats { + pdsDistribution: { key: string; value: number }[]; + history: HistoryGraphPoint[]; + pdsShare: PdsShare; } -export const getPdsDistribution = query(async () => { - const raw = await grab('/v0/pds/distribution'); - const data = new Map(); +export const getStats = query(async () => { + const raw = await fetchRawStats(); + const distribution = new Map(); - for (const point of raw) { - const url = URL.parse(point.pds); + for (const [pds, count] of Object.entries(raw.pdsDistribution)) { + const url = URL.parse(pds); const key = url?.hostname.endsWith('bsky.network') ? 'Bluesky PBC' - : url && url.protocol === 'https:' && point.count > 100 + : url && url.protocol === 'https:' && count > 100 ? url.hostname : 'Other'; - const count = data.getOrInsert(key, 0); - data.set(key, count + point.count); + const total = distribution.getOrInsert(key, 0); + distribution.set(key, total + count); } - return data - .entries() - .map((d) => ({ key: d[0], value: d[1] })) - .toArray(); + return { + pdsShare: raw.pdsShare, + pdsDistribution: distribution + .entries() + .map((d) => ({ key: d[0], value: d[1] })) + .toArray(), + history: Object.entries(raw.history) + .map(([date, d]) => ({ + date: new Date(date), + ...d, + })) + .sort((a, b) => a.date.getTime() - b.date.getTime()) + .reduce((acc, d) => { + const prev = acc.at(-1); + acc.push({ + date: d.date, + dids: d.dids + (prev?.dids ?? 0), + operations: d.operations + (prev?.operations ?? 0), + }); + return acc; + }, []), + }; +}); + +export interface Counts { + lastOperation: string; + operations: number; + dids: number; +} + +export const getCounts = query.live(async function* () { + while (true) { + const data = await fetchRawStats(); + + yield { + dids: data.dids, + operations: data.operations, + lastOperation: data.lastOperation, + }; + + await new Promise((resolve) => setTimeout(resolve, 10_000)); + } }); diff --git a/uplc/src/api.rs b/uplc/src/api.rs index e539799..2d0dba8 100644 --- a/uplc/src/api.rs +++ b/uplc/src/api.rs @@ -1,14 +1,14 @@ -use crate::db::DB; +use crate::{db::DB, stats::Stats}; use axum::{ Json, Router, http, response::{IntoResponse, Response}, routing::get, }; use color_eyre::eyre::Result; -use serde_json::{Value, json}; -use std::{env, fmt}; -use tokio::net::TcpListener; +use std::{env, fmt, sync::Arc}; +use tokio::{net::TcpListener, sync::RwLock}; +#[allow(unused)] struct AppError(color_eyre::eyre::Error); impl IntoResponse for AppError { @@ -24,167 +24,15 @@ impl + fmt::Debug> From for AppError { } } -fn get_counts(db: &DB) -> Result> { - Ok(db - .handle()? - .prepare( - r#"SELECT - count(DISTINCT did) as dids, - count(did) as operations, - strftime(max(created_at), '%Y-%m-%dT%H:%M:%SZ') as last_operation - FROM operations"#, - )? - .query_one([], |row| { - Ok(Json(json!({ - "dids": row.get::<_, i64>(0)?, - "operations": row.get::<_, i64>(1)?, - "lastOperation": row.get::<_, String>(2)?, - }))) - })?) -} - -fn get_counts_history(db: &DB) -> Result>> { - let conn = db.handle()?; - let mut stmt = conn.prepare( - r#" - WITH - -- Bucket every operation by week; flag the first occurrence of each DID - bucketed AS ( - SELECT - DATE_TRUNC('week', created_at) AS bucket_time, - ROW_NUMBER() OVER (PARTITION BY did ORDER BY created_at) = 1 AS is_new_did - FROM operations - ), - -- Per-bucket deltas for both metrics in one pass - bucket_stats AS ( - SELECT - bucket_time, - COUNT(*) FILTER (WHERE is_new_did) AS new_dids, - COUNT(*) AS new_ops - FROM bucketed - GROUP BY ALL - ) - SELECT - strftime(bucket_time::TIMESTAMP, '%Y-%m-%dT%H:%M:%SZ') AS date, - SUM(new_dids) OVER ( - ORDER BY bucket_time - ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW - ) AS dids, - SUM(new_ops) OVER ( - ORDER BY bucket_time - ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW - ) AS operations - FROM bucket_stats - ORDER BY bucket_time - "#, - )?; - - let rows = stmt.query_map([], |row| { - let date: String = row.get(0)?; - let dids: i64 = row.get(1)?; - let operations: i64 = row.get(2)?; - Ok((date, dids, operations)) - })?; - - let mut results = Vec::new(); - for row in rows { - let (date, dids, operations) = row?; - results.push(json!({ - "date": date, - "dids": dids, - "operations": operations, - })); - } - - Ok(Json(results)) -} - -fn dids_pds_share(db: &DB) -> Result> { - Ok(db - .handle()? - .prepare( - r#"WITH latest_ops AS ( - SELECT DISTINCT ON (did) - did, - service, - tombstone - FROM operations - ORDER BY did, seq DESC - ) - SELECT - SUM(CASE WHEN tombstone THEN 1 ELSE 0 END) AS deactivated_count, - SUM(CASE WHEN NOT tombstone AND service IS NOT NULL THEN 1 ELSE 0 END) AS with_pds_count, - SUM(CASE WHEN NOT tombstone AND service IS NULL THEN 1 ELSE 0 END) AS without_pds_count - FROM latest_ops"#, - )? - .query_one([], |row| { - Ok(Json(json!({ - "deactivated": row.get::<_, i64>(0)?, - "withPds": row.get::<_, i64>(1)?, - "withoutPds": row.get::<_, i64>(2)?, - }))) - })?) -} - -fn pds_distribution(db: &DB) -> Result> { - let services: Vec = db - .handle()? - .prepare( - r#"WITH latest_ops AS ( - SELECT DISTINCT ON (did) - did, - service - FROM operations - ORDER BY did, seq DESC - ) - SELECT - service, - count(*) as did_count - FROM latest_ops - WHERE service IS NOT NULL - GROUP BY service - ORDER BY did_count DESC"#, - )? - .query_map([], |row| { - Ok(json!({ - "pds": row.get::<_, String>(0)?, - "count": row.get::<_, i64>(1)?, - })) - })? - .collect::, _>>()?; - - Ok(Json(json!(services))) -} - -pub async fn run_api(db: DB) -> Result<()> { +#[allow(unused)] +pub async fn run_api(db: DB, stats: Arc>) -> Result<()> { let app = Router::new() .route("/alive", get(async || "I am alive")) .route( - "/v0/counts", - get({ - let db = db.clone(); - move || async move { get_counts(&db).map_err(AppError::from) } - }), - ) - .route( - "/v0/counts/history", - get({ - let db = db.clone(); - move || async move { get_counts_history(&db).map_err(AppError::from) } - }), - ) - .route( - "/v0/pds/share", - get({ - let db = db.clone(); - move || async move { dids_pds_share(&db).map_err(AppError::from) } - }), - ) - .route( - "/v0/pds/distribution", + "/v0/stats", get({ - let db = db.clone(); - move || async move { pds_distribution(&db).map_err(AppError::from) } + let stats = stats.clone(); + move || async move { Json(stats.read().await.json()) } }), ); diff --git a/uplc/src/main.rs b/uplc/src/main.rs index 9ae2031..2dfe36c 100644 --- a/uplc/src/main.rs +++ b/uplc/src/main.rs @@ -1,23 +1,35 @@ -use crate::{api::run_api, db::DB, plc::run_plc_sync}; +use std::sync::Arc; + +use crate::{api::run_api, db::DB, plc::run_plc_sync, stats::Stats}; use color_eyre::eyre::Result; +use tokio::sync::RwLock; mod api; mod db; mod plc; +mod stats; #[tokio::main] async fn main() -> Result<()> { color_eyre::install()?; let db = DB::new("./.data/uplc.db".into())?; + println!("Loading initial stats..."); + let mut stats = Stats::default(); + stats.refresh(&db)?; + let stats = Arc::new(RwLock::new(stats)); + println!("Done! Starting sync & api"); + let plc_sync_handle = tokio::spawn({ let db = db.clone(); - async move { run_plc_sync(db).await } + let stats = stats.clone(); + async move { run_plc_sync(db, stats).await } }); let api_handle = tokio::spawn({ let db = db.clone(); - async move { run_api(db).await } + let stats = stats.clone(); + async move { run_api(db, stats).await } }); tokio::select! { diff --git a/uplc/src/plc.rs b/uplc/src/plc.rs index c66bf31..670b36a 100644 --- a/uplc/src/plc.rs +++ b/uplc/src/plc.rs @@ -1,7 +1,11 @@ -use crate::db::{DB, OperationLogEntry}; +use crate::{ + db::{DB, OperationLogEntry}, + stats::Stats, +}; use color_eyre::eyre::Result; use serde::Deserialize; -use std::{collections::HashMap, time::Duration}; +use std::{collections::HashMap, sync::Arc, time::Duration}; +use tokio::sync::RwLock; #[derive(Debug, Deserialize)] pub struct CreateOpV1 { @@ -72,7 +76,7 @@ impl From for OperationLogEntry { const USER_AGENT: &str = "change-me (+https://willow.sh)"; // todo look at switching to websocket after backfilled -pub async fn run_plc_sync(db: DB) -> Result<()> { +pub async fn run_plc_sync(db: DB, stats: Arc>) -> Result<()> { let client = reqwest::Client::builder().user_agent(USER_AGENT).build()?; let mut after = db.get_cursor()?.unwrap_or(0); @@ -104,6 +108,7 @@ pub async fn run_plc_sync(db: DB) -> Result<()> { after = ops.last().unwrap().seq; total += count; + stats.write().await.count_batch(&ops)?; db.upsert_operations(ops)?; if count < 200 { diff --git a/uplc/src/stats.rs b/uplc/src/stats.rs new file mode 100644 index 0000000..33ece6f --- /dev/null +++ b/uplc/src/stats.rs @@ -0,0 +1,201 @@ +use crate::db::{DB, OperationLogEntry}; +use chrono::{DateTime, Datelike, Duration}; +use color_eyre::eyre::{Result, eyre}; +use serde::Serialize; +use std::collections::HashMap; + +#[derive(Debug, Default, Serialize)] +pub struct CountsHistoryPoint { + pub dids: i64, + pub operations: i64, +} + +#[derive(Debug, Default, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct PdsShare { + pub deactivated: i64, + pub with_pds: i64, + pub without_pds: i64, +} + +#[derive(Debug, Default, Serialize)] +#[serde(rename_all = "camelCase")] +pub struct Stats { + pub operations: i64, + pub dids: i64, + pub last_operation: String, + pub history: HashMap, + pub pds_share: PdsShare, + pub pds_distribution: HashMap, +} + +impl Stats { + pub fn refresh(&mut self, db: &DB) -> Result<()> { + let conn = db.handle()?; + + let mut stmt = conn.prepare( + r#" + WITH bucketed AS ( + SELECT + DATE_TRUNC('week', created_at) AS bucket_time, + ROW_NUMBER() OVER (PARTITION BY did ORDER BY created_at) = 1 AS is_new_did + FROM operations + ) + SELECT + strftime(bucket_time::TIMESTAMP, '%Y-%m-%dT%H:%M:%SZ') AS date, + COUNT(*) FILTER (WHERE is_new_did) AS new_dids, + COUNT(*) AS new_ops + FROM bucketed + GROUP BY ALL + "#, + )?; + + self.history = stmt + .query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + CountsHistoryPoint { + dids: row.get(1)?, + operations: row.get(2)?, + }, + )) + })? + .collect::, _>>()?; + + // Merge the scalar counts, PDS share, and PDS distribution into one query. + // `latest_ops` is computed once; `count(*)` on it replaces `count(DISTINCT did)`. + let mut stmt = conn.prepare( + r#" + WITH + all_ops AS ( + SELECT + count(*) AS total_ops, + strftime(max(created_at), '%Y-%m-%dT%H:%M:%SZ') AS last_op + FROM operations + ), + latest_ops AS ( + SELECT DISTINCT ON (did) + did, + service, + tombstone + FROM operations + ORDER BY did, seq DESC + ), + pds_agg AS ( + SELECT + count(*) AS dids, + SUM(CASE WHEN tombstone THEN 1 ELSE 0 END) AS deactivated, + SUM(CASE WHEN NOT tombstone AND service IS NOT NULL THEN 1 ELSE 0 END) AS with_pds, + SUM(CASE WHEN NOT tombstone AND service IS NULL THEN 1 ELSE 0 END) AS without_pds + FROM latest_ops + ) + SELECT + a.total_ops, + a.last_op, + p.dids, + p.deactivated, + p.with_pds, + p.without_pds, + d.service, + d.count + FROM all_ops a + CROSS JOIN pds_agg p + LEFT JOIN ( + SELECT service, count(*) AS count + FROM latest_ops + WHERE service IS NOT NULL + GROUP BY service + ORDER BY count DESC + ) d ON TRUE + ORDER BY d.count DESC + "#, + )?; + + let rows = stmt + .query_map([], |row| { + Ok(( + row.get::<_, i64>(0)?, // total_ops + row.get::<_, String>(1)?, // last_op + row.get::<_, i64>(2)?, // dids + row.get::<_, i64>(3)?, // deactivated + row.get::<_, i64>(4)?, // with_pds + row.get::<_, i64>(5)?, // without_pds + row.get::<_, Option>(6)?, // service + row.get::<_, Option>(7)?, // count + )) + })? + .collect::, _>>()?; + + if let Some((total_ops, last_op, dids, deactivated, with_pds, without_pds, _, _)) = + rows.first() + { + self.operations = *total_ops; + self.last_operation = last_op.clone(); + self.dids = *dids; + self.pds_share = PdsShare { + deactivated: *deactivated, + with_pds: *with_pds, + without_pds: *without_pds, + }; + } + + self.pds_distribution = rows + .into_iter() + .filter_map(|(_, _, _, _, _, _, service, count)| { + service.map(|pds| (pds, count.unwrap_or(0))) + }) + .collect(); + + Ok(()) + } + + pub fn count_batch(&mut self, batch: &Vec) -> Result<()> { + self.operations += batch.len() as i64; + + for op in batch { + if op.prev.is_none() { + self.dids += 1; + } + + if op.tombstone { + self.pds_share.deactivated += 1; + } + + if let Some(service) = &op.service { + self.pds_share.with_pds += 1; + + if let Some(count) = self.pds_distribution.get_mut(service) { + *count += 1; + } else { + self.pds_distribution.insert(service.clone(), 1); + } + } else { + self.pds_share.without_pds += 1; + } + + let entry = self.history.entry(week_key(&op.created_at)?).or_default(); + entry.operations += 1; + if op.prev.is_none() { + entry.dids += 1; + } + } + + if let Some(last) = batch.last() { + self.last_operation = last.created_at.clone(); + } + + Ok(()) + } + + pub fn json(&self) -> serde_json::Value { + serde_json::json!(self) + } +} + +fn week_key(created_at: &str) -> Result { + let dt = DateTime::parse_from_rfc3339(created_at) + .map_err(|e| eyre!("invalid created_at {created_at:?}: {e}"))?; + let date = dt.date_naive(); + let monday = date - Duration::days(date.weekday().num_days_from_monday() as i64); + Ok(format!("{}T00:00:00Z", monday.format("%Y-%m-%d"))) +} -- 2.51.2