diff --git a/crates/agent/src/jobs.rs b/crates/agent/src/jobs.rs index 2f0070217eb..7587caf68b8 100644 --- a/crates/agent/src/jobs.rs +++ b/crates/agent/src/jobs.rs @@ -182,7 +182,7 @@ where } #[allow(dead_code)] -async fn read_stdout(mut reader: async_process::ChildStdio) -> Result, Error> { +async fn read_stdout(mut reader: async_process::ChildOutput) -> Result, Error> { let mut buffer = Vec::new(); let _ = tokio::io::copy(&mut reader, &mut buffer) .await diff --git a/crates/async-process/src/lib.rs b/crates/async-process/src/lib.rs index b8b5682c126..7df8d0bccb3 100644 --- a/crates/async-process/src/lib.rs +++ b/crates/async-process/src/lib.rs @@ -1,18 +1,16 @@ pub use std::process::{Command, Output, Stdio}; use tokio::io::{AsyncReadExt, AsyncWriteExt}; -#[cfg(unix)] -use std::os::fd::OwnedFd as OwnedImpl; -#[cfg(windows)] -use std::os::fd::OwnedHandle as OwnedImpl; +#[cfg(not(unix))] +compile_error!("async-process supports only unix targets"); pub struct Child { pid: libc::pid_t, status: Status, - pub stdin: Option, - pub stdout: Option, - pub stderr: Option, + pub stdin: Option, + pub stdout: Option, + pub stderr: Option, } /// State machine for reaping the child process. @@ -32,13 +30,23 @@ enum Status { Abandoned, } -pub type ChildStdio = tokio::fs::File; +/// A child's stdin pipe. +/// +/// Child pipes are non-blocking and driven by the IO driver of the runtime +/// that was current at `Child::from`, and must be polled while it's alive. +/// +/// `flush` and `shutdown` are no-ops, as writes go directly to the pipe. +/// Dropping a `ChildStdin` is what closes it, sending EOF to the child. +pub type ChildStdin = tokio::net::unix::pipe::Sender; + +/// A child's stdout or stderr pipe. See [`ChildStdin`]. +pub type ChildOutput = tokio::net::unix::pipe::Receiver; impl From for Child { fn from(mut inner: std::process::Child) -> Self { - let stdin = map_stdio(inner.stdin.take()); - let stdout = map_stdio(inner.stdout.take()); - let stderr = map_stdio(inner.stderr.take()); + let stdin = map_stdin(inner.stdin.take()); + let stdout = map_output(inner.stdout.take()); + let stderr = map_output(inner.stderr.take()); let pid = inner.id() as libc::pid_t; let status = Status::Reaping(tokio::runtime::Handle::current().spawn_blocking(move || { @@ -198,13 +206,12 @@ pub async fn input_output(cmd: &mut Command, input: &[u8]) -> std::io::Result(f: Option) -> Option -where - F: Into, -{ - let f: Option = f.map(Into::into); - let f: Option = f.map(Into::into); - f.map(Into::into) +fn map_stdin(f: Option) -> Option { + Some(ChildStdin::from_owned_fd(f?.into()).expect("piped child stdin is a pipe")) +} + +fn map_output>(f: Option) -> Option { + Some(ChildOutput::from_owned_fd(f?.into()).expect("piped child output is a pipe")) } #[cfg(test)] diff --git a/crates/connector-init/src/rpc.rs b/crates/connector-init/src/rpc.rs index 241f18cf76c..5fec50c6c19 100644 --- a/crates/connector-init/src/rpc.rs +++ b/crates/connector-init/src/rpc.rs @@ -202,7 +202,7 @@ where /// the connector is already gone. Logging it at `warn` puts pure downstream noise beside /// the causal error, so it's logged at `debug`. Any other error is unexpected and remains /// a `warn`. -async fn write_stdin(stdin: &mut async_process::ChildStdio, buffer: &[u8]) { +async fn write_stdin(stdin: &mut async_process::ChildStdin, buffer: &[u8]) { let Err(error) = stdin.write_all(buffer).await else { return; }; diff --git a/crates/control-plane-api/src/jobs.rs b/crates/control-plane-api/src/jobs.rs index 2f0070217eb..7587caf68b8 100644 --- a/crates/control-plane-api/src/jobs.rs +++ b/crates/control-plane-api/src/jobs.rs @@ -182,7 +182,7 @@ where } #[allow(dead_code)] -async fn read_stdout(mut reader: async_process::ChildStdio) -> Result, Error> { +async fn read_stdout(mut reader: async_process::ChildOutput) -> Result, Error> { let mut buffer = Vec::new(); let _ = tokio::io::copy(&mut reader, &mut buffer) .await diff --git a/crates/json/src/validator/mod.rs b/crates/json/src/validator/mod.rs index 77369fd0934..9a25d1340e6 100644 --- a/crates/json/src/validator/mod.rs +++ b/crates/json/src/validator/mod.rs @@ -71,6 +71,9 @@ where keywords: &'s [Keyword], // Bit flags of this Frame. flags: u8, + // KIND_* classes of `keywords`, which let per-node passes skip Frames + // having no keyword of interest without scanning them. + kinds: u8, // Counter retains: // * The current index into Keyword::Properties (when an object). // * The number of valid Keyword::Contains applications (when an array). @@ -93,6 +96,35 @@ const FLAG_VALID_ANY_OF: u8 = 0x08; // FLAG_VALID_ONE_OF is set if a Keyword::OneOf in-place application validated. const FLAG_VALID_ONE_OF: u8 = 0x10; +// Keyword classes of a Frame's schema. Validation visits each document node +// with every active Frame, and most Frames (in-place applications especially) +// have no keyword relevant to a given visit: these classes let a visit skip +// them outright. Frame retains these eight as a u8 (keeping Frame within a +// cache line); the u16 classes beyond them only matter while winding the +// Frame, and their type keeps them from being tested against Frame::kinds. +// +// Keywords applied to object properties. +const KIND_PROPERTIES: u8 = 0x01; +// Keyword::PropertyNames. +const KIND_PROPERTY_NAMES: u8 = 0x02; +// Keywords applied to array items. +const KIND_ITEMS: u8 = 0x04; +// Keywords checked against string nodes. +const KIND_STRING: u8 = 0x08; +// Keywords checked against numeric nodes. +const KIND_NUMBER: u8 = 0x10; +// Keywords checked against a whole object or array, after its children. +// Includes Keyword::Properties, which checks for missing required properties. +const KIND_CONTAINER: u8 = 0x20; +// Keywords checked against every node. +const KIND_NODE: u8 = 0x40; +// Keyword::AnyOf or OneOf, which are checked as the Frame unwinds. +const KIND_ANY_OR_ONE_OF: u8 = 0x80; +// Keywords having in-place applications. +const KIND_IN_PLACE: u16 = 0x0100; +// Keyword::UnevaluatedItems or UnevaluatedProperties. +const KIND_UNEVALUATED: u16 = 0x0200; + struct FrameSpeculative<'s, A> where A: Annotation, @@ -158,6 +190,7 @@ where parent_keyword: None, keywords: &[], flags: 0, + kinds: 0, outcomes: Vec::new(), counter: 0, speculative: None, diff --git a/crates/json/src/validator/validation.rs b/crates/json/src/validator/validation.rs index 7e09bef5303..88bcb6b27be 100644 --- a/crates/json/src/validator/validation.rs +++ b/crates/json/src/validator/validation.rs @@ -1,6 +1,8 @@ use super::{ FLAG_INVALID, FLAG_VALID_ANY_OF, FLAG_VALID_IF_ELSE, FLAG_VALID_IF_THEN, FLAG_VALID_ONE_OF, - Frame, FrameSpeculative, Outcome, ScopedOutcome, Validation, + Frame, FrameSpeculative, KIND_ANY_OR_ONE_OF, KIND_CONTAINER, KIND_IN_PLACE, KIND_ITEMS, + KIND_NODE, KIND_NUMBER, KIND_PROPERTIES, KIND_PROPERTY_NAMES, KIND_STRING, KIND_UNEVALUATED, + Outcome, ScopedOutcome, Validation, }; use crate::{ AsNode, Node, @@ -109,6 +111,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in active.start..active.end { + if self.stack[frame].kinds & KIND_ITEMS == 0 { + continue; + } let mut matched = false; for kw in self.stack[frame].keywords { @@ -153,6 +158,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in active.start..active.end { + if self.stack[frame].kinds & KIND_PROPERTY_NAMES == 0 { + continue; + } for kw in self.stack[frame].keywords { if let Keyword::PropertyNames { property_names } = &kw { self.wind_frame( @@ -187,6 +195,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in active.start..active.end { + if self.stack[frame].kinds & KIND_PROPERTIES == 0 { + continue; + } let mut matched = false; for kw in self.stack[frame].keywords { @@ -275,6 +286,9 @@ where ) { let frame = self.stack.len(); let keywords = &*schema.keywords; + let kinds = keywords + .iter() + .fold(0, |kinds, kw| kinds | keyword_kinds(kw)); // Ensure we have capacity for another Frame. if self.stack.len() == self.stack.capacity() { @@ -302,6 +316,7 @@ where parent_keyword, keywords, flags: 0, + kinds: kinds as u8, outcomes: Vec::new(), counter: 0, speculative: None, @@ -318,27 +333,19 @@ where // * Or, if this Frame is an in-place application (e.g. $ref) of a // parent that tracks evaluated children, then it too must track // evaluated children (but does not track speculative validations). - for kw in keywords { - match kw { - Keyword::UnevaluatedItems { .. } => { - if let Node::Array(items) = node.as_node() { - self.stack[frame].speculative = Some(Box::new(FrameSpeculative { - evaluated: vec![0u32; (items.len() + 31) / 32].into(), - invalid_unevaluated: Some(vec![0u32; (items.len() + 31) / 32].into()), - outcomes_unevaluated: Vec::new(), - })) - } - } - Keyword::UnevaluatedProperties { .. } => { - if let Node::Object(fields) = node.as_node() { - self.stack[frame].speculative = Some(Box::new(FrameSpeculative { - evaluated: vec![0u32; (fields.len() + 31) / 32].into(), - invalid_unevaluated: Some(vec![0u32; (fields.len() + 31) / 32].into()), - outcomes_unevaluated: Vec::new(), - })) - } - } - _ => (), + if kinds & KIND_UNEVALUATED != 0 { + for kw in keywords { + let len = match (kw, node.as_node()) { + (Keyword::UnevaluatedItems { .. }, Node::Array(items)) => items.len(), + (Keyword::UnevaluatedProperties { .. }, Node::Object(fields)) => fields.len(), + _ => continue, + }; + let words = len.div_ceil(32); + self.stack[frame].speculative = Some(Box::new(FrameSpeculative { + evaluated: vec![0u32; words].into(), + invalid_unevaluated: Some(vec![0u32; words].into()), + outcomes_unevaluated: Vec::new(), + })); } } if in_place { @@ -353,6 +360,9 @@ where } } + if kinds & KIND_IN_PLACE == 0 { + return; + } for kw in keywords { match kw { Keyword::AllOf { all_of } => { @@ -432,6 +442,9 @@ where } for frame in &mut self.stack[active] { + if frame.kinds & KIND_CONTAINER == 0 { + continue; + } for kw in frame.keywords { let invalid: Outcome<'s, A> = match kw { Keyword::MaxItems { max_items } => { @@ -476,6 +489,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in &mut self.stack[active] { + if frame.kinds & KIND_CONTAINER == 0 { + continue; + } for kw in frame.keywords { let invalid: Outcome<'s, A> = match kw { Keyword::MaxProperties { max_properties } => { @@ -519,6 +535,9 @@ where let mut as_number: Option> = None; for frame in &mut self.stack[active] { + if frame.kinds & KIND_STRING == 0 { + continue; + } for kw in frame.keywords { let invalid: Outcome<'s, A> = match kw { Keyword::Format { format } => { @@ -601,6 +620,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in &mut self.stack[active] { + if frame.kinds & KIND_NUMBER == 0 { + continue; + } for kw in frame.keywords { let invalid: Outcome<'s, A> = match kw { Keyword::MinimumPosInt { minimum } => { @@ -715,6 +737,9 @@ where let active = *self.active.last().unwrap() as usize..self.stack.len(); for frame in &mut self.stack[active] { + if frame.kinds & KIND_NODE == 0 { + continue; + } for kw in frame.keywords { let invalid: Outcome<'s, A> = match kw { Keyword::False => Outcome::False, @@ -758,29 +783,48 @@ where #[inline] fn unwind_frame(&mut self, child_index: usize, tape_index: i32) { - let mut frame = self.stack.pop().unwrap(); - let parent = &mut self.stack[frame.parent_frame as usize]; + // Unwind the top Frame in place, and only then drop it. Popping it by + // value was an observed hot-spot: the Frame was typically just written + // by wind_frame, and moving it whole stalls on those pending stores. + let top = self.stack.len() - 1; + let (parents, frame) = self.stack.split_at_mut(top); + let frame = &mut frame[0]; + let parent = &mut parents[frame.parent_frame as usize]; + + Self::unwind_into(&self.filter, frame, parent, child_index, tape_index); + self.stack.truncate(top); + } - for kw in frame.keywords { - let invalid = match kw { - Keyword::AnyOf { .. } => { - if frame.flags & FLAG_VALID_ANY_OF != 0 { - continue; + #[inline] + fn unwind_into( + filter: &F, + frame: &mut Frame<'s, A>, + parent: &mut Frame<'s, A>, + child_index: usize, + tape_index: i32, + ) { + if frame.kinds & KIND_ANY_OR_ONE_OF != 0 { + for kw in frame.keywords { + let invalid = match kw { + Keyword::AnyOf { .. } => { + if frame.flags & FLAG_VALID_ANY_OF != 0 { + continue; + } + Outcome::AnyOfNotMatched } - Outcome::AnyOfNotMatched - } - Keyword::OneOf { .. } => { - if frame.flags & FLAG_VALID_ONE_OF != 0 { - continue; + Keyword::OneOf { .. } => { + if frame.flags & FLAG_VALID_ONE_OF != 0 { + continue; + } + Outcome::OneOfNotMatched } - Outcome::OneOfNotMatched - } - _ => continue, - }; - invalidate(&self.filter, &mut frame, invalid, tape_index); + _ => continue, + }; + invalidate(filter, frame, invalid, tape_index); + } } if frame.speculative.is_some() { - unwind_speculative(&mut frame); + unwind_speculative(frame); } let Some(parent_keyword) = frame.parent_keyword else { @@ -816,12 +860,7 @@ where if is_invalid(frame.flags) { return; } else if parent.flags & FLAG_VALID_ONE_OF != 0 { - invalidate( - &self.filter, - parent, - Outcome::OneOfMultipleMatched, - tape_index, - ); + invalidate(filter, parent, Outcome::OneOfMultipleMatched, tape_index); } else { parent.flags |= FLAG_VALID_ONE_OF; } @@ -851,7 +890,7 @@ where frame.outcomes.clear(); if is_valid(frame.flags) { - if let Some(outcome) = (self.filter)(Outcome::NotIsValid) { + if let Some(outcome) = filter(Outcome::NotIsValid) { frame.outcomes.push(ScopedOutcome { outcome, schema_curi: schema::get_curi(&frame.keywords), @@ -879,7 +918,12 @@ where // Unevaluated applications have bespoke handling. Keyword::UnevaluatedItems { .. } | Keyword::UnevaluatedProperties { .. } => { - unwind_unevaluated(child_index as u32, frame.flags, frame.outcomes, parent); + unwind_unevaluated( + child_index as u32, + frame.flags, + std::mem::take(&mut frame.outcomes), + parent, + ); return; } @@ -889,7 +933,7 @@ where if parent.outcomes.is_empty() { std::mem::swap(&mut parent.outcomes, &mut frame.outcomes); } else { - parent.outcomes.extend(frame.outcomes.into_iter()); + parent.outcomes.append(&mut frame.outcomes); } if is_invalid(frame.flags) { @@ -919,6 +963,77 @@ where } } +/// The KIND_* classes of a Keyword, by the validation passes which examine it. +/// Classes retained by Frame use only the low eight bits. +#[inline(always)] +fn keyword_kinds(kw: &Keyword) -> u16 { + match kw { + Keyword::Annotation { .. } + | Keyword::Const { .. } + | Keyword::Enum { .. } + | Keyword::False + | Keyword::Type { .. } => KIND_NODE as u16, + + Keyword::Properties { .. } => (KIND_PROPERTIES | KIND_CONTAINER) as u16, + Keyword::AdditionalProperties { .. } | Keyword::PatternProperties { .. } => { + KIND_PROPERTIES as u16 + } + Keyword::UnevaluatedProperties { .. } => KIND_PROPERTIES as u16 | KIND_UNEVALUATED, + Keyword::PropertyNames { .. } => KIND_PROPERTY_NAMES as u16, + + Keyword::Items { .. } | Keyword::PrefixItems { .. } => KIND_ITEMS as u16, + Keyword::Contains { .. } => KIND_ITEMS as u16, + Keyword::UnevaluatedItems { .. } => KIND_ITEMS as u16 | KIND_UNEVALUATED, + + Keyword::MaxContains { .. } + | Keyword::MaxItems { .. } + | Keyword::MaxProperties { .. } + | Keyword::MinContains { .. } + | Keyword::MinItems { .. } + | Keyword::MinProperties { .. } + | Keyword::UniqueItems {} => KIND_CONTAINER as u16, + + Keyword::Format { .. } + | Keyword::MaxLength { .. } + | Keyword::MinLength { .. } + | Keyword::Pattern { .. } + | Keyword::XStrMaximum { .. } + | Keyword::XStrMinimum { .. } => KIND_STRING as u16, + + Keyword::ExclusiveMaximumFloat { .. } + | Keyword::ExclusiveMaximumNegInt { .. } + | Keyword::ExclusiveMaximumPosInt { .. } + | Keyword::ExclusiveMinimumFloat { .. } + | Keyword::ExclusiveMinimumNegInt { .. } + | Keyword::ExclusiveMinimumPosInt { .. } + | Keyword::MaximumFloat { .. } + | Keyword::MaximumNegInt { .. } + | Keyword::MaximumPosInt { .. } + | Keyword::MinimumFloat { .. } + | Keyword::MinimumNegInt { .. } + | Keyword::MinimumPosInt { .. } + | Keyword::MultipleOfFloat { .. } + | Keyword::MultipleOfNegInt { .. } + | Keyword::MultipleOfPosInt { .. } => KIND_NUMBER as u16, + + Keyword::AnyOf { .. } | Keyword::OneOf { .. } => KIND_IN_PLACE | KIND_ANY_OR_ONE_OF as u16, + Keyword::AllOf { .. } + | Keyword::DependentSchemas { .. } + | Keyword::DynamicRef { .. } + | Keyword::Else { .. } + | Keyword::If { .. } + | Keyword::Not { .. } + | Keyword::Ref { .. } + | Keyword::Then { .. } => KIND_IN_PLACE, + + Keyword::Anchor { .. } + | Keyword::Definitions { .. } + | Keyword::Defs { .. } + | Keyword::DynamicAnchor { .. } + | Keyword::Id { .. } => 0, + } +} + #[inline(never)] fn unwind_unevaluated<'s, A: Annotation>( child_index: u32, diff --git a/crates/materialize-consistency/src/shim.rs b/crates/materialize-consistency/src/shim.rs index f446e261481..607091cbb68 100644 --- a/crates/materialize-consistency/src/shim.rs +++ b/crates/materialize-consistency/src/shim.rs @@ -123,8 +123,8 @@ pub struct Shim { /// A connector process, and the pipes we speak to it over. struct Instance { child: async_process::Child, - stdin: async_process::ChildStdio, - stdout: Option, + stdin: async_process::ChildStdin, + stdout: Option, } /// The value `materialize-boilerplate` expects in `FLOW_RUNTIME_CODEC`. @@ -451,11 +451,9 @@ impl Zombie { let _ = self.instance.stdin.write_all(&msg).await; } let _ = self.instance.stdin.flush().await; - // `shutdown` on a `ChildStdio` — a `tokio::fs::File` — flushes; it does not close - // the pipe, so it does not end the zombie's session. What ends the session is the - // fence refusing its commit; the pipe closes later, when the request pump returns - // and drops the file. - let _ = self.instance.stdin.shutdown().await; + // Replaying the queue doesn't end the zombie's session. What ends it is the fence + // refusing its commit; the pipe closes later, when the request pump returns and + // drops it. self.frozen = false; self.sealed = true; } @@ -549,7 +547,7 @@ async fn fire_faults( async fn pump_requests( shim: Arc, mut from_runtime: R, - mut to_connector: async_process::ChildStdio, + mut to_connector: async_process::ChildStdin, command: Vec, live_pid: libc::pid_t, ) -> anyhow::Result<()> @@ -565,15 +563,8 @@ where buffer.reserve(1); } if from_runtime.read_buf(&mut buffer).await? == 0 { - // The runtime closed our stdin. Flush what the connector has been given so it can - // finish its session — the pipe itself closes when this function returns and the - // stdio handles drop, which is what the connector sees as EOF. - let _ = to_connector.shutdown().await; - if let Some(z) = &mut zombie { - if !z.frozen { - let _ = z.instance.stdin.shutdown().await; - } - } + // The runtime closed our stdin. The connector's pipe closes when this function + // returns and the stdio handles drop, which is what the connector sees as EOF. return Ok(()); } @@ -694,7 +685,7 @@ where async fn pump_responses( shim: Arc, - mut from_connector: async_process::ChildStdio, + mut from_connector: async_process::ChildOutput, live_pid: libc::pid_t, ) -> anyhow::Result<()> { let mut buffer = Vec::with_capacity(32 * 1024); diff --git a/crates/runtime-next/src/leader/capture/task.rs b/crates/runtime-next/src/leader/capture/task.rs index 7ea9d6e31ac..e7eed861a14 100644 --- a/crates/runtime-next/src/leader/capture/task.rs +++ b/crates/runtime-next/src/leader/capture/task.rs @@ -220,8 +220,8 @@ impl Task { }; let mut close_policy = close_policy::Policy::new(min_txn_duration, max_txn_duration); - // Cap combiner usage at 64MB to favor small transactions. - close_policy.combiner_usage_bytes = 0..(64 * 1024 * 1024); + // Bound to avoid combiner spills and to favor small transactions. + close_policy.read_bytes = 0..(64 * 1024 * 1024); Ok(Self { bindings, diff --git a/crates/service-kit/README.md b/crates/service-kit/README.md index fbad342eb95..f8fa98d2113 100644 --- a/crates/service-kit/README.md +++ b/crates/service-kit/README.md @@ -39,9 +39,13 @@ become visible on the dashboard. - **No auth.** Bind admin on loopback only. - **Trace-override is additive** — it raises verbosity for one handler but - never suppresses what the base filter would keep. Cost when no override - is set is one extra `enabled()` check per disabled callsite (atomic load, - short scope walk only inside a handler span). + never suppresses what the base filter would keep. With no override set, + disabled callsites cost nothing extra: `OverrideFilter` reports + `Interest::never` and an `INFO` level hint. Setting the first override (or + clearing the last) rebuilds the process's callsite interest cache, and + while any is set every callsite below the base level is checked + dynamically (a scope walk per event or span — costly under h2's per-frame + spans), so overrides are for debugging sessions, not standing config. - **Handler spans must always be created.** `OverrideFilter` short-circuits to `true` for the `service_kit::handler` target — that's where override state is hung. Don't filter the target out at the base. diff --git a/crates/service-kit/src/handlers.rs b/crates/service-kit/src/handlers.rs index 7da9fb221fd..fdfc59ed636 100644 --- a/crates/service-kit/src/handlers.rs +++ b/crates/service-kit/src/handlers.rs @@ -159,7 +159,13 @@ impl Registry { return false; }; let value = level.map(|l| level_to_u8(&l)).unwrap_or(0); - slot.trace_override.store(value, Ordering::Relaxed); + let prev = slot.trace_override.swap(value, Ordering::Relaxed); + let rebuild = crate::trace::count_override_change(prev, value); + drop(inner); + + if rebuild { + crate::trace::rebuild_interest(); + } true } @@ -224,11 +230,19 @@ impl Registry { fn finish(&self, id: u64, view: FinishedView) { let mut inner = self.0.lock().unwrap(); - inner.live.remove(&id); + // A finished handler's override no longer admits anything. + let rebuild = inner.live.remove(&id).is_some_and(|slot| { + crate::trace::count_override_change(slot.trace_override.swap(0, Ordering::Relaxed), 0) + }); if inner.recent.len() == RECENT_CAPACITY { inner.recent.pop_front(); } inner.recent.push_back(view); + drop(inner); + + if rebuild { + crate::trace::rebuild_interest(); + } } } diff --git a/crates/service-kit/src/trace.rs b/crates/service-kit/src/trace.rs index b6cf3550d3c..3bc20445f39 100644 --- a/crates/service-kit/src/trace.rs +++ b/crates/service-kit/src/trace.rs @@ -15,14 +15,60 @@ //! override admits the event's level, it passes. The override is *additive* — //! it never suppresses an event the base filter would keep. //! -//! Cost when no override is active: [`OverrideFilter`]'s `max_level_hint` is -//! `TRACE`, so disabled `trace!`/`debug!` callsites do one extra `enabled()` -//! check (an atomic load, plus — only inside a handler span — a short scope -//! walk) rather than being statically skipped. +//! Cost when no override is active: none. [`OverrideFilter`] then reports +//! `Interest::never` and an `INFO` level hint, so disabled `trace!`/`debug!` +//! callsites are skipped statically, as under the base filter alone. Setting +//! the first override (or clearing the last) rebuilds the process's callsite +//! interest cache, and while any override is active every callsite below the +//! base level is checked dynamically: a scope walk per event or span. +//! Hot paths such as h2's per-frame spans make that dynamic check expensive, +//! which is why it's confined to the (rare) periods an override is set. use crate::Registry; use std::sync::Arc; -use std::sync::atomic::{AtomicU8, Ordering}; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering}; + +/// Number of live handlers with a trace override set, process-wide: callsite +/// interest is itself process-global. +static ACTIVE_OVERRIDES: AtomicUsize = AtomicUsize::new(0); + +/// Serializes [`rebuild_interest`]. Held only around the rebuild, never while +/// reading `ACTIVE_OVERRIDES`: the rebuild calls back into [`OverrideFilter`]. +static REBUILD: Mutex<()> = Mutex::new(()); + +/// Count a handler's override moving from `prev` to `next` (0 = none), +/// returning whether the first override was set or the last cleared, in which +/// case the caller must [`rebuild_interest`] once it has released its locks. +/// +/// Callers hold the [`Registry`] lock under which they swapped the override, +/// so each handler's set is counted before its clear and the count never +/// transiently reads zero (or wraps) while an override is live. +pub(crate) fn count_override_change(prev: u8, next: u8) -> bool { + match (prev != 0, next != 0) { + (false, true) => ACTIVE_OVERRIDES.fetch_add(1, Ordering::SeqCst) == 0, + (true, false) => ACTIVE_OVERRIDES.fetch_sub(1, Ordering::SeqCst) == 1, + _ => false, + } +} + +/// Rebuild the process's callsite interest cache and max level from the +/// current `ACTIVE_OVERRIDES`. +/// +/// tracing-core samples the max level hint when a rebuild starts and applies +/// it when it ends, and stores per-callsite interest unsynchronized, so two +/// overlapping rebuilds can finish out of order and leave the older answer in +/// place (e.g. an `INFO` max level while an override is set). Serialized, the +/// last rebuild starts after every count change that requested one, and so +/// reflects the final count. +pub(crate) fn rebuild_interest() { + let _guard = REBUILD.lock().unwrap(); + tracing::callsite::rebuild_interest_cache(); +} + +fn any_override_active() -> bool { + ACTIVE_OVERRIDES.load(Ordering::Relaxed) != 0 +} /// Compose `base` with an [`OverrideFilter`] over `registry`, yielding a filter /// to attach to a `fmt` (or other) layer via `Layer::with_filter`. Events pass @@ -61,6 +107,9 @@ where if meta.target() == crate::handlers::HANDLER_SPAN_TARGET { return true; } + if !any_override_active() { + return false; + } let want = crate::handlers::level_to_u8(meta.level()); let Some(span) = cx.lookup_current() else { return false; @@ -78,15 +127,24 @@ where ) -> tracing::subscriber::Interest { if meta.target() == crate::handlers::HANDLER_SPAN_TARGET { tracing::subscriber::Interest::always() - } else { - // An override set later may admit this callsite, so we can't cache - // a static decision: ask `enabled` per event. + } else if any_override_active() { + // An active override may admit this callsite within its handler + // span, and not elsewhere: ask `enabled` per event. tracing::subscriber::Interest::sometimes() + } else { + // Re-evaluated when an override is set (`rebuild_interest`). + tracing::subscriber::Interest::never() } } fn max_level_hint(&self) -> Option { - Some(tracing_subscriber::filter::LevelFilter::TRACE) + if any_override_active() { + Some(tracing_subscriber::filter::LevelFilter::TRACE) + } else { + // Handler spans are `info_span!`s, which the base filter's hint + // (composed by `or`) must cover for them to be created. + Some(tracing_subscriber::filter::LevelFilter::INFO) + } } fn on_new_span( @@ -198,4 +256,74 @@ mod tests { assert_eq!(n2(), 0); }); } + + #[test] + fn overrides_rebuild_cached_callsite_interest() { + let registry = Registry::new(); + let count = Arc::new(AtomicUsize::new(0)); + let n = || count.load(Ordering::Relaxed); + let active = || ACTIVE_OVERRIDES.load(Ordering::SeqCst); + let max_level = tracing::level_filters::LevelFilter::current; + + let subscriber = + tracing_subscriber::registry().with(CountLayer(count.clone()).with_filter( + layer_filter(tracing_subscriber::EnvFilter::new("info"), registry.clone()), + )); + + tracing::subscriber::with_default(subscriber, || { + // A single callsite, hit throughout: its cached interest (and the + // max level gating it) must follow overrides as they come and go. + let probe = || tracing::trace!("probe"); + + let first = registry.register("test.kind"); + let second = registry.register("test.kind"); + let (first_id, second_id) = { + let live = registry.snapshot().live; + (live[0].id, live[1].id) + }; + + first.span().in_scope(probe); + assert_eq!((n(), active()), (0, 0)); + + assert!(registry.set_trace_override(first_id, Some(tracing::Level::TRACE))); + first.span().in_scope(probe); + probe(); // Outside any handler span. + assert_eq!( + (n(), active(), max_level()), + (1, 1, tracing::level_filters::LevelFilter::TRACE) + ); + + // Clearing the last override returns the callsite to `never`... + assert!(registry.set_trace_override(first_id, None)); + first.span().in_scope(probe); + assert_eq!( + (n(), active(), max_level()), + (1, 0, tracing::level_filters::LevelFilter::INFO) + ); + + // ...and setting one again flips it back. + assert!(registry.set_trace_override(first_id, Some(tracing::Level::TRACE))); + first.span().in_scope(probe); + assert_eq!((n(), active()), (2, 1)); + + // Finishing a handler with its override set releases its count, + // without disturbing another handler's override. + assert!(registry.set_trace_override(second_id, Some(tracing::Level::TRACE))); + assert_eq!(active(), 2); + drop(first); + assert_eq!(active(), 1); + second.span().in_scope(probe); + assert_eq!(n(), 3); + + drop(second); + assert_eq!( + (active(), max_level()), + (0, tracing::level_filters::LevelFilter::INFO) + ); + + let third = registry.register("test.kind"); + third.span().in_scope(probe); + assert_eq!(n(), 3); + }); + } }