saluki_core/data_model/event/trace_stats/
mod.rs

1//! Trace stats.
2
3use stringtheory::MetaString;
4
5use crate::data_model::tags::TagSet;
6
7/// Trace statistics output from the APM Stats transform.
8///
9/// Contains pre-aggregated trace statistics grouped by client/tracer. The encoder wraps this
10/// in a `StatsPayload` protobuf and adds agent-level metadata (`agentHostname`, `agentEnv`,
11/// `agentVersion`, `clientComputed`, `splitPayload`) from ADP configuration.
12#[derive(Clone, Debug, PartialEq, Default)]
13pub struct TraceStats {
14    /// Multiple client payloads, one per PayloadAggregationKey (hostname/env/version/container).
15    stats: Vec<ClientStatsPayload>,
16}
17
18impl TraceStats {
19    /// Creates a new `TraceStats` with the given client stats payloads.
20    pub fn new(stats: Vec<ClientStatsPayload>) -> Self {
21        Self { stats }
22    }
23
24    /// Returns a reference to the client stats payloads.
25    pub fn stats(&self) -> &[ClientStatsPayload] {
26        &self.stats
27    }
28
29    /// Returns a mutable reference to the client stats payloads.
30    pub fn stats_mut(&mut self) -> &mut Vec<ClientStatsPayload> {
31        &mut self.stats
32    }
33}
34
35/// Tracer-level stats payload.
36///
37/// Groups stats by tracer/container identity (hostname, env, version, `container_id`).
38#[derive(Clone, Debug, PartialEq, Default)]
39pub struct ClientStatsPayload {
40    hostname: MetaString,
41    env: MetaString,
42    version: MetaString,
43    stats: Vec<ClientStatsBucket>,
44    lang: MetaString,
45    tracer_version: MetaString,
46    runtime_id: MetaString,
47    sequence: u64,
48    agent_aggregation: MetaString,
49    service: MetaString,
50    container_id: MetaString,
51    tags: TagSet,
52    git_commit_sha: MetaString,
53    image_tag: MetaString,
54    process_tags_hash: u64,
55    process_tags: MetaString,
56}
57
58impl ClientStatsPayload {
59    /// Creates a new `ClientStatsPayload` with the required identity fields.
60    pub fn new(hostname: impl Into<MetaString>, env: impl Into<MetaString>, version: impl Into<MetaString>) -> Self {
61        Self {
62            hostname: hostname.into(),
63            env: env.into(),
64            version: version.into(),
65            ..Self::default()
66        }
67    }
68
69    /// Sets the stats buckets.
70    pub fn with_stats(mut self, stats: Vec<ClientStatsBucket>) -> Self {
71        self.stats = stats;
72        self
73    }
74
75    /// Sets the tracer language.
76    pub fn with_lang(mut self, lang: impl Into<MetaString>) -> Self {
77        self.lang = lang.into();
78        self
79    }
80
81    /// Sets the tracer version.
82    pub fn with_tracer_version(mut self, tracer_version: impl Into<MetaString>) -> Self {
83        self.tracer_version = tracer_version.into();
84        self
85    }
86
87    /// Sets the runtime identifier.
88    pub fn with_runtime_id(mut self, runtime_id: impl Into<MetaString>) -> Self {
89        self.runtime_id = runtime_id.into();
90        self
91    }
92
93    /// Sets the message sequence number.
94    pub fn with_sequence(mut self, sequence: u64) -> Self {
95        self.sequence = sequence;
96        self
97    }
98
99    /// Sets the agent aggregation key.
100    pub fn with_agent_aggregation(mut self, agent_aggregation: impl Into<MetaString>) -> Self {
101        self.agent_aggregation = agent_aggregation.into();
102        self
103    }
104
105    /// Sets the main service name.
106    pub fn with_service(mut self, service: impl Into<MetaString>) -> Self {
107        self.service = service.into();
108        self
109    }
110
111    /// Sets the container identifier.
112    pub fn with_container_id(mut self, container_id: impl Into<MetaString>) -> Self {
113        self.container_id = container_id.into();
114        self
115    }
116
117    /// Sets the orchestrator tags.
118    pub fn with_tags(mut self, tags: impl Into<TagSet>) -> Self {
119        self.tags = tags.into();
120        self
121    }
122
123    /// Sets the git commit SHA.
124    pub fn with_git_commit_sha(mut self, git_commit_sha: impl Into<MetaString>) -> Self {
125        self.git_commit_sha = git_commit_sha.into();
126        self
127    }
128
129    /// Sets the container image tag.
130    pub fn with_image_tag(mut self, image_tag: impl Into<MetaString>) -> Self {
131        self.image_tag = image_tag.into();
132        self
133    }
134
135    /// Sets the process tags hash.
136    pub fn with_process_tags_hash(mut self, process_tags_hash: u64) -> Self {
137        self.process_tags_hash = process_tags_hash;
138        self
139    }
140
141    /// Sets the process tags.
142    pub fn with_process_tags(mut self, process_tags: impl Into<MetaString>) -> Self {
143        self.process_tags = process_tags.into();
144        self
145    }
146
147    /// Returns the hostname.
148    pub fn hostname(&self) -> &str {
149        &self.hostname
150    }
151
152    /// Returns the environment.
153    pub fn env(&self) -> &str {
154        &self.env
155    }
156
157    /// Returns the version.
158    pub fn version(&self) -> &str {
159        &self.version
160    }
161
162    /// Returns the stats buckets.
163    pub fn stats(&self) -> &[ClientStatsBucket] {
164        &self.stats
165    }
166
167    /// Returns the tracer language.
168    pub fn lang(&self) -> &str {
169        &self.lang
170    }
171
172    /// Returns the tracer version.
173    pub fn tracer_version(&self) -> &str {
174        &self.tracer_version
175    }
176
177    /// Returns the runtime identifier.
178    pub fn runtime_id(&self) -> &str {
179        &self.runtime_id
180    }
181
182    /// Returns the message sequence number.
183    pub fn sequence(&self) -> u64 {
184        self.sequence
185    }
186
187    /// Returns the agent aggregation key.
188    pub fn agent_aggregation(&self) -> &str {
189        &self.agent_aggregation
190    }
191
192    /// Returns the main service name.
193    pub fn service(&self) -> &str {
194        &self.service
195    }
196
197    /// Returns the container identifier.
198    pub fn container_id(&self) -> &str {
199        &self.container_id
200    }
201
202    /// Returns the orchestrator tags.
203    pub fn tags(&self) -> &TagSet {
204        &self.tags
205    }
206
207    /// Returns the git commit SHA.
208    pub fn git_commit_sha(&self) -> &str {
209        &self.git_commit_sha
210    }
211
212    /// Returns the container image tag.
213    pub fn image_tag(&self) -> &str {
214        &self.image_tag
215    }
216
217    /// Returns the process tags hash.
218    pub fn process_tags_hash(&self) -> u64 {
219        self.process_tags_hash
220    }
221
222    /// Returns the process tags.
223    pub fn process_tags(&self) -> &str {
224        &self.process_tags
225    }
226
227    /// Adds a new client statistics bucket to this payload.
228    pub fn add_stats(&mut self, stats: ClientStatsBucket) {
229        self.stats.push(stats);
230    }
231
232    /// Consumes the statistics buckets and returns them.
233    ///
234    /// No statistics buckets will remain in `self`.
235    pub fn take_stats(&mut self) -> Vec<ClientStatsBucket> {
236        std::mem::take(&mut self.stats)
237    }
238}
239
240/// A time bucket containing aggregated stats.
241///
242/// Stats are grouped into fixed-duration buckets (typically 10 seconds).
243#[derive(Clone, Debug, PartialEq, Default)]
244pub struct ClientStatsBucket {
245    /// Bucket start timestamp in nanoseconds since Unix epoch.
246    start: u64,
247    /// Bucket duration in nanoseconds.
248    duration: u64,
249    /// Grouped stats within this bucket.
250    stats: Vec<ClientGroupedStats>,
251    /// Time shift applied by the agent.
252    agent_time_shift: i64,
253}
254
255impl ClientStatsBucket {
256    /// Creates a new `ClientStatsBucket` with the given time range and stats.
257    pub fn new(start: u64, duration: u64, stats: Vec<ClientGroupedStats>) -> Self {
258        Self {
259            start,
260            duration,
261            stats,
262            agent_time_shift: 0,
263        }
264    }
265
266    /// Sets the grouped stats.
267    pub fn with_stats(mut self, stats: Vec<ClientGroupedStats>) -> Self {
268        self.stats = stats;
269        self
270    }
271
272    /// Sets the agent time shift.
273    pub fn with_agent_time_shift(mut self, agent_time_shift: i64) -> Self {
274        self.agent_time_shift = agent_time_shift;
275        self
276    }
277
278    /// Returns the bucket start timestamp in nanoseconds.
279    pub fn start(&self) -> u64 {
280        self.start
281    }
282
283    /// Returns the bucket duration in nanoseconds.
284    pub fn duration(&self) -> u64 {
285        self.duration
286    }
287
288    /// Returns the grouped stats within this bucket.
289    pub fn stats(&self) -> &[ClientGroupedStats] {
290        &self.stats
291    }
292
293    /// Returns a mutable reference to the grouped stats within this bucket.
294    pub fn stats_mut(&mut self) -> &mut Vec<ClientGroupedStats> {
295        &mut self.stats
296    }
297
298    /// Returns the agent time shift.
299    pub fn agent_time_shift(&self) -> i64 {
300        self.agent_time_shift
301    }
302
303    /// Consumes the grouped statistics and returns them.
304    ///
305    /// No statistics groups will remain in `self`.
306    pub fn take_stats(&mut self) -> Vec<ClientGroupedStats> {
307        std::mem::take(&mut self.stats)
308    }
309}
310
311/// Aggregated stats for spans grouped by aggregation key.
312///
313/// Contains both the aggregation key fields (service, name, resource, etc.) and
314/// the aggregated values (hits, errors, duration, latency distributions).
315#[derive(Clone, Debug, PartialEq, Default)]
316pub struct ClientGroupedStats {
317    // Aggregation key fields
318    service: MetaString,
319    name: MetaString,
320    resource: MetaString,
321    http_status_code: u32,
322    span_type: MetaString,
323    db_type: MetaString,
324    span_kind: MetaString,
325    peer_tags: Vec<MetaString>,
326    is_trace_root: Option<bool>,
327    grpc_status_code: MetaString,
328    http_method: MetaString,
329    http_endpoint: MetaString,
330
331    // Aggregated values
332    hits: u64,
333    errors: u64,
334    duration: u64,
335    ok_summary: Vec<u8>,
336    error_summary: Vec<u8>,
337    synthetics: bool,
338    top_level_hits: u64,
339}
340
341impl ClientGroupedStats {
342    /// Creates a new `ClientGroupedStats` with the required aggregation key fields.
343    pub fn new(service: impl Into<MetaString>, name: impl Into<MetaString>, resource: impl Into<MetaString>) -> Self {
344        Self {
345            service: service.into(),
346            name: name.into(),
347            resource: resource.into(),
348            ..Self::default()
349        }
350    }
351
352    // Builder methods for aggregation key fields
353
354    /// Sets the HTTP status code.
355    pub fn with_http_status_code(mut self, http_status_code: u32) -> Self {
356        self.http_status_code = http_status_code;
357        self
358    }
359
360    /// Sets the span type.
361    pub fn with_span_type(mut self, span_type: impl Into<MetaString>) -> Self {
362        self.span_type = span_type.into();
363        self
364    }
365
366    /// Sets the database type.
367    pub fn with_db_type(mut self, db_type: impl Into<MetaString>) -> Self {
368        self.db_type = db_type.into();
369        self
370    }
371
372    /// Sets the span kind.
373    pub fn with_span_kind(mut self, span_kind: impl Into<MetaString>) -> Self {
374        self.span_kind = span_kind.into();
375        self
376    }
377
378    /// Sets the peer tags.
379    pub fn with_peer_tags(mut self, peer_tags: Vec<MetaString>) -> Self {
380        self.peer_tags = peer_tags;
381        self
382    }
383
384    /// Sets whether this is a trace root.
385    pub fn with_is_trace_root(mut self, is_trace_root: Option<bool>) -> Self {
386        self.is_trace_root = is_trace_root;
387        self
388    }
389
390    /// Sets the gRPC status code.
391    pub fn with_grpc_status_code(mut self, grpc_status_code: impl Into<MetaString>) -> Self {
392        self.grpc_status_code = grpc_status_code.into();
393        self
394    }
395
396    /// Sets the HTTP method.
397    pub fn with_http_method(mut self, http_method: impl Into<MetaString>) -> Self {
398        self.http_method = http_method.into();
399        self
400    }
401
402    /// Sets the HTTP endpoint.
403    pub fn with_http_endpoint(mut self, http_endpoint: impl Into<MetaString>) -> Self {
404        self.http_endpoint = http_endpoint.into();
405        self
406    }
407
408    // Builder methods for aggregated values
409
410    /// Sets the hit count.
411    pub fn with_hits(mut self, hits: u64) -> Self {
412        self.hits = hits;
413        self
414    }
415
416    /// Sets the error count.
417    pub fn with_errors(mut self, errors: u64) -> Self {
418        self.errors = errors;
419        self
420    }
421
422    /// Sets the total duration in nanoseconds.
423    pub fn with_duration(mut self, duration: u64) -> Self {
424        self.duration = duration;
425        self
426    }
427
428    /// Sets the DDSketch summary for successful spans.
429    pub fn with_ok_summary(mut self, ok_summary: Vec<u8>) -> Self {
430        self.ok_summary = ok_summary;
431        self
432    }
433
434    /// Sets the DDSketch summary for error spans.
435    pub fn with_error_summary(mut self, error_summary: Vec<u8>) -> Self {
436        self.error_summary = error_summary;
437        self
438    }
439
440    /// Sets the synthetics traffic flag.
441    pub fn with_synthetics(mut self, synthetics: bool) -> Self {
442        self.synthetics = synthetics;
443        self
444    }
445
446    /// Sets the top-level hit count.
447    pub fn with_top_level_hits(mut self, top_level_hits: u64) -> Self {
448        self.top_level_hits = top_level_hits;
449        self
450    }
451
452    // Getters for aggregation key fields
453
454    /// Returns the service name.
455    pub fn service(&self) -> &str {
456        &self.service
457    }
458
459    /// Returns the operation name.
460    pub fn name(&self) -> &str {
461        &self.name
462    }
463
464    /// Returns the resource name.
465    pub fn resource(&self) -> &str {
466        &self.resource
467    }
468
469    /// Returns the HTTP status code.
470    pub fn http_status_code(&self) -> u32 {
471        self.http_status_code
472    }
473
474    /// Returns the span type.
475    pub fn span_type(&self) -> &str {
476        &self.span_type
477    }
478
479    /// Returns the database type.
480    pub fn db_type(&self) -> &str {
481        &self.db_type
482    }
483
484    /// Returns the span kind.
485    pub fn span_kind(&self) -> &str {
486        &self.span_kind
487    }
488
489    /// Returns the peer tags.
490    pub fn peer_tags(&self) -> &[MetaString] {
491        &self.peer_tags
492    }
493
494    /// Returns whether this is a trace root.
495    pub fn is_trace_root(&self) -> Option<bool> {
496        self.is_trace_root
497    }
498
499    /// Returns the gRPC status code.
500    pub fn grpc_status_code(&self) -> &str {
501        &self.grpc_status_code
502    }
503
504    /// Returns the HTTP method.
505    pub fn http_method(&self) -> &str {
506        &self.http_method
507    }
508
509    /// Returns the HTTP endpoint.
510    pub fn http_endpoint(&self) -> &str {
511        &self.http_endpoint
512    }
513
514    // Getters for aggregated values
515
516    /// Returns the hit count.
517    pub fn hits(&self) -> u64 {
518        self.hits
519    }
520
521    /// Returns the error count.
522    pub fn errors(&self) -> u64 {
523        self.errors
524    }
525
526    /// Returns the total duration in nanoseconds.
527    pub fn duration(&self) -> u64 {
528        self.duration
529    }
530
531    /// Returns the DDSketch summary for successful spans.
532    pub fn ok_summary(&self) -> &[u8] {
533        &self.ok_summary
534    }
535
536    /// Returns the DDSketch summary for error spans.
537    pub fn error_summary(&self) -> &[u8] {
538        &self.error_summary
539    }
540
541    /// Returns the synthetics traffic flag.
542    pub fn synthetics(&self) -> bool {
543        self.synthetics
544    }
545
546    /// Returns the top-level hit count.
547    pub fn top_level_hits(&self) -> u64 {
548        self.top_level_hits
549    }
550}