Skip to main content

glean_core/upload/
mod.rs

1// This Source Code Form is subject to the terms of the Mozilla Public
2// License, v. 2.0. If a copy of the MPL was not distributed with this
3// file, You can obtain one at https://mozilla.org/MPL/2.0/.
4
5//! Manages the pending pings queue and directory.
6//!
7//! * Keeps track of pending pings, loading any unsent ping from disk on startup;
8//! * Exposes [`get_upload_task`](PingUploadManager::get_upload_task) API for
9//!   the platform layer to request next upload task;
10//! * Exposes
11//!   [`process_ping_upload_response`](PingUploadManager::process_ping_upload_response)
12//!   API to check the HTTP response from the ping upload and either delete the
13//!   corresponding ping from disk or re-enqueue it for sending.
14
15use 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; // in milliseconds
46
47#[derive(Debug, MallocSizeOf)]
48struct RateLimiter {
49    /// The instant the current interval has started.
50    started: Option<Instant>,
51    /// The count for the current interval.
52    count: u32,
53    /// The duration of each interval.
54    interval: Duration,
55    /// The maximum count per interval.
56    max_count: u32,
57}
58
59/// An enum to represent the current state of the RateLimiter.
60#[derive(PartialEq)]
61enum RateLimiterState {
62    /// The RateLimiter has not reached the maximum count and is still incrementing.
63    Incrementing,
64    /// The RateLimiter has reached the maximum count for the  current interval.
65    ///
66    /// This variant contains the remaining time (in milliseconds)
67    /// until the rate limiter is not throttled anymore.
68    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    // The counter should reset if
91    //
92    // 1. It has never started;
93    // 2. It has been started more than the interval time ago;
94    // 3. Something goes wrong while trying to calculate the elapsed time since the last reset.
95    fn should_reset(&self) -> bool {
96        if self.started.is_none() {
97            return true;
98        }
99
100        // Safe unwrap, we already stated that `self.started` is not `None` above.
101        if self.elapsed() > self.interval {
102            return true;
103        }
104
105        false
106    }
107
108    /// Tries to increment the internal counter.
109    ///
110    /// # Returns
111    ///
112    /// The current state of the RateLimiter.
113    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            // Note that `remining` can't be a negative number because we just called `reset`,
120            // which will check if it is and reset if so.
121            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/// An enum representing the possible upload tasks to be performed by an uploader.
135///
136/// When asking for the next ping request to upload,
137/// the requester may receive one out of three possible tasks.
138#[derive(PartialEq, Eq, Debug)]
139pub enum PingUploadTask {
140    /// An upload task
141    Upload {
142        /// The ping request for upload
143        /// See [`PingRequest`](struct.PingRequest.html) for more information.
144        request: PingRequest,
145    },
146
147    /// A flag signaling that the pending pings directories are not done being processed,
148    /// thus the requester should wait and come back later.
149    Wait {
150        /// The time in milliseconds
151        /// the requester should wait before requesting a new task.
152        time: u64,
153    },
154
155    /// A flag signaling that requester doesn't need to request any more upload tasks at this moment.
156    ///
157    /// There are three possibilities for this scenario:
158    /// * Pending pings queue is empty, no more pings to request;
159    /// * Requester has gotten more than MAX_WAIT_ATTEMPTS (3, by default) `PingUploadTask::Wait` responses in a row;
160    /// * Requester has reported more than MAX_RECOVERABLE_FAILURES_PER_UPLOADING_WINDOW
161    ///   recoverable upload failures on the same uploading window (see below)
162    ///   and should stop requesting at this moment.
163    ///
164    /// An "uploading window" starts when a requester gets a new
165    /// `PingUploadTask::Upload(PingRequest)` response and finishes when they
166    /// finally get a `PingUploadTask::Done` or `PingUploadTask::Wait` response.
167    Done {
168        #[doc(hidden)]
169        /// Unused field. Required because UniFFI can't handle variants without fields.
170        unused: i8,
171    },
172}
173
174impl PingUploadTask {
175    /// Whether the current task is an upload task.
176    pub fn is_upload(&self) -> bool {
177        matches!(self, PingUploadTask::Upload { .. })
178    }
179
180    /// Whether the current task is wait task.
181    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/// Manages the pending pings queue and directory.
191#[derive(Debug)]
192pub struct PingUploadManager {
193    /// A FIFO queue storing a `PingRequest` for each pending ping.
194    queue: RwLock<VecDeque<PingRequest>>,
195    /// A manager for the pending pings directories.
196    directory_manager: PingDirectoryManager,
197    /// A flag signaling if we are done processing the pending pings directories.
198    processed_pending_pings: Arc<AtomicBool>,
199    /// A vector to store the pending pings processed off-thread.
200    cached_pings: Arc<RwLock<PingPayloadsByDirectory>>,
201    /// The number of upload failures for the current uploading window.
202    recoverable_failure_count: AtomicU32,
203    /// The number or times in a row a user has received a `PingUploadTask::Wait` response.
204    wait_attempt_count: AtomicU32,
205    /// A ping counter to help rate limit the ping uploads.
206    ///
207    /// To keep resource usage in check,
208    /// we may want to limit the amount of pings sent in a given interval.
209    rate_limiter: Option<RwLock<RateLimiter>>,
210    /// The name of the programming language used by the binding creating this instance of PingUploadManager.
211    ///
212    /// This will be used to build the value User-Agent header for each ping request.
213    language_binding_name: String,
214    /// Metrics related to ping uploading.
215    upload_metrics: UploadMetrics,
216    /// Policies for ping storage, uploading and requests.
217    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                    // SAFETY: The front element is a valid interior pointer and thus valid to pass
229                    // to an external function.
230                    unsafe { ops.malloc_enclosing_size_of(front) }
231                } else {
232                    // This assumes that no memory is allocated when the VecDeque is empty.
233                    0
234                }
235            } else {
236                // If `ops` can't estimate the size of a pointer,
237                // we can estimate the allocation size by the size of each element and the
238                // allocated capacity.
239                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>() // Allocated inside the `self.processed_pending_pings` `Arc`.
246            + 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    /// Creates a new PingUploadManager.
264    ///
265    /// # Arguments
266    ///
267    /// * `data_path` - Path to the pending pings directory.
268    /// * `language_binding_name` - The name of the language binding calling this managers instance.
269    ///
270    /// # Panics
271    ///
272    /// Will panic if unable to spawn a new thread.
273    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    /// Spawns a new thread and processes the pending pings directories,
290    /// filling up the queue with whatever pings are in there.
291    ///
292    /// # Returns
293    ///
294    /// The `JoinHandle` to the spawned thread
295    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                // Be sure to drop local_cached_pings lock before triggering upload.
305                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    /// Creates a new upload manager with no limitations, for tests.
328    #[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        // Disable all policies for tests, if necessary individuals tests can re-enable them.
333        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        // When building for tests, always scan the pending pings directories and do it sync.
342        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    /// Attempts to build a ping request from a ping file payload.
363    ///
364    /// Returns the `PingRequest` or `None` if unable to build,
365    /// in which case it will delete the ping file and record an error.
366    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                // Record the error.
398                // Currently the only possible error is PingBodyOverflow.
399                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    /// Enqueue a ping for upload.
411    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        // Checks if a ping with this `document_id` is already enqueued.
423        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    /// Enqueues pings that might have been cached.
457    ///
458    /// The size of the PENDING_PINGS_DIRECTORY directory will be calculated
459    /// (by accumulating each ping's size in that directory)
460    /// and in case we exceed the quota, defined by the `quota` arg,
461    /// outstanding pings get deleted and are not enqueued.
462    ///
463    /// The size of the DELETION_REQUEST_PINGS_DIRECTORY will not be calculated
464    /// and no deletion-request pings will be deleted. Deletion request pings
465    /// are not very common and usually don't contain any data,
466    /// we don't expect that directory to ever reach quota.
467    /// Most importantly, we don't want to ever delete deletion-request pings.
468    ///
469    /// # Arguments
470    ///
471    /// * `glean` - The Glean object holding the database.
472    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            // The pending pings vector is sorted by date in ascending order (oldest -> newest).
498            // We need to calculate the size of the pending pings directory
499            // and delete the **oldest** pings in case quota is reached.
500            // Thus, we reverse the order of the pending pings vector,
501            // so that we iterate in descending order (newest -> oldest).
502            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                // We don't want to spam the log for every ping over the quota.
508                // Size is checked first; if both limits are exceeded simultaneously,
509                // size_quota takes precedence as the recorded reason.
510                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                // Once we reach the number of allowed pings we start deleting,
520                // no matter what size.
521                // We already log this before the loop.
522                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            // After calculating the size of the pending pings directory,
543            // we record the calculated number and reverse the pings array back for enqueueing.
544            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            // Enqueue the remaining pending pings and
550            // enqueue all deletion-request pings.
551            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    /// Adds rate limiting capability to this upload manager.
563    ///
564    /// The rate limiter will limit the amount of calls to `get_upload_task` per interval.
565    ///
566    /// Setting this will restart count and timer in case there was a previous rate limiter set
567    /// (e.g. if we have reached the current limit and call this function, we start counting again
568    /// and the caller is allowed to asks for tasks).
569    ///
570    /// # Arguments
571    ///
572    /// * `interval` - the amount of seconds in each rate limiting window.
573    /// * `max_tasks` - the maximum amount of task requests allowed per interval.
574    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    /// Reads a ping file, creates a `PingRequest` and adds it to the queue.
590    ///
591    /// Duplicate requests won't be added.
592    ///
593    /// # Arguments
594    ///
595    /// * `glean` - The Glean object holding the database.
596    /// * `document_id` - The UUID of the ping in question.
597    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    /// Clears the pending pings queue, leaves the deletion-request pings.
604    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        // Helper to decide whether to return PingUploadTask::Wait or PingUploadTask::Done.
621        //
622        // We want to limit the amount of PingUploadTask::Wait returned in a row,
623        // in case we reach MAX_WAIT_ATTEMPTS we want to actually return PingUploadTask::Done.
624        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        // This is a no-op in case there are no cached pings.
641        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                    // Synchronous timer starts.
684                    // We're in the uploader thread anyway.
685                    // But also: No data is stored on disk.
686                    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                // Adding the `Date` header just before actual upload happens.
695                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    /// Gets the next `PingUploadTask`.
709    ///
710    /// # Arguments
711    ///
712    /// * `glean` - The Glean object holding the database.
713    /// * `log_ping` - Whether to log the ping before returning.
714    ///
715    /// # Returns
716    ///
717    /// The next [`PingUploadTask`](enum.PingUploadTask.html).
718    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    /// Processes the response from an attempt to upload a ping.
733    ///
734    /// Based on the HTTP status of said response,
735    /// the possible outcomes are:
736    ///
737    /// * **200 - 299 Success**
738    ///   Any status on the 2XX range is considered a succesful upload,
739    ///   which means the corresponding ping file can be deleted.
740    ///   _Known 2XX status:_
741    ///   * 200 - OK. Request accepted into the pipeline.
742    ///
743    /// * **400 - 499 Unrecoverable error**
744    ///   Any status on the 4XX range means something our client did is not correct.
745    ///   It is unlikely that the client is going to recover from this by retrying,
746    ///   so in this case the corresponding ping file can also be deleted.
747    ///   _Known 4XX status:_
748    ///   * 404 - not found - POST/PUT to an unknown namespace
749    ///   * 405 - wrong request type (anything other than POST/PUT)
750    ///   * 411 - missing content-length header
751    ///   * 413 - request body too large Note that if we have badly-behaved clients that
752    ///           retry on 4XX, we should send back 202 on body/path too long).
753    ///   * 414 - request path too long (See above)
754    ///
755    /// * **Any other error**
756    ///   For any other error, a warning is logged and the ping is re-enqueued.
757    ///   _Known other errors:_
758    ///   * 500 - internal error
759    ///
760    /// # Note
761    ///
762    /// The disk I/O performed by this function is not done off-thread,
763    /// as it is expected to be called off-thread by the platform.
764    ///
765    /// # Arguments
766    ///
767    /// * `glean` - The Glean object holding the database.
768    /// * `document_id` - The UUID of the ping in question.
769    /// * `status` - The HTTP status of the response.
770    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/// Splits log message into chunks on Android.
861#[cfg(target_os = "android")]
862pub fn chunked_log_info(path: &str, payload: &str) {
863    // Since the logcat ring buffer size is configurable, but it's 'max payload' size is not,
864    // we must break apart long pings into chunks no larger than the max payload size of 4076b.
865    // We leave some head space for our prefix.
866    const MAX_LOG_PAYLOAD_SIZE_BYTES: usize = 4000;
867
868    // If the length of the ping will fit within one logcat payload, then we can
869    // short-circuit here and avoid some overhead, otherwise we must split up the
870    // message so that we don't truncate it.
871    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    // Otherwise we break it apart into chunks of smaller size,
877    // prefixing it with the path and a counter.
878    let mut start = 0;
879    let mut end = MAX_LOG_PAYLOAD_SIZE_BYTES;
880    let mut chunk_idx = 1;
881    // Might be off by 1 on edge cases, but do we really care?
882    let total_chunks = payload.len() / MAX_LOG_PAYLOAD_SIZE_BYTES + 1;
883
884    while end < payload.len() {
885        // Find char boundary from the end.
886        // It's UTF-8, so it is within 4 bytes from here.
887        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        // Move on with the string
903        start = end;
904        end = end + MAX_LOG_PAYLOAD_SIZE_BYTES;
905        chunk_idx += 1;
906    }
907
908    // Print any suffix left
909    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/// Logs payload in one go (all other OS).
921#[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        // Try and get the next request.
942        // Verify request was not returned
943        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        // Enqueue a ping
953        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        // Try and get the next request.
967        // Verify request was returned
968        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        // Enqueue a ping multiple times
979        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        // Verify a request is returned for each submitted ping
996        for _ in 0..n {
997            let task = upload_manager.get_upload_task(&glean, false);
998            assert!(task.is_upload());
999        }
1000
1001        // Verify that after all requests are returned, none are left
1002        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        // Add a rate limiter to the upload mangager with max of 10 pings every 3 seconds.
1015        let max_pings_per_interval = 10;
1016        upload_manager.set_rate_limiter(3, 10);
1017
1018        // Enqueue the max number of pings allowed per uploading window
1019        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        // Verify a request is returned for each submitted ping
1035        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        // Enqueue just one more ping
1041        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        // Verify that we are indeed told to wait because we are at capacity
1055        match upload_manager.get_upload_task(&glean, false) {
1056            PingUploadTask::Wait { time } => {
1057                // Wait for the uploading window to reset
1058                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        // Enqueue a ping multiple times
1074        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        // Clear the queue
1090        drop(upload_manager.clear_ping_queue());
1091
1092        // Verify there really isn't any ping in the queue
1093        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        // Register a ping for testing
1104        let ping_type = PingType::new(
1105            "test",
1106            true,
1107            /* send_if_empty */ 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        // Submit the ping multiple times
1119        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        // Clear the queue
1130        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        // Verify there really isn't any other pings in the queue
1139        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        // Register a ping for testing
1147        let ping_type = PingType::new(
1148            "test",
1149            true,
1150            /* send_if_empty */ 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        // Submit the ping multiple times
1162        let n = 10;
1163        for _ in 0..n {
1164            ping_type.submit_sync(&glean, None);
1165        }
1166
1167        // Create a new upload manager pointing to the same data_path as the glean instance.
1168        let upload_manager = PingUploadManager::no_policy(dir.path());
1169
1170        // Verify the requests were properly enqueued
1171        for _ in 0..n {
1172            let task = upload_manager.get_upload_task(&glean, false);
1173            assert!(task.is_upload());
1174        }
1175
1176        // Verify that after all requests are returned, none are left
1177        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        // Register a ping for testing
1188        let ping_type = PingType::new(
1189            "test",
1190            true,
1191            /* send_if_empty */ 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        // Submit a ping
1203        ping_type.submit_sync(&glean, None);
1204
1205        // Get the pending ping directory path
1206        let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1207
1208        // Get the submitted PingRequest
1209        match glean.get_upload_task() {
1210            PingUploadTask::Upload { request } => {
1211                // Simulate the processing of a sucessfull request
1212                let document_id = request.document_id;
1213                glean.process_ping_upload_response(&document_id, UploadResult::http_status(200));
1214                // Verify file was deleted
1215                assert!(!pending_pings_dir.join(document_id).exists());
1216            }
1217            _ => panic!("Expected upload manager to return the next request!"),
1218        }
1219
1220        // Verify that after request is returned, none are left
1221        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        // Register a ping for testing
1229        let ping_type = PingType::new(
1230            "test",
1231            true,
1232            /* send_if_empty */ 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        // Submit a ping
1244        ping_type.submit_sync(&glean, None);
1245
1246        // Get the pending ping directory path
1247        let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1248
1249        // Get the submitted PingRequest
1250        match glean.get_upload_task() {
1251            PingUploadTask::Upload { request } => {
1252                // Simulate the processing of a client error
1253                let document_id = request.document_id;
1254                glean.process_ping_upload_response(&document_id, UploadResult::http_status(404));
1255                // Verify file was deleted
1256                assert!(!pending_pings_dir.join(document_id).exists());
1257            }
1258            _ => panic!("Expected upload manager to return the next request!"),
1259        }
1260
1261        // Verify that after request is returned, none are left
1262        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        // Register a ping for testing
1270        let ping_type = PingType::new(
1271            "test",
1272            true,
1273            /* send_if_empty */ 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        // Submit a ping
1285        ping_type.submit_sync(&glean, None);
1286
1287        // Get the submitted PingRequest
1288        match glean.get_upload_task() {
1289            PingUploadTask::Upload { request } => {
1290                // Simulate the processing of a client error
1291                let document_id = request.document_id;
1292                glean.process_ping_upload_response(&document_id, UploadResult::http_status(500));
1293                // Verify this ping was indeed re-enqueued
1294                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        // Verify that after request is returned, none are left
1305        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        // Register a ping for testing
1313        let ping_type = PingType::new(
1314            "test",
1315            true,
1316            /* send_if_empty */ 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        // Submit a ping
1328        ping_type.submit_sync(&glean, None);
1329
1330        // Get the pending ping directory path
1331        let pending_pings_dir = dir.path().join(PENDING_PINGS_DIRECTORY);
1332
1333        // Get the submitted PingRequest
1334        match glean.get_upload_task() {
1335            PingUploadTask::Upload { request } => {
1336                // Simulate the processing of a client error
1337                let document_id = request.document_id;
1338                glean.process_ping_upload_response(
1339                    &document_id,
1340                    UploadResult::unrecoverable_failure(),
1341                );
1342                // Verify file was deleted
1343                assert!(!pending_pings_dir.join(document_id).exists());
1344            }
1345            _ => panic!("Expected upload manager to return the next request!"),
1346        }
1347
1348        // Verify that after request is returned, none are left
1349        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        // Enqueue a ping
1365        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        // Try and get the first request.
1379        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        // Schedule the next one while the first one is "in progress"
1386        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        // Mark as processed
1400        upload_manager.process_ping_upload_response(
1401            &glean,
1402            &req.document_id,
1403            UploadResult::http_status(200),
1404        );
1405
1406        // Get the second request.
1407        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        // Mark as processed
1414        upload_manager.process_ping_upload_response(
1415            &glean,
1416            &req.document_id,
1417            UploadResult::http_status(200),
1418        );
1419
1420        // ... and then we're done.
1421        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        // Register a ping for testing
1434        let ping_type = PingType::new(
1435            "test",
1436            true,
1437            /* send_if_empty */ 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        // Submit a ping
1449        ping_type.submit_sync(&glean, None);
1450
1451        // Get the submitted PingRequest
1452        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        // Create a new upload manager so that we have access to its functions directly,
1465        // make it synchronous so we don't have to manually wait for the scanning to finish.
1466        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        // Try to enqueue a ping with the same doc_id twice
1472        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        // Get a task once
1498        let task = upload_manager.get_upload_task(&glean, false);
1499        assert!(task.is_upload());
1500
1501        // There should be no more queued tasks
1502        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        // Register a ping for testing
1513        let ping_type = PingType::new(
1514            "test",
1515            true,
1516            /* send_if_empty */ 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        // Submit the ping multiple times
1528        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        // Set a policy for max recoverable failures, this is usually disabled for tests.
1536        let max_recoverable_failures = 3;
1537        upload_manager
1538            .policy
1539            .set_max_recoverable_failures(Some(max_recoverable_failures));
1540
1541        // Return the max recoverable error failures in a row
1542        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        // Verify that after returning the max amount of recoverable failures,
1556        // we are done even though we haven't gotten all the enqueued requests.
1557        assert_eq!(
1558            upload_manager.get_upload_task(&glean, false),
1559            PingUploadTask::done()
1560        );
1561
1562        // Verify all requests are returned when we try again.
1563        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        // Register a ping for testing
1574        let ping_type = PingType::new(
1575            "test",
1576            true,
1577            /* send_if_empty */ 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        // Submit the ping multiple times
1589        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        // The pending pings array is sorted by date in ascending order,
1597        // the newest element is the last one.
1598        let (_, newest_ping) = &pending_pings.last().unwrap();
1599        let PingPayload {
1600            document_id: newest_ping_id,
1601            ..
1602        } = &newest_ping;
1603
1604        // Create a new upload manager pointing to the same data_path as the glean instance.
1605        let mut upload_manager = PingUploadManager::no_policy(dir.path());
1606
1607        // Set the quota to just a little over the size on an empty ping file.
1608        // This way we can check that one ping is kept and all others are deleted.
1609        //
1610        // From manual testing I figured out an empty ping file is 324bytes,
1611        // I am setting this a little over just so that minor changes to the ping structure
1612        // don't immediatelly break this.
1613        upload_manager
1614            .policy
1615            .set_max_pending_pings_directory_size(Some(500));
1616
1617        // Get a task once
1618        // One ping should have been enqueued.
1619        // Make sure it is the newest ping.
1620        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        // Verify that no other requests were returned,
1626        // they should all have been deleted because pending pings quota was hit.
1627        assert_eq!(
1628            upload_manager.get_upload_task(&glean, false),
1629            PingUploadTask::done()
1630        );
1631
1632        // Verify that the correct number of deleted pings was recorded
1633        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        // Register a ping for testing
1656        let ping_type = PingType::new(
1657            "test",
1658            true,
1659            /* send_if_empty */ 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        // How many pings we allow at maximum
1671        let count_quota = 3;
1672        // The number of pings we fill the pending pings directory with.
1673        let n = 10;
1674
1675        // Submit the ping multiple times
1676        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        // The pending pings array is sorted by date in ascending order,
1683        // the newest element is the last one.
1684        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        // Create a new upload manager pointing to the same data_path as the glean instance.
1692        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        // Get a task once
1699        // One ping should have been enqueued.
1700        // Make sure it is the newest ping.
1701        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        // Verify that no other requests were returned,
1709        // they should all have been deleted because pending pings quota was hit.
1710        assert_eq!(
1711            upload_manager.get_upload_task(&glean, false),
1712            PingUploadTask::done()
1713        );
1714
1715        // Verify that the correct number of deleted pings was recorded
1716        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        // Register a ping for testing
1739        let ping_type = PingType::new(
1740            "test",
1741            true,
1742            /* send_if_empty */ 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        // The number of pings we fill the pending pings directory with.
1755        let n = 10;
1756
1757        // Submit the ping multiple times
1758        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        // The pending pings array is sorted by date in ascending order,
1765        // the newest element is the last one.
1766        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        // Create a new upload manager pointing to the same data_path as the glean instance.
1774        let mut upload_manager = PingUploadManager::no_policy(dir.path());
1775
1776        // From manual testing we figured out a basically empty ping file is 399 bytes,
1777        // so this allows 3 pings with some headroom in case of future changes.
1778        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        // Get a task once
1784        // One ping should have been enqueued.
1785        // Make sure it is the newest ping.
1786        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        // Verify that no other requests were returned,
1794        // they should all have been deleted because pending pings quota was hit.
1795        assert_eq!(
1796            upload_manager.get_upload_task(&glean, false),
1797            PingUploadTask::done()
1798        );
1799
1800        // Verify that the correct number of deleted pings was recorded
1801        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        // Verify the labeled deletion counter attributes deletions to size_quota
1818        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        // Register a ping for testing
1840        let ping_type = PingType::new(
1841            "test",
1842            true,
1843            /* send_if_empty */ 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        // The number of pings we fill the pending pings directory with.
1856        let n = 10;
1857
1858        // Submit the ping multiple times
1859        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        // The pending pings array is sorted by date in ascending order,
1866        // the newest element is the last one.
1867        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        // Create a new upload manager pointing to the same data_path as the glean instance.
1875        let mut upload_manager = PingUploadManager::no_policy(dir.path());
1876
1877        // Set a large enough size quota so it never triggers before the count quota does.
1878        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        // Get a task once
1884        // One ping should have been enqueued.
1885        // Make sure it is the newest ping.
1886        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        // Verify that no other requests were returned,
1894        // they should all have been deleted because pending pings quota was hit.
1895        assert_eq!(
1896            upload_manager.get_upload_task(&glean, false),
1897            PingUploadTask::done()
1898        );
1899
1900        // Verify that the correct number of deleted pings was recorded
1901        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        // Verify the labeled deletion counter attributes deletions to count_quota
1918        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            /* send_if_empty */ 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        // Submit fewer pings than any quota.
1954        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        // Define a max_wait_attemps policy, this is disabled for tests by default.
2008        let max_wait_attempts = 3;
2009        upload_manager
2010            .policy
2011            .set_max_wait_attempts(Some(max_wait_attempts));
2012
2013        // Add a rate limiter to the upload mangager with max of 1 ping 5secs.
2014        //
2015        // We arbitrarily set the maximum pings per interval to a very low number,
2016        // when the rate limiter reaches it's limit get_upload_task returns a PingUploadTask::Wait,
2017        // which will allow us to test the limitations around returning too many of those in a row.
2018        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        // Enqueue two pings
2023        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        // Get the first ping, it should be returned normally.
2049        match upload_manager.get_upload_task(&glean, false) {
2050            PingUploadTask::Upload { .. } => {}
2051            _ => panic!("Expected upload manager to return the next request!"),
2052        }
2053
2054        // Try to get the next ping,
2055        // we should be throttled and thus get a PingUploadTask::Wait.
2056        // Check that we are indeed allowed to get this response as many times as expected.
2057        for _ in 0..max_wait_attempts {
2058            let task = upload_manager.get_upload_task(&glean, false);
2059            assert!(task.is_wait());
2060        }
2061
2062        // Check that after we get PingUploadTask::Wait the allowed number of times,
2063        // we then get PingUploadTask::Done.
2064        assert_eq!(
2065            upload_manager.get_upload_task(&glean, false),
2066            PingUploadTask::done()
2067        );
2068
2069        // Wait for the rate limiter to allow upload tasks again.
2070        thread::sleep(Duration::from_secs(secs_per_interval));
2071
2072        // Check that we are allowed again to get pings.
2073        let task = upload_manager.get_upload_task(&glean, false);
2074        assert!(task.is_upload());
2075
2076        // And once we are done we don't need to wait anymore.
2077        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        // Enqueue a ping and start processing it
2102        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        // Attempt to re-enqueue the same ping
2116        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        // No new pings should have been enqueued so the upload task is Done.
2128        assert_eq!(
2129            upload_manager.get_upload_task(&glean, false),
2130            PingUploadTask::done()
2131        );
2132
2133        // Process the upload response
2134        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        // Register a ping for testing
2147        let ping_type = PingType::new(
2148            "test",
2149            true,
2150            /* send_if_empty */ 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        // Submit a ping
2162        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        // Get the submitted PingRequest
2171        match glean.get_upload_task() {
2172            PingUploadTask::Upload { request } => {
2173                // Simulate the processing of a sucessful request
2174                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        // Verify that after request is returned, none are left
2187        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        // Register a ping for testing
2196        let ping_type = PingType::new(
2197            "test",
2198            true,
2199            /* send_if_empty */ 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        // Submit a ping
2211        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        // Get the submitted PingRequest
2220        match glean.get_upload_task() {
2221            PingUploadTask::Upload { request } => {
2222                // Simulate the processing of a sucessful request
2223                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        // Verify that after request is returned, none are left
2237        assert_eq!(glean.get_upload_task(), PingUploadTask::done());
2238    }
2239}