From 7fed2b3cbc14f65c114bb9338c26ba34df970f01 Mon Sep 17 00:00:00 2001 From: Aria Date: Thu, 11 Jun 2026 18:36:54 +0100 Subject: [PATCH] wip: most of a tracing subscriber for binted services, going to a central registry --- Cargo.lock | 180 ++++++++++++++++++++++++++++++++++ crates/lib/Cargo.toml | 1 + crates/lib/src/bint.rs | 67 +++++++++++++ crates/lib/src/lib.rs | 93 ++++++++++-------- crates/lib/src/messagepack.rs | 20 ++++ crates/lib/src/subscriber.rs | 87 ++++++++++++++++ 6 files changed, 408 insertions(+), 40 deletions(-) create mode 100644 crates/lib/src/bint.rs create mode 100644 crates/lib/src/messagepack.rs create mode 100644 crates/lib/src/subscriber.rs diff --git a/Cargo.lock b/Cargo.lock index ef3fae0..fa432e0 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -77,6 +77,7 @@ dependencies = [ "tokio", "tower", "tracing", + "ulid", ] [[package]] @@ -91,6 +92,12 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "bumpalo" +version = "3.20.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "72f5acc6cb2ba439de613abc23857ec3d78374d8ed5ac84e9d11336e87da8649" + [[package]] name = "bytes" version = "1.11.1" @@ -145,6 +152,18 @@ dependencies = [ "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" @@ -231,6 +250,17 @@ version = "1.0.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" +[[package]] +name = "js-sys" +version = "0.3.100" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2025f20d7a4fa7785846e7b63d10a76d3f1cee98ee5cb79ea59703f95e42162" +dependencies = [ + "cfg-if", + "futures-util", + "wasm-bindgen", +] + [[package]] name = "lazy_static" version = "1.5.0" @@ -334,6 +364,15 @@ version = "0.2.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a89322df9ebe1c1578d689c92318e070967d1042b512afbe49518723f4e6d5cd" +[[package]] +name = "ppv-lite86" +version = "0.2.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "85eae3c4ed2f50dcfe72643da4befc30deadb458a9b590d720cde2f2b1e97da9" +dependencies = [ + "zerocopy", +] + [[package]] name = "proc-macro2" version = "1.0.106" @@ -352,6 +391,41 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "r-efi" +version = "5.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "69cdb34c158ceb288df11e18b4bd39de994f6657d83847bdffdbd7f346754b0f" + +[[package]] +name = "rand" +version = "0.9.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "44c5af06bb1b7d3216d91932aed5265164bf384dc89cd6ba05cf59a35f5f76ea" +dependencies = [ + "rand_chacha", + "rand_core", +] + +[[package]] +name = "rand_chacha" +version = "0.9.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb" +dependencies = [ + "ppv-lite86", + "rand_core", +] + +[[package]] +name = "rand_core" +version = "0.9.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c" +dependencies = [ + "getrandom", +] + [[package]] name = "regex-automata" version = "0.4.14" @@ -369,6 +443,12 @@ version = "0.8.9" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a96887878f22d7bad8a3b6dc5b7440e0ada9a245242924394987b21cf2210a4c" +[[package]] +name = "rustversion" +version = "1.0.22" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" + [[package]] name = "ryu" version = "1.0.23" @@ -622,6 +702,16 @@ dependencies = [ "tracing-log", ] +[[package]] +name = "ulid" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "470dbf6591da1b39d43c14523b2b469c86879a53e8b758c8e090a470fe7b1fbe" +dependencies = [ + "rand", + "web-time", +] + [[package]] name = "unicode-ident" version = "1.0.24" @@ -640,6 +730,70 @@ 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" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "20064672db26d7cdc89c7798c48a0fdfac8213434a1186e5ef29fd560ae223d6" +dependencies = [ + "wit-bindgen", +] + +[[package]] +name = "wasm-bindgen" +version = "0.2.123" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a254a4b10c19a76f09a27640e7ffbf9bc30bf67e16a3bf28aaefa4920fe81563" +dependencies = [ + "cfg-if", + "once_cell", + "rustversion", + "wasm-bindgen-macro", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-macro" +version = "0.2.123" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "24a40fc75b0ec6f3746ceb10d36f53a93dcd68a93b11b6445983945d79eba0dc" +dependencies = [ + "quote", + "wasm-bindgen-macro-support", +] + +[[package]] +name = "wasm-bindgen-macro-support" +version = "0.2.123" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "908f34bd9b9ce3d4caf07b72dfab63d61504d156856c6bd3cd87fa350cf3985b" +dependencies = [ + "bumpalo", + "proc-macro2", + "quote", + "syn", + "wasm-bindgen-shared", +] + +[[package]] +name = "wasm-bindgen-shared" +version = "0.2.123" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7acbf7616c27b194bbb550bf77ed0c2c3e5b7fd1260a93082b95fb7f47959b92" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "web-time" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5a6580f308b1fad9207618087a65c04e7a10bc77e02c8e84e9b00dd4b12fa0bb" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + [[package]] name = "windows-link" version = "0.2.1" @@ -655,6 +809,32 @@ dependencies = [ "windows-link", ] +[[package]] +name = "wit-bindgen" +version = "0.57.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e" + +[[package]] +name = "zerocopy" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ce1022995ff5ff5d841ad7d994facc23098cd40152f2c1d11cd607c6f530653f" +dependencies = [ + "zerocopy-derive", +] + +[[package]] +name = "zerocopy-derive" +version = "0.8.52" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1ae7f38b72ec2a254e2b87ef277cf2cd4fb97cbebf944faa6f33354da0867930" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "zmij" version = "1.0.21" diff --git a/crates/lib/Cargo.toml b/crates/lib/Cargo.toml index f3de1b4..7241170 100644 --- a/crates/lib/Cargo.toml +++ b/crates/lib/Cargo.toml @@ -8,3 +8,4 @@ axum = { workspace = true } tokio = { workspace = true } tower = { workspace = true } tracing = { workspace = true } +ulid = "1.2.1" diff --git a/crates/lib/src/bint.rs b/crates/lib/src/bint.rs new file mode 100644 index 0000000..00c72c2 --- /dev/null +++ b/crates/lib/src/bint.rs @@ -0,0 +1,67 @@ +//! 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; +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/lib.rs b/crates/lib/src/lib.rs index df08cbc..15f3e10 100644 --- a/crates/lib/src/lib.rs +++ b/crates/lib/src/lib.rs @@ -1,55 +1,68 @@ -use std::{marker::PhantomData, pin::Pin}; +use std::{ + num::NonZeroU64, + sync::{Arc, Mutex, RwLock}, +}; -use tower::Service; +use tracing::Subscriber; +use ulid::{Generator, Ulid}; -pub trait ServiceExt: Service { - fn bint(self) -> Binted - where - Self: Sized; +mod bint; +#[doc(inline)] +pub use bint::*; + +mod subscriber; +use subscriber::*; + +use crate::messagepack::MessagePackBytes; + +mod messagepack; + +#[derive(Default)] +struct Registry { + ulid_generator: Mutex, + + in_mem: RwLock)>>, } -impl ServiceExt for S -where - S: Service, -{ - fn bint(self) -> Binted { - Binted::new(self) +impl Registry { + pub(crate) 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(RequestEntry::new(inp)))); + + RegistrySubscriber::new(this, ulid) + } + + pub(crate) 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); + } } } -pub struct Binted> { - inner: S, - _pd: PhantomData, +struct RequestEntry { + inp: String, + log: Vec, } -impl> Binted { - pub fn new(inner: S) -> Self { +impl RequestEntry { + fn new(inp: String) -> Self { Self { - inner, - _pd: PhantomData, + inp, + log: Vec::with_capacity(8), } } } -impl> Service for Binted -where - S::Future: Send + 'static, -{ - 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 { - Box::pin(self.inner.call(req)) - } +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/messagepack.rs b/crates/lib/src/messagepack.rs new file mode 100644 index 0000000..2bd9f0a --- /dev/null +++ b/crates/lib/src/messagepack.rs @@ -0,0 +1,20 @@ +use tracing::{Event, span}; + +pub struct MessagePackBytes(Vec); + +impl From<&span::Attributes<'_>> for MessagePackBytes { + fn from(value: &span::Attributes<'_>) -> Self { + todo!() + } +} + +impl From<&span::Record<'_>> for MessagePackBytes { + fn from(value: &span::Record<'_>) -> Self { + todo!() + } +} +impl From<&Event<'_>> for MessagePackBytes { + fn from(value: &Event<'_>) -> Self { + todo!() + } +} diff --git a/crates/lib/src/subscriber.rs b/crates/lib/src/subscriber.rs new file mode 100644 index 0000000..7eb2f9a --- /dev/null +++ b/crates/lib/src/subscriber.rs @@ -0,0 +1,87 @@ +use std::{ + num::NonZero, + sync::{ + Arc, + atomic::{AtomicU64, Ordering}, + }, +}; + +use tracing::{Event, Metadata, Subscriber, span}; +use ulid::Ulid; + +use crate::{LogEntry, Registry, messagepack::MessagePackBytes}; + +/// receives tracing events / spans and bridges them to the registry +pub(crate) struct RegistrySubscriber { + registry: Arc, + next_id: AtomicU64, + ulid: Ulid, +} + +impl RegistrySubscriber { + pub fn new(registry: Arc, ulid: Ulid) -> Self { + Self { + registry, + next_id: AtomicU64::new(1), + ulid, + } + } +} + +impl Subscriber for RegistrySubscriber { + fn enabled(&self, metadata: &Metadata<'_>) -> bool { + true // filering is provided by layers above our subscriber + } + + 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, + self.ulid, + LogEntry::AddSpan(NonZero::new(id).unwrap(), MessagePackBytes::from(span)), + ); + + span::Id::from_u64(id) + } + + fn record(&self, span: &span::Id, values: &span::Record<'_>) { + Registry::add_log_entry( + &self.registry, + 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, + 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, + self.ulid, + LogEntry::Event(MessagePackBytes::from(event)), + ); + } + + fn enter(&self, span: &span::Id) { + Registry::add_log_entry( + &self.registry, + self.ulid, + LogEntry::EnterSpan(span.into_non_zero_u64()), + ); + } + + fn exit(&self, span: &span::Id) { + Registry::add_log_entry( + &self.registry, + self.ulid, + LogEntry::ExitSpan(span.into_non_zero_u64()), + ); + } +} -- 2.51.2