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")))
+}