Skip to main content

relay_server/services/buffer/envelope_stack/
sqlite.rs

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/// An error returned when doing an operation on [`SqliteEnvelopeStack`].
19#[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)]
28/// An [`EnvelopeStack`] that is implemented on an SQLite database.
29///
30/// For efficiency reasons, the implementation has an in-memory buffer that is periodically spooled
31/// to disk in a batched way.
32pub struct SqliteEnvelopeStack {
33    /// Shared SQLite database pool which will be used to read and write from disk.
34    envelope_store: SqliteEnvelopeStore,
35    /// Maximum number of bytes in the in-memory cache before we write to disk.
36    batch_size_bytes: NonZeroUsize,
37    /// The project key of the project to which all the envelopes belong.
38    own_key: ProjectKey,
39    /// The project key of the root project of the trace to which all the envelopes belong.
40    sampling_key: ProjectKey,
41    /// In-memory stack containing a batch of envelopes that either have not been written to disk yet, or have been read from disk recently.
42    batch: Vec<DatabaseEnvelope>,
43    /// Boolean representing whether calls to `push()` and `peek()` check disk in case not enough
44    /// elements are available in the `batches_buffer`.
45    check_disk: bool,
46    /// The tag value of this partition which is used for reporting purposes.
47    partition_tag: String,
48    /// Time after which to flush the buffer to disk, regardless of batch size.
49    flush_timeout: Option<Duration>,
50    /// Time of last flush to disk (or creation time of the envelope stack).
51    last_flush: Instant,
52    /// Optional throttle which paces how quickly envelopes are unspooled from disk.
53    unspool_throttle: Option<Arc<Throttle>>,
54}
55
56/// Configuration for a [`SqliteEnvelopeStack`].
57#[derive(Debug, Clone)]
58pub struct SqliteEnvelopeStackConfig {
59    /// The partition of the envelope buffer to which the stack belongs.
60    pub partition_id: u8,
61    /// Maximum number of bytes in the in-memory cache before writing to disk.
62    pub batch_size_bytes: usize,
63    /// Time after which to flush the buffer to disk.
64    pub flush_timeout: Option<Duration>,
65    /// Optional throttle which paces how quickly envelopes are unspooled from disk.
66    pub unspool_throttle: Option<Arc<Throttle>>,
67}
68
69impl SqliteEnvelopeStack {
70    /// Creates a new empty [`SqliteEnvelopeStack`].
71    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    /// Threshold above which the [`SqliteEnvelopeStack`] will spool data from the `buffer` to disk.
94    ///
95    /// Returns the reason why the batch should be spooled.
96    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    /// Spools to disk a batch of envelopes from the `batch`.
112    ///
113    /// In case there is a failure while writing envelopes, all the envelopes that were enqueued
114    /// to be written to disk are lost. The explanation for this behavior can be found in the body
115    /// of the method.
116    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        // When early return here, we are acknowledging that the elements that we popped from
123        // the buffer are lost in case of failure. We are doing this on purpose, since if we were
124        // to have a database corruption during runtime, and we were to put the values back into
125        // the buffer we will end up with an infinite cycle.
126        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        // If we successfully spooled to disk, we know that data should be there.
139        self.check_disk = true;
140        self.last_flush = Instant::now();
141
142        Ok(())
143    }
144
145    /// Unspools from disk a batch of envelopes and appends them to the `batch`.
146    ///
147    /// In case there is a failure while deleting envelopes, the envelopes will be lost.
148    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    /// Validates that the incoming [`Envelope`] has the same project keys at the
175    /// [`SqliteEnvelopeStack`].
176    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    /// Helper function to calculate the total size of a slice of envelopes after compression
246    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        // Create envelopes first so we can calculate actual size
283        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        // We push the 4 envelopes without errors because they are below the threshold.
300        for envelope in envelopes.clone() {
301            assert!(stack.push(envelope).await.is_ok());
302        }
303
304        // We push 1 more envelope which results in spooling, which fails because of a database
305        // problem.
306        let envelope = mock_envelope(Utc::now());
307        assert!(matches!(
308            stack.push(envelope).await,
309            Err(SqliteEnvelopeStackError::EnvelopeStoreError(_))
310        ));
311
312        // The stack now contains the last of the 1 elements that were added. If we add a new one
313        // we will end up with 2.
314        let envelope = mock_envelope(Utc::now());
315        assert!(stack.push(envelope.clone()).await.is_ok());
316        assert_eq!(stack.batch.len(), 1);
317
318        // We pop the remaining elements, expecting the last added envelope to be on top.
319        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        // We pop with an invalid db.
345        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        // We pop with no elements.
369        // We pop with no elements.
370        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        // We push 5 envelopes.
393        for envelope in envelopes.clone() {
394            assert!(stack.push(envelope).await.is_ok());
395        }
396        assert_eq!(stack.batch.len(), 5);
397
398        // We peek the top element.
399        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        // We pop 5 envelopes.
406        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        // First push: no spool
437        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        // Second push: no spool
443        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        // Third push (after timeout): spool
448        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        // Fourth push: no spool
454        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        // Envelopes from memory and disk are still returned in stack order.
459        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        // Create envelopes first so we can calculate actual size
472        let envelopes = mock_envelopes(7);
473        let threshold_size = calculate_compressed_size(&envelopes[..5]) - 1;
474
475        // Create stack with threshold just below the size of first 5 envelopes
476        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        // We push 7 envelopes.
490        for envelope in envelopes.clone() {
491            assert!(stack.push(envelope).await.is_ok());
492        }
493        assert_eq!(stack.batch.len(), 2);
494
495        // We peek the top element.
496        let peeked = stack.peek().await.unwrap().unwrap();
497        assert_eq!(
498            peeked.timestamp_millis(),
499            envelopes[6].received_at().timestamp_millis()
500        );
501
502        // We pop envelopes, and we expect that the last 2 are in memory, since the first 5
503        // should have been spooled to disk.
504        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        // We peek the top element, which since the buffer is empty should result in a disk load.
514        let peeked = stack.peek().await.unwrap().unwrap();
515        assert_eq!(
516            peeked.timestamp_millis(),
517            envelopes[4].received_at().timestamp_millis()
518        );
519
520        // We insert a new envelope, to test the load from disk happening during `peek()` gives
521        // priority to this envelope in the stack.
522        let envelope = mock_envelope(Utc::now());
523        assert!(stack.push(envelope.clone()).await.is_ok());
524
525        // We pop and expect the newly inserted element.
526        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        // We pop 5 envelopes, which should not result in a disk load since `peek()` already should
533        // have caused it.
534        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        // We push 5 envelopes and check that there is nothing on disk.
564        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        // We drain the stack and make sure everything was spooled to disk.
571        stack.flush().await;
572        assert_eq!(envelope_store.total_count().await.unwrap(), 5);
573    }
574}