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#[derive(Debug)]
36#[allow(private_interfaces)]
37pub enum PolymorphicEnvelopeBuffer {
38 InMemory(EnvelopeBuffer<MemoryStackProvider>),
40 Sqlite(EnvelopeBuffer<SqliteStackProvider>),
42}
43
44impl PolymorphicEnvelopeBuffer {
45 pub fn is_memory(&self) -> bool {
47 match self {
48 Self::InMemory(_) => true,
49 Self::Sqlite(_) => false,
50 }
51 }
52
53 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 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 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 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 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 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 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 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 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 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 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 pub async fn shutdown(&mut self) -> bool {
195 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 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#[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#[derive(Debug)]
243struct EnvelopeBuffer<P: StackProvider> {
244 priority_queue: priority_queue::PriorityQueue<QueueItem<ProjectKeyPair, P::Stack>, Priority>,
246 stacks_by_project: hashbrown::HashMap<ProjectKey, BTreeSet<ProjectKeyPair>>,
248 stack_provider: P,
253 total_count: i64,
260 tracked_count: u64,
266 total_count_initialized: bool,
271 partition_tag: String,
273}
274
275impl EnvelopeBuffer<MemoryStackProvider> {
276 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 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 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 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 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 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 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 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 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 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 stack.next_project_fetch = Instant::now() + next_fetch;
483 });
484 }
485
486 pub fn is_empty(&self) -> bool {
488 self.priority_queue.is_empty()
489 }
490
491 pub fn has_capacity(&self) -> bool {
493 self.stack_provider.has_store_capacity()
494 }
495
496 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 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 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 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 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 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
609pub 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 (true, true) => self.received_at.cmp(&other.received_at),
687 (true, false) => Ordering::Greater,
688 (false, true) => Ordering::Less,
689 (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 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 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 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 assert_eq!(peek_received_at(&mut buffer).await, time1);
848
849 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 assert_eq!(peek_received_at(&mut buffer).await, time1);
859
860 buffer.mark_ready(&project_key1, true);
862 assert_eq!(peek_received_at(&mut buffer).await, time1);
863
864 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 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 assert_eq!(
932 buffer.peek().await.unwrap().last_received_at().unwrap(),
933 time1
934 );
935
936 buffer.mark_ready(&project_key2, true);
938 assert_eq!(
939 buffer.peek().await.unwrap().last_received_at().unwrap(),
940 time2
941 );
942
943 buffer.mark_ready(&project_key2, false);
945 assert_eq!(
946 buffer.peek().await.unwrap().last_received_at().unwrap(),
947 time1
948 );
949
950 buffer.mark_ready(&project_key1, true);
952 assert_eq!(
953 buffer.peek().await.unwrap().last_received_at().unwrap(),
954 time1
955 );
956
957 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 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 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 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 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 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, ¤t_config)
1102 .await
1103 .unwrap();
1104 let mut buffer = EnvelopeBuffer::<SqliteStackProvider>::new(0, ¤t_config, None)
1105 .await
1106 .unwrap();
1107
1108 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 assert!(buffer.priority_queue.is_empty());
1127 assert!(buffer.stacks_by_project.is_empty());
1128
1129 buffer.initialize().await;
1130
1131 assert_eq!(buffer.priority_queue.len(), 1);
1134 assert_eq!(buffer.stacks_by_project.len(), 2);
1137 }
1138}