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