Repository navigation
feat(observed): expose sink interest query - #787
Sander Saares (sandersaares) wants to merge 4 commits into
Conversation
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## main #787 +/- ##
========================================
Coverage 100.0% 100.0%
========================================
Files 738 736 -2
Lines 98474 98588 +114
========================================
+ Hits 98474 98588 +114
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
There was a problem hiding this comment.
Copilot review overview
🟢 Approval recommended
The API reuses existing aggregation logic, has focused contract tests, and has no unresolved findings.
Review effort: Balanced
Findings: None
What changed in this PR
This PR exposes the existing sink-interest check so foreign telemetry adapters can decide whether to construct an event.
Changes:
- Makes
Sink::is_interestedpublic and documents its current-interest, not delivery-guarantee, contract. - Adds public API tests for leaf and composite sinks, initialization changes, and sampling.
- Updates the design and feature documentation.
| File | Description |
|---|---|
| crates/observed/tests/sink_interest.rs | Tests the public interest-query contract. |
| crates/observed/src/sink/core.rs | Exposes and documents the existing query. |
| crates/observed/src/processing/processor.rs | Clarifies processor-interest documentation. |
| crates/observed/FEATURES.md | Lists the sink-interest query. |
| crates/observed/docs/implementation.md | Explains interest aggregation. |
| crates/observed/DESIGN.md | Documents the query’s scope and limits. |
| crates/observed/Cargo.toml | Adds a test dependency. |
| Cargo.lock | Records the dependency change. |
💡 Configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
[Copilot speaking] Performance goal and measurementsThe goal is to avoid constructing and formatting events that no downstream processor wants, while minimizing the extra work for accepted events. End-to-end adapter pathsA downstream tracing-adapter microbenchmark gives these instruction counts with the optimized observed dispatch. Percentages are relative to the original adapter path.
These rejection savings require the adapter to use the new query for early admission; exposing the query alone does not change tracing's filtering. Subscriber-wide admission can prevent caller field evaluation. Per-layer admission preserves independent consumers but may still evaluate caller expressions. Accepted events still pay for dynamic admission. Isolated improvement from this PR's dispatch optimizationHolding the adapter, resolved dependency versions and toolchain fixed, the core optimization saves 55 instructions per accepted event:
Rejected TRACE stays at 194 and 360 respectively. Inclusive Callgrind attribution places the entire saving in The isolated core comparison is 9d3ece6 to a6c2de5. Measurement scopeOne warmed emission, a trivial numeric inline field and one synthetic INFO-interested processor, without sampling. Subscriber construction, warm-up, teardown, exporter I/O, message formatting and expensive caller computation are excluded. Every instruction count repeated unchanged. These are instruction-count measurements, not production CPU, latency, throughput or allocation measurements; sampled and composite performance is not represented. Environment: WSL Linux x86_64, Rust 1.95.0, Valgrind 3.22.0, Gungraun runner 0.19.4, |
|
Posted by an AI agent Human review required at |
Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
|
[Copilot speaking] Addressed the recipient-selection feedback in dee5ed3 by clarifying the supported contract and adding a deterministic regression. Single-pass decisions are retained intentionally: interest checks are advisory observations, not an atomic snapshot, and earlier decisions need not be reconsidered within the same emission. The leaf/composite regression covers the described initialization sequence: the earlier rejected recipient stays skipped for that dynamic emission and is selected on the next. Decisions are not cached across emissions. Typed lazy construction and sampled recipient rechecks remain unchanged. Runtime code is unchanged, preserving the measured accepted-event saving. The corresponding inline review thread is resolved. |
| /// Routes an already-built event without a lazy-construction interest pass. | ||
| pub(crate) fn emit_dynamic(&self, event: &dyn DynEvent) { |
|
At head 1. New public methodimpl observed::Sink {
#[must_use]
#[inline]
pub fn is_interested(
&self,
description: &observed::metadata::EventDescription,
) -> bool;
}
This replaces the previously private 2. Updated public trait contractThe signature is unchanged: observed::processing::EventProcessor::is_interested(
&self,
description: &EventDescription,
) -> bool;Its documentation now explicitly states:
The existing restriction remains: interest depends on metadata and state that changes at most once; per-call sampling/rate limiting belongs in 3. Changed behavior behind an unchanged public functionobserved::interop::emit_dyn_event(
sink: &Sink,
event: &dyn DynEvent,
);Dynamic events now bypass the lazy typed-event machinery:
Leaf-local timestamps, enrichment, sampling, and processor order remain preserved. No public items are removed, no existing public signatures change, and no public types, fields, variants, re-exports, or feature flags are added or changed. The removed event-state types and new dispatch helpers are crate-private/private. |
There was a problem hiding this comment.
I think this document can be cut down. For example "Public API tests use frozen clocks" is just a common good practice and mentioning this in a document doesn't add any value
| let timestamp = self.clock.system_time(); | ||
| let view = EventView::new(event, self.enrichment.current(), self.isolated_enrichment, self.id, timestamp); | ||
| first.process(&view); | ||
| // Evaluate each remaining recipient after earlier processors have run, | ||
| // so their initialization effects are visible without caching interest. | ||
| for processor in interested { | ||
| processor.process(&view); | ||
| } |
There was a problem hiding this comment.
The simplification in event_state.rs is worth keeping. We can reduce duplication in sink/core.rs without reintroducing the dynamic-event wrapper: extract only view construction and recipient delivery.
Concretely, add use std::iter; and use std::time::SystemTime;, then use this structure inside impl SingleSinkState:
#[inline]
fn deliver<'p>(
&self,
event: &dyn DynEvent,
timestamp: SystemTime,
recipients: impl Iterator<Item = &'p Arc<dyn EventProcessor>>,
) {
let view = EventView::new(
event,
self.enrichment.current(),
self.isolated_enrichment,
self.id,
timestamp,
);
for processor in recipients {
processor.process(&view);
}
}
fn dispatch(&self, event: &dyn DynEvent, description: &EventDescription) {
let timestamp = self.clock.system_time();
if let Some(sampler) = &self.sampler
&& sampler.sample(&EventSamplingContext::new(
description,
self.id,
timestamp,
)) == EventSamplingDecision::Drop
{
return;
}
self.deliver(
event,
timestamp,
self.processors
.iter()
.filter(|processor| processor.is_interested(description)),
);
}
fn prepare_dynamic_dispatch<'a>(
&'a self,
event: &'a dyn DynEvent,
description: &'a EventDescription,
) -> Option<impl FnOnce() + 'a> {
let mut interested = self.processors
.iter()
.filter(move |processor| processor.is_interested(description));
let first = interested.next()?;
Some(move || {
if self.sampler.is_some() {
// Sampling can change interest; recheck every processor afterward.
self.dispatch(event, description);
return;
}
self.deliver(
event,
self.clock.system_time(),
iter::once(first).chain(interested),
);
})
}This preserves the PR's ordering: unsampled dynamic delivery retains the first selected recipient; later interest checks remain lazy and observe preceding processor callbacks; sampled delivery rechecks every processor after sampling; timestamps are read before sampling and enrichment is captured afterward. Admission and reentrancy-guard placement remain unchanged.
The helper uses a monomorphized iterator, with no iterator boxing, recipient collection, added allocation, or new runtime mode branch. #[inline] makes eliminating the helper call straightforward, although an exact zero-overhead claim would require generated-code or benchmark evidence.
I would share this delivery stage, not force construction, admission, sampling, and guard acquisition into one universal pipeline.
|
AI review preparation has started for #787 at head The submitted GitHub review will contain the result. |
martintmk
left a comment
There was a problem hiding this comment.
Posted by an AI agent
Warning: Incomplete review
I reviewed the sink-interest API and public contracts from source, dynamic-dispatch correctness, tests, performance-related code, naming, telemetry, and documentation at dee5ed39baaf8e89f9db6afe929626cd0b3c9728. No recovery mechanisms changed. I could not check:
- Exported API comparison: matching base/head API captures were unavailable, and the credentialed review environment did not permit generating them.
No new findings were identified in the reviewed areas. The existing request to disclose the dynamic-delivery semantics change remains applicable; it is not duplicated here.
I did not run local builds, tests, or benchmarks. Supplied Linux ARM CI evidence shows all 19 added integration tests passing twice on verified synthetic merge 5d88263263236159f4827b2dda15d92bf70df184, containing this head and later main 27a8b19174fb99633555d9c2194e7efd1782b087. This is not a locally run base/head reproduction.
No overall verdict is given.
martintmk
left a comment
There was a problem hiding this comment.
Posted by an AI agent
Warning: Incomplete review
I reviewed public contracts, correctness, tests, performance, naming, telemetry, and documentation. Resilience had no relevant behavior. I could not check:
- Exported public API: exact base/head outputs were unavailable, and building PR code on the shared credentialed host was not permitted.
No local builds, tests, or benchmarks were run. The earlier request to disclose changed dynamic-delivery semantics remains applicable and is not duplicated here.
No overall verdict is given.
| pub(crate) struct IntermediateEvent<'a, F> { | ||
| inner: Inner<'a, F>, | ||
| /// Holds a typed event builder until processor interest permits construction. | ||
| pub(crate) struct IntermediateEvent<F> { |
There was a problem hiding this comment.
Posted by an AI agent · Non-blocking
IntermediateEvent retains a redundant construction layer
Problem
The dynamic variant was removed, leaving IntermediateEvent to store only the typed builder and source location that Sink::emit already has; emit_impl then evaluates this wrapper immediately before dispatch. It no longer encodes a choice or enforces an invariant.
Why this matters
The extra internal type and forwarding method obscure where construction happens and keep a state-machine layer with only one state.
Suggested fix
Keep EvaluatedEvent as the adapter, but perform the interest check and build in Sink::emit before acquiring the recursion guard, then remove IntermediateEvent and update the implementation guide's construction-state wording.


[Copilot speaking]
Foreign telemetry adapters need to check downstream interest before constructing events.
Expose
Sink::is_interestedto report whether any downstream processor is interested in the described event, using the existing processor and composite-sink aggregation. LikeEventProcessor::is_interested, the query inspects the event description rather than event field values. Its result reflects current interest, not a lifetime-cacheable decision or a delivery guarantee.