From a5b48c5a059bc4178d6c93b08a7b85e7ba74ceb1 Mon Sep 17 00:00:00 2001 From: Aria Date: Fri, 26 Jun 2026 21:39:39 +0100 Subject: [PATCH] more shit --- Cargo.lock | 381 ++++++----------- Cargo.toml | 16 +- crates/demo/Cargo.toml | 6 + crates/demo/src/main.rs | 14 +- crates/lib/Cargo.toml | 22 +- crates/lib/src/aggregator/mod.rs | 98 +++++ .../{registry => aggregator}/subscriber.rs | 50 ++- crates/lib/src/bint.rs | 67 --- crates/lib/src/destination.rs | 146 +++++++ crates/lib/src/lib.rs | 12 +- crates/lib/src/logs.rs | 396 ++++++++++++++++++ crates/lib/src/messagepack.rs | 249 +++++++++-- crates/lib/src/registry/mod.rs | 75 ---- crates/lib/src/router.rs | 85 ++++ 14 files changed, 1148 insertions(+), 469 deletions(-) create mode 100644 crates/lib/src/aggregator/mod.rs rename crates/lib/src/{registry => aggregator}/subscriber.rs (60%) delete mode 100644 crates/lib/src/bint.rs create mode 100644 crates/lib/src/destination.rs create mode 100644 crates/lib/src/logs.rs delete mode 100644 crates/lib/src/registry/mod.rs create mode 100644 crates/lib/src/router.rs diff --git a/Cargo.lock b/Cargo.lock index 65c37e9..50b73c2 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11,77 +11,24 @@ dependencies = [ "memchr", ] -[[package]] -name = "atomic-waker" -version = "1.1.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" - [[package]] name = "autocfg" version = "1.5.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" -[[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 = "bogos-binted" version = "0.1.0" dependencies = [ - "axum", + "futures", + "redb", "rmp", + "rmp-serde", + "serde", + "thiserror", "tokio", + "tokio-stream", "tower", "tracing", "ulid", @@ -105,12 +52,6 @@ version = "3.20.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" -[[package]] -name = "bytes" -version = "1.11.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1e748733b7cbc798e1434b6ac524f0c1ff2ab456fe201501e6497c8417a4fc33" - [[package]] name = "cfg-if" version = "1.0.4" @@ -118,12 +59,18 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9330f8b2ff13f34540b44e946ef35111825727b38d33286ef986142615121801" [[package]] -name = "form_urlencoded" -version = "1.2.2" +name = "futures" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "cb4cb245038516f5f85277875cdaa4f7d2c9a0fa0468de06ed190163b1581fcf" +checksum = "8b147ee9d1f6d097cef9ce628cd2ee62288d963e16fb287bd9286455b241382d" dependencies = [ - "percent-encoding", + "futures-channel", + "futures-core", + "futures-executor", + "futures-io", + "futures-sink", + "futures-task", + "futures-util", ] [[package]] @@ -133,6 +80,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "07bbe89c50d7a535e539b8c17bc0b49bdb77747034daa8087407d655f3f7cc1d" dependencies = [ "futures-core", + "futures-sink", ] [[package]] @@ -142,121 +90,74 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7e3450815272ef58cec6d564423f6e755e25379b217b0bc688e295ba24df6b1d" [[package]] -name = "futures-task" -version = "0.3.32" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "037711b3d59c33004d3856fbdc83b99d4ff37a24768fa1be9ce3538a1cde4393" - -[[package]] -name = "futures-util" +name = "futures-executor" version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" +checksum = "baf29c38818342a3b26b5b923639e7b1f4a61fc5e76102d4b1981c6dc7a7579d" dependencies = [ "futures-core", "futures-task", - "pin-project-lite", - "slab", -] - -[[package]] -name = "getrandom" -version = "0.3.4" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" -dependencies = [ - "cfg-if", - "libc", - "r-efi", - "wasip2", -] - -[[package]] -name = "http" -version = "1.4.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6970f50e31d6fc17d3fa27329444bfa74e196cf62e95052a3f6fee181dba6425" -dependencies = [ - "bytes", - "itoa", + "futures-util", ] [[package]] -name = "http-body" -version = "1.0.1" +name = "futures-io" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1efedce1fb8e6913f23e0c92de8e62cd5b772a67e7b3946df930a62566c93184" -dependencies = [ - "bytes", - "http", -] +checksum = "cecba35d7ad927e23624b22ad55235f2239cfa44fd10428eecbeba6d6a717718" [[package]] -name = "http-body-util" -version = "0.1.3" +name = "futures-macro" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b021d93e26becf5dc7e1b75b1bed1fd93124b374ceb73f43d4d4eafec896a64a" +checksum = "e835b70203e41293343137df5c0664546da5745f82ec9b84d40be8336958447b" dependencies = [ - "bytes", - "futures-core", - "http", - "http-body", - "pin-project-lite", + "proc-macro2", + "quote", + "syn 2.0.117", ] [[package]] -name = "httparse" -version = "1.10.1" +name = "futures-sink" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6dbf3de79e51f3d586ab4cb9d5c3e2c14aa28ed23d180cf89b4df0454a69cc87" +checksum = "c39754e157331b013978ec91992bde1ac089843443c49cbc7f46150b0fad0893" [[package]] -name = "httpdate" -version = "1.0.3" +name = "futures-task" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "df3b46402a9d5adb4c86a0cf463f42e19994e3ee891101b1841f30a545cb49a9" +checksum = "037711b3d59c33004d3856fbdc83b99d4ff37a24768fa1be9ce3538a1cde4393" [[package]] -name = "hyper" -version = "1.10.1" +name = "futures-util" +version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "55281c53a1894c864990125767da440a4e630446785086f52523b20033b74498" +checksum = "389ca41296e6190b48053de0321d02a77f32f8a5d2461dd38762c0593805c6d6" dependencies = [ - "atomic-waker", - "bytes", "futures-channel", "futures-core", - "http", - "http-body", - "httparse", - "httpdate", - "itoa", + "futures-io", + "futures-macro", + "futures-sink", + "futures-task", + "memchr", "pin-project-lite", - "smallvec", - "tokio", + "slab", ] [[package]] -name = "hyper-util" -version = "0.1.20" +name = "getrandom" +version = "0.3.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "96547c2556ec9d12fb1578c4eaf448b04993e7fb79cbaad930a656880a6bdfa0" +checksum = "899def5c37c4fd7b2664648c28120ecec138e4d395b459e5ca34f9cce2dd77fd" dependencies = [ - "bytes", - "http", - "http-body", - "hyper", - "pin-project-lite", - "tokio", - "tower-service", + "cfg-if", + "libc", + "r-efi", + "wasip2", ] -[[package]] -name = "itoa" -version = "1.0.18" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" - [[package]] name = "js-sys" version = "0.3.100" @@ -295,35 +196,12 @@ dependencies = [ "regex-automata", ] -[[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.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6b947ae49db0d222b1dbc6b113ce7248a3fc3a6ca21b696717bfc000ba4484d8" -[[package]] -name = "mime" -version = "0.3.17" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" - -[[package]] -name = "mio" -version = "1.2.1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "02bd0af71c67b473010cbbc60715ee815645a4dc942899111f494b4b737d6fda" -dependencies = [ - "libc", - "wasi", - "windows-sys", -] - [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -348,12 +226,6 @@ version = "1.21.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9f7c3e4beb33f85d45ae3e3a1792185706c8e16d043238c593331cc7cd313b50" -[[package]] -name = "percent-encoding" -version = "2.3.2" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9b4f627cb1b25917193a259e49bdad08f671f8d9708acfd5fe0a8c1455d87220" - [[package]] name = "pin-project" version = "1.1.13" @@ -371,7 +243,7 @@ checksum = "c96395f0a926bc13b1c17622aaddda1ecb55d49c8f1bf9777e4d877800a43f8b" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -442,6 +314,15 @@ dependencies = [ "getrandom", ] +[[package]] +name = "redb" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e925444704b5f17d32bf42f5b6e2df050bceebc3dcd6e71cc73dafe8092e839" +dependencies = [ + "libc", +] + [[package]] name = "regex-automata" version = "0.4.14" @@ -469,80 +350,49 @@ dependencies = [ ] [[package]] -name = "rustversion" -version = "1.0.22" +name = "rmp-serde" +version = "1.3.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" +checksum = "72f81bee8c8ef9b577d1681a70ebbc962c232461e397b22c208c43c04b67a155" +dependencies = [ + "rmp", + "serde", +] [[package]] -name = "ryu" -version = "1.0.23" +name = "rustversion" +version = "1.0.22" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" [[package]] name = "serde" -version = "1.0.228" +version = "1.0.229" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +checksum = "4148590afebada386688f18773da617792bf2ef03ffc1e4cbd2b1d45b023e0ba" dependencies = [ "serde_core", + "serde_derive", ] [[package]] name = "serde_core" -version = "1.0.228" +version = "1.0.229" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" +checksum = "67dca2c9c51e58a4791a4b1ed58308b39c64224d349a935ab5039aa360942a48" dependencies = [ "serde_derive", ] [[package]] name = "serde_derive" -version = "1.0.228" +version = "1.0.229" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" +checksum = "e7a5d71263a5a7d47b41f6b3f06ba276f10cc18b0931f1799f710578e2309348" dependencies = [ "proc-macro2", "quote", - "syn", -] - -[[package]] -name = "serde_json" -version = "1.0.150" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e8014e44b4736ed0538adeecded0fce2a272f22dc9578a7eb6b2d9993c74cfb9" -dependencies = [ - "itoa", - "memchr", - "serde", - "serde_core", - "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", + "syn 3.0.3", ] [[package]] @@ -567,20 +417,21 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "67b1b7a3b5fe4f1376887184045fcf45c69e92af734b7aaddc05fb777b6fbd03" [[package]] -name = "socket2" -version = "0.6.4" +name = "syn" +version = "2.0.117" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" +checksum = "e665b8803e7b1d2a727f4023456bbbbe74da67099c585258af0ad9c5013b9b99" dependencies = [ - "libc", - "windows-sys", + "proc-macro2", + "quote", + "unicode-ident", ] [[package]] name = "syn" -version = "2.0.117" +version = "3.0.3" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "e665b8803e7b1d2a727f4023456bbbbe74da67099c585258af0ad9c5013b9b99" +checksum = "53e9bae58849f64dfa4f5d5ae372c8341f7305f82a3868709269343628b659a3" dependencies = [ "proc-macro2", "quote", @@ -593,6 +444,26 @@ version = "1.0.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0bf256ce5efdfa370213c1dabab5935a12e49f2c58d15e9eac2870d3b4f27263" +[[package]] +name = "thiserror" +version = "2.0.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "09a43598840e33d5b0331f38c5e30d13bb11c11210a4b58f0d9b18a5a5eefcd9" +dependencies = [ + "thiserror-impl", +] + +[[package]] +name = "thiserror-impl" +version = "2.0.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "43cbfe0cf76104d42a574802844187e84a305e531ed54455f11fbde0f10541cd" +dependencies = [ + "proc-macro2", + "quote", + "syn 3.0.3", +] + [[package]] name = "thread_local" version = "1.1.9" @@ -608,12 +479,8 @@ version = "1.52.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8fc7f01b389ac15039e4dc9531aa973a135d7a4135281b12d7c1bc79fd57fffe" dependencies = [ - "libc", - "mio", "pin-project-lite", - "socket2", "tokio-macros", - "windows-sys", ] [[package]] @@ -624,7 +491,18 @@ checksum = "385a6cb71ab9ab790c5fe8d67f1645e6c450a7ce006a33de03daa956cf70a496" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", +] + +[[package]] +name = "tokio-stream" +version = "0.1.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "32da49809aab5c3bc678af03902d4ccddea2a87d028d86392a4b1560c6906c70" +dependencies = [ + "futures-core", + "pin-project-lite", + "tokio", ] [[package]] @@ -637,10 +515,8 @@ dependencies = [ "futures-util", "pin-project-lite", "sync_wrapper", - "tokio", "tower-layer", "tower-service", - "tracing", ] [[package]] @@ -661,7 +537,6 @@ version = "0.1.44" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "63e71662fa4b2a2c3a26f570f037eb95bb1f85397f3cd8076caed2f026a6d100" dependencies = [ - "log", "pin-project-lite", "tracing-attributes", "tracing-core", @@ -675,7 +550,7 @@ checksum = "7490cfa5ec963746568740651ac6781f701c9c5ea257c58e057f3ba8cf69e8da" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] [[package]] @@ -749,12 +624,6 @@ version = "0.1.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ba73ea9cf16a25df0c8caa16c51acb937d5712a8429db78a3ee29d5dcacd3a65" -[[package]] -name = "wasi" -version = "0.11.1+wasi-snapshot-preview1" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "ccf3ec651a847eb01de73ccad15eb7d99f80485de043efb2f370cd654f4ea44b" - [[package]] name = "wasip2" version = "1.0.3+wasi-0.2.9" @@ -796,7 +665,7 @@ dependencies = [ "bumpalo", "proc-macro2", "quote", - "syn", + "syn 2.0.117", "wasm-bindgen-shared", ] @@ -857,11 +726,5 @@ checksum = "1ae7f38b72ec2a254e2b87ef277cf2cd4fb97cbebf944faa6f33354da0867930" dependencies = [ "proc-macro2", "quote", - "syn", + "syn 2.0.117", ] - -[[package]] -name = "zmij" -version = "1.0.21" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "b8848ee67ecc8aedbaf3e4122217aff892639231befc6a1b58d29fff4c2cabaa" diff --git a/Cargo.toml b/Cargo.toml index d2855a3..ee84071 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -7,11 +7,21 @@ members = [ [workspace.dependencies] -axum = "0.8.9" +# async stuff tokio = "1.52.3" +tokio-stream = "0.1.18" +futures = "0.3.32" + +# service trait tower = "0.5.3" + +# logging tracing = "0.1.44" tracing-futures = "0.2.5" tracing-subscriber = "0.3.23" -rmp = "0.8.15" -rmpv = "1.3.1" + +# serialisation +serde = "1.0.229" + +# utilities +thiserror = "2.0.19" diff --git a/crates/demo/Cargo.toml b/crates/demo/Cargo.toml index 1f5462b..cf1d270 100644 --- a/crates/demo/Cargo.toml +++ b/crates/demo/Cargo.toml @@ -4,9 +4,15 @@ version = "0.1.0" edition = "2024" [dependencies] +# async stuff tokio = { workspace = true, features = ["macros", "rt", "time"] } + +# service trait tower = { workspace = true, features = ["util"] } + +# logging tracing = { workspace = true } tracing-futures = { workspace = true } tracing-subscriber = { workspace = true, features = ["env-filter"] } + bogos-binted = { path = "../lib" } diff --git a/crates/demo/src/main.rs b/crates/demo/src/main.rs index e345f1f..6a313e6 100644 --- a/crates/demo/src/main.rs +++ b/crates/demo/src/main.rs @@ -1,10 +1,10 @@ -use std::{pin::Pin, task::Poll, time::Duration}; - -use bogos_binted::ServiceExt as _; +use bogos_binted::{ServiceExt as _, TraceDestination, in_memory_database}; +use std::{pin::Pin, sync::Arc, task::Poll, time::Duration}; use tower::{Service, ServiceExt as _}; use tracing::{info, info_span}; use tracing_futures::Instrument; +/// Our fancy service that does something useful struct UppercaseService; async fn crunch_numbers() { @@ -39,7 +39,11 @@ async fn main() { tracing_subscriber::fmt::init(); info!("initialising"); - let mut svc = UppercaseService.bint(); + let db = Arc::new(in_memory_database()); + let mut dest = TraceDestination::new(db.clone()); + let mut svc = UppercaseService.route_logs_to("uppercase_service", &mut dest); + + tokio::spawn(dest.process()); dbg!( svc.ready() @@ -49,4 +53,6 @@ async fn main() { .await .unwrap() ); + + // TODO: consume the produced logs through a nice api (yet to be created) } diff --git a/crates/lib/Cargo.toml b/crates/lib/Cargo.toml index 55a1fe8..d9ff805 100644 --- a/crates/lib/Cargo.toml +++ b/crates/lib/Cargo.toml @@ -4,9 +4,27 @@ version = "0.1.0" edition = "2024" [dependencies] -axum = { workspace = true } +# async stuff tokio = { workspace = true } +tokio-stream = { workspace = true } +futures = { workspace = true } + +# service trait tower = { workspace = true } + +# logs tracing = { workspace = true } -rmp = { workspace = true } + +# internal identifiers ulid = "1.2.1" + + # database +redb = "4.1.0" + +# serialisation +serde = { workspace = true, features = ["derive"] } +rmp = "0.8.15" +rmp-serde = "1.3.1" + +# utilities +thiserror = { workspace = true } diff --git a/crates/lib/src/aggregator/mod.rs b/crates/lib/src/aggregator/mod.rs new file mode 100644 index 0000000..b304ce6 --- /dev/null +++ b/crates/lib/src/aggregator/mod.rs @@ -0,0 +1,98 @@ +use std::{ + num::NonZeroU64, + sync::{Arc, Mutex, RwLock}, +}; + +use tokio::sync::mpsc::Sender; +use tracing::{Subscriber, debug}; +use ulid::{Generator, Ulid}; + +use crate::messagepack::MessagePackBytes; + +mod subscriber; +#[doc(inline)] +pub use subscriber::*; + +pub(crate) type FinishedLogsPayload = (Ulid, Entry); +pub(crate) type FinishedLogsDest = Sender; + +/// Aggregates in-progress logs from multiple subscribers quickly, in-memory, and then yeets it down a channel once the future finishes. +pub struct Aggregator { + ulid_generator: Mutex, + in_mem: RwLock)>>, + dest: FinishedLogsDest, +} + +impl Aggregator { + pub fn new(dest: FinishedLogsDest) -> Self { + Aggregator { + ulid_generator: Mutex::default(), + in_mem: RwLock::default(), + dest, + } + } + + /// Get a subscriber to use for a future. `inp` should be the input to the service in debug form, so it can be logged. + pub fn subscriber(this: Arc, inp: String) -> impl Subscriber { + let ulid = this.ulid_generator.lock().unwrap().generate().unwrap(); + debug!(message = "new subscriber created", inp, ?ulid); + + this.in_mem + .write() + .unwrap() + .push((ulid, Mutex::new(Entry::new(inp)))); + + AggregatorSubscriber::new(this, ulid) + } + + /// Append to the log for a given future. + /// If [`Self::finished`] has already been called, or `ulid` is not for this aggregator, this has no effect. + // NOTE: Called from within subscriber, can't log anything in here. + fn add_log_entry(this: &Arc, ulid: Ulid, log: LogEntry) { + if let Some((_, entry)) = this.in_mem.read().unwrap().iter().find(|(u, _)| *u == ulid) { + entry.lock().unwrap().log.push(log); + } + } + + /// Called once a future is done to finalise its logs + // This might be called in a `Drop` implementation so should not block for long. + // NOTE: Called from within subscriber drop, can't log anything in here + fn finished(this: &Arc, ulid: Ulid) { + let mut in_mem = this.in_mem.write().unwrap(); + if let Some((pos, _)) = in_mem.iter().enumerate().find(|(_, (u, _))| *u == ulid) { + let (ulid, entry) = in_mem.swap_remove(pos); + let entry = entry.into_inner().unwrap(); + + // if this fails, then drop the log silently. this sucks, but not as bad as + // holding up the drop implementation does + let _ = this.dest.try_send((ulid, entry)); + } + } +} + +/// A log, possibly in-progress, recorded from a [`AggregatorSubscriber`] +#[derive(Debug, Clone)] +pub struct Entry { + pub inp: String, + pub log: Vec, +} + +impl Entry { + fn new(inp: String) -> Self { + Self { + inp, + log: Vec::with_capacity(8), + } + } +} + +/// Entry in the flat log that [`Aggregator`] uses internally +#[derive(Debug, Clone)] +pub enum LogEntry { + AddSpan(NonZeroU64, MessagePackBytes), + RecordSpan(NonZeroU64, MessagePackBytes), + SpanFollows(NonZeroU64, NonZeroU64), // 1 follows 0 + EnterSpan(NonZeroU64), + Event(MessagePackBytes), + ExitSpan(NonZeroU64), +} diff --git a/crates/lib/src/registry/subscriber.rs b/crates/lib/src/aggregator/subscriber.rs similarity index 60% rename from crates/lib/src/registry/subscriber.rs rename to crates/lib/src/aggregator/subscriber.rs index 8db868a..23e19ac 100644 --- a/crates/lib/src/registry/subscriber.rs +++ b/crates/lib/src/aggregator/subscriber.rs @@ -1,4 +1,7 @@ -use super::{LogEntry, Registry}; +//! Implements a tracer subscriber that sends things to a [`Aggregator`]. +//! Logs are written as a series of [`LogEntry`]s appended to a vec. +//! Parsing this into a useful structure is done later, in [`::crate::destination`]. +use super::{Aggregator, LogEntry}; use crate::messagepack::MessagePackBytes; use std::{ num::NonZero, @@ -10,18 +13,18 @@ use std::{ use tracing::{Event, Metadata, Subscriber, span}; use ulid::Ulid; -/// Receives tracing events / spans and bridges them to the registry -pub struct RegistrySubscriber { - registry: Arc, +/// Receives tracing events / spans and bridges them to the aggregator +pub struct AggregatorSubscriber { + aggregator: Arc, next_id: AtomicU64, ulid: Ulid, } // Start of lifecycle -impl RegistrySubscriber { - pub fn new(registry: Arc, ulid: Ulid) -> Self { +impl AggregatorSubscriber { + pub fn new(aggregator: Arc, ulid: Ulid) -> Self { Self { - registry, + aggregator, next_id: AtomicU64::new(1), ulid, } @@ -29,16 +32,17 @@ impl RegistrySubscriber { } // Gets called by tracing during the future -impl Subscriber for RegistrySubscriber { +impl Subscriber for AggregatorSubscriber { fn enabled(&self, _metadata: &Metadata<'_>) -> bool { - true // filering is provided by layers above our subscriber + // TODO: level filtering + true } fn new_span(&self, span: &span::Attributes<'_>) -> span::Id { let id = self.next_id.fetch_add(1, Ordering::Relaxed); - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::AddSpan(NonZero::new(id).unwrap(), MessagePackBytes::from(span)), ); @@ -47,40 +51,40 @@ impl Subscriber for RegistrySubscriber { } fn record(&self, span: &span::Id, values: &span::Record<'_>) { - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::RecordSpan(span.into_non_zero_u64(), MessagePackBytes::from(values)), ); } fn record_follows_from(&self, span: &span::Id, follows: &span::Id) { - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::SpanFollows(follows.into_non_zero_u64(), span.into_non_zero_u64()), ); } fn event(&self, event: &Event<'_>) { - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::Event(MessagePackBytes::from(event)), ); } fn enter(&self, span: &span::Id) { - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::EnterSpan(span.into_non_zero_u64()), ); } fn exit(&self, span: &span::Id) { - Registry::add_log_entry( - &self.registry, + Aggregator::add_log_entry( + &self.aggregator, self.ulid, LogEntry::ExitSpan(span.into_non_zero_u64()), ); @@ -88,8 +92,8 @@ impl Subscriber for RegistrySubscriber { } // After the future is done, the subscriber gets dropped -impl Drop for RegistrySubscriber { +impl Drop for AggregatorSubscriber { fn drop(&mut self) { - Registry::finished(&self.registry, self.ulid); + Aggregator::finished(&self.aggregator, self.ulid); } } diff --git a/crates/lib/src/bint.rs b/crates/lib/src/bint.rs deleted file mode 100644 index 66dbe9f..0000000 --- a/crates/lib/src/bint.rs +++ /dev/null @@ -1,67 +0,0 @@ -//! thing that wraps a service, and ext trait for convenience -use std::{fmt::Debug, marker::PhantomData, pin::Pin, sync::Arc}; -use tower::Service; -use tracing::instrument::WithSubscriber; - -use crate::registry::Registry; -pub struct Binted> { - inner: S, - registry: Arc, - _pd: PhantomData, -} - -impl> Binted { - pub fn new(inner: S) -> Self { - Self { - inner, - registry: Arc::default(), - _pd: PhantomData, - } - } -} - -impl> Service for Binted -where - S::Future: Send + 'static, - I: Debug, -{ - type Response = S::Response; - - type Error = S::Error; - - type Future = - Pin> + Send + 'static>>; - - fn poll_ready( - &mut self, - cx: &mut std::task::Context<'_>, - ) -> std::task::Poll> { - self.inner.poll_ready(cx) - } - - fn call(&mut self, req: I) -> Self::Future { - let inp = format!("{:?}", req); - Box::pin( - self.inner - .call(req) - .with_subscriber(Registry::subscriber(self.registry.clone(), inp)), - ) - } -} - -// convenience .bint method -pub trait ServiceExt: Service { - fn bint(self) -> Binted - where - Self: Sized; -} - -impl ServiceExt for S -where - S: Service, - I: Debug, -{ - fn bint(self) -> Binted { - Binted::new(self) - } -} diff --git a/crates/lib/src/destination.rs b/crates/lib/src/destination.rs new file mode 100644 index 0000000..6276ab6 --- /dev/null +++ b/crates/lib/src/destination.rs @@ -0,0 +1,146 @@ +//! Code for writing logs out to somewhere they can be read. +use std::ops::Deref; + +use redb::{ + Database, StorageError, TableDefinition, TableError, TableHandle, TransactionError, + backends::InMemoryBackend, +}; +use rmp::decode::bytes::BytesReadError; +use thiserror::Error; +use tokio::sync::mpsc::{self, Receiver}; +use tokio_stream::{StreamExt, wrappers::ReceiverStream}; +use tracing::{Level, debug, instrument, warn}; +use ulid::Ulid; + +use crate::{ + aggregator::{Entry as AggregatorEntry, FinishedLogsDest, FinishedLogsPayload}, + logs::ServiceCallLogs, +}; + +const CHANNEL_SIZE: usize = 8; + +/// Receives finished logs and writes them out. +pub struct TraceDestination { + receivers: Vec, + db: D, +} + +/// A place that logs come from. Usually, there's a [`Aggregator`] on the other end. +struct LogSource { + /// A static identifier for this source + slug: &'static str, + + /// The table to write logs to + dest_table: TableDefinition<'static, u128, Vec>, + + /// The pipe finished logs are sent down. + recv: Receiver<(Ulid, AggregatorEntry)>, +} + +impl> TraceDestination { + /// Write logs to the given [`Database`]. This may be a shared reference, so you can query the database in other places. + pub fn new(db: D) -> Self { + Self { + receivers: vec![], + db, + } + } + + /// Get a receiver with the given identifier. + pub(crate) fn receiver_for(&mut self, slug: &'static str) -> FinishedLogsDest { + let (send, recv) = mpsc::channel(CHANNEL_SIZE); + let table = TableDefinition::new(slug); + self.receivers.push(LogSource { + slug, + recv, + dest_table: table, + }); + send + } + + /// Continuously process all received logs and store them in the database. + /// Loops until all sources are closed. + #[instrument(level = Level::DEBUG)] + pub async fn process(self) { + let streams = self + .receivers + .into_iter() + .map(|ls| ReceiverStream::new(ls.recv).map(move |payload| (ls.dest_table, payload))); + let mut selector = futures::prelude::stream::select_all(streams); + + while let Some((table, payload)) = selector.next().await { + if let Err(e) = Self::process_one(&*self.db, table, payload) { + warn!("error when writing log to database: {e}"); + // TODO: dump the whole log so it doesn't get lost entirely? + } + } + } + + #[instrument(level = Level::DEBUG, skip_all, fields(table = table.name(), ulid = ?id, len = entry.log.len()))] + fn process_one( + db: &Database, + table: TableDefinition<'static, u128, Vec>, + (id, entry): FinishedLogsPayload, + ) -> Result<(), ProcessPayloadError> { + let stream: ServiceCallLogs = entry + .try_into() + .map_err(ProcessPayloadError::LogSerialisingError)?; + debug!( + message = "serialised logs", + spans = stream.num_spans(), + events = stream.num_events() + ); + + { + let write = db.begin_write()?; + + let key = id.0; + let value = rmp_serde::encode::to_vec(&stream)?; + debug!(message = "converted to message pack", len = value.len()); + + write.open_table(table)?.insert(key, value)?; + + write.commit().unwrap(); + debug!("written"); + + Ok(()) + } + } +} + +#[derive(Debug, Error)] +enum ProcessPayloadError { + #[error("table error: {0}")] + TableError(#[from] TableError), + #[error("storage error: {0}")] + StorageError(#[from] StorageError), + #[error("transaction error: {0}")] + TransactionError(#[from] TransactionError), + #[error("encoding error: {0}")] + EncodeError(#[from] rmp_serde::encode::Error), + #[error("error serialising logs")] + LogSerialisingError(rmp::decode::ValueReadError), +} + +impl std::fmt::Debug for TraceDestination { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("TraceDestination") + .field("receivers", &self.receivers) + .finish_non_exhaustive() + } +} + +impl std::fmt::Debug for LogSource { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("LogSink") + .field("slug", &self.slug) + .field("recv", &self.recv) + .finish_non_exhaustive() + } +} + +pub fn in_memory_database() -> Database { + Database::builder() + .create_with_backend(InMemoryBackend::new()) + .unwrap() +} diff --git a/crates/lib/src/lib.rs b/crates/lib/src/lib.rs index 9807b22..e7694ce 100644 --- a/crates/lib/src/lib.rs +++ b/crates/lib/src/lib.rs @@ -1,6 +1,12 @@ -mod bint; +mod router; #[doc(inline)] -pub use bint::*; +pub use router::*; +mod destination; +#[doc(inline)] +pub use destination::*; + +pub mod logs; + +mod aggregator; mod messagepack; -mod registry; diff --git a/crates/lib/src/logs.rs b/crates/lib/src/logs.rs new file mode 100644 index 0000000..94ce233 --- /dev/null +++ b/crates/lib/src/logs.rs @@ -0,0 +1,396 @@ +//! Convenient representation of a finished set of logs. +use crate::{ + aggregator::{self, LogEntry}, + messagepack::{INTERNAL_MESSAGEPACK_ERROR, MessagePackBytes}, +}; +use serde::{Deserialize, Serialize}; + +/// A finished set of logs from a service getting called, in a nice storage format. +#[derive(Debug, Clone, PartialEq, Deserialize, Serialize)] +pub struct ServiceCallLogs { + span_attributes: Vec, + relationships: Vec, + events: Vec, + called_with: String, +} + +/// Attributes attached to a span, or event +pub type Attributes = Vec<(String, AttributeValue)>; + +/// A log event, or message. +#[derive(Debug, Clone, PartialEq, Deserialize, Serialize)] +pub struct Event { + span_idx: usize, + attributes: Attributes, +} + +/// What an attribute is stored as. This may be an alternative (ie Debug) representation of the actual type. +#[derive(Debug, Clone, PartialEq, Deserialize, Serialize)] +pub enum AttributeValue { + F64(f64), + U64(u64), + I64(i64), + U128(u128), + I128(i128), + Bool(bool), + String(String), + Bytes(Box<[u8]>), +} + +impl ServiceCallLogs { + /// The number of individual spans, including the root one that everything is a child of. + pub fn num_spans(&self) -> usize { + self.span_attributes.len() + } + + /// The number of log events (messages) + pub fn num_events(&self) -> usize { + self.events.len() + } +} + +/// Turns a stream of log entries into a structured representation of them. +impl TryFrom for ServiceCallLogs { + type Error = <&'static MessagePackBytes as TryInto>::Error; + + fn try_from(entry: aggregator::Entry) -> Result { + let log = entry.log; + let mut spans = vec![Attributes::default()]; + let mut relationships = vec![]; + let mut events = vec![]; + let mut span_stack = vec![0usize]; + for entry in log.iter() { + match entry { + LogEntry::Event(message_pack_bytes) => { + let attributes = message_pack_bytes + .try_into() + .expect(INTERNAL_MESSAGEPACK_ERROR); + + events.push(Event { + span_idx: *span_stack.last().unwrap(), + attributes, + }); + } + + LogEntry::AddSpan(non_zero, message_pack_bytes) => { + assert!( + non_zero.get() as usize == spans.len(), + "log entry added a span out of order" + ); + + let attrs = message_pack_bytes.try_into()?; + spans.push(attrs); + relationships.push(SpanRelationship::SpanChild( + non_zero.get() as usize, + *span_stack.last().unwrap(), + )); + } + LogEntry::RecordSpan(non_zero, message_pack_bytes) => { + assert!( + (non_zero.get() as usize) < spans.len(), + "recordspan on not yet found span" + ); + let attrs: Attributes = message_pack_bytes.try_into()?; + spans[non_zero.get() as usize].extend(attrs.into_iter()); + } + + LogEntry::EnterSpan(non_zero) => { + assert!( + (non_zero.get() as usize) < spans.len(), + "enterspan on not yet found span" + ); + span_stack.push(non_zero.get() as usize); + // TODO: record entries & exits probably + } + LogEntry::ExitSpan(non_zero) => { + assert!( + (non_zero.get() as usize) < spans.len(), + "exitspan on not yet found span" + ); + span_stack + .pop_if(|c| *c == non_zero.get() as usize) + .expect("exitspan referred to not currently active span"); + } + LogEntry::SpanFollows(non_zero, non_zero1) => { + assert!( + (non_zero.get() as usize) < spans.len(), + "spanfollows on not yet found span" + ); + assert!( + (non_zero1.get() as usize) < spans.len(), + "spanfollows on not yet found span" + ); + relationships.push(SpanRelationship::SpanFollows( + non_zero.get() as usize, + non_zero1.get() as usize, + )); + } + } + } + + Ok(ServiceCallLogs { + span_attributes: spans, + relationships, + events, + called_with: entry.inp, + }) + } +} + +/// Defines a relationship between two spans, or a log message attached to one. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub enum SpanRelationship { + SpanChild(usize, usize), // 0 is a child of 1 + SpanFollows(usize, usize), // 1 follows 0 +} + +#[cfg(test)] +mod tests { + + use crate::{ + aggregator::{Entry, LogEntry}, + logs::{AttributeValue, Attributes, Event, ServiceCallLogs, SpanRelationship}, + messagepack::{MessagePackBytes, Visitor}, + }; + + fn empty_attributes() -> Attributes { + vec![] + } + fn empty_messagepack_attributes() -> MessagePackBytes { + Visitor::new(0).into() + } + + fn simple_attributes(thing: bool) -> Attributes { + vec![("thing".to_string(), AttributeValue::Bool(thing))] + } + fn simple_messagepack_attributes(thing: bool) -> MessagePackBytes { + let mut visitor = Visitor::new(1); + visitor.raw_record_bool("thing", thing); + visitor.into() + } + + #[test] + fn test_serialise_no_events() { + let inp = Entry { + inp: "".to_string(), + log: vec![], + }; + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![vec![]], + relationships: vec![], + events: vec![], + called_with: "".to_string() + } + ); + } + + #[test] + fn test_serialise_one_simple_event() { + let inp = Entry { + inp: "".to_string(), + log: vec![LogEntry::Event(simple_messagepack_attributes(true))], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![Attributes::default()], + called_with: "".to_string(), + relationships: vec![], + events: vec![Event { + span_idx: 0, + attributes: simple_attributes(true) + }] + } + ); + } + + #[test] + fn test_serialise_multiple_simple_events() { + let inp = Entry { + inp: "".to_string(), + log: vec![ + LogEntry::Event(simple_messagepack_attributes(true)), + LogEntry::Event(simple_messagepack_attributes(false)), + LogEntry::Event(simple_messagepack_attributes(true)), + ], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![Attributes::default()], + relationships: vec![], + called_with: "".to_string(), + events: vec![ + Event { + span_idx: 0, + attributes: simple_attributes(true) + }, + Event { + span_idx: 0, + attributes: simple_attributes(false) + }, + Event { + span_idx: 0, + attributes: simple_attributes(true) + }, + ] + } + ); + } + + #[test] + fn test_serialise_add_unused_span() { + let span_id = 1.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![LogEntry::AddSpan(span_id, empty_messagepack_attributes())], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![empty_attributes(), empty_attributes()], + called_with: "".to_string(), + relationships: vec![SpanRelationship::SpanChild(1, 0)], + events: vec![], + } + ); + } + + #[test] + fn test_serialise_add_enter_leave_empty_span() { + let span_id = 1.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![ + LogEntry::AddSpan(span_id, empty_messagepack_attributes()), + LogEntry::EnterSpan(span_id), + LogEntry::ExitSpan(span_id), + ], + }; + + // TODO: probably we should track span entry/leave even if messages aren't logged? + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![empty_attributes(), empty_attributes()], + called_with: "".to_string(), + relationships: vec![SpanRelationship::SpanChild(1, 0),], + events: vec![] + } + ); + } + + #[test] + fn test_serialise_add_unused_simple_span() { + let span_id = 1.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![LogEntry::AddSpan( + span_id, + simple_messagepack_attributes(true), + )], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![empty_attributes(), simple_attributes(true)], + called_with: "".to_string(), + relationships: vec![SpanRelationship::SpanChild(1, 0)], + events: vec![] + } + ); + } + + #[test] + fn test_serialise_add_enter_leave_simple_span() { + let span_id = 1.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![ + LogEntry::AddSpan(span_id, simple_messagepack_attributes(true)), + LogEntry::EnterSpan(span_id), + LogEntry::ExitSpan(span_id), + ], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![empty_attributes(), simple_attributes(true)], + called_with: "".to_string(), + relationships: vec![SpanRelationship::SpanChild(1, 0)], + events: vec![] + } + ); + } + + #[test] + fn test_serialise_nested_empty_spans() { + let span_id = 1.try_into().unwrap(); + let span_id_2 = 2.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![ + LogEntry::AddSpan(span_id, empty_messagepack_attributes()), + LogEntry::EnterSpan(span_id), + LogEntry::AddSpan(span_id_2, empty_messagepack_attributes()), + LogEntry::EnterSpan(span_id_2), + LogEntry::ExitSpan(span_id_2), + LogEntry::ExitSpan(span_id), + ], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + span_attributes: vec![empty_attributes(), empty_attributes(), empty_attributes()], + called_with: "".to_string(), + relationships: vec![ + SpanRelationship::SpanChild(1, 0), + SpanRelationship::SpanChild(2, 1) + ], + events: vec![] + } + ); + } + + #[test] + fn test_serialise_nested_simple_spans() { + let span_id = 1.try_into().unwrap(); + let span_id_2 = 2.try_into().unwrap(); + let inp = Entry { + inp: "".to_string(), + log: vec![ + LogEntry::AddSpan(span_id, simple_messagepack_attributes(true)), + LogEntry::EnterSpan(span_id), + LogEntry::AddSpan(span_id_2, simple_messagepack_attributes(false)), + LogEntry::EnterSpan(span_id_2), + LogEntry::ExitSpan(span_id_2), + LogEntry::ExitSpan(span_id), + ], + }; + + assert_eq!( + TryInto::::try_into(inp).unwrap(), + ServiceCallLogs { + called_with: "".to_string(), + span_attributes: vec![ + empty_attributes(), + simple_attributes(true), + simple_attributes(false) + ], + relationships: vec![ + SpanRelationship::SpanChild(1, 0), + SpanRelationship::SpanChild(2, 1) + ], + events: vec![] + } + ); + } +} diff --git a/crates/lib/src/messagepack.rs b/crates/lib/src/messagepack.rs index b82e1db..0da09d5 100644 --- a/crates/lib/src/messagepack.rs +++ b/crates/lib/src/messagepack.rs @@ -1,14 +1,54 @@ -use rmp::encode::ByteBuf; +use rmp::{ + decode::{Bytes, MarkerReadError, RmpRead, bytes::BytesReadError}, + encode::{ByteBuf, RmpWrite}, +}; use tracing::{Event, field::Visit, span}; +use crate::logs::{AttributeValue, Attributes}; + +#[derive(Debug, Clone)] pub struct MessagePackBytes(Vec); +impl MessagePackBytes { + fn bytes(&self) -> Bytes<'_> { + Bytes::new(&self.0[..]) + } +} + impl From<&span::Attributes<'_>> for MessagePackBytes { fn from(value: &span::Attributes<'_>) -> Self { - let mut visitor = Visitor::new(value.fields().len()); + let metadata = value.metadata(); + let mut extra_attrs = 3; // name, target, level + if let Some(_) = metadata.module_path() { + extra_attrs += 1; + } + if let Some(_) = metadata.file() { + extra_attrs += 1; + } + if let Some(_) = metadata.line() { + extra_attrs += 1; + } + + let mut visitor = Visitor::new(value.fields().len() + extra_attrs); + + // TODO: maybe allow controlling which of these actually get recorded somehow? + visitor.raw_record_str("name", metadata.name()); + visitor.raw_record_str("target", metadata.target()); + + visitor.raw_record_str("level", metadata.level().as_str()); // TODO: as int? + if let Some(module_path) = metadata.module_path() { + visitor.raw_record_str("module_path", module_path); + } + if let Some(file) = metadata.file() { + visitor.raw_record_str("file", file); + } + if let Some(line) = metadata.line() { + visitor.raw_record_u64("line", line as u64); + } + value.record(&mut visitor); - Self(visitor.0.into_vec()) + visitor.into() } } @@ -17,7 +57,7 @@ impl From<&span::Record<'_>> for MessagePackBytes { let mut visitor = Visitor::new(value.len()); value.record(&mut visitor); - Self(visitor.0.into_vec()) + visitor.into() } } impl From<&Event<'_>> for MessagePackBytes { @@ -25,7 +65,7 @@ impl From<&Event<'_>> for MessagePackBytes { let mut visitor = Visitor::new(value.fields().count()); value.record(&mut visitor); - Self(visitor.0.into_vec()) + visitor.into() } } @@ -33,71 +73,214 @@ const SIGNED_128: i8 = 27; const UNSIGNED_128: i8 = 34; #[derive(Debug, Default)] -struct Visitor(ByteBuf); +pub(crate) struct Visitor(ByteBuf); impl Visitor { - fn new(len: usize) -> Self { + pub(crate) fn new(len: usize) -> Self { let mut this = Self(ByteBuf::new()); rmp::encode::write_map_len(&mut this.0, len.try_into().unwrap()).unwrap(); this } } +// the actual bulk of the record bits. these are pub(crate) so we can use them when testing - +// tracing spans are a pain / to construct impl Visitor { - fn record_field_name(&mut self, field: &tracing::field::Field) { - rmp::encode::write_str_len(&mut self.0, field.name().len().try_into().unwrap()).unwrap(); - rmp::encode::write_str(&mut self.0, field.name()).unwrap(); + pub(crate) fn raw_record_field_name_str(&mut self, field: &'static str) { + rmp::encode::write_str(&mut self.0, field).unwrap(); + } + + pub(crate) fn raw_record_debug(&mut self, field: &'static str, value: &dyn core::fmt::Debug) { + self.raw_record_field_name_str(field); + let s = format!("{:?}", value); + rmp::encode::write_str(&mut self.0, &s).unwrap(); + } + + pub(crate) fn raw_record_f64(&mut self, field: &'static str, value: f64) { + self.raw_record_field_name_str(field); + rmp::encode::write_f64(&mut self.0, value).unwrap(); + } + + pub(crate) fn raw_record_i64(&mut self, field: &'static str, value: i64) { + self.raw_record_field_name_str(field); + rmp::encode::write_i64(&mut self.0, value).unwrap(); + } + + pub(crate) fn raw_record_u64(&mut self, field: &'static str, value: u64) { + self.raw_record_field_name_str(field); + rmp::encode::write_u64(&mut self.0, value).unwrap(); + } + + pub(crate) fn raw_record_i128(&mut self, field: &'static str, value: i128) { + self.raw_record_field_name_str(field); + rmp::encode::write_ext_meta(&mut self.0, 16, SIGNED_128).unwrap(); + self.0.write_bytes(&value.to_be_bytes()).unwrap(); + } + + pub(crate) fn raw_record_u128(&mut self, field: &'static str, value: u128) { + self.raw_record_field_name_str(field); + rmp::encode::write_ext_meta(&mut self.0, 16, UNSIGNED_128).unwrap(); + self.0.write_bytes(&value.to_be_bytes()).unwrap(); + } + + pub(crate) fn raw_record_bool(&mut self, field: &'static str, value: bool) { + self.raw_record_field_name_str(field); + rmp::encode::write_bool(&mut self.0, value).unwrap(); + } + + pub(crate) fn raw_record_str(&mut self, field: &'static str, value: &str) { + self.raw_record_field_name_str(field); + rmp::encode::write_str(&mut self.0, value).unwrap(); + } + + pub(crate) fn raw_record_bytes(&mut self, field: &'static str, value: &[u8]) { + self.raw_record_field_name_str(field); + rmp::encode::write_bin(&mut self.0, value).unwrap(); } } impl Visit for Visitor { fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn core::fmt::Debug) { - self.record_field_name(field); - let s = format!("{:?}", value); - rmp::encode::write_str_len(&mut self.0, s.len().try_into().unwrap()).unwrap(); - rmp::encode::write_str(&mut self.0, &s).unwrap(); + self.raw_record_debug(field.name(), value) } fn record_f64(&mut self, field: &tracing::field::Field, value: f64) { - self.record_field_name(field); - rmp::encode::write_f64(&mut self.0, value).unwrap(); + self.raw_record_f64(field.name(), value) } fn record_i64(&mut self, field: &tracing::field::Field, value: i64) { - self.record_field_name(field); - rmp::encode::write_i64(&mut self.0, value).unwrap(); + self.raw_record_i64(field.name(), value) } fn record_u64(&mut self, field: &tracing::field::Field, value: u64) { - self.record_field_name(field); - rmp::encode::write_u64(&mut self.0, value).unwrap(); + self.raw_record_u64(field.name(), value) } fn record_i128(&mut self, field: &tracing::field::Field, value: i128) { - self.record_field_name(field); - rmp::encode::write_ext_meta(&mut self.0, 16, SIGNED_128).unwrap(); - rmp::encode::write_bin(&mut self.0, &value.to_be_bytes()).unwrap(); + self.raw_record_i128(field.name(), value) } fn record_u128(&mut self, field: &tracing::field::Field, value: u128) { - self.record_field_name(field); - rmp::encode::write_ext_meta(&mut self.0, 16, UNSIGNED_128).unwrap(); - rmp::encode::write_bin(&mut self.0, &value.to_be_bytes()).unwrap(); + self.raw_record_u128(field.name(), value) } fn record_bool(&mut self, field: &tracing::field::Field, value: bool) { - self.record_field_name(field); - rmp::encode::write_bool(&mut self.0, value).unwrap(); + self.raw_record_bool(field.name(), value) } fn record_str(&mut self, field: &tracing::field::Field, value: &str) { - self.record_field_name(field); - rmp::encode::write_str(&mut self.0, value).unwrap(); + self.raw_record_str(field.name(), value) } fn record_bytes(&mut self, field: &tracing::field::Field, value: &[u8]) { - self.record_field_name(field); - rmp::encode::write_bin_len(&mut self.0, value.len().try_into().unwrap()).unwrap(); - rmp::encode::write_bin(&mut self.0, value).unwrap(); + self.raw_record_bytes(field.name(), value) + } +} + +impl Into for Visitor { + fn into(self) -> MessagePackBytes { + MessagePackBytes(self.0.into_vec()) + } +} + +impl Visitor {} + +pub(crate) const INTERNAL_MESSAGEPACK_ERROR: &'static str = + "invalid messagepack data received internally"; + +impl TryInto for &MessagePackBytes { + type Error = rmp::decode::ValueReadError; + + fn try_into(self) -> Result { + let mut rd = self.bytes(); + let mut attrs = Vec::new(); + + // note: the .clone means that we use a temporary cursor, ie we peek rather than consuming + // this is needed because the read_ methods expect to consume the marker, but we need + // to know what it is first (and, in some cases, allocate a buffer of the correct size). + let len = rmp::decode::read_map_len(&mut rd)?; + for _ in 0..len { + let name_len = rmp::decode::read_str_len(&mut rd.clone())?; + let mut name = vec![0; name_len.try_into().unwrap()]; + rmp::decode::read_str(&mut rd, &mut name).expect(INTERNAL_MESSAGEPACK_ERROR); + + let mut after_marker = rd.clone(); + let value = match rmp::decode::read_marker(&mut after_marker)? { + rmp::Marker::F64 => AttributeValue::F64(rmp::decode::read_f64(&mut rd)?), + rmp::Marker::U64 => AttributeValue::U64(rmp::decode::read_u64(&mut rd)?), + rmp::Marker::I64 => AttributeValue::U64(rmp::decode::read_u64(&mut rd)?), + + rmp::Marker::False => AttributeValue::Bool(false), + rmp::Marker::True => AttributeValue::Bool(true), + + rmp::Marker::FixStr(_) + | rmp::Marker::Str8 + | rmp::Marker::Str16 + | rmp::Marker::Str32 => { + let len = rmp::decode::read_str_len(&mut rd.clone())? as usize; + + let mut buf = vec![0; len]; + rmp::decode::read_str(&mut rd, &mut buf).expect(INTERNAL_MESSAGEPACK_ERROR); + + AttributeValue::String( + String::from_utf8(buf).expect(INTERNAL_MESSAGEPACK_ERROR), + ) + } + + rmp::Marker::FixExt16 => { + let signed = match rmp::decode::read_i8(&mut after_marker)? { + SIGNED_128 => true, + UNSIGNED_128 => false, + _ => unreachable!(), + }; + let mut buf = [0; 16]; + rd.read_exact_buf(&mut buf).map_err(MarkerReadError)?; + if signed { + AttributeValue::U128(u128::from_be_bytes(buf)) + } else { + AttributeValue::I128(i128::from_be_bytes(buf)) + } + } + + rmp::Marker::Bin8 | rmp::Marker::Bin16 | rmp::Marker::Bin32 => { + let len = rmp::decode::read_bin_len(&mut rd)? as usize; + let mut buf = vec![0; len]; + rd.read_exact_buf(&mut buf).map_err(MarkerReadError)?; + AttributeValue::Bytes(buf.into_boxed_slice()) + } + + rmp::Marker::FixPos(_) + | rmp::Marker::FixMap(_) + | rmp::Marker::FixArray(_) + | rmp::Marker::Null + | rmp::Marker::Reserved + | rmp::Marker::Ext8 + | rmp::Marker::Ext16 + | rmp::Marker::Ext32 + | rmp::Marker::F32 + | rmp::Marker::U8 + | rmp::Marker::U16 + | rmp::Marker::U32 + | rmp::Marker::I8 + | rmp::Marker::I16 + | rmp::Marker::I32 + | rmp::Marker::FixExt1 + | rmp::Marker::FixExt2 + | rmp::Marker::FixExt4 + | rmp::Marker::FixExt8 + | rmp::Marker::Array16 + | rmp::Marker::Array32 + | rmp::Marker::Map16 + | rmp::Marker::Map32 + | rmp::Marker::FixNeg(_) => unreachable!(), + }; + + attrs.push(( + String::from_utf8(name).expect(INTERNAL_MESSAGEPACK_ERROR), + value, + )); + } + + Ok(attrs) } } diff --git a/crates/lib/src/registry/mod.rs b/crates/lib/src/registry/mod.rs deleted file mode 100644 index 8d9edc6..0000000 --- a/crates/lib/src/registry/mod.rs +++ /dev/null @@ -1,75 +0,0 @@ -use std::{ - num::NonZeroU64, - sync::{Arc, Mutex, RwLock}, -}; - -use tracing::Subscriber; -use ulid::{Generator, Ulid}; - -use crate::messagepack::MessagePackBytes; - -mod subscriber; -#[doc(inline)] -pub use subscriber::*; - -/// Provides a way of getting subscribers to track futures, stores the collected results -/// quickly in-memory, and then yeets it down a channel once the future finishes. -#[derive(Default)] -pub struct Registry { - ulid_generator: Mutex, - in_mem: RwLock)>>, -} - -impl Registry { - /// Get a subscriber to use for a future. `inp` should be the input to the registry in debug form, so it can be logged. - pub fn subscriber(this: Arc, inp: String) -> impl Subscriber { - // TODO: level filter can go here - let ulid = this.ulid_generator.lock().unwrap().generate().unwrap(); - - this.in_mem - .write() - .unwrap() - .push((ulid, Mutex::new(RegistryEntry::new(inp)))); - - RegistrySubscriber::new(this, ulid) - } - - /// Append to the log for a given future. - /// If [`Self::finished`] has already been called, or `ulid` is not for this registry, this has no effect. - fn add_log_entry(this: &Arc, ulid: Ulid, log: LogEntry) { - if let Some((_, entry)) = this.in_mem.read().unwrap().iter().find(|(u, _)| *u == ulid) { - entry.lock().unwrap().log.push(log); - } - } - - /// Called once a future is done to finalise its logs - // This might be called in a `Drop` implementation so should not block for long. - fn finished(this: &Arc, ulid: Ulid) { - todo!() - } -} - -/// A log, possibly in-progress, recorded from a [`RegistrySubscriber`] -pub struct RegistryEntry { - pub inp: String, - pub log: Vec, -} - -impl RegistryEntry { - fn new(inp: String) -> Self { - Self { - inp, - log: Vec::with_capacity(8), - } - } -} - -/// Entry in the flat log that [`Registry`] uses internally -pub enum LogEntry { - AddSpan(NonZeroU64, MessagePackBytes), - RecordSpan(NonZeroU64, MessagePackBytes), - SpanFollows(NonZeroU64, NonZeroU64), // 1 follows 0 - EnterSpan(NonZeroU64), - Event(MessagePackBytes), - ExitSpan(NonZeroU64), -} diff --git a/crates/lib/src/router.rs b/crates/lib/src/router.rs new file mode 100644 index 0000000..9fecf54 --- /dev/null +++ b/crates/lib/src/router.rs @@ -0,0 +1,85 @@ +//! Thing that wraps a [`tower::Service`], and an ext trait to create it (for convenience) +use redb::Database; +use std::{fmt::Debug, marker::PhantomData, ops::Deref, pin::Pin, sync::Arc}; +use tower::Service; +use tracing::instrument::WithSubscriber; + +use crate::{aggregator::Aggregator, destination::TraceDestination}; + +/// Instruments each call to a service with a subscriber that writes to an aggregator +pub struct LogsRouted> { + inner: S, + aggregator: Arc, + _pd: PhantomData, +} + +impl> LogsRouted { + /// Create a new service, wrapping `inner`. Logs will be routed to the given `dest`, using `slug` as an identifier. + pub fn new>( + inner: S, + slug: &'static str, + dest: &mut TraceDestination, + ) -> Self { + Self { + inner, + aggregator: Arc::new(Aggregator::new(dest.receiver_for(slug))), + + _pd: PhantomData, + } + } +} + +impl> Service for LogsRouted +where + S::Future: Send + 'static, + I: Debug, +{ + type Response = S::Response; + + type Error = S::Error; + + type Future = + Pin> + Send + 'static>>; + + fn poll_ready( + &mut self, + cx: &mut std::task::Context<'_>, + ) -> std::task::Poll> { + self.inner.poll_ready(cx) + } + + fn call(&mut self, req: I) -> Self::Future { + let inp = format!("{:?}", req); + Box::pin( + self.inner + .call(req) + .with_subscriber(Aggregator::subscriber(self.aggregator.clone(), inp)), + ) + } +} + +/// Adds [`ServiceExt::route_logs_to`] to [`Service`]s. +pub trait ServiceExt: Service { + /// Instrument each call to this service so that logs end up written to `dest` + fn route_logs_to>( + self, + slug: &'static str, + dest: &mut TraceDestination, + ) -> LogsRouted + where + Self: Sized; +} + +impl ServiceExt for S +where + S: Service, + I: Debug, +{ + fn route_logs_to>( + self, + slug: &'static str, + dest: &mut TraceDestination, + ) -> LogsRouted { + LogsRouted::new(self, slug, dest) + } +} -- 2.51.2