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 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 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 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 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 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 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 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 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}