From 94021f4b78fb948a649116a279a10fd6dc0ecfbd Mon Sep 17 00:00:00 2001 From: "Willow (GHOST)" Date: Sat, 25 Jul 2026 14:07:39 +0100 Subject: [PATCH] feat(uplc): setup http server don't love the code, but some hard negotiating took place to get us here --- uplc/Cargo.lock | 113 +++++++++++++++++++++++++++++++++++++++++++++++ uplc/Cargo.toml | 4 +- uplc/src/db.rs | 18 +++++--- uplc/src/main.rs | 62 ++++++++++++++++++++++++++ 4 files changed, 189 insertions(+), 8 deletions(-) diff --git a/uplc/Cargo.lock b/uplc/Cargo.lock index db256b9..56bc3aa 100644 --- a/uplc/Cargo.lock +++ b/uplc/Cargo.lock @@ -265,6 +265,58 @@ dependencies = [ "pkg-config", ] +[[package]] +name = "axum" +version = "0.8.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "31b698c5f9a010f6573133b09e0de5408834d0c82f8d7475a89fc1867a71cd90" +dependencies = [ + "axum-core", + "bytes", + "form_urlencoded", + "futures-util", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-util", + "itoa", + "matchit", + "memchr", + "mime", + "percent-encoding", + "pin-project-lite", + "serde_core", + "serde_json", + "serde_path_to_error", + "serde_urlencoded", + "sync_wrapper", + "tokio", + "tower", + "tower-layer", + "tower-service", + "tracing", +] + +[[package]] +name = "axum-core" +version = "0.5.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "08c78f31d7b1291f7ee735c1c6780ccde7785daae9a9206026862dab7d8792d1" +dependencies = [ + "bytes", + "futures-core", + "http", + "http-body", + "http-body-util", + "mime", + "pin-project-lite", + "sync_wrapper", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "backtrace" version = "0.3.76" @@ -541,6 +593,7 @@ dependencies = [ "hashlink", "libduckdb-sys", "num-integer", + "r2d2", "strum", ] @@ -835,6 +888,12 @@ version = "1.10.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" +[[package]] +name = "httpdate" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" + [[package]] name = "hyper" version = "1.11.0" @@ -849,6 +908,7 @@ dependencies = [ "http", "http-body", "httparse", + "httpdate", "itoa", "pin-project-lite", "smallvec", @@ -1251,6 +1311,12 @@ version = "0.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "112b39cec0b298b6c1999fee3e31427f74f676e4cb9879ed1a121b43661a4154" +[[package]] +name = "matchit" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "47e1ffaa40ddd1f3ed91f717a33c8c0ee23fff369e3aa8772b9605cc1d22f4c3" + [[package]] name = "memchr" version = "2.8.3" @@ -1486,6 +1552,17 @@ version = "6.0.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f8dcc9c7d52a811697d2151c701e0d08956f92b0e24136cf4cf27b57a6a0d9bf" +[[package]] +name = "r2d2" +version = "0.8.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "51de85fb3fb6524929c8a2eb85e6b6d363de4e8c48f9e2c2eac4944abc181c93" +dependencies = [ + "log", + "parking_lot", + "scheduled-thread-pool", +] + [[package]] name = "rand" version = "0.10.2" @@ -1756,6 +1833,15 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "scheduled-thread-pool" +version = "0.2.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3cbc66816425a074528352f5789333ecff06ca41b36b0b0efdfbb29edc391a19" +dependencies = [ + "parking_lot", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -1834,6 +1920,29 @@ dependencies = [ "zmij", ] +[[package]] +name = "serde_path_to_error" +version = "0.1.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "10a9ff822e371bb5403e391ecd83e182e0e77ba7f6fe0160b795797109d1b457" +dependencies = [ + "itoa", + "serde", + "serde_core", +] + +[[package]] +name = "serde_urlencoded" +version = "0.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3491c14715ca2294c4d6a88f15e84739788c1d030eed8c110436aafdaa2f3fd" +dependencies = [ + "form_urlencoded", + "itoa", + "ryu", + "serde", +] + [[package]] name = "sharded-slab" version = "0.1.7" @@ -2138,6 +2247,7 @@ dependencies = [ "tokio", "tower-layer", "tower-service", + "tracing", ] [[package]] @@ -2176,6 +2286,7 @@ version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ + "log", "pin-project-lite", "tracing-core", ] @@ -2245,8 +2356,10 @@ checksum = "8ecb6da28b8a351d773b68d5825ac39017e680750f980f3a1a85cd8dd28a47c1" name = "uplc" version = "0.1.0" dependencies = [ + "axum", "color-eyre", "duckdb", + "r2d2", "reqwest", "serde", "serde_json", diff --git a/uplc/Cargo.toml b/uplc/Cargo.toml index 9cdc65e..360bab0 100644 --- a/uplc/Cargo.toml +++ b/uplc/Cargo.toml @@ -5,8 +5,10 @@ edition = "2024" [dependencies] color-eyre = "0.6.5" -duckdb = { version = "1.10505.0" } +duckdb = { version = "1.10505.0", features = ["r2d2"] } reqwest = "0.13.4" serde = { version = "1.0.229", features = ["derive"] } serde_json = "1.0.151" tokio = { version = "1.53.1", features = ["full"] } +axum = "0.8.1" +r2d2 = "0.8.10" diff --git a/uplc/src/db.rs b/uplc/src/db.rs index ebc7e5d..c5f505b 100644 --- a/uplc/src/db.rs +++ b/uplc/src/db.rs @@ -1,13 +1,16 @@ use color_eyre::eyre::Result; -use duckdb::Connection; +use duckdb::DuckdbConnectionManager; +#[derive(Clone)] pub struct DB { - connection: Connection, + pool: r2d2::Pool, } impl DB { pub fn new(path: String) -> Result { - let connection = duckdb::Connection::open(path)?; + let manager = DuckdbConnectionManager::file(path)?; + let pool = r2d2::Pool::new(manager)?; + let connection = pool.get()?; connection.execute( r#"CREATE TABLE IF NOT EXISTS operations ( @@ -31,11 +34,11 @@ impl DB { [], )?; - Ok(DB { connection }) + Ok(DB { pool }) } pub fn upsert_operation(&self, op: OperationLogEntry) -> Result<()> { - self.connection.execute( + self.pool.get()?.execute( r#"INSERT INTO operations (did, cid, prev, op_type, service, created_at, nullified) VALUES (?, ?, ?, ?, ?, ?, ?) @@ -56,7 +59,7 @@ impl DB { } pub fn set_cursor(&self, cursor: &str) -> Result<()> { - self.connection.execute( + self.pool.get()?.execute( r#"INSERT INTO kv (key, value) VALUES ('cursor', ?) ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value"#, @@ -68,7 +71,8 @@ impl DB { pub fn get_cursor(&self) -> Result> { let result = self - .connection + .pool + .get()? .prepare("SELECT value FROM kv WHERE key = 'cursor'")? .query_one([], |row| row.get::<_, String>(0)); diff --git a/uplc/src/main.rs b/uplc/src/main.rs index 27788cc..62ed2df 100644 --- a/uplc/src/main.rs +++ b/uplc/src/main.rs @@ -1,19 +1,81 @@ use crate::{db::DB, plc::ExportedOp}; +use axum::{ + Router, + http::StatusCode, + response::{IntoResponse, Response}, + routing::get, +}; use color_eyre::eyre::{Result, eyre}; +use tokio::net::TcpListener; mod db; mod plc; const USER_AGENT: &str = "change-me (+https://willow.sh)"; +struct AppError(color_eyre::eyre::Error); + +impl IntoResponse for AppError { + fn into_response(self) -> Response { + (StatusCode::INTERNAL_SERVER_ERROR, self.0.to_string()).into_response() + } +} + +impl> From for AppError { + fn from(err: E) -> Self { + AppError(err.into()) + } +} + #[tokio::main] async fn main() -> Result<()> { + color_eyre::install()?; let db = DB::new("./.data/uplc.db".into())?; + + let app = Router::new().route( + "/", + get({ + let db = db.clone(); + move || async move { + db.get_cursor() + .map(|c| c.unwrap_or(String::from("1970-01-01T00:00:00.000Z"))) + .map_err(AppError::from) + } + }), + ); + + let listener = TcpListener::bind("127.0.0.1:7712").await?; + println!("HTTP server listening on http://127.0.0.1:7712"); + let server = axum::serve(listener, app); + + let sync_handle = tokio::spawn({ + let db = db.clone(); + async move { sync_loop(db).await } + }); + + tokio::select! { + result = server => { + if let Err(e) = result { + eprintln!("HTTP server error: {}", e); + } + } + result = sync_handle => { + if let Err(e) = result { + eprintln!("Sync loop panicked: {}", e); + } + } + } + + Ok(()) +} + +async fn sync_loop(db: DB) -> Result<()> { let client = reqwest::Client::builder().user_agent(USER_AGENT).build()?; let mut after = db .get_cursor()? .unwrap_or(String::from("1970-01-01T00:00:00.000Z")); + let mut total = 0; let mut page = 0; -- 2.51.2