Skip to main content

relay_server/services/buffer/envelope_buffer/
mod.rs

1use std::cmp::Ordering;
2use std::collections::BTreeSet;
3use std::convert::Infallible;
4use std::error::Error;
5use std::mem;
6use std::sync::Arc;
7use std::time::Duration;
8
9use chrono::{DateTime, Utc};
10use hashbrown::HashSet;
11use relay_base_schema::project::ProjectKey;
12use relay_config::ConfigSnapshot;
13use tokio::time::{Instant, timeout};
14
15use crate::envelope::Envelope;
16use crate::envelope::Item;
17use crate::services::buffer::common::ProjectKeyPair;
18use crate::services::buffer::envelope_stack::EnvelopeStack;
19use crate::services::buffer::envelope_stack::sqlite::SqliteEnvelopeStackError;
20use crate::services::buffer::envelope_store::sqlite::SqliteEnvelopeStoreError;
21use crate::services::buffer::stack_provider::memory::MemoryStackProvider;
22use crate::services::buffer::stack_provider::sqlite::SqliteStackProvider;
23use crate::services::buffer::stack_provider::{StackCreationType, StackProvider};
24use crate::services::buffer::throttle::Throttle;
25use crate::statsd::{RelayDistributions, RelayGauges, RelayTimers};
26use crate::utils::MemoryChecker;
27
28/// Polymorphic envelope buffering interface.
29///
30/// The underlying buffer can either be disk-based or memory-based,
31/// depending on the given configuration.
32///
33/// NOTE: This is implemented as an enum because a trait object with async methods would not be
34/// object safe.
35#[derive(Debug)]
36#[allow(private_interfaces)]
37pub enum PolymorphicEnvelopeBuffer {
38    /// An enveloper buffer that uses in-memory envelopes stacks.
39    InMemory(EnvelopeBuffer<MemoryStackProvider>),
40    /// An enveloper buffer that uses sqlite envelopes stacks.
41    Sqlite(EnvelopeBuffer<SqliteStackProvider>),
42}
43
44impl PolymorphicEnvelopeBuffer {
45    /// Returns true if the implementation stores all envelopes in RAM.
46    pub fn is_memory(&self) -> bool {
47        match self {
48            Self::InMemory(_) => true,
49            Self::Sqlite(_) => false,
50        }
51    }
52
53    /// Creates either a memory-based or a disk-based envelope buffer,
54    /// depending on the given configuration.
55    pub async fn from_config(
56        partition_id: u8,
57        config: &ConfigSnapshot,
58        memory_checker: MemoryChecker,
59        throttle: Option<Arc<Throttle>>,
60    ) -> Result<Self, EnvelopeBufferError> {
61        let buffer = if config.spool_envelopes_path(partition_id).is_some() {
62            relay_log::trace!("PolymorphicEnvelopeBuffer: initializing sqlite envelope buffer");
63            let buffer =
64                EnvelopeBuffer::<SqliteStackProvider>::new(partition_id, config, throttle).await?;
65            Self::Sqlite(buffer)
66        } else {
67            relay_log::trace!("PolymorphicEnvelopeBuffer: initializing memory envelope buffer");
68            let buffer = EnvelopeBuffer::<MemoryStackProvider>::new(partition_id, memory_checker);
69            Self::InMemory(buffer)
70        };
71
72        Ok(buffer)
73    }
74
75    /// Initializes the envelope buffer.
76    pub async fn initialize(&mut self) {
77        match self {
78            PolymorphicEnvelopeBuffer::InMemory(buffer) => buffer.initialize().await,
79            PolymorphicEnvelopeBuffer::Sqlite(buffer) => buffer.initialize().await,
80        }
81    }
82
83    /// Adds an envelope to the buffer.
84    pub async fn push(&mut self, envelope: Box<Envelope>) -> Result<(), EnvelopeBufferError> {
85        relay_statsd::metric!(
86            distribution(RelayDistributions::BufferEnvelopeBodySize) =
87                envelope.items().map(Item::len).sum::<usize>() as u64,
88            partition_id = self.partition_tag()
89        );
90
91        relay_statsd::metric!(
92            timer(RelayTimers::BufferPush),
93            partition_id = self.partition_tag(),
94            {
95                match self {
96                    Self::Sqlite(buffer) => buffer.push(envelope).await,
97                    Self::InMemory(buffer) => buffer.push(envelope).await,
98                }
99            }
100        )
101    }
102
103    /// Returns a reference to the next-in-line envelope.
104    pub async fn peek(&mut self) -> Result<Peek, EnvelopeBufferError> {
105        relay_statsd::metric!(
106            timer(RelayTimers::BufferPeek),
107            partition_id = self.partition_tag(),
108            {
109                match self {
110                    Self::Sqlite(buffer) => buffer.peek().await,
111                    Self::InMemory(buffer) => buffer.peek().await,
112                }
113            }
114        )
115    }
116
117    /// Pops the next-in-line envelope.
118    pub async fn pop(&mut self) -> Result<Option<Box<Envelope>>, EnvelopeBufferError> {
119        relay_statsd::metric!(
120            timer(RelayTimers::BufferPop),
121            partition_id = self.partition_tag(),
122            {
123                match self {
124                    Self::Sqlite(buffer) => buffer.pop().await,
125                    Self::InMemory(buffer) => buffer.pop().await,
126                }
127            }
128        )
129    }
130
131    /// Marks a project as ready or not ready.
132    ///
133    /// The buffer re-prioritizes its envelopes based on this information.
134    /// Returns `true` if at least one priority was changed.
135    pub fn mark_ready(&mut self, project: &ProjectKey, is_ready: bool) -> bool {
136        relay_log::trace!(
137            project_key = project.as_str(),
138            "buffer marked {}",
139            if is_ready { "ready" } else { "not ready" }
140        );
141        match self {
142            Self::Sqlite(buffer) => buffer.mark_ready(project, is_ready),
143            Self::InMemory(buffer) => buffer.mark_ready(project, is_ready),
144        }
145    }
146
147    /// Marks a stack as seen.
148    ///
149    /// Non-ready stacks are deprioritized when they are marked as seen, such that
150    /// the next call to `.peek()` will look at a different stack. This prevents
151    /// head-of-line blocking.
152    pub fn mark_seen(&mut self, project_key_pair: &ProjectKeyPair, next_fetch: Duration) {
153        match self {
154            Self::Sqlite(buffer) => buffer.mark_seen(project_key_pair, next_fetch),
155            Self::InMemory(buffer) => buffer.mark_seen(project_key_pair, next_fetch),
156        }
157    }
158
159    /// Returns `true` whether the buffer has capacity to accept new [`Envelope`]s.
160    pub fn has_capacity(&self) -> bool {
161        match self {
162            Self::Sqlite(buffer) => buffer.has_capacity(),
163            Self::InMemory(buffer) => buffer.has_capacity(),
164        }
165    }
166
167    /// Returns `true` if the buffer contains no envelopes, in memory or on disk.
168    pub fn is_empty(&self) -> bool {
169        match self {
170            Self::Sqlite(buffer) => buffer.is_empty(),
171            Self::InMemory(buffer) => buffer.is_empty(),
172        }
173    }
174
175    /// Returns the total number of envelopes that have been spooled since the startup. It does
176    /// not include the count that existed in a persistent spooler before.
177    pub fn item_count(&self) -> u64 {
178        match self {
179            Self::Sqlite(buffer) => buffer.tracked_count,
180            Self::InMemory(buffer) => buffer.tracked_count,
181        }
182    }
183
184    /// Returns the total number of bytes that the spooler storage uses or `None` if the number
185    /// cannot be reliably determined.
186    pub fn total_size(&self) -> Option<u64> {
187        match self {
188            Self::Sqlite(buffer) => buffer.stack_provider.total_size(),
189            Self::InMemory(buffer) => buffer.stack_provider.total_size(),
190        }
191    }
192
193    /// Shuts down the [`PolymorphicEnvelopeBuffer`].
194    pub async fn shutdown(&mut self) -> bool {
195        // Currently, we want to flush the buffer only for disk, since the in memory implementation
196        // tries to not do anything and pop as many elements as possible within the shutdown
197        // timeout.
198        match self {
199            Self::Sqlite(buffer) if !buffer.stack_provider.ephemeral() => {
200                buffer.flush().await;
201                true
202            }
203            _ => {
204                relay_log::trace!("shutdown procedure not needed");
205                false
206            }
207        }
208    }
209
210    /// Returns the partition tag for this [`PolymorphicEnvelopeBuffer`].
211    fn partition_tag(&self) -> &str {
212        match self {
213            PolymorphicEnvelopeBuffer::InMemory(buffer) => &buffer.partition_tag,
214            PolymorphicEnvelopeBuffer::Sqlite(buffer) => &buffer.partition_tag,
215        }
216    }
217}
218
219/// Error that occurs while interacting with the envelope buffer.
220#[derive(Debug, thiserror::Error)]
221pub enum EnvelopeBufferError {
222    #[error("sqlite")]
223    SqliteStore(#[from] SqliteEnvelopeStoreError),
224
225    #[error("sqlite")]
226    SqliteStack(#[from] SqliteEnvelopeStackError),
227
228    #[error("failed to push envelope to the buffer")]
229    PushFailed,
230}
231
232impl From<Infallible> for EnvelopeBufferError {
233    fn from(value: Infallible) -> Self {
234        match value {}
235    }
236}
237
238/// An envelope buffer that holds an individual stack for each project/sampling project combination.
239///
240/// Envelope stacks are organized in a priority queue, and are re-prioritized every time an envelope
241/// is pushed, popped, or when a project becomes ready.
242#[derive(Debug)]
243struct EnvelopeBuffer<P: StackProvider> {
244    /// The central priority queue.
245    priority_queue: priority_queue::PriorityQueue<QueueItem<ProjectKeyPair, P::Stack>, Priority>,
246    /// A lookup table to find all stacks involving a project.
247    stacks_by_project: hashbrown::HashMap<ProjectKey, BTreeSet<ProjectKeyPair>>,
248    /// A provider of stacks that provides utilities to create stacks, check their capacity...
249    ///
250    /// This indirection is needed because different stack implementations might need different
251    /// initialization (e.g. a database connection).
252    stack_provider: P,
253    /// The total count of envelopes that the buffer is working with.
254    ///
255    /// Note that this count is not meant to be perfectly accurate since the initialization of the
256    /// count might not succeed if it takes more than a set timeout. For example, if we load the
257    /// count of all envelopes from disk, and it takes more than the time we set, we will mark the
258    /// initial count as 0 and just count incoming and outgoing envelopes from the buffer.
259    total_count: i64,
260    /// The total count of envelopes that the buffer is working with ignoring envelopes that
261    /// were previously stored on disk.
262    ///
263    /// On startup this will always be 0 and will only count incoming envelopes. If a reliable
264    /// count of currently buffered envelopes is required, prefer this over `total_count`
265    tracked_count: u64,
266    /// Whether the count initialization succeeded or not.
267    ///
268    /// This boolean is just used for tagging the metric that tracks the total count of envelopes
269    /// in the buffer.
270    total_count_initialized: bool,
271    /// The tag value of this partition which is used for reporting purposes.
272    partition_tag: String,
273}
274
275impl EnvelopeBuffer<MemoryStackProvider> {
276    /// Creates an empty memory-based buffer.
277    pub fn new(partition_id: u8, memory_checker: MemoryChecker) -> Self {
278        Self {
279            stacks_by_project: Default::default(),
280            priority_queue: Default::default(),
281            stack_provider: MemoryStackProvider::new(memory_checker),
282            total_count: 0,
283            tracked_count: 0,
284            total_count_initialized: false,
285            partition_tag: partition_id.to_string(),
286        }
287    }
288}
289
290#[allow(dead_code)]
291impl EnvelopeBuffer<SqliteStackProvider> {
292    /// Creates an empty sqlite-based buffer.
293    pub async fn new(
294        partition_id: u8,
295        config: &ConfigSnapshot,
296        throttle: Option<Arc<Throttle>>,
297    ) -> Result<Self, EnvelopeBufferError> {
298        Ok(Self {
299            stacks_by_project: Default::default(),
300            priority_queue: Default::default(),
301            stack_provider: SqliteStackProvider::new(partition_id, config, throttle).await?,
302            total_count: 0,
303            tracked_count: 0,
304            total_count_initialized: false,
305            partition_tag: partition_id.to_string(),
306        })
307    }
308}
309
310impl<P: StackProvider> EnvelopeBuffer<P>
311where
312    EnvelopeBufferError: From<<P::Stack as EnvelopeStack>::Error>,
313{
314    /// Initializes the [`EnvelopeBuffer`] given the initialization state from the
315    /// [`StackProvider`].
316    pub async fn initialize(&mut self) {
317        relay_statsd::metric!(
318            timer(RelayTimers::BufferInitialization),
319            partition_id = &self.partition_tag,
320            {
321                let initialization_state = self.stack_provider.initialize().await;
322                self.load_stacks(initialization_state.project_key_pairs)
323                    .await;
324                self.load_store_total_count().await;
325            }
326        );
327    }
328
329    /// Pushes an envelope to the appropriate envelope stack and re-prioritizes the stack.
330    ///
331    /// If the envelope stack does not exist, a new stack is pushed to the priority queue.
332    /// The priority of the stack is updated with the envelope's received_at time.
333    pub async fn push(&mut self, envelope: Box<Envelope>) -> Result<(), EnvelopeBufferError> {
334        let received_at = envelope.received_at();
335
336        let project_key_pair = ProjectKeyPair::from_envelope(&envelope);
337        if let Some((
338            QueueItem {
339                key: _,
340                value: stack,
341            },
342            _,
343        )) = self.priority_queue.get_mut(&project_key_pair)
344        {
345            stack.push(envelope).await?;
346        } else {
347            // Since we have initialization code that creates all the necessary stacks, we assume
348            // that any new stack that is added during the envelope buffer's lifecycle, is recreated.
349            self.push_stack(
350                StackCreationType::New,
351                ProjectKeyPair::from_envelope(&envelope),
352                Some(envelope),
353            )
354            .await?;
355        }
356        self.priority_queue
357            .change_priority_by(&project_key_pair, |prio| {
358                prio.received_at = received_at;
359            });
360
361        self.total_count += 1;
362        self.tracked_count += 1;
363        self.track_total_count();
364
365        Ok(())
366    }
367
368    /// Returns a reference to the next-in-line envelope, if one exists.
369    pub async fn peek(&mut self) -> Result<Peek, EnvelopeBufferError> {
370        let Some((
371            QueueItem {
372                key: project_key_pair,
373                value: stack,
374            },
375            Priority {
376                readiness,
377                next_project_fetch,
378                ..
379            },
380        )) = self.priority_queue.peek_mut()
381        else {
382            return Ok(Peek::Empty);
383        };
384
385        let ready = readiness.ready();
386
387        Ok(match (stack.peek().await?, ready) {
388            (None, _) => Peek::Empty,
389            (Some(last_received_at), true) => Peek::Ready {
390                project_key_pair: *project_key_pair,
391                last_received_at,
392            },
393            (Some(last_received_at), false) => Peek::NotReady {
394                project_key_pair: *project_key_pair,
395                next_project_fetch: *next_project_fetch,
396                last_received_at,
397            },
398        })
399    }
400
401    /// Returns the next-in-line envelope, if one exists.
402    ///
403    /// The priority of the envelope's stack is updated with the next envelope's received_at
404    /// time. If the stack is empty after popping, it is removed from the priority queue.
405    pub async fn pop(&mut self) -> Result<Option<Box<Envelope>>, EnvelopeBufferError> {
406        let Some((QueueItem { key, value: stack }, _)) = self.priority_queue.peek_mut() else {
407            return Ok(None);
408        };
409        let project_key_pair = *key;
410        let envelope = stack.pop().await?.expect("found an empty stack");
411
412        let last_received_at = stack.peek().await?;
413
414        match last_received_at {
415            None => {
416                self.pop_stack(project_key_pair);
417            }
418            Some(last_received_at) => {
419                self.priority_queue
420                    .change_priority_by(&project_key_pair, |prio| {
421                        prio.received_at = last_received_at;
422                    });
423            }
424        }
425
426        // We are fine with the count going negative, since it represents that more data was popped,
427        // than it was initially counted, meaning that we had a wrong total count from
428        // initialization.
429        self.total_count -= 1;
430        self.tracked_count = self.tracked_count.saturating_sub(1);
431        self.track_total_count();
432
433        Ok(Some(envelope))
434    }
435
436    /// Re-prioritizes all stacks that involve the given project key by setting it to "ready".
437    ///
438    /// Returns `true` if at least one priority was changed.
439    pub fn mark_ready(&mut self, project: &ProjectKey, is_ready: bool) -> bool {
440        let mut changed = false;
441        if let Some(project_key_pairs) = self.stacks_by_project.get(project) {
442            for project_key_pair in project_key_pairs {
443                self.priority_queue
444                    .change_priority_by(project_key_pair, |stack| {
445                        let mut found = false;
446                        for (subkey, readiness) in [
447                            (
448                                project_key_pair.own_key,
449                                &mut stack.readiness.own_project_ready,
450                            ),
451                            (
452                                project_key_pair.sampling_key,
453                                &mut stack.readiness.sampling_project_ready,
454                            ),
455                        ] {
456                            if subkey == *project {
457                                found = true;
458                                if *readiness != is_ready {
459                                    changed = true;
460                                    *readiness = is_ready;
461                                }
462                            }
463                        }
464                        debug_assert!(found);
465                    });
466            }
467        }
468
469        changed
470    }
471
472    /// Marks a stack as seen.
473    ///
474    /// Non-ready stacks are deprioritized when they are marked as seen, such that
475    /// the next call to `.peek()` will look at a different stack. This prevents
476    /// head-of-line blocking.
477    pub fn mark_seen(&mut self, project_key_pair: &ProjectKeyPair, next_fetch: Duration) {
478        self.priority_queue
479            .change_priority_by(project_key_pair, |stack| {
480                // We use the next project fetch to debounce project fetching and avoid head of
481                // line blocking of non-ready stacks.
482                stack.next_project_fetch = Instant::now() + next_fetch;
483            });
484    }
485
486    /// Returns `true` if the buffer contains no envelopes, in memory or on disk.
487    pub fn is_empty(&self) -> bool {
488        self.priority_queue.is_empty()
489    }
490
491    /// Returns `true` if the underlying storage has the capacity to store more envelopes.
492    pub fn has_capacity(&self) -> bool {
493        self.stack_provider.has_store_capacity()
494    }
495
496    /// Flushes the envelope buffer.
497    pub async fn flush(&mut self) {
498        let priority_queue = mem::take(&mut self.priority_queue);
499        self.stack_provider
500            .flush(priority_queue.into_iter().map(|(q, _)| q.value))
501            .await;
502    }
503
504    /// Pushes a new [`EnvelopeStack`] with the given [`Envelope`] inserted.
505    async fn push_stack(
506        &mut self,
507        stack_creation_type: StackCreationType,
508        project_key_pair: ProjectKeyPair,
509        envelope: Option<Box<Envelope>>,
510    ) -> Result<(), EnvelopeBufferError> {
511        let received_at = envelope.as_ref().map_or(Utc::now(), |e| e.received_at());
512
513        let mut stack = self
514            .stack_provider
515            .create_stack(stack_creation_type, project_key_pair);
516        if let Some(envelope) = envelope {
517            stack.push(envelope).await?;
518        }
519
520        let previous_entry = self.priority_queue.push(
521            QueueItem {
522                key: project_key_pair,
523                value: stack,
524            },
525            Priority::new(received_at),
526        );
527        debug_assert!(previous_entry.is_none());
528        for project_key in project_key_pair.iter() {
529            self.stacks_by_project
530                .entry(project_key)
531                .or_default()
532                .insert(project_key_pair);
533        }
534        relay_statsd::metric!(
535            gauge(RelayGauges::BufferStackCount) = self.priority_queue.len() as u64,
536            partition_id = &self.partition_tag
537        );
538
539        Ok(())
540    }
541
542    /// Pops an [`EnvelopeStack`] with the supplied [`EnvelopeBufferError`].
543    fn pop_stack(&mut self, project_key_pair: ProjectKeyPair) {
544        for project_key in project_key_pair.iter() {
545            self.stacks_by_project
546                .get_mut(&project_key)
547                .expect("project_key is missing from lookup")
548                .remove(&project_key_pair);
549        }
550        self.priority_queue.remove(&project_key_pair);
551
552        relay_statsd::metric!(
553            gauge(RelayGauges::BufferStackCount) = self.priority_queue.len() as u64,
554            partition_id = &self.partition_tag
555        );
556    }
557
558    /// Creates all the [`EnvelopeStack`]s with no data given a set of [`ProjectKeyPair`].
559    async fn load_stacks(&mut self, project_key_pairs: HashSet<ProjectKeyPair>) {
560        for project_key_pair in project_key_pairs {
561            self.push_stack(StackCreationType::Initialization, project_key_pair, None)
562                .await
563                .expect("Pushing an empty stack raised an error");
564        }
565    }
566
567    /// Loads the total count from the store if it takes less than a specified duration.
568    ///
569    /// The total count returned by the store is related to the count of elements that the buffer
570    /// will process, besides the count of elements that will be added and removed during its
571    /// lifecycle
572    async fn load_store_total_count(&mut self) {
573        let total_count = timeout(Duration::from_secs(1), async {
574            self.stack_provider.store_total_count().await
575        })
576        .await;
577        match total_count {
578            Ok(total_count) => {
579                self.total_count = total_count as i64;
580                self.total_count_initialized = true;
581            }
582            Err(error) => {
583                self.total_count_initialized = false;
584                relay_log::error!(
585                    error = &error as &dyn Error,
586                    "failed to load the total envelope count of the store",
587                );
588            }
589        };
590        self.track_total_count();
591    }
592
593    /// Emits a metric to track the total count of envelopes that are in the envelope buffer.
594    fn track_total_count(&self) {
595        let total_count = self.total_count as f64;
596        let initialized = match self.total_count_initialized {
597            true => "true",
598            false => "false",
599        };
600        relay_statsd::metric!(
601            gauge(RelayGauges::BufferEnvelopesCount) = total_count,
602            initialized = initialized,
603            stack_type = self.stack_provider.stack_type(),
604            partition_id = &self.partition_tag
605        );
606    }
607}
608
609/// Contains the state of the first element in the buffer.
610pub enum Peek {
611    Empty,
612    Ready {
613        project_key_pair: ProjectKeyPair,
614        last_received_at: DateTime<Utc>,
615    },
616    NotReady {
617        project_key_pair: ProjectKeyPair,
618        next_project_fetch: Instant,
619        last_received_at: DateTime<Utc>,
620    },
621}
622
623impl Peek {
624    pub fn last_received_at(&self) -> Option<DateTime<Utc>> {
625        match self {
626            Self::Empty => None,
627            Self::Ready {
628                last_received_at, ..
629            }
630            | Self::NotReady {
631                last_received_at, ..
632            } => Some(*last_received_at),
633        }
634    }
635}
636
637#[derive(Debug)]
638struct QueueItem<K, V> {
639    key: K,
640    value: V,
641}
642
643impl<K, V> std::borrow::Borrow<K> for QueueItem<K, V> {
644    fn borrow(&self) -> &K {
645        &self.key
646    }
647}
648
649impl<K: std::hash::Hash, V> std::hash::Hash for QueueItem<K, V> {
650    fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
651        self.key.hash(state);
652    }
653}
654
655impl<K: PartialEq, V> PartialEq for QueueItem<K, V> {
656    fn eq(&self, other: &Self) -> bool {
657        self.key == other.key
658    }
659}
660
661impl<K: PartialEq, V> Eq for QueueItem<K, V> {}
662
663#[derive(Debug, Clone)]
664struct Priority {
665    readiness: Readiness,
666    received_at: DateTime<Utc>,
667    next_project_fetch: Instant,
668}
669
670impl Priority {
671    fn new(received_at: DateTime<Utc>) -> Self {
672        Self {
673            readiness: Readiness::new(),
674            received_at,
675            next_project_fetch: Instant::now(),
676        }
677    }
678}
679
680impl Ord for Priority {
681    fn cmp(&self, other: &Self) -> Ordering {
682        match (self.readiness.ready(), other.readiness.ready()) {
683            // Assuming that two priorities differ only w.r.t. the `last_peek`, we want to prioritize
684            // stacks that were the least recently peeked. The rationale behind this is that we want
685            // to keep cycling through different stacks while peeking.
686            (true, true) => self.received_at.cmp(&other.received_at),
687            (true, false) => Ordering::Greater,
688            (false, true) => Ordering::Less,
689            // For non-ready stacks, we invert the priority, such that projects that are not
690            // ready and did not receive envelopes recently can be evicted.
691            (false, false) => self
692                .next_project_fetch
693                .cmp(&other.next_project_fetch)
694                .reverse()
695                .then(self.received_at.cmp(&other.received_at).reverse()),
696        }
697    }
698}
699
700impl PartialOrd for Priority {
701    fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
702        Some(self.cmp(other))
703    }
704}
705
706impl PartialEq for Priority {
707    fn eq(&self, other: &Self) -> bool {
708        self.cmp(other).is_eq()
709    }
710}
711
712impl Eq for Priority {}
713
714#[derive(Debug, Clone, Copy)]
715struct Readiness {
716    own_project_ready: bool,
717    sampling_project_ready: bool,
718}
719
720impl Readiness {
721    fn new() -> Self {
722        // Optimistically set ready state to true.
723        // The large majority of stack creations are re-creations after a stack was emptied.
724        Self {
725            own_project_ready: true,
726            sampling_project_ready: true,
727        }
728    }
729
730    fn ready(&self) -> bool {
731        self.own_project_ready && self.sampling_project_ready
732    }
733}
734
735#[cfg(test)]
736mod tests {
737    use relay_base_schema::project::ProjectId;
738    use relay_common::Dsn;
739    use relay_config::Config;
740    use relay_event_schema::protocol::EventId;
741    use relay_sampling::DynamicSamplingContext;
742    use std::str::FromStr;
743    use std::sync::Arc;
744    use uuid::Uuid;
745
746    use crate::SqliteEnvelopeStore;
747    use crate::envelope::{Item, ItemType};
748    use crate::extractors::RequestMeta;
749    use crate::services::buffer::common::ProjectKeyPair;
750    use crate::services::buffer::envelope_store::sqlite::DatabaseEnvelope;
751    use crate::services::buffer::testutils::utils::mock_envelopes;
752    use crate::utils::MemoryStat;
753
754    use super::*;
755
756    impl Peek {
757        fn is_empty(&self) -> bool {
758            matches!(self, Peek::Empty)
759        }
760    }
761
762    fn new_envelope(
763        own_key: ProjectKey,
764        sampling_key: Option<ProjectKey>,
765        event_id: Option<EventId>,
766    ) -> Box<Envelope> {
767        let mut envelope = Envelope::from_request(
768            None,
769            RequestMeta::new(Dsn::from_str(&format!("http://{own_key}@localhost/1")).unwrap()),
770        );
771        if let Some(sampling_key) = sampling_key {
772            envelope.set_dsc(DynamicSamplingContext {
773                public_key: sampling_key,
774                project_id: Some(ProjectId::new(42)),
775                trace_id: "67e5504410b1426f9247bb680e5fe0c8".parse().unwrap(),
776                release: None,
777                user: Default::default(),
778                replay_id: None,
779                environment: None,
780                transaction: None,
781                sample_rate: None,
782                sampled: None,
783                other: Default::default(),
784            });
785            envelope.add_item(Item::new(ItemType::Transaction));
786        }
787        if let Some(event_id) = event_id {
788            envelope.set_event_id(event_id);
789        }
790        envelope
791    }
792
793    fn mock_config(path: &str) -> Arc<Config> {
794        Config::from_json_value(serde_json::json!({
795            "spool": {
796                "envelopes": {
797                    "path": path
798                }
799            }
800        }))
801        .unwrap()
802        .into()
803    }
804
805    fn mock_memory_checker() -> MemoryChecker {
806        MemoryChecker::new(MemoryStat::default(), mock_config("my/db/path").clone())
807    }
808
809    async fn peek_received_at(buffer: &mut EnvelopeBuffer<MemoryStackProvider>) -> DateTime<Utc> {
810        buffer.peek().await.unwrap().last_received_at().unwrap()
811    }
812
813    #[tokio::test]
814    async fn test_insert_pop() {
815        let mut buffer = EnvelopeBuffer::<MemoryStackProvider>::new(0, mock_memory_checker());
816
817        let project_key1 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fed").unwrap();
818        let project_key2 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fee").unwrap();
819        let project_key3 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fef").unwrap();
820
821        assert!(buffer.pop().await.unwrap().is_none());
822        assert!(buffer.peek().await.unwrap().is_empty());
823
824        let envelope1 = new_envelope(project_key1, None, None);
825        let time1 = envelope1.meta().received_at();
826        buffer.push(envelope1).await.unwrap();
827
828        let envelope2 = new_envelope(project_key2, None, None);
829        let time2 = envelope2.meta().received_at();
830        buffer.push(envelope2).await.unwrap();
831
832        // Both projects are ready, so project 2 is on top (has the newest envelopes):
833        assert_eq!(peek_received_at(&mut buffer).await, time2);
834
835        buffer.mark_ready(&project_key1, false);
836        buffer.mark_ready(&project_key2, false);
837
838        // Both projects are not ready, so project 1 is on top (has the oldest envelopes):
839        assert_eq!(peek_received_at(&mut buffer).await, time1);
840
841        let envelope3 = new_envelope(project_key3, None, None);
842        let time3 = envelope3.meta().received_at();
843        buffer.push(envelope3).await.unwrap();
844        buffer.mark_ready(&project_key3, false);
845
846        // All projects are not ready, so project 1 is on top (has the oldest envelopes):
847        assert_eq!(peek_received_at(&mut buffer).await, time1);
848
849        // After marking a project ready, it goes to the top:
850        buffer.mark_ready(&project_key3, true);
851        assert_eq!(peek_received_at(&mut buffer).await, time3);
852        assert_eq!(
853            buffer.pop().await.unwrap().unwrap().meta().public_key(),
854            project_key3
855        );
856
857        // After popping, project 1 is on top again:
858        assert_eq!(peek_received_at(&mut buffer).await, time1);
859
860        // Mark project 1 as ready (still on top):
861        buffer.mark_ready(&project_key1, true);
862        assert_eq!(peek_received_at(&mut buffer).await, time1);
863
864        // Mark project 2 as ready as well (now on top because most recent):
865        buffer.mark_ready(&project_key2, true);
866        assert_eq!(peek_received_at(&mut buffer).await, time2);
867        assert_eq!(
868            buffer.pop().await.unwrap().unwrap().meta().public_key(),
869            project_key2
870        );
871
872        // Pop last element:
873        assert_eq!(
874            buffer.pop().await.unwrap().unwrap().meta().public_key(),
875            project_key1
876        );
877        assert!(buffer.pop().await.unwrap().is_none());
878        assert!(buffer.peek().await.unwrap().is_empty());
879    }
880
881    #[tokio::test]
882    async fn test_project_internal_order() {
883        let mut buffer = EnvelopeBuffer::<MemoryStackProvider>::new(0, mock_memory_checker());
884
885        let project_key = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fed").unwrap();
886
887        let envelope1 = new_envelope(project_key, None, None);
888        let time1 = envelope1.meta().received_at();
889        let envelope2 = new_envelope(project_key, None, None);
890        let time2 = envelope2.meta().received_at();
891
892        assert!(time2 > time1);
893
894        buffer.push(envelope1).await.unwrap();
895        buffer.push(envelope2).await.unwrap();
896
897        assert_eq!(
898            buffer.pop().await.unwrap().unwrap().meta().received_at(),
899            time2
900        );
901        assert_eq!(
902            buffer.pop().await.unwrap().unwrap().meta().received_at(),
903            time1
904        );
905        assert!(buffer.pop().await.unwrap().is_none());
906    }
907
908    #[tokio::test]
909    async fn test_sampling_projects() {
910        let mut buffer = EnvelopeBuffer::<MemoryStackProvider>::new(0, mock_memory_checker());
911
912        let project_key1 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fed").unwrap();
913        let project_key2 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fef").unwrap();
914
915        let envelope1 = new_envelope(project_key1, None, None);
916        let time1 = envelope1.received_at();
917        buffer.push(envelope1).await.unwrap();
918
919        let envelope2 = new_envelope(project_key2, None, None);
920        let time2 = envelope2.received_at();
921        buffer.push(envelope2).await.unwrap();
922
923        let envelope3 = new_envelope(project_key1, Some(project_key2), None);
924        let time3 = envelope3.meta().received_at();
925        buffer.push(envelope3).await.unwrap();
926
927        buffer.mark_ready(&project_key1, false);
928        buffer.mark_ready(&project_key2, false);
929
930        // Nothing is ready, instant1 is on top:
931        assert_eq!(
932            buffer.peek().await.unwrap().last_received_at().unwrap(),
933            time1
934        );
935
936        // Mark project 2 ready, gets on top:
937        buffer.mark_ready(&project_key2, true);
938        assert_eq!(
939            buffer.peek().await.unwrap().last_received_at().unwrap(),
940            time2
941        );
942
943        // Revert
944        buffer.mark_ready(&project_key2, false);
945        assert_eq!(
946            buffer.peek().await.unwrap().last_received_at().unwrap(),
947            time1
948        );
949
950        // Project 1 ready:
951        buffer.mark_ready(&project_key1, true);
952        assert_eq!(
953            buffer.peek().await.unwrap().last_received_at().unwrap(),
954            time1
955        );
956
957        // when both projects are ready, event no 3 ends up on top:
958        buffer.mark_ready(&project_key2, true);
959        assert_eq!(
960            buffer.pop().await.unwrap().unwrap().meta().received_at(),
961            time3
962        );
963        assert_eq!(
964            buffer.peek().await.unwrap().last_received_at().unwrap(),
965            time2
966        );
967
968        buffer.mark_ready(&project_key2, false);
969        assert_eq!(buffer.pop().await.unwrap().unwrap().received_at(), time1);
970        assert_eq!(buffer.pop().await.unwrap().unwrap().received_at(), time2);
971
972        assert!(buffer.pop().await.unwrap().is_none());
973    }
974
975    #[tokio::test]
976    async fn test_project_keys_distinct() {
977        let project_key1 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fed").unwrap();
978        let project_key2 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fef").unwrap();
979
980        let project_key_pair1 = ProjectKeyPair::new(project_key1, project_key2);
981        let project_key_pair2 = ProjectKeyPair::new(project_key2, project_key1);
982
983        assert_ne!(project_key_pair1, project_key_pair2);
984
985        let mut buffer = EnvelopeBuffer::<MemoryStackProvider>::new(0, mock_memory_checker());
986        buffer
987            .push(new_envelope(project_key1, Some(project_key2), None))
988            .await
989            .unwrap();
990        buffer
991            .push(new_envelope(project_key2, Some(project_key1), None))
992            .await
993            .unwrap();
994        assert_eq!(buffer.priority_queue.len(), 2);
995    }
996
997    #[test]
998    fn test_total_order() {
999        let p1 = Priority {
1000            readiness: Readiness {
1001                own_project_ready: true,
1002                sampling_project_ready: true,
1003            },
1004            received_at: Utc::now(),
1005            next_project_fetch: Instant::now(),
1006        };
1007        let mut p2 = p1.clone();
1008        p2.next_project_fetch += Duration::from_millis(1);
1009
1010        // Last peek does not matter because project is ready:
1011        assert_eq!(p1.cmp(&p2), Ordering::Equal);
1012        assert_eq!(p1, p2);
1013    }
1014
1015    #[tokio::test]
1016    async fn test_last_peek_internal_order() {
1017        let mut buffer = EnvelopeBuffer::<MemoryStackProvider>::new(0, mock_memory_checker());
1018
1019        let project_key_1 = ProjectKey::parse("a94ae32be2584e0bbd7a4cbb95971fed").unwrap();
1020        let event_id_1 = EventId::new();
1021        let envelope1 = new_envelope(project_key_1, None, Some(event_id_1));
1022        let time1 = envelope1.received_at();
1023
1024        let project_key_2 = ProjectKey::parse("b56ae32be2584e0bbd7a4cbb95971fed").unwrap();
1025        let event_id_2 = EventId::new();
1026        let envelope2 = new_envelope(project_key_2, None, Some(event_id_2));
1027        let time2 = envelope2.received_at();
1028
1029        buffer.push(envelope1).await.unwrap();
1030        buffer.push(envelope2).await.unwrap();
1031
1032        buffer.mark_ready(&project_key_1, false);
1033        buffer.mark_ready(&project_key_2, false);
1034
1035        // event_id_1 is first element:
1036        let Peek::NotReady {
1037            last_received_at, ..
1038        } = buffer.peek().await.unwrap()
1039        else {
1040            panic!();
1041        };
1042        assert_eq!(last_received_at, time1);
1043
1044        // Second peek returns same element:
1045        let Peek::NotReady {
1046            last_received_at,
1047            project_key_pair,
1048            ..
1049        } = buffer.peek().await.unwrap()
1050        else {
1051            panic!();
1052        };
1053        assert_eq!(last_received_at, time1);
1054        assert_ne!(last_received_at, time2);
1055
1056        buffer.mark_seen(&project_key_pair, Duration::ZERO);
1057
1058        // After mark_seen, event 2 is on top:
1059        let Peek::NotReady {
1060            last_received_at, ..
1061        } = buffer.peek().await.unwrap()
1062        else {
1063            panic!();
1064        };
1065        assert_eq!(last_received_at, time2);
1066        assert_ne!(last_received_at, time1);
1067
1068        let Peek::NotReady {
1069            last_received_at,
1070            project_key_pair,
1071            ..
1072        } = buffer.peek().await.unwrap()
1073        else {
1074            panic!();
1075        };
1076        assert_eq!(last_received_at, time2);
1077        assert_ne!(last_received_at, time1);
1078
1079        buffer.mark_seen(&project_key_pair, Duration::ZERO);
1080
1081        // After another mark_seen, cycle back to event 1:
1082        let Peek::NotReady {
1083            last_received_at, ..
1084        } = buffer.peek().await.unwrap()
1085        else {
1086            panic!();
1087        };
1088        assert_eq!(last_received_at, time1);
1089        assert_ne!(last_received_at, time2);
1090    }
1091
1092    #[tokio::test]
1093    async fn test_initialize_buffer() {
1094        let path = std::env::temp_dir()
1095            .join(Uuid::new_v4().to_string())
1096            .into_os_string()
1097            .into_string()
1098            .unwrap();
1099        let config = mock_config(&path);
1100        let current_config = config.current();
1101        let mut store = SqliteEnvelopeStore::prepare(0, &current_config)
1102            .await
1103            .unwrap();
1104        let mut buffer = EnvelopeBuffer::<SqliteStackProvider>::new(0, &current_config, None)
1105            .await
1106            .unwrap();
1107
1108        // We write 5 envelopes to disk so that we can check if they are loaded. These envelopes
1109        // belong to the same project keys, so they belong to the same envelope stack.
1110        let envelopes = mock_envelopes(10);
1111        assert!(
1112            store
1113                .insert_batch(
1114                    envelopes
1115                        .into_iter()
1116                        .map(|e| DatabaseEnvelope::try_from(e.as_ref()).unwrap())
1117                        .collect::<Vec<_>>()
1118                        .try_into()
1119                        .unwrap()
1120                )
1121                .await
1122                .is_ok()
1123        );
1124
1125        // We assume that the buffer is empty.
1126        assert!(buffer.priority_queue.is_empty());
1127        assert!(buffer.stacks_by_project.is_empty());
1128
1129        buffer.initialize().await;
1130
1131        // We assume that we loaded only 1 envelope stack, because of the project keys combinations
1132        // of the envelopes we inserted above.
1133        assert_eq!(buffer.priority_queue.len(), 1);
1134        // We expect to have an entry per project key, since we have 1 pair, the total entries
1135        // should be 2.
1136        assert_eq!(buffer.stacks_by_project.len(), 2);
1137    }
1138}