diff --git a/Cargo.lock b/Cargo.lock index 4a64855..39e8bd0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -34,9 +34,9 @@ dependencies = [ "async-compression", "chrono", "clap", - "env_logger", "futures", "log", + "poem", "reqwest", "reqwest-middleware", "reqwest-retry", @@ -47,6 +47,22 @@ dependencies = [ "tokio-postgres", "tokio-stream", "tokio-util", + "tracing-subscriber", +] + +[[package]] +name = "alloc-no-stdlib" +version = "2.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc7bb162ec39d46ab1ca8c77bf72e890535becd1751bb45f64c597edb4c8c6b3" + +[[package]] +name = "alloc-stdlib" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94fb8275041c72129eb51b7d0322c29b8387a0386127718b096429201a5d6ece" +dependencies = [ + "alloc-no-stdlib", ] [[package]] @@ -193,6 +209,27 @@ dependencies = [ "generic-array", ] +[[package]] +name = "brotli" +version = "8.0.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9991eea70ea4f293524138648e41ee89b0b2b12ddef3b255effa43c8056e0e0d" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", + "brotli-decompressor", +] + +[[package]] +name = "brotli-decompressor" +version = "5.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "874bb8112abecc98cbd6d81ea4fa7e94fb9449648c93cc89aa40c81c24d7de03" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", +] + [[package]] name = "bumpalo" version = "3.19.0" @@ -218,6 +255,8 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "65193589c6404eb80b450d618eaf9a2cafaaafd57ecce47370519ef674a7bd44" dependencies = [ "find-msvc-tools", + "jobserver", + "libc", "shlex", ] @@ -227,6 +266,12 @@ version = "1.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2fd1289c04a9ea8cb22300a459a72a385d7c73d3259e2ed7dcb2af674838cfa9" +[[package]] +name = "cfg_aliases" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "613afe47fcd5fac7ccf1db93babcb082c5994d996f20b8b159f2ad1658eb5724" + [[package]] name = "chrono" version = "0.4.42" @@ -293,9 +338,12 @@ version = "0.4.30" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "485abf41ac0c8047c07c87c72c8fb3eb5197f6e9d7ded615dfd1a00ae00a0f64" dependencies = [ + "brotli", "compression-core", "flate2", "memchr", + "zstd", + "zstd-safe", ] [[package]] @@ -379,29 +427,6 @@ dependencies = [ "cfg-if", ] -[[package]] -name = "env_filter" -version = "0.1.3" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "186e05a59d4c50738528153b83b0b0194d3a29507dfec16eccd4b342903397d0" -dependencies = [ - "log", - "regex", -] - -[[package]] -name = "env_logger" -version = "0.11.8" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "13c863f0904021b108aa8b2f55046443e6b1ebde8fd4a15c399893aae4fa069f" -dependencies = [ - "anstream", - "anstyle", - "env_filter", - "jiff", - "log", -] - [[package]] name = "equivalent" version = "1.0.2" @@ -631,6 +656,30 @@ version = "0.15.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9229cfe53dfd69f0609a49f65461bd93001ea1ef889cd5529dd176593f5338a1" +[[package]] +name = "headers" +version = "0.4.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b3314d5adb5d94bcdf56771f2e50dbbc80bb4bdf88967526706205ac9eff24eb" +dependencies = [ + "base64", + "bytes", + "headers-core", + "http", + "httpdate", + "mime", + "sha1", +] + +[[package]] +name = "headers-core" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "54b4a22553d4242c49fddb9ba998a99962b5cc6f22cb5a3482bec22522403ce4" +dependencies = [ + "http", +] + [[package]] name = "heck" version = "0.5.0" @@ -686,6 +735,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.7.0" @@ -700,6 +755,7 @@ dependencies = [ "http", "http-body", "httparse", + "httpdate", "itoa", "pin-project-lite", "pin-utils", @@ -899,9 +955,9 @@ dependencies = [ [[package]] name = "indexmap" -version = "2.11.1" +version = "2.11.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "206a8042aec68fa4a62e8d3f7aa4ceb508177d9324faf261e1959e495b7a1921" +checksum = "4b0f83760fb341a774ed326568e19f5a863af4a952def8c39f9ab92fd95b88e5" dependencies = [ "equivalent", "hashbrown", @@ -959,27 +1015,13 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4a5f13b858c8d314ee3e8f639011f7ccefe71f97f96e50151fb991f267928e2c" [[package]] -name = "jiff" -version = "0.2.15" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "be1f93b8b1eb69c77f24bbb0afdf66f54b632ee39af40ca21c4365a1d7347e49" -dependencies = [ - "jiff-static", - "log", - "portable-atomic", - "portable-atomic-util", - "serde", -] - -[[package]] -name = "jiff-static" -version = "0.2.15" +name = "jobserver" +version = "0.1.34" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "03343451ff899767262ec32146f6d559dd759fdadf42ff0e227c7c48f72594b4" +checksum = "9afb3de4395d6b3e67a780b6de64b51c978ecf11cb9a462c66be7d4ca9039d33" dependencies = [ - "proc-macro2", - "quote", - "syn", + "getrandom 0.3.3", + "libc", ] [[package]] @@ -992,6 +1034,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "lazy_static" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bbd2bcb4c963f2ddae06a2efc7e9f3591312473c50c6685e1f298068316e66fe" + [[package]] name = "libc" version = "0.2.175" @@ -1037,6 +1085,15 @@ version = "0.4.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" +[[package]] +name = "matchers" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d1525a2a28c7f4fa0fc98bb91ae755d1e2d1505079e05539e35bc876b5d65ae9" +dependencies = [ + "regex-automata", +] + [[package]] name = "md-5" version = "0.10.6" @@ -1096,6 +1153,27 @@ dependencies = [ "tempfile", ] +[[package]] +name = "nix" +version = "0.30.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74523f3a35e05aba87a1d978330aef40f67b0304ac79c1c00b294c9830543db6" +dependencies = [ + "bitflags 2.9.4", + "cfg-if", + "cfg_aliases", + "libc", +] + +[[package]] +name = "nu-ansi-term" +version = "0.50.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d4a28e057d01f97e61255210fcff094d74ed0466038633e95017f5beb68e4399" +dependencies = [ + "windows-sys 0.52.0", +] + [[package]] name = "num-traits" version = "0.2.19" @@ -1261,18 +1339,49 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" [[package]] -name = "portable-atomic" -version = "1.11.1" +name = "poem" +version = "3.1.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f84267b20a16ea918e43c6a88433c2d54fa145c92a811b5b047ccbe153674483" +checksum = "9f977080932c87287147dca052951c3e2696f8759863f6b4e4c0c9ffe7a4cc8b" +dependencies = [ + "async-compression", + "bytes", + "futures-util", + "headers", + "http", + "http-body-util", + "hyper", + "hyper-util", + "mime", + "nix", + "parking_lot 0.12.4", + "percent-encoding", + "pin-project-lite", + "poem-derive", + "regex", + "rfc7239", + "serde", + "serde_json", + "serde_urlencoded", + "smallvec", + "sync_wrapper", + "thiserror 2.0.16", + "tokio", + "tokio-util", + "tracing", + "wildmatch", +] [[package]] -name = "portable-atomic-util" -version = "0.2.4" +name = "poem-derive" +version = "3.1.12" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d8a2f0d8d040d7848a709caf78912debcc3f33ee4b3cac47d73d1e1069e83507" +checksum = "056e2fea6de1cb240ffe23cfc4fc370b629f8be83b5f27e16b7acd5231a72de4" dependencies = [ - "portable-atomic", + "proc-macro-crate", + "proc-macro2", + "quote", + "syn", ] [[package]] @@ -1325,6 +1434,15 @@ dependencies = [ "zerocopy", ] +[[package]] +name = "proc-macro-crate" +version = "3.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "219cb19e96be00ab2e37d6e299658a0cfa83e52429179969b0f0121b4ac46983" +dependencies = [ + "toml_edit", +] + [[package]] name = "proc-macro2" version = "1.0.101" @@ -1544,6 +1662,15 @@ dependencies = [ "rand 0.8.5", ] +[[package]] +name = "rfc7239" +version = "0.1.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4a82f1d1e38e9a85bb58ffcfadf22ed6f2c94e8cd8581ec2b0f80a2a6858350f" +dependencies = [ + "uncased", +] + [[package]] name = "ring" version = "0.17.14" @@ -1662,18 +1789,28 @@ dependencies = [ [[package]] name = "serde" -version = "1.0.219" +version = "1.0.226" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5f0e2c6ed6606019b4e29e69dbaba95b11854410e5347d525002456dbbb786b6" +checksum = "0dca6411025b24b60bfa7ec1fe1f8e710ac09782dca409ee8237ba74b51295fd" +dependencies = [ + "serde_core", + "serde_derive", +] + +[[package]] +name = "serde_core" +version = "1.0.226" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba2ba63999edb9dac981fb34b3e5c0d111a69b0924e253ed29d83f7c99e966a4" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.219" +version = "1.0.226" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "5b0276cf7f2c73365f7157c8123c21cd9a50fbbd844757af28ca1f5925fc2a00" +checksum = "8db53ae22f34573731bafa1db20f04027b2d25e02d8205921b569171699cdb33" dependencies = [ "proc-macro2", "quote", @@ -1704,6 +1841,17 @@ dependencies = [ "serde", ] +[[package]] +name = "sha1" +version = "0.10.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e3bf829a2d51ab4a5ddf1352d8470c140cadc8301b2ae1789db023f01cedd6ba" +dependencies = [ + "cfg-if", + "cpufeatures", + "digest", +] + [[package]] name = "sha2" version = "0.10.9" @@ -1715,6 +1863,15 @@ dependencies = [ "digest", ] +[[package]] +name = "sharded-slab" +version = "0.1.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f40ca3c46823713e0d4209592e8d6e826aa57e928f09752619fc696c499637f6" +dependencies = [ + "lazy_static", +] + [[package]] name = "shlex" version = "1.3.0" @@ -1902,6 +2059,15 @@ dependencies = [ "syn", ] +[[package]] +name = "thread_local" +version = "1.1.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f60246a4944f24f6e018aa17cdeffb7818b76356965d03b07d6a9886e8962185" +dependencies = [ + "cfg-if", +] + [[package]] name = "tinystr" version = "0.8.1" @@ -2029,6 +2195,36 @@ dependencies = [ "tokio", ] +[[package]] +name = "toml_datetime" +version = "0.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32f1085dec27c2b6632b04c80b3bb1b4300d6495d1e129693bdda7d91e72eec1" +dependencies = [ + "serde_core", +] + +[[package]] +name = "toml_edit" +version = "0.23.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3effe7c0e86fdff4f69cdd2ccc1b96f933e24811c5441d44904e8683e27184b" +dependencies = [ + "indexmap", + "toml_datetime", + "toml_parser", + "winnow", +] + +[[package]] +name = "toml_parser" +version = "1.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4cf893c33be71572e0e9aa6dd15e6677937abd686b066eac3f8cd3531688a627" +dependencies = [ + "winnow", +] + [[package]] name = "tower" version = "0.5.2" @@ -2103,6 +2299,36 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b9d12581f227e93f094d3af2ae690a574abb8a2b9b7a96e7cfe9647b2b617678" dependencies = [ "once_cell", + "valuable", +] + +[[package]] +name = "tracing-log" +version = "0.2.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ee855f1f400bd0e5c02d150ae5de3840039a3f54b025156404e34c23c03f47c3" +dependencies = [ + "log", + "once_cell", + "tracing-core", +] + +[[package]] +name = "tracing-subscriber" +version = "0.3.20" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2054a14f5307d601f88daf0553e1cbf472acc4f2c51afab632431cdcd72124d5" +dependencies = [ + "matchers", + "nu-ansi-term", + "once_cell", + "regex-automata", + "sharded-slab", + "smallvec", + "thread_local", + "tracing", + "tracing-core", + "tracing-log", ] [[package]] @@ -2117,6 +2343,15 @@ version = "1.18.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1dccffe3ce07af9386bfd29e80c0ab1a8205a2fc34e4bcd40364df902cfa8f3f" +[[package]] +name = "uncased" +version = "0.9.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e1b88fcfe09e89d3866a5c11019378088af2d24c3fbd4f0543f96b479ec90697" +dependencies = [ + "version_check", +] + [[package]] name = "unicode-bidi" version = "0.3.18" @@ -2174,6 +2409,12 @@ version = "0.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "06abde3611657adf66d383f00b093d7faecc7fa57071cce2578660c9f1010821" +[[package]] +name = "valuable" +version = "0.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" + [[package]] name = "vcpkg" version = "0.2.15" @@ -2346,6 +2587,12 @@ dependencies = [ "web-sys", ] +[[package]] +name = "wildmatch" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "39b7d07a236abaef6607536ccfaf19b396dbe3f5110ddb73d39f4562902ed382" + [[package]] name = "winapi" version = "0.3.9" @@ -2609,6 +2856,15 @@ version = "0.53.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "271414315aff87387382ec3d271b52d7ae78726f5d44ac98b4f4030c91880486" +[[package]] +name = "winnow" +version = "0.7.13" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "21a0236b59786fed61e2a80582dd500fe61f18b5dca67a4a067d0bc9039339cf" +dependencies = [ + "memchr", +] + [[package]] name = "wit-bindgen" version = "0.45.1" @@ -2724,3 +2980,31 @@ dependencies = [ "quote", "syn", ] + +[[package]] +name = "zstd" +version = "0.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e91ee311a569c327171651566e07972200e76fcfe2242a4fa446149a3881c08a" +dependencies = [ + "zstd-safe", +] + +[[package]] +name = "zstd-safe" +version = "7.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f49c4d5f0abb602a93fb8736af2a4f4dd9512e36f7f570d66e65ff867ed3b9d" +dependencies = [ + "zstd-sys", +] + +[[package]] +name = "zstd-sys" +version = "2.0.16+zstd.1.5.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91e19ebc2adc8f83e43039e79776e3fda8ca919132d68a1fed6a5faca2683748" +dependencies = [ + "cc", + "pkg-config", +] diff --git a/Cargo.toml b/Cargo.toml index 08873af..4d0750b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "allegedly" -description = "public ledger tools and services (for the PLC)" +description = "public ledger server tools and services (for the PLC)" license = "MIT OR Apache-2.0" version = "0.1.0" edition = "2024" @@ -11,9 +11,9 @@ anyhow = "1.0.99" async-compression = { version = "0.4.30", features = ["futures-io", "tokio", "gzip"] } chrono = { version = "0.4.42", features = ["serde"] } clap = { version = "4.5.47", features = ["derive", "env"] } -env_logger = "0.11.8" futures = "0.3.31" log = "0.4.28" +poem = { version = "3.1.12", features = ["compression"] } reqwest = { version = "0.12.23", features = ["stream"] } reqwest-middleware = "0.4.2" reqwest-retry = "0.7.0" @@ -24,3 +24,4 @@ tokio = { version = "1.47.1", features = ["full"] } tokio-postgres = { version = "0.7.13", features = ["with-chrono-0_4", "with-serde_json-1"] } tokio-stream = { version = "0.1.17", features = ["io-util"] } tokio-util = { version = "0.7.16", features = ["compat"] } +tracing-subscriber = { version = "0.3.20", features = ["env-filter"] } diff --git a/readme.md b/readme.md index e9853f8..4c143bd 100644 --- a/readme.md +++ b/readme.md @@ -1,12 +1,20 @@ # Allegedly -Some [public ledger](https://github.com/did-method-plc/did-method-plc) tools and services +Some [public ledger](https://github.com/did-method-plc/did-method-plc) server tools and services Allegedly can - Tail PLC ops to stdout: `allegedly tail | jq` - Export PLC ops to weekly gzipped bundles: `allegdly bundle --dest ./some-folder` - Dump bundled ops to stdout FAST: `allegedly backfill --source-workers 6 | pv -l > /ops-unordered.jsonl` +- Wrap the reference PLC server and run it as a mirror: + + ```bash + allegedly mirror \ + --bind 0.0.0.0:8000 \ + --wrap http://127.0.0.1:3000 \ + --wrap-pg "postgresql://postgres:postgres@localhost:5432/postgres" + ``` (add `--help` to any command for more info about it) diff --git a/src/bin/allegedly.rs b/src/bin/allegedly.rs index c45071a..478f5f0 100644 --- a/src/bin/allegedly.rs +++ b/src/bin/allegedly.rs @@ -1,10 +1,10 @@ use allegedly::{ Db, Dt, ExportPage, FolderSource, HttpSource, PageBoundaryState, backfill, backfill_to_pg, - bin_init, pages_to_pg, pages_to_weeks, poll_upstream, + bin_init, pages_to_pg, pages_to_weeks, poll_upstream, serve, }; -use clap::{Parser, Subcommand}; +use clap::{CommandFactory, Parser, Subcommand}; use reqwest::Url; -use std::{path::PathBuf, time::Instant}; +use std::{net::SocketAddr, path::PathBuf, time::Instant}; use tokio::sync::{mpsc, oneshot}; #[derive(Debug, Parser)] @@ -71,6 +71,19 @@ enum Commands { #[arg(long, action)] clobber: bool, }, + /// Wrap a did-method-plc server, syncing upstream and blocking op submits + Mirror { + /// the wrapped did-method-plc server + #[arg(long, env)] + wrap: Url, + /// the wrapped did-method-plc server's database (write access required) + #[arg(long, env)] + wrap_pg: Url, + /// wrapping server listen address + #[arg(short, long, env)] + #[clap(default_value = "127.0.0.1:8000")] + bind: SocketAddr, + }, /// Poll an upstream PLC server and log new ops to stdout Tail { /// Begin tailing from a specific timestamp for replay or wait-until @@ -121,9 +134,10 @@ fn full_pages(mut rx: mpsc::Receiver) -> mpsc::Receiver #[tokio::main] async fn main() { - bin_init("main"); - let args = Cli::parse(); + let matches = Cli::command().get_matches(); + let name = matches.subcommand().map(|(name, _)| name).unwrap_or("???"); + bin_init(name); let t0 = Instant::now(); match args.command { @@ -206,6 +220,37 @@ async fn main() { std::fs::create_dir_all(&dest).unwrap(); pages_to_weeks(rx, dest, clobber).await.unwrap(); } + Commands::Mirror { + wrap, + wrap_pg, + bind, + } => { + let db = Db::new(wrap_pg.as_str()).await.unwrap(); + let latest = db + .get_latest() + .await + .unwrap() + .expect("there to be at least one op in the db. did you backfill?"); + + let (tx, rx) = mpsc::channel(2); + // upstream poller + tokio::task::spawn(async move { + log::info!("starting poll reader..."); + let mut url = args.upstream; + url.set_path("/export"); + tokio::task::spawn( + async move { poll_upstream(Some(latest), url, tx).await.unwrap() }, + ); + }); + // db writer + let poll_db = db.clone(); + tokio::task::spawn(async move { + log::info!("starting db writer..."); + pages_to_pg(poll_db, rx).await.unwrap(); + }); + + serve(wrap, bind).await.unwrap(); + } Commands::Tail { after } => { let mut url = args.upstream; url.set_path("/export"); diff --git a/src/lib.rs b/src/lib.rs index ec3c2c2..5f5f43f 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -2,12 +2,14 @@ use serde::Deserialize; mod backfill; mod client; +mod mirror; mod plc_pg; mod poll; mod weekly; pub use backfill::backfill; pub use client::CLIENT; +pub use mirror::serve; pub use plc_pg::{Db, backfill_to_pg, pages_to_pg}; pub use poll::{PageBoundaryState, get_page, poll_upstream}; pub use weekly::{BundleSource, FolderSource, HttpSource, Week, pages_to_weeks, week_to_pages}; @@ -58,11 +60,8 @@ impl From<&Op<'_>> for OpKey { } } -pub fn bin_init(name: &str) { - use env_logger::{Builder, Env}; - Builder::from_env(Env::new().filter_or("RUST_LOG", "info")).init(); - - log::info!( +pub fn logo(name: &str) -> String { + format!( r" \ | | | | @@ -70,6 +69,19 @@ pub fn bin_init(name: &str) { _/ _\ _| _| \___| \__, | \___| \__,_| _| \_, | (v{}) ____| __/ ", - env!("CARGO_PKG_VERSION") - ); + env!("CARGO_PKG_VERSION"), + ) +} + +pub fn bin_init(name: &str) { + if std::env::var_os("RUST_LOG").is_none() { + unsafe { std::env::set_var("RUST_LOG", "info") }; + } + let filter = tracing_subscriber::EnvFilter::from_default_env(); + tracing_subscriber::fmt() + .with_writer(std::io::stderr) + .with_env_filter(filter) + .init(); + + log::info!("{}", logo(name)); } diff --git a/src/mirror.rs b/src/mirror.rs new file mode 100644 index 0000000..ab2c104 --- /dev/null +++ b/src/mirror.rs @@ -0,0 +1,70 @@ +use crate::logo; +use poem::{ + EndpointExt, Error, IntoResponse, Request, Response, Result, Route, Server, get, handler, + http::{StatusCode, Uri}, + listener::TcpListener, + middleware::{AddData, CatchPanic, Compression, Cors, Tracing}, + web::Data, +}; +use reqwest::{Client, Url}; +use std::net::SocketAddr; +use std::time::Duration; + +#[derive(Debug, Clone)] +struct State { + client: Client, + plc: Url, +} + +#[handler] +fn hello() -> String { + logo("mirror") +} + +#[handler] +async fn proxy(req: &Request, Data(state): Data<&State>) -> Result { + let mut target = state.plc.clone(); + target.set_path(req.uri().path()); + let upstream_res = state + .client + .get(target) + .headers(req.headers().clone()) + .send() + .await + .map_err(|e| { + log::error!("upstream req fail: {e}"); + Error::from_string("request to plc server failed", StatusCode::BAD_GATEWAY) + })?; + let mut res = Response::default(); + upstream_res.headers().iter().for_each(|(k, v)| { + res.headers_mut().insert(k, v.to_owned()); + }); + res.set_status(upstream_res.status()); + res.set_version(upstream_res.version()); + res.set_body(upstream_res.bytes().await.unwrap()); + Ok(res) +} + +#[handler] +async fn nope(uri: &Uri) -> Result { + log::info!("ha nope, {uri:?}"); + Ok(()) +} + +pub async fn serve(plc: Url, bind: SocketAddr) -> std::io::Result<()> { + let client = Client::builder() + .timeout(Duration::from_secs(3)) + .build() + .unwrap(); + let state = State { client, plc }; + + let app = Route::new() + .at("/", get(hello)) + .at("/:any", get(proxy).post(nope)) + .with(AddData::new(state)) + .with(Cors::new().allow_credentials(false)) + .with(Compression::new()) + .with(CatchPanic::new()) + .with(Tracing); + Server::new(TcpListener::bind(bind)).run(app).await +} diff --git a/src/plc_pg.rs b/src/plc_pg.rs index 5e9a877..983d46b 100644 --- a/src/plc_pg.rs +++ b/src/plc_pg.rs @@ -70,6 +70,21 @@ impl Db { Ok(client) } + + pub async fn get_latest(&self) -> Result, PgError> { + let client = self.connect().await?; + let dt: Option
= client + .query_opt( + r#"SELECT "createdAt" + FROM operations + ORDER BY "createdAt" DESC + LIMIT 1"#, + &[], + ) + .await? + .map(|row| row.get(0)); + Ok(dt) + } } pub async fn pages_to_pg(db: Db, mut pages: mpsc::Receiver) -> Result<(), PgError> { @@ -78,7 +93,8 @@ pub async fn pages_to_pg(db: Db, mut pages: mpsc::Receiver) -> Resul let ops_stmt = client .prepare( r#"INSERT INTO operations (did, operation, cid, nullified, "createdAt") - VALUES ($1, $2, $3, $4, $5)"#, + VALUES ($1, $2, $3, $4, $5) + ON CONFLICT do nothing"#, ) .await?; let did_stmt = client