Skip to main content

vector_buffers/
internal_events.rs

1use std::time::Duration;
2
3use metrics::Histogram;
4use vector_common::NamedInternalEvent;
5use vector_common::{
6    counter, gauge, histogram,
7    internal_event::{CounterName, GaugeName, HistogramName, InternalEvent, error_type},
8    registered_event,
9};
10
11#[derive(NamedInternalEvent)]
12pub struct BufferCreated {
13    pub buffer_id: String,
14    pub idx: usize,
15    pub max_size_events: usize,
16    pub max_size_bytes: u64,
17}
18
19impl InternalEvent for BufferCreated {
20    #[expect(clippy::cast_precision_loss)]
21    fn emit(self) {
22        let stage = self.idx.to_string();
23        if self.max_size_events != 0 {
24            gauge!(
25                GaugeName::BufferMaxSizeEvents,
26                "buffer_id" => self.buffer_id.clone(),
27                "stage" => stage.clone(),
28            )
29            .set(self.max_size_events as f64);
30            // DEPRECATED: buffer-bytes-events-metrics
31            gauge!(
32                GaugeName::BufferMaxEventSize,
33                "buffer_id" => self.buffer_id.clone(),
34                "stage" => stage.clone(),
35            )
36            .set(self.max_size_events as f64);
37        }
38        if self.max_size_bytes != 0 {
39            gauge!(
40                GaugeName::BufferMaxSizeBytes,
41                "buffer_id" => self.buffer_id.clone(),
42                "stage" => stage.clone(),
43            )
44            .set(self.max_size_bytes as f64);
45            // DEPRECATED: buffer-bytes-events-metrics
46            gauge!(
47                GaugeName::BufferMaxByteSize,
48                "buffer_id" => self.buffer_id,
49                "stage" => stage,
50            )
51            .set(self.max_size_bytes as f64);
52        }
53    }
54}
55
56#[derive(NamedInternalEvent)]
57pub struct BufferEventsReceived {
58    pub buffer_id: String,
59    pub idx: usize,
60    pub count: u64,
61    pub byte_size: u64,
62    pub total_count: u64,
63    pub total_byte_size: u64,
64}
65
66impl InternalEvent for BufferEventsReceived {
67    #[expect(clippy::cast_precision_loss)]
68    fn emit(self) {
69        counter!(
70            CounterName::BufferReceivedEventsTotal,
71            "buffer_id" => self.buffer_id.clone(),
72            "stage" => self.idx.to_string()
73        )
74        .increment(self.count);
75
76        counter!(
77            CounterName::BufferReceivedBytesTotal,
78            "buffer_id" => self.buffer_id.clone(),
79            "stage" => self.idx.to_string()
80        )
81        .increment(self.byte_size);
82        // DEPRECATED: buffer-bytes-events-metrics
83        gauge!(
84            GaugeName::BufferEvents,
85            "buffer_id" => self.buffer_id.clone(),
86            "stage" => self.idx.to_string()
87        )
88        .set(self.total_count as f64);
89        gauge!(
90            GaugeName::BufferSizeEvents,
91            "buffer_id" => self.buffer_id.clone(),
92            "stage" => self.idx.to_string()
93        )
94        .set(self.total_count as f64);
95        gauge!(
96            GaugeName::BufferSizeBytes,
97            "buffer_id" => self.buffer_id.clone(),
98            "stage" => self.idx.to_string()
99        )
100        .set(self.total_byte_size as f64);
101        // DEPRECATED: buffer-bytes-events-metrics
102        gauge!(
103            GaugeName::BufferByteSize,
104            "buffer_id" => self.buffer_id,
105            "stage" => self.idx.to_string()
106        )
107        .set(self.total_byte_size as f64);
108    }
109}
110
111#[derive(NamedInternalEvent)]
112pub struct BufferEventsSent {
113    pub buffer_id: String,
114    pub idx: usize,
115    pub count: u64,
116    pub byte_size: u64,
117    pub total_count: u64,
118    pub total_byte_size: u64,
119}
120
121impl InternalEvent for BufferEventsSent {
122    #[expect(clippy::cast_precision_loss)]
123    fn emit(self) {
124        counter!(
125            CounterName::BufferSentEventsTotal,
126            "buffer_id" => self.buffer_id.clone(),
127            "stage" => self.idx.to_string()
128        )
129        .increment(self.count);
130        counter!(
131            CounterName::BufferSentBytesTotal,
132            "buffer_id" => self.buffer_id.clone(),
133            "stage" => self.idx.to_string()
134        )
135        .increment(self.byte_size);
136        // DEPRECATED: buffer-bytes-events-metrics
137        gauge!(
138            GaugeName::BufferEvents,
139            "buffer_id" => self.buffer_id.clone(),
140            "stage" => self.idx.to_string()
141        )
142        .set(self.total_count as f64);
143        gauge!(
144            GaugeName::BufferSizeEvents,
145            "buffer_id" => self.buffer_id.clone(),
146            "stage" => self.idx.to_string()
147        )
148        .set(self.total_count as f64);
149        gauge!(
150            GaugeName::BufferSizeBytes,
151            "buffer_id" => self.buffer_id.clone(),
152            "stage" => self.idx.to_string()
153        )
154        .set(self.total_byte_size as f64);
155        // DEPRECATED: buffer-bytes-events-metrics
156        gauge!(
157            GaugeName::BufferByteSize,
158            "buffer_id" => self.buffer_id,
159            "stage" => self.idx.to_string()
160        )
161        .set(self.total_byte_size as f64);
162    }
163}
164
165#[derive(NamedInternalEvent)]
166pub struct BufferEventsDropped {
167    pub buffer_id: String,
168    pub idx: usize,
169    pub count: u64,
170    pub byte_size: u64,
171    pub total_count: u64,
172    pub total_byte_size: u64,
173    pub intentional: bool,
174    pub reason: &'static str,
175}
176
177impl InternalEvent for BufferEventsDropped {
178    #[expect(clippy::cast_precision_loss)]
179    fn emit(self) {
180        let intentional_str = if self.intentional { "true" } else { "false" };
181        if self.intentional {
182            debug!(
183                message = "Events dropped.",
184                count = %self.count,
185                byte_size = %self.byte_size,
186                intentional = %intentional_str,
187                reason = %self.reason,
188                buffer_id = %self.buffer_id,
189                stage = %self.idx,
190            );
191        } else {
192            error!(
193                message = "Events dropped.",
194                count = %self.count,
195                byte_size = %self.byte_size,
196                intentional = %intentional_str,
197                reason = %self.reason,
198                buffer_id = %self.buffer_id,
199                stage = %self.idx,
200            );
201        }
202
203        counter!(
204            CounterName::BufferDiscardedEventsTotal,
205            "buffer_id" => self.buffer_id.clone(),
206            "stage" => self.idx.to_string(),
207            "intentional" => intentional_str,
208        )
209        .increment(self.count);
210        counter!(
211            CounterName::BufferDiscardedBytesTotal,
212            "buffer_id" => self.buffer_id.clone(),
213            "stage" => self.idx.to_string(),
214            "intentional" => intentional_str,
215        )
216        .increment(self.byte_size);
217        // DEPRECATED: buffer-bytes-events-metrics
218        gauge!(
219            GaugeName::BufferEvents,
220            "buffer_id" => self.buffer_id.clone(),
221            "stage" => self.idx.to_string()
222        )
223        .set(self.total_count as f64);
224        gauge!(
225            GaugeName::BufferSizeEvents,
226            "buffer_id" => self.buffer_id.clone(),
227            "stage" => self.idx.to_string()
228        )
229        .set(self.total_count as f64);
230        gauge!(
231            GaugeName::BufferSizeBytes,
232            "buffer_id" => self.buffer_id.clone(),
233            "stage" => self.idx.to_string()
234        )
235        .set(self.total_byte_size as f64);
236        // DEPRECATED: buffer-bytes-events-metrics
237        gauge!(
238            GaugeName::BufferByteSize,
239            "buffer_id" => self.buffer_id,
240            "stage" => self.idx.to_string()
241        )
242        .set(self.total_byte_size as f64);
243    }
244}
245
246#[derive(NamedInternalEvent)]
247pub struct BufferReadError {
248    pub error_code: &'static str,
249    pub error: String,
250}
251
252impl InternalEvent for BufferReadError {
253    fn emit(self) {
254        error!(
255            message = "Error encountered during buffer read.",
256            error = %self.error,
257            error_code = self.error_code,
258            error_type = error_type::READER_FAILED,
259            stage = "processing",
260        );
261        counter!(
262            CounterName::BufferErrorsTotal, "error_code" => self.error_code,
263            "error_type" => "reader_failed",
264            "stage" => "processing",
265        )
266        .increment(1);
267    }
268}
269
270registered_event! {
271    BufferSendDuration {
272        stage: usize,
273    } => {
274        send_duration: Histogram = histogram!(HistogramName::BufferSendDurationSeconds, "stage" => self.stage.to_string()),
275    }
276
277    fn emit(&self, duration: Duration) {
278        self.send_duration.record(duration);
279    }
280}