1use std::collections::HashMap;
16use std::collections::VecDeque;
17use std::mem;
18use std::path::PathBuf;
19use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
20use std::sync::{Arc, RwLock, RwLockWriteGuard};
21use std::time::{Duration, Instant};
22
23use chrono::Utc;
24use malloc_size_of::MallocSizeOf;
25use malloc_size_of_derive::MallocSizeOf;
26
27use crate::error::ErrorKind;
28use crate::TimerId;
29use crate::{internal_metrics::UploadMetrics, Glean};
30pub use directory::process_metadata;
31use directory::{PingDirectoryManager, PingPayloadsByDirectory};
32use policy::Policy;
33use request::create_date_header_value;
34
35pub use directory::{PingMetadata, PingPayload};
36pub use request::{HeaderMap, PingRequest};
37pub use result::{UploadResult, UploadTaskAction};
38
39mod directory;
40mod policy;
41mod request;
42mod result;
43
44const WAIT_TIME_FOR_PING_PROCESSING: u64 = 1000; #[derive(Debug, MallocSizeOf)]
47struct RateLimiter {
48 started: Option<Instant>,
50 count: u32,
52 interval: Duration,
54 max_count: u32,
56}
57
58#[derive(PartialEq)]
60enum RateLimiterState {
61 Incrementing,
63 Throttled(u64),
68}
69
70impl RateLimiter {
71 pub fn new(interval: Duration, max_count: u32) -> Self {
72 Self {
73 started: None,
74 count: 0,
75 interval,
76 max_count,
77 }
78 }
79
80 fn reset(&mut self) {
81 self.started = Some(Instant::now());
82 self.count = 0;
83 }
84
85 fn elapsed(&self) -> Duration {
86 self.started.unwrap().elapsed()
87 }
88
89 fn should_reset(&self) -> bool {
95 if self.started.is_none() {
96 return true;
97 }
98
99 if self.elapsed() > self.interval {
101 return true;
102 }
103
104 false
105 }
106
107 pub fn get_state(&mut self) -> RateLimiterState {
113 if self.should_reset() {
114 self.reset();
115 }
116
117 if self.count == self.max_count {
118 let remaining = self.interval.as_millis() - self.elapsed().as_millis();
121 return RateLimiterState::Throttled(
122 remaining
123 .try_into()
124 .unwrap_or(self.interval.as_secs() * 1000),
125 );
126 }
127
128 self.count += 1;
129 RateLimiterState::Incrementing
130 }
131}
132
133#[derive(PartialEq, Eq, Debug)]
138pub enum PingUploadTask {
139 Upload {
141 request: PingRequest,
144 },
145
146 Wait {
149 time: u64,
152 },
153
154 Done {
167 #[doc(hidden)]
168 unused: i8,
170 },
171}
172
173impl PingUploadTask {
174 pub fn is_upload(&self) -> bool {
176 matches!(self, PingUploadTask::Upload { .. })
177 }
178
179 pub fn is_wait(&self) -> bool {
181 matches!(self, PingUploadTask::Wait { .. })
182 }
183
184 pub(crate) fn done() -> Self {
185 PingUploadTask::Done { unused: 0 }
186 }
187}
188
189#[derive(Debug)]
191pub struct PingUploadManager {
192 queue: RwLock<VecDeque<PingRequest>>,
194 directory_manager: PingDirectoryManager,
196 processed_pending_pings: Arc<AtomicBool>,
198 cached_pings: Arc<RwLock<PingPayloadsByDirectory>>,
200 recoverable_failure_count: AtomicU32,
202 wait_attempt_count: AtomicU32,
204 rate_limiter: Option<RwLock<RateLimiter>>,
209 language_binding_name: String,
213 upload_metrics: UploadMetrics,
215 policy: Policy,
217
218 in_flight: RwLock<HashMap<String, (TimerId, TimerId)>>,
219}
220
221impl MallocSizeOf for PingUploadManager {
222 fn size_of(&self, ops: &mut malloc_size_of::MallocSizeOfOps) -> usize {
223 let shallow_size = {
224 let queue = self.queue.read().unwrap();
225 if ops.has_malloc_enclosing_size_of() {
226 if let Some(front) = queue.front() {
227 unsafe { ops.malloc_enclosing_size_of(front) }
230 } else {
231 0
233 }
234 } else {
235 queue.capacity() * mem::size_of::<PingRequest>()
239 }
240 };
241
242 let mut n = shallow_size
243 + self.directory_manager.size_of(ops)
244 + mem::size_of::<AtomicBool>() + self.cached_pings.read().unwrap().size_of(ops)
246 + self.rate_limiter.as_ref().map(|rl| {
247 let lock = rl.read().unwrap();
248 (*lock).size_of(ops)
249 }).unwrap_or(0)
250 + self.language_binding_name.size_of(ops)
251 + self.upload_metrics.size_of(ops)
252 + self.policy.size_of(ops);
253
254 let in_flight = self.in_flight.read().unwrap();
255 n += in_flight.size_of(ops);
256
257 n
258 }
259}
260
261impl PingUploadManager {
262 pub fn new<P: Into<PathBuf>>(data_path: P, language_binding_name: &str) -> Self {
273 Self {
274 queue: RwLock::new(VecDeque::new()),
275 directory_manager: PingDirectoryManager::new(data_path),
276 processed_pending_pings: Arc::new(AtomicBool::new(false)),
277 cached_pings: Arc::new(RwLock::new(PingPayloadsByDirectory::default())),
278 recoverable_failure_count: AtomicU32::new(0),
279 wait_attempt_count: AtomicU32::new(0),
280 rate_limiter: None,
281 language_binding_name: language_binding_name.into(),
282 upload_metrics: UploadMetrics::new(),
283 policy: Policy::default(),
284 in_flight: RwLock::new(HashMap::default()),
285 }
286 }
287
288 pub fn scan_pending_pings_directories(
295 &self,
296 trigger_upload: bool,
297 ) -> std::thread::JoinHandle<()> {
298 let local_manager = self.directory_manager.clone();
299 let local_cached_pings = self.cached_pings.clone();
300 let local_flag = self.processed_pending_pings.clone();
301 crate::thread::spawn("glean.ping_directory_manager.process_dir", move || {
302 {
303 let mut local_cached_pings = local_cached_pings
305 .write()
306 .expect("Can't write to pending pings cache.");
307 local_cached_pings.extend(local_manager.process_dirs());
308 local_flag.store(true, Ordering::SeqCst);
309 }
310 if trigger_upload {
311 crate::dispatcher::launch(|| {
312 if let Some(state) = crate::maybe_global_state().and_then(|s| s.lock().ok()) {
313 if let Err(e) = state.callbacks.trigger_upload() {
314 log::error!(
315 "Triggering upload after pending ping scan failed. Error: {}",
316 e
317 );
318 }
319 }
320 });
321 }
322 })
323 .expect("Unable to spawn thread to process pings directories.")
324 }
325
326 #[cfg(test)]
328 pub fn no_policy<P: Into<PathBuf>>(data_path: P) -> Self {
329 let mut upload_manager = Self::new(data_path, "Test");
330
331 upload_manager.policy.set_max_recoverable_failures(None);
333 upload_manager.policy.set_max_wait_attempts(None);
334 upload_manager.policy.set_max_ping_body_size(None);
335 upload_manager
336 .policy
337 .set_max_pending_pings_directory_size(None);
338 upload_manager.policy.set_max_pending_pings_count(None);
339
340 upload_manager
342 .scan_pending_pings_directories(false)
343 .join()
344 .unwrap();
345
346 upload_manager
347 }
348
349 fn processed_pending_pings(&self) -> bool {
350 self.processed_pending_pings.load(Ordering::SeqCst)
351 }
352
353 fn recoverable_failure_count(&self) -> u32 {
354 self.recoverable_failure_count.load(Ordering::SeqCst)
355 }
356
357 fn wait_attempt_count(&self) -> u32 {
358 self.wait_attempt_count.load(Ordering::SeqCst)
359 }
360
361 fn build_ping_request(&self, glean: &Glean, ping: PingPayload) -> Option<PingRequest> {
366 let PingPayload {
367 document_id,
368 upload_path: path,
369 json_body: body,
370 headers,
371 body_has_info_sections,
372 ping_name,
373 uploader_capabilities,
374 } = ping;
375 let mut request = PingRequest::builder(
376 &self.language_binding_name,
377 self.policy.max_ping_body_size(),
378 )
379 .document_id(&document_id)
380 .path(path)
381 .body(body)
382 .body_has_info_sections(body_has_info_sections)
383 .ping_name(ping_name)
384 .uploader_capabilities(uploader_capabilities);
385
386 if let Some(headers) = headers {
387 request = request.headers(headers);
388 }
389
390 match request.build() {
391 Ok(request) => Some(request),
392 Err(e) => {
393 log::warn!("Error trying to build ping request: {}", e);
394 self.directory_manager.delete_file(&document_id);
395
396 if let ErrorKind::PingBodyOverflow(s) = e.kind() {
399 self.upload_metrics
400 .discarded_exceeding_pings_size
401 .accumulate_sync(glean, *s as i64 / 1024);
402 }
403
404 None
405 }
406 }
407 }
408
409 pub fn enqueue_ping(&self, glean: &Glean, ping: PingPayload) {
411 let mut queue = self
412 .queue
413 .write()
414 .expect("Can't write to pending pings queue.");
415
416 let PingPayload {
417 ref document_id,
418 upload_path: ref path,
419 ..
420 } = ping;
421 if queue
423 .iter()
424 .any(|request| request.document_id.as_str() == document_id)
425 {
426 log::warn!(
427 "Attempted to enqueue a duplicate ping {} at {}.",
428 document_id,
429 path
430 );
431 return;
432 }
433
434 {
435 let in_flight = self.in_flight.read().unwrap();
436 if in_flight.contains_key(document_id) {
437 log::warn!(
438 "Attempted to enqueue an in-flight ping {} at {}.",
439 document_id,
440 path
441 );
442 self.upload_metrics
443 .in_flight_pings_dropped
444 .add_sync(glean, 0);
445 return;
446 }
447 }
448
449 log::trace!("Enqueuing ping {} at {}", document_id, path);
450 if let Some(request) = self.build_ping_request(glean, ping) {
451 queue.push_back(request)
452 }
453 }
454
455 fn enqueue_cached_pings(&self, glean: &Glean) {
472 let mut cached_pings = self
473 .cached_pings
474 .write()
475 .expect("Can't write to pending pings cache.");
476
477 if cached_pings.len() > 0 {
478 let mut pending_pings_directory_size: u64 = 0;
479 let mut pending_pings_count = 0;
480 let mut deleting = false;
481 let mut delete_reason: Option<&'static str> = None;
482
483 let total = cached_pings.pending_pings.len() as u64;
484 self.upload_metrics
485 .pending_pings
486 .add_sync(glean, total.try_into().unwrap_or(0));
487
488 if total > self.policy.max_pending_pings_count() {
489 log::warn!(
490 "More than {} pending pings in the directory, will delete {} old pings.",
491 self.policy.max_pending_pings_count(),
492 total - self.policy.max_pending_pings_count()
493 );
494 }
495
496 cached_pings.pending_pings.reverse();
502 cached_pings.pending_pings.retain(|(file_size, PingPayload {document_id, ..})| {
503 pending_pings_count += 1;
504 pending_pings_directory_size += file_size;
505
506 if !deleting && pending_pings_directory_size > self.policy.max_pending_pings_directory_size() {
510 log::warn!(
511 "Pending pings directory has reached the size quota of {} bytes, outstanding pings will be deleted.",
512 self.policy.max_pending_pings_directory_size()
513 );
514 deleting = true;
515 delete_reason = Some("size_quota");
516 }
517
518 if !deleting && pending_pings_count > self.policy.max_pending_pings_count() {
522 deleting = true;
523 delete_reason = Some("count_quota");
524 }
525
526 if deleting && self.directory_manager.delete_file(document_id) {
527 self.upload_metrics
528 .deleted_pings_after_quota_hit
529 .add_sync(glean, 1);
530 if let Some(reason) = delete_reason {
531 self.upload_metrics
532 .pending_pings_deleted
533 .get(reason)
534 .add_sync(glean, 1);
535 }
536 return false;
537 }
538
539 true
540 });
541 cached_pings.pending_pings.reverse();
544 self.upload_metrics
545 .pending_pings_directory_size
546 .accumulate_sync(glean, pending_pings_directory_size as i64 / 1024);
547
548 cached_pings
551 .deletion_request_pings
552 .drain(..)
553 .for_each(|(_, ping)| self.enqueue_ping(glean, ping));
554 cached_pings
555 .pending_pings
556 .drain(..)
557 .for_each(|(_, ping)| self.enqueue_ping(glean, ping));
558 }
559 }
560
561 pub fn set_rate_limiter(&mut self, interval: u64, max_tasks: u32) {
574 self.rate_limiter = Some(RwLock::new(RateLimiter::new(
575 Duration::from_secs(interval),
576 max_tasks,
577 )));
578 }
579
580 pub(crate) fn set_max_pending_pings_count(&mut self, n: u64) {
581 self.policy.set_max_pending_pings_count(Some(n));
582 }
583
584 pub(crate) fn set_max_pending_pings_directory_size(&mut self, n: u64) {
585 self.policy.set_max_pending_pings_directory_size(Some(n));
586 }
587
588 pub fn enqueue_ping_from_file(&self, glean: &Glean, document_id: &str) {
597 if let Some(ping) = self.directory_manager.process_file(document_id) {
598 self.enqueue_ping(glean, ping);
599 }
600 }
601
602 pub fn clear_ping_queue(&self) -> RwLockWriteGuard<'_, VecDeque<PingRequest>> {
604 log::trace!("Clearing ping queue");
605 let mut queue = self
606 .queue
607 .write()
608 .expect("Can't write to pending pings queue.");
609
610 queue.retain(|ping| ping.is_deletion_request());
611 log::trace!(
612 "{} pings left in the queue (only deletion-request expected)",
613 queue.len()
614 );
615 queue
616 }
617
618 fn get_upload_task_internal(&self, glean: &Glean, log_ping: bool) -> PingUploadTask {
619 let wait_or_done = |time: u64| {
624 self.wait_attempt_count.fetch_add(1, Ordering::SeqCst);
625 if self.wait_attempt_count() > self.policy.max_wait_attempts() {
626 PingUploadTask::done()
627 } else {
628 PingUploadTask::Wait { time }
629 }
630 };
631
632 if !self.processed_pending_pings() {
633 log::info!(
634 "Tried getting an upload task, but processing is ongoing. Will come back later."
635 );
636 return wait_or_done(WAIT_TIME_FOR_PING_PROCESSING);
637 }
638
639 self.enqueue_cached_pings(glean);
641
642 if self.recoverable_failure_count() >= self.policy.max_recoverable_failures() {
643 log::warn!(
644 "Reached maximum recoverable failures for the current uploading window. You are done."
645 );
646 return PingUploadTask::done();
647 }
648
649 let mut queue = self
650 .queue
651 .write()
652 .expect("Can't write to pending pings queue.");
653 match queue.front() {
654 Some(request) => {
655 if let Some(rate_limiter) = &self.rate_limiter {
656 let mut rate_limiter = rate_limiter
657 .write()
658 .expect("Can't write to the rate limiter.");
659 if let RateLimiterState::Throttled(remaining) = rate_limiter.get_state() {
660 log::info!(
661 "Tried getting an upload task, but we are throttled at the moment."
662 );
663 return wait_or_done(remaining);
664 }
665 }
666
667 log::info!(
668 "New upload task with id {} (path: {})",
669 request.document_id,
670 request.path
671 );
672
673 if log_ping {
674 if let Some(body) = request.pretty_body() {
675 chunked_log_info(&request.path, &body);
676 } else {
677 chunked_log_info(&request.path, "<invalid ping payload>");
678 }
679 }
680
681 {
682 let mut in_flight = self.in_flight.write().unwrap();
686 let success_id = self.upload_metrics.send_success.start_sync();
687 let failure_id = self.upload_metrics.send_failure.start_sync();
688 in_flight.insert(request.document_id.clone(), (success_id, failure_id));
689 }
690
691 let mut request = queue.pop_front().unwrap();
692
693 request
695 .headers
696 .insert("Date".to_string(), create_date_header_value(Utc::now()));
697
698 PingUploadTask::Upload { request }
699 }
700 None => {
701 log::info!("No more pings to upload! You are done.");
702 PingUploadTask::done()
703 }
704 }
705 }
706
707 pub fn get_upload_task(&self, glean: &Glean, log_ping: bool) -> PingUploadTask {
718 let task = self.get_upload_task_internal(glean, log_ping);
719
720 if !task.is_wait() && self.wait_attempt_count() > 0 {
721 self.wait_attempt_count.store(0, Ordering::SeqCst);
722 }
723
724 if !task.is_upload() && self.recoverable_failure_count() > 0 {
725 self.recoverable_failure_count.store(0, Ordering::SeqCst);
726 }
727
728 task
729 }
730
731 pub fn process_ping_upload_response(
770 &self,
771 glean: &Glean,
772 document_id: &str,
773 status: UploadResult,
774 ) -> UploadTaskAction {
775 use UploadResult::*;
776
777 let stop_time = zeitstempel::now_awake();
778
779 if let Some(label) = status.get_label() {
780 let metric = self.upload_metrics.ping_upload_failure.get(label);
781 metric.add_sync(glean, 1);
782 }
783
784 let send_ids = {
785 let mut lock = self.in_flight.write().unwrap();
786 lock.remove(document_id)
787 };
788
789 if send_ids.is_none() {
790 self.upload_metrics.missing_send_ids.add_sync(glean, 1);
791 }
792
793 match status {
794 HttpStatus { code } if (200..=299).contains(&code) => {
795 log::info!("Ping {} successfully sent {}.", document_id, code);
796 if let Some((success_id, failure_id)) = send_ids {
797 self.upload_metrics
798 .send_success
799 .set_stop_and_accumulate(glean, success_id, stop_time);
800 self.upload_metrics.send_failure.cancel_sync(failure_id);
801 }
802 #[cfg(feature = "sqlite")]
803 if glean.store_submitted_pings_enabled {
804 glean
805 .storage()
806 .mark_ping_as_uploaded(document_id, Utc::now());
807 }
808 self.directory_manager.delete_file(document_id);
809 }
810
811 UnrecoverableFailure { .. } | HttpStatus { code: 400..=499 } | Incapable { .. } => {
812 log::warn!(
813 "Unrecoverable upload failure while attempting to send ping {}. Error was {:?}",
814 document_id,
815 status
816 );
817 if let Some((success_id, failure_id)) = send_ids {
818 self.upload_metrics.send_success.cancel_sync(success_id);
819 self.upload_metrics
820 .send_failure
821 .set_stop_and_accumulate(glean, failure_id, stop_time);
822 }
823 #[cfg(feature = "sqlite")]
824 if glean.store_submitted_pings_enabled {
825 glean.storage().mark_ping_as_upload_failed(document_id);
826 }
827 self.directory_manager.delete_file(document_id);
828 }
829
830 RecoverableFailure { .. } | HttpStatus { .. } => {
831 log::warn!(
832 "Recoverable upload failure while attempting to send ping {}, will retry. Error was {:?}",
833 document_id,
834 status
835 );
836 if let Some((success_id, failure_id)) = send_ids {
837 self.upload_metrics.send_success.cancel_sync(success_id);
838 self.upload_metrics
839 .send_failure
840 .set_stop_and_accumulate(glean, failure_id, stop_time);
841 }
842 self.enqueue_ping_from_file(glean, document_id);
843 self.recoverable_failure_count
844 .fetch_add(1, Ordering::SeqCst);
845 }
846
847 Done { .. } => {
848 log::debug!("Uploader signaled Done. Exiting.");
849 if let Some((success_id, failure_id)) = send_ids {
850 self.upload_metrics.send_success.cancel_sync(success_id);
851 self.upload_metrics.send_failure.cancel_sync(failure_id);
852 }
853 return UploadTaskAction::End;
854 }
855 };
856
857 UploadTaskAction::Next
858 }
859}
860
861#[cfg(target_os = "android")]
863pub fn chunked_log_info(path: &str, payload: &str) {
864 const MAX_LOG_PAYLOAD_SIZE_BYTES: usize = 4000;
868
869 if path.len() + payload.len() <= MAX_LOG_PAYLOAD_SIZE_BYTES {
873 log::info!("Glean ping to URL: {}\n{}", path, payload);
874 return;
875 }
876
877 let mut start = 0;
880 let mut end = MAX_LOG_PAYLOAD_SIZE_BYTES;
881 let mut chunk_idx = 1;
882 let total_chunks = payload.len() / MAX_LOG_PAYLOAD_SIZE_BYTES + 1;
884
885 while end < payload.len() {
886 for _ in 0..4 {
889 if payload.is_char_boundary(end) {
890 break;
891 }
892 end -= 1;
893 }
894
895 log::info!(
896 "Glean ping to URL: {} [Part {} of {}]\n{}",
897 path,
898 chunk_idx,
899 total_chunks,
900 &payload[start..end]
901 );
902
903 start = end;
905 end = end + MAX_LOG_PAYLOAD_SIZE_BYTES;
906 chunk_idx += 1;
907 }
908
909 if start < payload.len() {
911 log::info!(
912 "Glean ping to URL: {} [Part {} of {}]\n{}",
913 path,
914 chunk_idx,
915 total_chunks,
916 &payload[start..]
917 );
918 }
919}
920
921#[cfg(not(target_os = "android"))]
923pub fn chunked_log_info(_path: &str, payload: &str) {
924 log::info!("{}", payload)
925}
926
927#[cfg(test)]
928mod test {
929 use std::thread;
930 use uuid::Uuid;
931
932 use super::*;
933 use crate::metrics::PingType;
934 use crate::{tests::new_glean, PENDING_PINGS_DIRECTORY};
935
936 const PATH: &str = "/submit/app_id/ping_name/schema_version/doc_id";
937
938 #[test]
939 fn doesnt_error_when_there_are_no_pending_pings() {
940 let (glean, _t) = new_glean(None);
941
942 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
945 }
946
947 #[test]
948 fn returns_ping_request_when_there_is_one() {
949 let (glean, dir) = new_glean(None);
950
951 let upload_manager = PingUploadManager::no_policy(dir.path());
952
953 upload_manager.enqueue_ping(
955 &glean,
956 PingPayload {
957 document_id: Uuid::new_v4().to_string(),
958 upload_path: PATH.into(),
959 json_body: "".into(),
960 headers: None,
961 body_has_info_sections: true,
962 ping_name: "ping-name".into(),
963 uploader_capabilities: vec![],
964 },
965 );
966
967 let task = upload_manager.get_upload_task(&glean, false);
970 assert!(task.is_upload());
971 }
972
973 #[test]
974 fn returns_as_many_ping_requests_as_there_are() {
975 let (glean, dir) = new_glean(None);
976
977 let upload_manager = PingUploadManager::no_policy(dir.path());
978
979 let n = 10;
981 for _ in 0..n {
982 upload_manager.enqueue_ping(
983 &glean,
984 PingPayload {
985 document_id: Uuid::new_v4().to_string(),
986 upload_path: PATH.into(),
987 json_body: "".into(),
988 headers: None,
989 body_has_info_sections: true,
990 ping_name: "ping-name".into(),
991 uploader_capabilities: vec![],
992 },
993 );
994 }
995
996 for _ in 0..n {
998 let task = upload_manager.get_upload_task(&glean, false);
999 assert!(task.is_upload());
1000 }
1001
1002 assert_eq!(
1004 upload_manager.get_upload_task(&glean, false),
1005 PingUploadTask::done()
1006 );
1007 }
1008
1009 #[test]
1010 fn limits_the_number_of_pings_when_there_is_rate_limiting() {
1011 let (glean, dir) = new_glean(None);
1012
1013 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1014
1015 let max_pings_per_interval = 10;
1017 upload_manager.set_rate_limiter(3, 10);
1018
1019 for _ in 0..max_pings_per_interval {
1021 upload_manager.enqueue_ping(
1022 &glean,
1023 PingPayload {
1024 document_id: Uuid::new_v4().to_string(),
1025 upload_path: PATH.into(),
1026 json_body: "".into(),
1027 headers: None,
1028 body_has_info_sections: true,
1029 ping_name: "ping-name".into(),
1030 uploader_capabilities: vec![],
1031 },
1032 );
1033 }
1034
1035 for _ in 0..max_pings_per_interval {
1037 let task = upload_manager.get_upload_task(&glean, false);
1038 assert!(task.is_upload());
1039 }
1040
1041 upload_manager.enqueue_ping(
1043 &glean,
1044 PingPayload {
1045 document_id: Uuid::new_v4().to_string(),
1046 upload_path: PATH.into(),
1047 json_body: "".into(),
1048 headers: None,
1049 body_has_info_sections: true,
1050 ping_name: "ping-name".into(),
1051 uploader_capabilities: vec![],
1052 },
1053 );
1054
1055 match upload_manager.get_upload_task(&glean, false) {
1057 PingUploadTask::Wait { time } => {
1058 thread::sleep(Duration::from_millis(time));
1060 }
1061 _ => panic!("Expected upload manager to return a wait task!"),
1062 };
1063
1064 let task = upload_manager.get_upload_task(&glean, false);
1065 assert!(task.is_upload());
1066 }
1067
1068 #[test]
1069 fn clearing_the_queue_works_correctly() {
1070 let (glean, dir) = new_glean(None);
1071
1072 let upload_manager = PingUploadManager::no_policy(dir.path());
1073
1074 for _ in 0..10 {
1076 upload_manager.enqueue_ping(
1077 &glean,
1078 PingPayload {
1079 document_id: Uuid::new_v4().to_string(),
1080 upload_path: PATH.into(),
1081 json_body: "".into(),
1082 headers: None,
1083 body_has_info_sections: true,
1084 ping_name: "ping-name".into(),
1085 uploader_capabilities: vec![],
1086 },
1087 );
1088 }
1089
1090 drop(upload_manager.clear_ping_queue());
1092
1093 assert_eq!(
1095 upload_manager.get_upload_task(&glean, false),
1096 PingUploadTask::done()
1097 );
1098 }
1099
1100 #[test]
1101 fn clearing_the_queue_doesnt_clear_deletion_request_pings() {
1102 let (mut glean, _t) = new_glean(None);
1103
1104 let ping_type = PingType::new(
1106 "test",
1107 true,
1108 true,
1109 true,
1110 true,
1111 true,
1112 vec![],
1113 vec![],
1114 true,
1115 vec![],
1116 );
1117 glean.register_ping_type(&ping_type);
1118
1119 let n = 10;
1121 for _ in 0..n {
1122 ping_type.submit_sync(&glean, None);
1123 }
1124
1125 glean
1126 .internal_pings
1127 .deletion_request
1128 .submit_sync(&glean, None);
1129
1130 drop(glean.upload_manager.clear_ping_queue());
1132
1133 let upload_task = glean.get_upload_task();
1134 match upload_task {
1135 PingUploadTask::Upload { request } => assert!(request.is_deletion_request()),
1136 _ => panic!("Expected upload manager to return the next request!"),
1137 }
1138
1139 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
1141 }
1142
1143 #[test]
1144 fn fills_up_queue_successfully_from_disk() {
1145 let (mut glean, dir) = new_glean(None);
1146
1147 let ping_type = PingType::new(
1149 "test",
1150 true,
1151 true,
1152 true,
1153 true,
1154 true,
1155 vec![],
1156 vec![],
1157 true,
1158 vec![],
1159 );
1160 glean.register_ping_type(&ping_type);
1161
1162 let n = 10;
1164 for _ in 0..n {
1165 ping_type.submit_sync(&glean, None);
1166 }
1167
1168 let upload_manager = PingUploadManager::no_policy(dir.path());
1170
1171 for _ in 0..n {
1173 let task = upload_manager.get_upload_task(&glean, false);
1174 assert!(task.is_upload());
1175 }
1176
1177 assert_eq!(
1179 upload_manager.get_upload_task(&glean, false),
1180 PingUploadTask::done()
1181 );
1182 }
1183
1184 #[test]
1185 fn processes_correctly_success_upload_response() {
1186 let (mut glean, dir) = new_glean(None);
1187
1188 let ping_type = PingType::new(
1190 "test",
1191 true,
1192 true,
1193 true,
1194 true,
1195 true,
1196 vec![],
1197 vec![],
1198 true,
1199 vec![],
1200 );
1201 glean.register_ping_type(&ping_type);
1202
1203 ping_type.submit_sync(&glean, None);
1205
1206 let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1208
1209 match glean.get_upload_task() {
1211 PingUploadTask::Upload { request } => {
1212 let document_id = request.document_id;
1214 glean.process_ping_upload_response(&document_id, UploadResult::http_status(200));
1215 assert!(!pending_pings_dir.join(document_id).exists());
1217 }
1218 _ => panic!("Expected upload manager to return the next request!"),
1219 }
1220
1221 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
1223 }
1224
1225 #[test]
1226 fn processes_correctly_client_error_upload_response() {
1227 let (mut glean, dir) = new_glean(None);
1228
1229 let ping_type = PingType::new(
1231 "test",
1232 true,
1233 true,
1234 true,
1235 true,
1236 true,
1237 vec![],
1238 vec![],
1239 true,
1240 vec![],
1241 );
1242 glean.register_ping_type(&ping_type);
1243
1244 ping_type.submit_sync(&glean, None);
1246
1247 let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1249
1250 match glean.get_upload_task() {
1252 PingUploadTask::Upload { request } => {
1253 let document_id = request.document_id;
1255 glean.process_ping_upload_response(&document_id, UploadResult::http_status(404));
1256 assert!(!pending_pings_dir.join(document_id).exists());
1258 }
1259 _ => panic!("Expected upload manager to return the next request!"),
1260 }
1261
1262 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
1264 }
1265
1266 #[test]
1267 fn processes_correctly_server_error_upload_response() {
1268 let (mut glean, _t) = new_glean(None);
1269
1270 let ping_type = PingType::new(
1272 "test",
1273 true,
1274 true,
1275 true,
1276 true,
1277 true,
1278 vec![],
1279 vec![],
1280 true,
1281 vec![],
1282 );
1283 glean.register_ping_type(&ping_type);
1284
1285 ping_type.submit_sync(&glean, None);
1287
1288 match glean.get_upload_task() {
1290 PingUploadTask::Upload { request } => {
1291 let document_id = request.document_id;
1293 glean.process_ping_upload_response(&document_id, UploadResult::http_status(500));
1294 match glean.get_upload_task() {
1296 PingUploadTask::Upload { request } => {
1297 assert_eq!(document_id, request.document_id);
1298 }
1299 _ => panic!("Expected upload manager to return the next request!"),
1300 }
1301 }
1302 _ => panic!("Expected upload manager to return the next request!"),
1303 }
1304
1305 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
1307 }
1308
1309 #[test]
1310 fn processes_correctly_unrecoverable_upload_response() {
1311 let (mut glean, dir) = new_glean(None);
1312
1313 let ping_type = PingType::new(
1315 "test",
1316 true,
1317 true,
1318 true,
1319 true,
1320 true,
1321 vec![],
1322 vec![],
1323 true,
1324 vec![],
1325 );
1326 glean.register_ping_type(&ping_type);
1327
1328 ping_type.submit_sync(&glean, None);
1330
1331 let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1333
1334 match glean.get_upload_task() {
1336 PingUploadTask::Upload { request } => {
1337 let document_id = request.document_id;
1339 glean.process_ping_upload_response(
1340 &document_id,
1341 UploadResult::unrecoverable_failure(),
1342 );
1343 assert!(!pending_pings_dir.join(document_id).exists());
1345 }
1346 _ => panic!("Expected upload manager to return the next request!"),
1347 }
1348
1349 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
1351 }
1352
1353 #[test]
1354 fn new_pings_are_added_while_upload_in_progress() {
1355 let (glean, dir) = new_glean(None);
1356
1357 let upload_manager = PingUploadManager::no_policy(dir.path());
1358
1359 let doc1 = Uuid::new_v4().to_string();
1360 let path1 = format!("/submit/app_id/test-ping/1/{}", doc1);
1361
1362 let doc2 = Uuid::new_v4().to_string();
1363 let path2 = format!("/submit/app_id/test-ping/1/{}", doc2);
1364
1365 upload_manager.enqueue_ping(
1367 &glean,
1368 PingPayload {
1369 document_id: doc1.clone(),
1370 upload_path: path1,
1371 json_body: "".into(),
1372 headers: None,
1373 body_has_info_sections: true,
1374 ping_name: "test-ping".into(),
1375 uploader_capabilities: vec![],
1376 },
1377 );
1378
1379 let req = match upload_manager.get_upload_task(&glean, false) {
1381 PingUploadTask::Upload { request } => request,
1382 _ => panic!("Expected upload manager to return the next request!"),
1383 };
1384 assert_eq!(doc1, req.document_id);
1385
1386 upload_manager.enqueue_ping(
1388 &glean,
1389 PingPayload {
1390 document_id: doc2.clone(),
1391 upload_path: path2,
1392 json_body: "".into(),
1393 headers: None,
1394 body_has_info_sections: true,
1395 ping_name: "test-ping".into(),
1396 uploader_capabilities: vec![],
1397 },
1398 );
1399
1400 upload_manager.process_ping_upload_response(
1402 &glean,
1403 &req.document_id,
1404 UploadResult::http_status(200),
1405 );
1406
1407 let req = match upload_manager.get_upload_task(&glean, false) {
1409 PingUploadTask::Upload { request } => request,
1410 _ => panic!("Expected upload manager to return the next request!"),
1411 };
1412 assert_eq!(doc2, req.document_id);
1413
1414 upload_manager.process_ping_upload_response(
1416 &glean,
1417 &req.document_id,
1418 UploadResult::http_status(200),
1419 );
1420
1421 assert_eq!(
1423 upload_manager.get_upload_task(&glean, false),
1424 PingUploadTask::done()
1425 );
1426 }
1427
1428 #[test]
1429 fn adds_debug_view_header_to_requests_when_tag_is_set() {
1430 let (mut glean, _t) = new_glean(None);
1431
1432 glean.set_debug_view_tag("valid-tag");
1433
1434 let ping_type = PingType::new(
1436 "test",
1437 true,
1438 true,
1439 true,
1440 true,
1441 true,
1442 vec![],
1443 vec![],
1444 true,
1445 vec![],
1446 );
1447 glean.register_ping_type(&ping_type);
1448
1449 ping_type.submit_sync(&glean, None);
1451
1452 match glean.get_upload_task() {
1454 PingUploadTask::Upload { request } => {
1455 assert_eq!(request.headers.get("X-Debug-ID").unwrap(), "valid-tag")
1456 }
1457 _ => panic!("Expected upload manager to return the next request!"),
1458 }
1459 }
1460
1461 #[test]
1462 fn duplicates_are_not_enqueued() {
1463 let (glean, dir) = new_glean(None);
1464
1465 let upload_manager = PingUploadManager::no_policy(dir.path());
1468
1469 let doc_id = Uuid::new_v4().to_string();
1470 let path = format!("/submit/app_id/test-ping/1/{}", doc_id);
1471
1472 upload_manager.enqueue_ping(
1474 &glean,
1475 PingPayload {
1476 document_id: doc_id.clone(),
1477 upload_path: path.clone(),
1478 json_body: "".into(),
1479 headers: None,
1480 body_has_info_sections: true,
1481 ping_name: "test-ping".into(),
1482 uploader_capabilities: vec![],
1483 },
1484 );
1485 upload_manager.enqueue_ping(
1486 &glean,
1487 PingPayload {
1488 document_id: doc_id,
1489 upload_path: path,
1490 json_body: "".into(),
1491 headers: None,
1492 body_has_info_sections: true,
1493 ping_name: "test-ping".into(),
1494 uploader_capabilities: vec![],
1495 },
1496 );
1497
1498 let task = upload_manager.get_upload_task(&glean, false);
1500 assert!(task.is_upload());
1501
1502 assert_eq!(
1504 upload_manager.get_upload_task(&glean, false),
1505 PingUploadTask::done()
1506 );
1507 }
1508
1509 #[test]
1510 fn maximum_of_recoverable_errors_is_enforced_for_uploading_window() {
1511 let (mut glean, dir) = new_glean(None);
1512
1513 let ping_type = PingType::new(
1515 "test",
1516 true,
1517 true,
1518 true,
1519 true,
1520 true,
1521 vec![],
1522 vec![],
1523 true,
1524 vec![],
1525 );
1526 glean.register_ping_type(&ping_type);
1527
1528 let n = 5;
1530 for _ in 0..n {
1531 ping_type.submit_sync(&glean, None);
1532 }
1533
1534 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1535
1536 let max_recoverable_failures = 3;
1538 upload_manager
1539 .policy
1540 .set_max_recoverable_failures(Some(max_recoverable_failures));
1541
1542 for _ in 0..max_recoverable_failures {
1544 match upload_manager.get_upload_task(&glean, false) {
1545 PingUploadTask::Upload { request } => {
1546 upload_manager.process_ping_upload_response(
1547 &glean,
1548 &request.document_id,
1549 UploadResult::recoverable_failure(),
1550 );
1551 }
1552 _ => panic!("Expected upload manager to return the next request!"),
1553 }
1554 }
1555
1556 assert_eq!(
1559 upload_manager.get_upload_task(&glean, false),
1560 PingUploadTask::done()
1561 );
1562
1563 for _ in 0..n {
1565 let task = upload_manager.get_upload_task(&glean, false);
1566 assert!(task.is_upload());
1567 }
1568 }
1569
1570 #[test]
1571 fn quota_is_enforced_when_enqueueing_cached_pings() {
1572 let (mut glean, dir) = new_glean(None);
1573
1574 let ping_type = PingType::new(
1576 "test",
1577 true,
1578 true,
1579 true,
1580 true,
1581 true,
1582 vec![],
1583 vec![],
1584 true,
1585 vec![],
1586 );
1587 glean.register_ping_type(&ping_type);
1588
1589 let n = 10;
1591 for _ in 0..n {
1592 ping_type.submit_sync(&glean, None);
1593 }
1594
1595 let directory_manager = PingDirectoryManager::new(dir.path());
1596 let pending_pings = directory_manager.process_dirs().pending_pings;
1597 let (_, newest_ping) = &pending_pings.last().unwrap();
1600 let PingPayload {
1601 document_id: newest_ping_id,
1602 ..
1603 } = &newest_ping;
1604
1605 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1607
1608 upload_manager
1615 .policy
1616 .set_max_pending_pings_directory_size(Some(500));
1617
1618 match upload_manager.get_upload_task(&glean, false) {
1622 PingUploadTask::Upload { request } => assert_eq!(&request.document_id, newest_ping_id),
1623 _ => panic!("Expected upload manager to return the next request!"),
1624 }
1625
1626 assert_eq!(
1629 upload_manager.get_upload_task(&glean, false),
1630 PingUploadTask::done()
1631 );
1632
1633 assert_eq!(
1635 n - 1,
1636 upload_manager
1637 .upload_metrics
1638 .deleted_pings_after_quota_hit
1639 .get_value(&glean, Some("metrics"))
1640 .unwrap()
1641 );
1642 assert_eq!(
1643 n,
1644 upload_manager
1645 .upload_metrics
1646 .pending_pings
1647 .get_value(&glean, Some("metrics"))
1648 .unwrap()
1649 );
1650 }
1651
1652 #[test]
1653 fn number_quota_is_enforced_when_enqueueing_cached_pings() {
1654 let (mut glean, dir) = new_glean(None);
1655
1656 let ping_type = PingType::new(
1658 "test",
1659 true,
1660 true,
1661 true,
1662 true,
1663 true,
1664 vec![],
1665 vec![],
1666 true,
1667 vec![],
1668 );
1669 glean.register_ping_type(&ping_type);
1670
1671 let count_quota = 3;
1673 let n = 10;
1675
1676 for _ in 0..n {
1678 ping_type.submit_sync(&glean, None);
1679 }
1680
1681 let directory_manager = PingDirectoryManager::new(dir.path());
1682 let pending_pings = directory_manager.process_dirs().pending_pings;
1683 let expected_pings = pending_pings
1686 .iter()
1687 .rev()
1688 .take(count_quota)
1689 .map(|(_, ping)| ping.document_id.clone())
1690 .collect::<Vec<_>>();
1691
1692 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1694
1695 upload_manager
1696 .policy
1697 .set_max_pending_pings_count(Some(count_quota as u64));
1698
1699 for ping_id in expected_pings.iter().rev() {
1703 match upload_manager.get_upload_task(&glean, false) {
1704 PingUploadTask::Upload { request } => assert_eq!(&request.document_id, ping_id),
1705 _ => panic!("Expected upload manager to return the next request!"),
1706 }
1707 }
1708
1709 assert_eq!(
1712 upload_manager.get_upload_task(&glean, false),
1713 PingUploadTask::done()
1714 );
1715
1716 assert_eq!(
1718 (n - count_quota) as i32,
1719 upload_manager
1720 .upload_metrics
1721 .deleted_pings_after_quota_hit
1722 .get_value(&glean, Some("metrics"))
1723 .unwrap()
1724 );
1725 assert_eq!(
1726 n as i32,
1727 upload_manager
1728 .upload_metrics
1729 .pending_pings
1730 .get_value(&glean, Some("metrics"))
1731 .unwrap()
1732 );
1733 }
1734
1735 #[test]
1736 fn size_and_count_quota_work_together_size_first() {
1737 let (mut glean, dir) = new_glean(None);
1738
1739 let ping_type = PingType::new(
1741 "test",
1742 true,
1743 true,
1744 true,
1745 true,
1746 true,
1747 vec![],
1748 vec![],
1749 true,
1750 vec![],
1751 );
1752 glean.register_ping_type(&ping_type);
1753
1754 let expected_number_of_pings = 3;
1755 let n = 10;
1757
1758 for _ in 0..n {
1760 ping_type.submit_sync(&glean, None);
1761 }
1762
1763 let directory_manager = PingDirectoryManager::new(dir.path());
1764 let pending_pings = directory_manager.process_dirs().pending_pings;
1765 let expected_pings = pending_pings
1768 .iter()
1769 .rev()
1770 .take(expected_number_of_pings)
1771 .map(|(_, ping)| ping.document_id.clone())
1772 .collect::<Vec<_>>();
1773
1774 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1776
1777 upload_manager
1780 .policy
1781 .set_max_pending_pings_directory_size(Some(1300));
1782 upload_manager.policy.set_max_pending_pings_count(Some(5));
1783
1784 for ping_id in expected_pings.iter().rev() {
1788 match upload_manager.get_upload_task(&glean, false) {
1789 PingUploadTask::Upload { request } => assert_eq!(&request.document_id, ping_id),
1790 _ => panic!("Expected upload manager to return the next request!"),
1791 }
1792 }
1793
1794 assert_eq!(
1797 upload_manager.get_upload_task(&glean, false),
1798 PingUploadTask::done()
1799 );
1800
1801 assert_eq!(
1803 (n - expected_number_of_pings) as i32,
1804 upload_manager
1805 .upload_metrics
1806 .deleted_pings_after_quota_hit
1807 .get_value(&glean, Some("metrics"))
1808 .unwrap()
1809 );
1810 assert_eq!(
1811 n as i32,
1812 upload_manager
1813 .upload_metrics
1814 .pending_pings
1815 .get_value(&glean, Some("metrics"))
1816 .unwrap()
1817 );
1818 assert_eq!(
1820 (n - expected_number_of_pings) as i32,
1821 upload_manager
1822 .upload_metrics
1823 .pending_pings_deleted
1824 .get("size_quota")
1825 .get_value(&glean, Some("health"))
1826 .unwrap()
1827 );
1828 assert!(upload_manager
1829 .upload_metrics
1830 .pending_pings_deleted
1831 .get("count_quota")
1832 .get_value(&glean, Some("health"))
1833 .is_none());
1834 }
1835
1836 #[test]
1837 fn size_and_count_quota_work_together_count_first() {
1838 let (mut glean, dir) = new_glean(None);
1839
1840 let ping_type = PingType::new(
1842 "test",
1843 true,
1844 true,
1845 true,
1846 true,
1847 true,
1848 vec![],
1849 vec![],
1850 true,
1851 vec![],
1852 );
1853 glean.register_ping_type(&ping_type);
1854
1855 let expected_number_of_pings = 2;
1856 let n = 10;
1858
1859 for _ in 0..n {
1861 ping_type.submit_sync(&glean, None);
1862 }
1863
1864 let directory_manager = PingDirectoryManager::new(dir.path());
1865 let pending_pings = directory_manager.process_dirs().pending_pings;
1866 let expected_pings = pending_pings
1869 .iter()
1870 .rev()
1871 .take(expected_number_of_pings)
1872 .map(|(_, ping)| ping.document_id.clone())
1873 .collect::<Vec<_>>();
1874
1875 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1877
1878 upload_manager
1880 .policy
1881 .set_max_pending_pings_directory_size(Some(100_000));
1882 upload_manager.policy.set_max_pending_pings_count(Some(2));
1883
1884 for ping_id in expected_pings.iter().rev() {
1888 match upload_manager.get_upload_task(&glean, false) {
1889 PingUploadTask::Upload { request } => assert_eq!(&request.document_id, ping_id),
1890 _ => panic!("Expected upload manager to return the next request!"),
1891 }
1892 }
1893
1894 assert_eq!(
1897 upload_manager.get_upload_task(&glean, false),
1898 PingUploadTask::done()
1899 );
1900
1901 assert_eq!(
1903 (n - expected_number_of_pings) as i32,
1904 upload_manager
1905 .upload_metrics
1906 .deleted_pings_after_quota_hit
1907 .get_value(&glean, Some("metrics"))
1908 .unwrap()
1909 );
1910 assert_eq!(
1911 n as i32,
1912 upload_manager
1913 .upload_metrics
1914 .pending_pings
1915 .get_value(&glean, Some("metrics"))
1916 .unwrap()
1917 );
1918 assert_eq!(
1920 (n - expected_number_of_pings) as i32,
1921 upload_manager
1922 .upload_metrics
1923 .pending_pings_deleted
1924 .get("count_quota")
1925 .get_value(&glean, Some("health"))
1926 .unwrap()
1927 );
1928 assert!(upload_manager
1929 .upload_metrics
1930 .pending_pings_deleted
1931 .get("size_quota")
1932 .get_value(&glean, Some("health"))
1933 .is_none());
1934 }
1935
1936 #[test]
1937 fn pending_pings_deleted_is_not_recorded_when_quota_not_hit() {
1938 let (mut glean, dir) = new_glean(None);
1939
1940 let ping_type = PingType::new(
1941 "test",
1942 true,
1943 true,
1944 true,
1945 true,
1946 true,
1947 vec![],
1948 vec![],
1949 true,
1950 vec![],
1951 );
1952 glean.register_ping_type(&ping_type);
1953
1954 for _ in 0..3 {
1956 ping_type.submit_sync(&glean, None);
1957 }
1958
1959 let mut upload_manager = PingUploadManager::no_policy(dir.path());
1960 upload_manager.policy.set_max_pending_pings_count(Some(10));
1961 upload_manager
1962 .policy
1963 .set_max_pending_pings_directory_size(Some(1024 * 1024));
1964
1965 upload_manager.get_upload_task(&glean, false);
1966
1967 assert!(upload_manager
1968 .upload_metrics
1969 .pending_pings_deleted
1970 .get("count_quota")
1971 .get_value(&glean, Some("health"))
1972 .is_none());
1973 assert!(upload_manager
1974 .upload_metrics
1975 .pending_pings_deleted
1976 .get("size_quota")
1977 .get_value(&glean, Some("health"))
1978 .is_none());
1979 }
1980
1981 #[test]
1982 fn pending_pings_config_overrides_are_applied() {
1983 let (_, dir) = new_glean(None);
1984
1985 let mut upload_manager = PingUploadManager::new(dir.path(), "test");
1986
1987 let custom_count: u64 = 42;
1988 let custom_size: u64 = 999_999;
1989 upload_manager.set_max_pending_pings_count(custom_count);
1990 upload_manager.set_max_pending_pings_directory_size(custom_size);
1991
1992 assert_eq!(
1993 custom_count,
1994 upload_manager.policy.max_pending_pings_count()
1995 );
1996 assert_eq!(
1997 custom_size,
1998 upload_manager.policy.max_pending_pings_directory_size()
1999 );
2000 }
2001
2002 #[test]
2003 fn maximum_wait_attemps_is_enforced() {
2004 let (glean, dir) = new_glean(None);
2005
2006 let mut upload_manager = PingUploadManager::no_policy(dir.path());
2007
2008 let max_wait_attempts = 3;
2010 upload_manager
2011 .policy
2012 .set_max_wait_attempts(Some(max_wait_attempts));
2013
2014 let secs_per_interval = 5;
2020 let max_pings_per_interval = 1;
2021 upload_manager.set_rate_limiter(secs_per_interval, max_pings_per_interval);
2022
2023 upload_manager.enqueue_ping(
2025 &glean,
2026 PingPayload {
2027 document_id: Uuid::new_v4().to_string(),
2028 upload_path: PATH.into(),
2029 json_body: "".into(),
2030 headers: None,
2031 body_has_info_sections: true,
2032 ping_name: "ping-name".into(),
2033 uploader_capabilities: vec![],
2034 },
2035 );
2036 upload_manager.enqueue_ping(
2037 &glean,
2038 PingPayload {
2039 document_id: Uuid::new_v4().to_string(),
2040 upload_path: PATH.into(),
2041 json_body: "".into(),
2042 headers: None,
2043 body_has_info_sections: true,
2044 ping_name: "ping-name".into(),
2045 uploader_capabilities: vec![],
2046 },
2047 );
2048
2049 match upload_manager.get_upload_task(&glean, false) {
2051 PingUploadTask::Upload { .. } => {}
2052 _ => panic!("Expected upload manager to return the next request!"),
2053 }
2054
2055 for _ in 0..max_wait_attempts {
2059 let task = upload_manager.get_upload_task(&glean, false);
2060 assert!(task.is_wait());
2061 }
2062
2063 assert_eq!(
2066 upload_manager.get_upload_task(&glean, false),
2067 PingUploadTask::done()
2068 );
2069
2070 thread::sleep(Duration::from_secs(secs_per_interval));
2072
2073 let task = upload_manager.get_upload_task(&glean, false);
2075 assert!(task.is_upload());
2076
2077 assert_eq!(
2079 upload_manager.get_upload_task(&glean, false),
2080 PingUploadTask::done()
2081 );
2082 }
2083
2084 #[test]
2085 fn wait_task_contains_expected_wait_time_when_pending_pings_dir_not_processed_yet() {
2086 let (glean, dir) = new_glean(None);
2087 let upload_manager = PingUploadManager::new(dir.path(), "test");
2088 match upload_manager.get_upload_task(&glean, false) {
2089 PingUploadTask::Wait { time } => {
2090 assert_eq!(time, WAIT_TIME_FOR_PING_PROCESSING);
2091 }
2092 _ => panic!("Expected upload manager to return a wait task!"),
2093 };
2094 }
2095
2096 #[test]
2097 fn cannot_enqueue_ping_while_its_being_processed() {
2098 let (glean, dir) = new_glean(None);
2099
2100 let upload_manager = PingUploadManager::no_policy(dir.path());
2101
2102 let identifier = &Uuid::new_v4();
2104 let ping = PingPayload {
2105 document_id: identifier.to_string(),
2106 upload_path: PATH.into(),
2107 json_body: "".into(),
2108 headers: None,
2109 body_has_info_sections: true,
2110 ping_name: "ping-name".into(),
2111 uploader_capabilities: vec![],
2112 };
2113 upload_manager.enqueue_ping(&glean, ping);
2114 assert!(upload_manager.get_upload_task(&glean, false).is_upload());
2115
2116 let ping = PingPayload {
2118 document_id: identifier.to_string(),
2119 upload_path: PATH.into(),
2120 json_body: "".into(),
2121 headers: None,
2122 body_has_info_sections: true,
2123 ping_name: "ping-name".into(),
2124 uploader_capabilities: vec![],
2125 };
2126 upload_manager.enqueue_ping(&glean, ping);
2127
2128 assert_eq!(
2130 upload_manager.get_upload_task(&glean, false),
2131 PingUploadTask::done()
2132 );
2133
2134 upload_manager.process_ping_upload_response(
2136 &glean,
2137 &identifier.to_string(),
2138 UploadResult::http_status(200),
2139 );
2140 }
2141
2142 #[cfg(feature = "sqlite")]
2143 #[test]
2144 fn stores_pings_during_submission_and_upload_if_enabled() {
2145 let (mut glean, _t) = new_glean(None);
2146 glean.set_store_submitted_pings_enabled(true);
2147
2148 let ping_type = PingType::new(
2150 "test",
2151 true,
2152 true,
2153 true,
2154 true,
2155 true,
2156 vec![],
2157 vec![],
2158 true,
2159 vec![],
2160 );
2161 glean.register_ping_type(&ping_type);
2162
2163 ping_type.submit_sync(&glean, None);
2165
2166 let pings = glean.storage().get_all_submitted_pings();
2167 assert_eq!(pings.len(), 1);
2168 let ping = pings.first().unwrap();
2169 assert!(ping.submitted_date.0 <= Utc::now());
2170 assert!(ping.uploaded_date.is_none());
2171
2172 match glean.get_upload_task() {
2174 PingUploadTask::Upload { request } => {
2175 let document_id = request.document_id;
2177 glean.process_ping_upload_response(&document_id, UploadResult::http_status(200));
2178 }
2179 _ => panic!("Expected upload manager to return the next request!"),
2180 }
2181
2182 let pings = glean.storage().get_all_submitted_pings();
2183 assert_eq!(pings.len(), 1);
2184 let ping = pings.first().unwrap();
2185 assert!(ping.submitted_date.0 <= Utc::now());
2186 assert!(ping.uploaded_date.is_some());
2187
2188 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
2190 }
2191
2192 #[cfg(feature = "sqlite")]
2193 #[test]
2194 fn stores_pings_during_submission_and_marks_as_upload_failed_when_appropriate() {
2195 let (mut glean, _t) = new_glean(None);
2196 glean.set_store_submitted_pings_enabled(true);
2197
2198 let ping_type = PingType::new(
2200 "test",
2201 true,
2202 true,
2203 true,
2204 true,
2205 true,
2206 vec![],
2207 vec![],
2208 true,
2209 vec![],
2210 );
2211 glean.register_ping_type(&ping_type);
2212
2213 ping_type.submit_sync(&glean, None);
2215
2216 let pings = glean.storage().get_all_submitted_pings();
2217 assert_eq!(pings.len(), 1);
2218 let ping = pings.first().unwrap();
2219 assert!(ping.submitted_date.0 <= Utc::now());
2220 assert!(ping.uploaded_date.is_none());
2221
2222 match glean.get_upload_task() {
2224 PingUploadTask::Upload { request } => {
2225 let document_id = request.document_id;
2227 glean.process_ping_upload_response(&document_id, UploadResult::http_status(400));
2228 }
2229 _ => panic!("Expected upload manager to return the next request!"),
2230 }
2231
2232 let pings = glean.storage().get_all_submitted_pings();
2233 assert_eq!(pings.len(), 1);
2234 let ping = pings.first().unwrap();
2235 assert!(ping.submitted_date.0 <= Utc::now());
2236 assert!(ping.upload_failed.is_some());
2237 assert!(ping.uploaded_date.is_none());
2238
2239 assert_eq!(glean.get_upload_task(), PingUploadTask::done());
2241 }
2242}