1use std::fmt::Debug;
2use std::num::NonZeroUsize;
3use std::sync::Arc;
4use std::time::{Duration, Instant};
5
6use chrono::{DateTime, Utc};
7use relay_base_schema::project::ProjectKey;
8
9use crate::envelope::Envelope;
10use crate::services::buffer::envelope_stack::EnvelopeStack;
11use crate::services::buffer::envelope_store::sqlite::{
12 DatabaseBatch, DatabaseEnvelope, InsertEnvelopeError, SqliteEnvelopeStore,
13 SqliteEnvelopeStoreError,
14};
15use crate::services::buffer::throttle::Throttle;
16use crate::statsd::RelayTimers;
17
18#[derive(Debug, thiserror::Error)]
20pub enum SqliteEnvelopeStackError {
21 #[error("envelope store error: {0}")]
22 EnvelopeStoreError(#[from] SqliteEnvelopeStoreError),
23 #[error("envelope encode error: {0}")]
24 Envelope(#[from] InsertEnvelopeError),
25}
26
27#[derive(Debug)]
28pub struct SqliteEnvelopeStack {
33 envelope_store: SqliteEnvelopeStore,
35 batch_size_bytes: NonZeroUsize,
37 own_key: ProjectKey,
39 sampling_key: ProjectKey,
41 batch: Vec<DatabaseEnvelope>,
43 check_disk: bool,
46 partition_tag: String,
48 flush_timeout: Option<Duration>,
50 last_flush: Instant,
52 unspool_throttle: Option<Arc<Throttle>>,
54}
55
56#[derive(Debug, Clone)]
58pub struct SqliteEnvelopeStackConfig {
59 pub partition_id: u8,
61 pub batch_size_bytes: usize,
63 pub flush_timeout: Option<Duration>,
65 pub unspool_throttle: Option<Arc<Throttle>>,
67}
68
69impl SqliteEnvelopeStack {
70 pub fn new(
72 config: SqliteEnvelopeStackConfig,
73 envelope_store: SqliteEnvelopeStore,
74 own_key: ProjectKey,
75 sampling_key: ProjectKey,
76 check_disk: bool,
77 ) -> Self {
78 Self {
79 envelope_store,
80 batch_size_bytes: NonZeroUsize::new(config.batch_size_bytes)
81 .expect("batch bytes should be > 0"),
82 own_key,
83 sampling_key,
84 batch: vec![],
85 check_disk,
86 partition_tag: config.partition_id.to_string(),
87 flush_timeout: config.flush_timeout,
88 last_flush: Instant::now(),
89 unspool_throttle: config.unspool_throttle,
90 }
91 }
92
93 fn should_spool_to_disk(&self) -> Option<&'static str> {
97 let batch_size = self.batch.iter().map(|e| e.len()).sum::<usize>();
98 if batch_size > self.batch_size_bytes.get() {
99 return Some("size");
100 }
101
102 if let Some(timeout) = self.flush_timeout
103 && self.last_flush.elapsed() > timeout
104 {
105 return Some("timeout");
106 }
107
108 None
109 }
110
111 async fn spool_to_disk(&mut self, why: &'static str) -> Result<(), SqliteEnvelopeStackError> {
117 let batch = std::mem::take(&mut self.batch);
118 let Ok(batch) = DatabaseBatch::try_from(batch) else {
119 return Ok(());
120 };
121
122 relay_statsd::metric!(
127 timer(RelayTimers::BufferSpool),
128 partition_id = &self.partition_tag,
129 reason = why,
130 {
131 self.envelope_store
132 .insert_batch(batch)
133 .await
134 .map_err(SqliteEnvelopeStackError::EnvelopeStoreError)?;
135 }
136 );
137
138 self.check_disk = true;
140 self.last_flush = Instant::now();
141
142 Ok(())
143 }
144
145 async fn unspool_from_disk(&mut self) -> Result<(), SqliteEnvelopeStackError> {
149 debug_assert!(self.batch.is_empty());
150 let batch = relay_statsd::metric!(
151 timer(RelayTimers::BufferUnspool),
152 partition_id = &self.partition_tag,
153 {
154 self.envelope_store
155 .delete_batch(self.own_key, self.sampling_key)
156 .await
157 .map_err(SqliteEnvelopeStackError::EnvelopeStoreError)?
158 }
159 );
160
161 match batch {
162 Some(batch) => {
163 self.batch = batch.into();
164 if let Some(throttle) = &self.unspool_throttle {
165 throttle.acquire(self.batch.len()).await;
166 }
167 }
168 None => self.check_disk = false,
169 }
170
171 Ok(())
172 }
173
174 fn validate_envelope(&self, envelope: &Envelope) -> bool {
177 let own_key = envelope.meta().public_key();
178 let sampling_key = envelope.sampling_key().unwrap_or(own_key);
179
180 self.own_key == own_key && self.sampling_key == sampling_key
181 }
182}
183
184impl EnvelopeStack for SqliteEnvelopeStack {
185 type Error = SqliteEnvelopeStackError;
186
187 async fn push(&mut self, envelope: Box<Envelope>) -> Result<(), Self::Error> {
188 debug_assert!(self.validate_envelope(&envelope));
189
190 if let Some(reason) = self.should_spool_to_disk() {
191 self.spool_to_disk(reason).await?;
192 }
193
194 let encoded_envelope = relay_statsd::metric!(
195 timer(RelayTimers::BufferEnvelopesSerialization),
196 partition_id = &self.partition_tag,
197 { DatabaseEnvelope::try_from(envelope.as_ref())? }
198 );
199 self.batch.push(encoded_envelope);
200
201 Ok(())
202 }
203
204 async fn peek(&mut self) -> Result<Option<DateTime<Utc>>, Self::Error> {
205 if self.batch.is_empty() && self.check_disk {
206 self.unspool_from_disk().await?
207 }
208
209 let Some(envelope) = self.batch.last() else {
210 return Ok(None);
211 };
212
213 Ok(Some(envelope.received_at()))
214 }
215
216 async fn pop(&mut self) -> Result<Option<Box<Envelope>>, Self::Error> {
217 if self.batch.is_empty() && self.check_disk {
218 self.unspool_from_disk().await?
219 }
220
221 let Some(envelope) = self.batch.pop() else {
222 return Ok(None);
223 };
224 let envelope = envelope.try_into()?;
225
226 Ok(Some(envelope))
227 }
228
229 async fn flush(mut self) {
230 if let Err(e) = self.spool_to_disk("flush").await {
231 relay_log::error!(error = &e as &dyn std::error::Error, "flush error");
232 }
233 }
234}
235
236#[cfg(test)]
237mod tests {
238 use chrono::Utc;
239 use relay_base_schema::project::ProjectKey;
240 use std::time::Duration;
241
242 use super::*;
243 use crate::services::buffer::testutils::utils::{mock_envelope, mock_envelopes, setup_db};
244
245 fn calculate_compressed_size(envelopes: &[Box<Envelope>]) -> usize {
247 envelopes
248 .iter()
249 .map(|e| DatabaseEnvelope::try_from(e.as_ref()).unwrap().len())
250 .sum()
251 }
252
253 #[tokio::test]
254 #[should_panic]
255 async fn test_push_with_mismatching_project_keys() {
256 let db = setup_db(false).await;
257 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
258 let mut stack = SqliteEnvelopeStack::new(
259 SqliteEnvelopeStackConfig {
260 partition_id: 0,
261 batch_size_bytes: 10,
262 flush_timeout: None,
263 unspool_throttle: None,
264 },
265 envelope_store,
266 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
267 ProjectKey::parse("c25ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
268 true,
269 );
270
271 let envelope = mock_envelope(Utc::now());
272 let _ = stack.push(envelope).await;
273 }
274
275 const COMPRESSED_ENVELOPE_SIZE: usize = 313;
276
277 #[tokio::test]
278 async fn test_push_when_db_is_not_valid() {
279 let db = setup_db(false).await;
280 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
281
282 let envelopes = mock_envelopes(4);
284 let threshold_size = calculate_compressed_size(&envelopes) - 1;
285
286 let mut stack = SqliteEnvelopeStack::new(
287 SqliteEnvelopeStackConfig {
288 partition_id: 0,
289 batch_size_bytes: threshold_size,
290 flush_timeout: None,
291 unspool_throttle: None,
292 },
293 envelope_store,
294 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
295 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
296 true,
297 );
298
299 for envelope in envelopes.clone() {
301 assert!(stack.push(envelope).await.is_ok());
302 }
303
304 let envelope = mock_envelope(Utc::now());
307 assert!(matches!(
308 stack.push(envelope).await,
309 Err(SqliteEnvelopeStackError::EnvelopeStoreError(_))
310 ));
311
312 let envelope = mock_envelope(Utc::now());
315 assert!(stack.push(envelope.clone()).await.is_ok());
316 assert_eq!(stack.batch.len(), 1);
317
318 let popped_envelope_1 = stack.pop().await.unwrap().unwrap();
320 assert_eq!(
321 popped_envelope_1.event_id().unwrap(),
322 envelope.event_id().unwrap()
323 );
324 assert_eq!(stack.batch.len(), 0);
325 }
326
327 #[tokio::test]
328 async fn test_pop_when_db_is_not_valid() {
329 let db = setup_db(false).await;
330 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
331 let mut stack = SqliteEnvelopeStack::new(
332 SqliteEnvelopeStackConfig {
333 partition_id: 0,
334 batch_size_bytes: 2,
335 flush_timeout: None,
336 unspool_throttle: None,
337 },
338 envelope_store,
339 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
340 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
341 true,
342 );
343
344 assert!(matches!(
346 stack.pop().await,
347 Err(SqliteEnvelopeStackError::EnvelopeStoreError(_))
348 ));
349 }
350
351 #[tokio::test]
352 async fn test_pop_when_stack_is_empty() {
353 let db = setup_db(true).await;
354 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
355 let mut stack = SqliteEnvelopeStack::new(
356 SqliteEnvelopeStackConfig {
357 partition_id: 0,
358 batch_size_bytes: 2,
359 flush_timeout: None,
360 unspool_throttle: None,
361 },
362 envelope_store,
363 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
364 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
365 true,
366 );
367
368 assert!(stack.pop().await.unwrap().is_none());
371 }
372
373 #[tokio::test]
374 async fn test_push_below_threshold_and_pop() {
375 let db = setup_db(true).await;
376 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
377 let mut stack = SqliteEnvelopeStack::new(
378 SqliteEnvelopeStackConfig {
379 partition_id: 0,
380 batch_size_bytes: 9999,
381 flush_timeout: None,
382 unspool_throttle: None,
383 },
384 envelope_store,
385 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
386 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
387 true,
388 );
389
390 let envelopes = mock_envelopes(5);
391
392 for envelope in envelopes.clone() {
394 assert!(stack.push(envelope).await.is_ok());
395 }
396 assert_eq!(stack.batch.len(), 5);
397
398 let peeked = stack.peek().await.unwrap().unwrap();
400 assert_eq!(
401 peeked.timestamp_millis(),
402 envelopes.clone()[4].received_at().timestamp_millis()
403 );
404
405 for envelope in envelopes.iter().rev() {
407 let popped_envelope = stack.pop().await.unwrap().unwrap();
408 assert_eq!(
409 popped_envelope.event_id().unwrap(),
410 envelope.event_id().unwrap()
411 );
412 }
413
414 assert_eq!(stack.batch.len(), 0);
415 }
416
417 #[tokio::test]
418 async fn test_push_with_flush_timeout() {
419 let db = setup_db(true).await;
420 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
421 let envelopes = mock_envelopes(4);
422 let timeout = Duration::from_secs(3600);
423 let mut stack = SqliteEnvelopeStack::new(
424 SqliteEnvelopeStackConfig {
425 partition_id: 0,
426 batch_size_bytes: calculate_compressed_size(&envelopes),
427 flush_timeout: Some(timeout),
428 unspool_throttle: None,
429 },
430 envelope_store.clone(),
431 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
432 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
433 false,
434 );
435
436 assert_eq!(stack.batch.len(), 0);
438 stack.push(envelopes[0].clone()).await.unwrap();
439 assert_eq!(stack.batch.len(), 1);
440 assert_eq!(envelope_store.total_count().await.unwrap(), 0);
441
442 stack.push(envelopes[1].clone()).await.unwrap();
444 assert_eq!(stack.batch.len(), 2);
445 assert_eq!(envelope_store.total_count().await.unwrap(), 0);
446
447 stack.last_flush = Instant::now() - timeout - Duration::from_secs(1);
449 stack.push(envelopes[2].clone()).await.unwrap();
450 assert_eq!(stack.batch.len(), 1);
451 assert_eq!(envelope_store.total_count().await.unwrap(), 2);
452
453 stack.push(envelopes[3].clone()).await.unwrap();
455 assert_eq!(stack.batch.len(), 2);
456 assert_eq!(envelope_store.total_count().await.unwrap(), 2);
457
458 for envelope in envelopes.iter().rev() {
460 let popped = stack.pop().await.unwrap().unwrap();
461 assert_eq!(popped.event_id().unwrap(), envelope.event_id().unwrap());
462 }
463 assert!(stack.pop().await.unwrap().is_none());
464 }
465
466 #[tokio::test]
467 async fn test_push_above_threshold_and_pop() {
468 let db = setup_db(true).await;
469 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
470
471 let envelopes = mock_envelopes(7);
473 let threshold_size = calculate_compressed_size(&envelopes[..5]) - 1;
474
475 let mut stack = SqliteEnvelopeStack::new(
477 SqliteEnvelopeStackConfig {
478 partition_id: 0,
479 batch_size_bytes: threshold_size,
480 flush_timeout: None,
481 unspool_throttle: None,
482 },
483 envelope_store,
484 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
485 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
486 true,
487 );
488
489 for envelope in envelopes.clone() {
491 assert!(stack.push(envelope).await.is_ok());
492 }
493 assert_eq!(stack.batch.len(), 2);
494
495 let peeked = stack.peek().await.unwrap().unwrap();
497 assert_eq!(
498 peeked.timestamp_millis(),
499 envelopes[6].received_at().timestamp_millis()
500 );
501
502 for envelope in envelopes[5..7].iter().rev() {
505 let popped_envelope = stack.pop().await.unwrap().unwrap();
506 assert_eq!(
507 popped_envelope.event_id().unwrap(),
508 envelope.event_id().unwrap()
509 );
510 }
511 assert_eq!(stack.batch.len(), 0);
512
513 let peeked = stack.peek().await.unwrap().unwrap();
515 assert_eq!(
516 peeked.timestamp_millis(),
517 envelopes[4].received_at().timestamp_millis()
518 );
519
520 let envelope = mock_envelope(Utc::now());
523 assert!(stack.push(envelope.clone()).await.is_ok());
524
525 let popped_envelope = stack.pop().await.unwrap().unwrap();
527 assert_eq!(
528 popped_envelope.event_id().unwrap(),
529 envelope.event_id().unwrap()
530 );
531
532 for envelope in envelopes[0..5].iter().rev() {
535 let popped_envelope = stack.pop().await.unwrap().unwrap();
536 assert_eq!(
537 popped_envelope.event_id().unwrap(),
538 envelope.event_id().unwrap()
539 );
540 }
541 assert_eq!(stack.batch.len(), 0);
542 }
543
544 #[tokio::test]
545 async fn test_drain() {
546 let db = setup_db(true).await;
547 let envelope_store = SqliteEnvelopeStore::new(0, db, Duration::from_millis(100));
548 let mut stack = SqliteEnvelopeStack::new(
549 SqliteEnvelopeStackConfig {
550 partition_id: 0,
551 batch_size_bytes: 10 * COMPRESSED_ENVELOPE_SIZE,
552 flush_timeout: None,
553 unspool_throttle: None,
554 },
555 envelope_store.clone(),
556 ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap(),
557 ProjectKey::parse("b81ae32be2584e0bbd7a4cbb95971fe1").unwrap(),
558 true,
559 );
560
561 let envelopes = mock_envelopes(5);
562
563 for envelope in envelopes.clone() {
565 assert!(stack.push(envelope).await.is_ok());
566 }
567 assert_eq!(stack.batch.len(), 5);
568 assert_eq!(envelope_store.total_count().await.unwrap(), 0);
569
570 stack.flush().await;
572 assert_eq!(envelope_store.total_count().await.unwrap(), 5);
573 }
574}