1use std::sync::Arc;
9use std::time::Duration;
10
11use bytes::Bytes;
12use rand::RngExt;
13use reqwest::{Method, StatusCode};
14use serde::Serialize;
15use serde::de::DeserializeOwned;
16use tokio::sync::Mutex;
17
18use crate::api::{ApiResponse, HumanVerificationCredential, ResponseCode};
19use crate::config::{API_CONTENT_TYPE, ProtonClientConfiguration, RetryPolicy};
20use crate::error::{ProtonApiError, ProtonError, Result};
21use crate::ids::SessionId;
22use crate::telemetry::{NoopTelemetry, Telemetry, TelemetryExt};
23
24const SESSION_ID_HEADER: &str = "x-pm-uid";
25const APP_VERSION_HEADER: &str = "x-pm-appversion";
26const STORAGE_TOKEN_HEADER: &str = "pm-storage-token";
27const HV_TOKEN_HEADER: &str = "x-pm-human-verification-token";
28const HV_TOKEN_TYPE_HEADER: &str = "x-pm-human-verification-token-type";
29
30#[derive(Debug, Clone)]
33pub struct Tokens {
34 pub access_token: String,
35 pub refresh_token: String,
36}
37
38#[derive(Clone)]
43pub struct ApiHttpClient {
44 inner: Arc<Inner>,
45 route_prefix: Arc<str>,
51}
52
53type TokensRefreshedCallback = Arc<dyn Fn(Tokens) + Send + Sync>;
55
56struct Inner {
57 http: reqwest::Client,
58 base_url: String,
59 config: ProtonClientConfiguration,
60 session_id: SessionId,
61 tokens: Mutex<Tokens>,
62 telemetry: std::sync::Mutex<Arc<dyn Telemetry>>,
68 on_tokens_refreshed: std::sync::Mutex<Option<TokensRefreshedCallback>>,
69}
70
71impl ApiHttpClient {
72 pub fn new(
74 config: ProtonClientConfiguration,
75 session_id: SessionId,
76 tokens: Tokens,
77 ) -> Result<Self> {
78 let http = reqwest::Client::builder()
84 .timeout(config.request_timeout)
85 .gzip(true)
86 .build()?;
87
88 let base_url = ensure_trailing_slash(&config.base_url);
89
90 Ok(Self {
91 inner: Arc::new(Inner {
92 http,
93 base_url,
94 config,
95 session_id,
96 tokens: Mutex::new(tokens),
97 telemetry: std::sync::Mutex::new(NoopTelemetry::shared()),
98 on_tokens_refreshed: std::sync::Mutex::new(None),
99 }),
100 route_prefix: Arc::from(""),
101 })
102 }
103
104 pub fn with_base_route(&self, route: impl Into<Arc<str>>) -> Self {
110 Self {
111 inner: Arc::clone(&self.inner),
112 route_prefix: route.into(),
113 }
114 }
115
116 pub async fn current_tokens(&self) -> Tokens {
118 self.inner.tokens.lock().await.clone()
119 }
120
121 pub fn set_telemetry(&self, telemetry: Arc<dyn Telemetry>) {
128 *self
129 .inner
130 .telemetry
131 .lock()
132 .expect("telemetry mutex poisoned") = telemetry;
133 }
134
135 pub fn set_on_tokens_refreshed(&self, callback: impl Fn(Tokens) + Send + Sync + 'static) {
138 *self
139 .inner
140 .on_tokens_refreshed
141 .lock()
142 .expect("on_tokens_refreshed mutex poisoned") = Some(Arc::new(callback));
143 }
144
145 fn telemetry(&self) -> Arc<dyn Telemetry> {
147 self.inner
148 .telemetry
149 .lock()
150 .expect("telemetry mutex poisoned")
151 .clone()
152 }
153
154 pub async fn get_storage_blob(&self, url: &str, token: &str) -> Result<Bytes> {
168 let mut timer = self.telemetry().start("storage_download");
169 let response = send_retrying(&self.inner.config.retry_policy, || {
170 let mut request = self
173 .inner
174 .http
175 .get(url)
176 .timeout(self.inner.config.storage_timeout)
177 .header(STORAGE_TOKEN_HEADER, token);
178 if !self.inner.config.user_agent.is_empty() {
179 request =
180 request.header(reqwest::header::USER_AGENT, &self.inner.config.user_agent);
181 }
182 request
183 })
184 .await?;
185 let status = response.status();
186 timer.attr("status", status.as_u16());
187 let bytes = response.bytes().await?;
188
189 if let Ok(envelope) = serde_json::from_slice::<ApiResponse>(&bytes) {
192 if !envelope.is_success() {
193 return Err(api_error(status, &bytes));
194 }
195 } else if !status.is_success() {
196 return Err(api_error(status, &bytes));
197 }
198
199 timer.success();
200 Ok(bytes)
201 }
202
203 pub async fn post_storage_blob(&self, url: &str, token: &str, blob: Bytes) -> Result<()> {
216 reqwest::multipart::Part::bytes(Vec::new())
219 .mime_str("application/octet-stream")
220 .map_err(ProtonError::from)?;
221
222 let blob_len = blob.len() as u64;
223 let mut timer = self.telemetry().start("storage_upload");
224 let response = send_retrying(&self.inner.config.retry_policy, || {
225 let body = reqwest::Body::from(blob.clone());
226 let part = reqwest::multipart::Part::stream_with_length(body, blob_len)
227 .file_name("blob")
228 .mime_str("application/octet-stream")
229 .expect("octet-stream is a valid MIME type");
230 let form = reqwest::multipart::Form::new().part("Block", part);
231
232 let mut request = self
233 .inner
234 .http
235 .post(url)
236 .timeout(self.inner.config.storage_timeout)
237 .header(STORAGE_TOKEN_HEADER, token)
238 .multipart(form);
239
240 if !self.inner.config.user_agent.is_empty() {
241 request =
242 request.header(reqwest::header::USER_AGENT, &self.inner.config.user_agent);
243 }
244 request
245 })
246 .await?;
247 let status = response.status();
248 timer.attr("status", status.as_u16());
249 let bytes = response.bytes().await?;
250
251 if let Ok(envelope) = serde_json::from_slice::<ApiResponse>(&bytes) {
252 if !envelope.is_success() {
253 return Err(api_error(status, &bytes));
254 }
255 } else if !status.is_success() {
256 return Err(api_error(status, &bytes));
257 }
258 timer.success();
259 Ok(())
260 }
261
262 pub fn session_id(&self) -> &SessionId {
263 &self.inner.session_id
264 }
265
266 pub async fn get<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
268 self.send::<(), T>(Method::GET, path, None).await
269 }
270
271 pub async fn post<B: Serialize, T: DeserializeOwned>(&self, path: &str, body: &B) -> Result<T> {
273 self.send::<B, T>(Method::POST, path, Some(body)).await
274 }
275
276 pub async fn post_multipart<B: Serialize, T: DeserializeOwned>(
279 &self,
280 path: &str,
281 metadata: &B,
282 binary_parts: &[(String, Vec<u8>)],
283 ) -> Result<T> {
284 let metadata = serde_json::to_vec(metadata)?;
285 let mut timer = self.telemetry().start("http_request");
286 timer.attr("method", "POST");
287 let rejected_token = self.inner.tokens.lock().await.access_token.clone();
288 let response = self
289 .send_multipart_with_token(path, &metadata, binary_parts, &rejected_token)
290 .await?;
291 let response = if response.status() == StatusCode::UNAUTHORIZED {
292 let bytes = response.bytes().await?;
293 if let Ok(envelope) = serde_json::from_slice::<ApiResponse>(&bytes)
294 && matches!(
295 envelope.code,
296 ResponseCode::AccountDeleted | ResponseCode::AccountDisabled
297 )
298 {
299 return Err(api_error(StatusCode::UNAUTHORIZED, &bytes));
300 }
301 let token = self.refresh_access_token(&rejected_token).await?;
302 self.send_multipart_with_token(path, &metadata, binary_parts, &token)
303 .await?
304 } else {
305 response
306 };
307 timer.attr("status", response.status().as_u16());
308 let parsed = parse_response(response).await?;
309 timer.success();
310 Ok(parsed)
311 }
312
313 pub async fn put<B: Serialize, T: DeserializeOwned>(&self, path: &str, body: &B) -> Result<T> {
315 self.send::<B, T>(Method::PUT, path, Some(body)).await
316 }
317
318 pub async fn delete<T: DeserializeOwned>(&self, path: &str) -> Result<T> {
320 self.send::<(), T>(Method::DELETE, path, None).await
321 }
322
323 async fn send<B: Serialize, T: DeserializeOwned>(
324 &self,
325 method: Method,
326 path: &str,
327 body: Option<&B>,
328 ) -> Result<T> {
329 let mut timer = self.telemetry().start("http_request");
330 timer.attr("method", method.as_str());
331
332 let access_token = self.inner.tokens.lock().await.access_token.clone();
333
334 let response = self
336 .send_with_token(method.clone(), path, body, &access_token)
337 .await?;
338
339 let response = if response.status() == StatusCode::UNAUTHORIZED {
340 self.handle_unauthorized(method, path, body, response, access_token)
341 .await?
342 } else {
343 response
344 };
345
346 timer.attr("status", response.status().as_u16());
347 let parsed = parse_response(response).await?;
348 timer.success();
349 Ok(parsed)
350 }
351
352 async fn handle_unauthorized<B: Serialize>(
353 &self,
354 method: Method,
355 path: &str,
356 body: Option<&B>,
357 response: reqwest::Response,
358 rejected_access_token: String,
359 ) -> Result<reqwest::Response> {
360 let bytes = response.bytes().await?;
362 if let Ok(envelope) = serde_json::from_slice::<ApiResponse>(&bytes)
363 && matches!(
364 envelope.code,
365 ResponseCode::AccountDeleted | ResponseCode::AccountDisabled
366 )
367 {
368 return Err(api_error(StatusCode::UNAUTHORIZED, &bytes));
369 }
370
371 let access_token = self.refresh_access_token(&rejected_access_token).await?;
372 self.send_with_token(method, path, body, &access_token)
373 .await
374 }
375
376 async fn send_with_token<B: Serialize>(
377 &self,
378 method: Method,
379 path: &str,
380 body: Option<&B>,
381 access_token: &str,
382 ) -> Result<reqwest::Response> {
383 let url = format!(
384 "{}{}{}",
385 self.inner.base_url,
386 self.route_prefix,
387 path.trim_start_matches('/')
388 );
389 send_retrying(&self.inner.config.retry_policy, || {
390 let mut request = self
391 .inner
392 .http
393 .request(method.clone(), &url)
394 .header(SESSION_ID_HEADER, self.inner.session_id.as_str())
395 .header(APP_VERSION_HEADER, &self.inner.config.app_version)
396 .header(reqwest::header::ACCEPT, API_CONTENT_TYPE)
397 .bearer_auth(access_token);
398
399 if !self.inner.config.user_agent.is_empty() {
400 request =
401 request.header(reqwest::header::USER_AGENT, &self.inner.config.user_agent);
402 }
403
404 if let Some(body) = body {
405 request = request.json(body);
406 }
407
408 request
409 })
410 .await
411 }
412
413 async fn send_multipart_with_token(
414 &self,
415 path: &str,
416 metadata: &[u8],
417 binary_parts: &[(String, Vec<u8>)],
418 access_token: &str,
419 ) -> Result<reqwest::Response> {
420 let url = format!(
421 "{}{}{}",
422 self.inner.base_url,
423 self.route_prefix,
424 path.trim_start_matches('/')
425 );
426 send_retrying(&self.inner.config.retry_policy, || {
427 let metadata_part = reqwest::multipart::Part::bytes(metadata.to_vec())
428 .file_name("Metadata")
429 .mime_str("application/json")
430 .expect("application/json is a valid MIME type");
431 let mut form = reqwest::multipart::Form::new().part("Metadata", metadata_part);
432 for (name, bytes) in binary_parts {
433 let part = reqwest::multipart::Part::bytes(bytes.clone())
434 .file_name(name.clone())
435 .mime_str("application/octet-stream")
436 .expect("octet-stream is a valid MIME type");
437 form = form.part(name.clone(), part);
438 }
439 let mut request = self
440 .inner
441 .http
442 .post(&url)
443 .timeout(self.inner.config.storage_timeout)
444 .header(SESSION_ID_HEADER, self.inner.session_id.as_str())
445 .header(APP_VERSION_HEADER, &self.inner.config.app_version)
446 .header(reqwest::header::ACCEPT, API_CONTENT_TYPE)
447 .bearer_auth(access_token)
448 .multipart(form);
449 if !self.inner.config.user_agent.is_empty() {
450 request =
451 request.header(reqwest::header::USER_AGENT, &self.inner.config.user_agent);
452 }
453 request
454 })
455 .await
456 }
457
458 async fn refresh_access_token(&self, rejected_access_token: &str) -> Result<String> {
462 let mut guard = self.inner.tokens.lock().await;
463
464 if guard.access_token != rejected_access_token {
465 return Ok(guard.access_token.clone());
466 }
467
468 let refreshed = self.request_refresh(&guard.refresh_token).await?;
469 *guard = refreshed.clone();
470
471 if let Some(ref cb) = *self
473 .inner
474 .on_tokens_refreshed
475 .lock()
476 .expect("on_tokens_refreshed mutex poisoned")
477 {
478 cb(refreshed.clone());
479 }
480
481 Ok(refreshed.access_token)
482 }
483
484 async fn request_refresh(&self, refresh_token: &str) -> Result<Tokens> {
485 let url = format!("{}auth/v4/refresh", self.inner.base_url);
486 let body = SessionRefreshRequest {
487 response_type: "token",
488 grant_type: "refresh_token",
489 refresh_token,
490 redirect_uri: &self.inner.config.refresh_redirect_uri,
491 };
492
493 let response = send_retrying(&self.inner.config.retry_policy, || {
496 self.inner
497 .http
498 .post(&url)
499 .header(SESSION_ID_HEADER, self.inner.session_id.as_str())
500 .header(APP_VERSION_HEADER, &self.inner.config.app_version)
501 .header(reqwest::header::ACCEPT, API_CONTENT_TYPE)
502 .json(&body)
503 })
504 .await?;
505
506 let refreshed: SessionRefreshResponse = parse_response(response).await?;
507 Ok(Tokens {
508 access_token: refreshed.access_token,
509 refresh_token: refreshed.refresh_token,
510 })
511 }
512}
513
514#[derive(Serialize)]
515struct SessionRefreshRequest<'a> {
516 #[serde(rename = "ResponseType")]
517 response_type: &'a str,
518 #[serde(rename = "GrantType")]
519 grant_type: &'a str,
520 #[serde(rename = "RefreshToken")]
521 refresh_token: &'a str,
522 #[serde(rename = "RedirectURI")]
523 redirect_uri: &'a str,
524}
525
526#[derive(serde::Deserialize)]
527struct SessionRefreshResponse {
528 #[serde(rename = "AccessToken")]
529 access_token: String,
530 #[serde(rename = "RefreshToken")]
531 refresh_token: String,
532}
533
534pub async fn get_unauthenticated<T: DeserializeOwned>(
540 config: &ProtonClientConfiguration,
541 path: &str,
542) -> Result<T> {
543 let http = reqwest::Client::builder()
544 .timeout(config.request_timeout)
545 .gzip(true)
546 .build()?;
547
548 let base_url = ensure_trailing_slash(&config.base_url);
549 let url = format!("{}{}", base_url, path.trim_start_matches('/'));
550
551 let response = send_retrying(&config.retry_policy, || {
552 let mut request = http
553 .get(&url)
554 .header(APP_VERSION_HEADER, &config.app_version)
555 .header(reqwest::header::ACCEPT, API_CONTENT_TYPE);
556 if !config.user_agent.is_empty() {
557 request = request.header(reqwest::header::USER_AGENT, &config.user_agent);
558 }
559 request
560 })
561 .await?;
562
563 parse_response(response).await
564}
565
566pub async fn post_unauthenticated<B: Serialize, T: DeserializeOwned>(
572 config: &ProtonClientConfiguration,
573 path: &str,
574 body: &B,
575) -> Result<T> {
576 post_unauthenticated_verified(config, path, body, None).await
577}
578
579pub async fn post_unauthenticated_verified<B: Serialize, T: DeserializeOwned>(
588 config: &ProtonClientConfiguration,
589 path: &str,
590 body: &B,
591 verification: Option<&HumanVerificationCredential>,
592) -> Result<T> {
593 let http = reqwest::Client::builder()
594 .timeout(config.request_timeout)
595 .gzip(true)
596 .build()?;
597
598 let base_url = ensure_trailing_slash(&config.base_url);
599 let url = format!("{}{}", base_url, path.trim_start_matches('/'));
600
601 let response = send_retrying(&config.retry_policy, || {
602 let mut request = http
603 .post(&url)
604 .header(APP_VERSION_HEADER, &config.app_version)
605 .header(reqwest::header::ACCEPT, API_CONTENT_TYPE)
606 .json(body);
607 if !config.user_agent.is_empty() {
608 request = request.header(reqwest::header::USER_AGENT, &config.user_agent);
609 }
610 if let Some(hv) = verification {
611 request = request
612 .header(HV_TOKEN_HEADER, &hv.token)
613 .header(HV_TOKEN_TYPE_HEADER, &hv.method);
614 }
615 request
616 })
617 .await?;
618
619 parse_response(response).await
620}
621
622async fn parse_response<T: DeserializeOwned>(response: reqwest::Response) -> Result<T> {
625 let status = response.status();
626 let bytes = response.bytes().await?;
627
628 if let Ok(envelope) = serde_json::from_slice::<ApiResponse>(&bytes) {
633 if !envelope.is_success() && envelope.code != ResponseCode::MultipleResponses {
634 return Err(api_error(status, &bytes));
635 }
636 } else if !status.is_success() {
637 return Err(api_error(status, &bytes));
638 }
639
640 Ok(serde_json::from_slice::<T>(&bytes)?)
641}
642
643fn api_error(status: StatusCode, bytes: &[u8]) -> ProtonError {
644 let envelope = serde_json::from_slice::<ApiResponse>(bytes).ok();
645 let code = envelope
646 .as_ref()
647 .map(|e| e.code)
648 .unwrap_or(ResponseCode::Unknown);
649 let details = envelope.as_ref().and_then(|e| e.details.clone());
650 let message = envelope.and_then(|e| e.error_message).unwrap_or_else(|| {
651 status
652 .canonical_reason()
653 .unwrap_or("unknown error")
654 .to_owned()
655 });
656
657 ProtonError::Api(ProtonApiError {
658 code,
659 http_status: status.as_u16(),
660 message,
661 details,
662 })
663}
664
665async fn send_retrying<F>(policy: &RetryPolicy, build: F) -> Result<reqwest::Response>
677where
678 F: Fn() -> reqwest::RequestBuilder,
679{
680 let mut attempt: u32 = 0;
681 loop {
682 match build().send().await {
683 Ok(response) => {
684 if attempt < policy.max_retries && is_retryable_status(response.status()) {
685 let delay = retry_after(&response).unwrap_or_else(|| backoff(policy, attempt));
686 tokio::time::sleep(delay).await;
687 attempt += 1;
688 continue;
689 }
690 return Ok(response);
691 }
692 Err(err) => {
693 if attempt < policy.max_retries && is_retryable_error(&err) {
694 tokio::time::sleep(backoff(policy, attempt)).await;
695 attempt += 1;
696 continue;
697 }
698 return Err(err.into());
699 }
700 }
701 }
702}
703
704fn is_retryable_status(status: StatusCode) -> bool {
707 matches!(status.as_u16(), 408 | 429 | 502 | 503 | 504)
708}
709
710fn is_retryable_error(err: &reqwest::Error) -> bool {
714 err.is_timeout() || err.is_connect()
715}
716
717fn retry_after(response: &reqwest::Response) -> Option<Duration> {
720 let value = response
721 .headers()
722 .get(reqwest::header::RETRY_AFTER)?
723 .to_str()
724 .ok()?;
725 parse_retry_after_secs(value)
726}
727
728fn parse_retry_after_secs(value: &str) -> Option<Duration> {
731 value.trim().parse().ok().map(Duration::from_secs)
732}
733
734fn backoff(policy: &RetryPolicy, attempt: u32) -> Duration {
737 let ceiling = policy
738 .base_delay
739 .saturating_mul(1u32.checked_shl(attempt).unwrap_or(u32::MAX))
740 .min(policy.max_delay);
741 let ceiling_ms = ceiling.as_millis() as u64;
742 if ceiling_ms == 0 {
743 return Duration::ZERO;
744 }
745 Duration::from_millis(rand::rng().random_range(0..=ceiling_ms))
746}
747
748fn ensure_trailing_slash(url: &str) -> String {
749 if url.ends_with('/') {
750 url.to_owned()
751 } else {
752 format!("{url}/")
753 }
754}
755
756#[cfg(test)]
757mod tests {
758 use super::*;
759
760 #[test]
761 fn retryable_statuses() {
762 for code in [408u16, 429, 502, 503, 504] {
763 assert!(is_retryable_status(StatusCode::from_u16(code).unwrap()));
764 }
765 for code in [200u16, 400, 401, 403, 404, 500] {
766 assert!(!is_retryable_status(StatusCode::from_u16(code).unwrap()));
767 }
768 }
769
770 #[test]
771 fn retry_after_parses_seconds_only() {
772 assert_eq!(parse_retry_after_secs("5"), Some(Duration::from_secs(5)));
773 assert_eq!(
774 parse_retry_after_secs(" 12 "),
775 Some(Duration::from_secs(12))
776 );
777 assert_eq!(parse_retry_after_secs("0"), Some(Duration::ZERO));
778 assert_eq!(
780 parse_retry_after_secs("Wed, 21 Oct 2015 07:28:00 GMT"),
781 None
782 );
783 assert_eq!(parse_retry_after_secs(""), None);
784 }
785
786 #[test]
787 fn backoff_grows_then_caps_within_jitter_bounds() {
788 let policy = RetryPolicy {
789 max_retries: 5,
790 base_delay: Duration::from_millis(100),
791 max_delay: Duration::from_millis(1000),
792 };
793 for attempt in 0..8u32 {
796 let ceiling = Duration::from_millis(100u64.saturating_mul(1 << attempt.min(20)))
797 .min(policy.max_delay);
798 for _ in 0..64 {
799 assert!(backoff(&policy, attempt) <= ceiling);
800 }
801 }
802 }
803
804 #[test]
805 fn backoff_handles_large_attempt_without_overflow() {
806 let policy = RetryPolicy::default();
807 assert!(backoff(&policy, 64) <= policy.max_delay);
809 }
810
811 #[test]
812 fn disabled_policy_has_no_retries() {
813 assert_eq!(RetryPolicy::disabled().max_retries, 0);
814 }
815
816 struct Capture(std::sync::Mutex<Vec<crate::telemetry::TelemetryEvent>>);
818
819 impl Telemetry for Capture {
820 fn record(&self, event: &crate::telemetry::TelemetryEvent) {
821 self.0.lock().unwrap().push(event.clone());
822 }
823 }
824
825 #[tokio::test]
826 async fn http_request_records_telemetry_event() {
827 use crate::telemetry::Outcome;
828 use tokio::io::{AsyncReadExt, AsyncWriteExt};
829
830 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
833 let addr = listener.local_addr().unwrap();
834 let server = tokio::spawn(async move {
835 let (mut sock, _) = listener.accept().await.unwrap();
836 let mut buf = [0u8; 2048];
837 let _ = sock.read(&mut buf).await.unwrap();
838 let body = br#"{"Code":1000}"#;
839 let head = format!(
840 "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
841 body.len()
842 );
843 sock.write_all(head.as_bytes()).await.unwrap();
844 sock.write_all(body).await.unwrap();
845 sock.flush().await.unwrap();
846 });
847
848 let config = ProtonClientConfiguration::new("test@1.0")
849 .with_base_url(format!("http://{addr}/"))
850 .with_retry_policy(RetryPolicy::disabled());
851 let client = ApiHttpClient::new(
852 config,
853 SessionId::from("test-session"),
854 Tokens {
855 access_token: "access".into(),
856 refresh_token: "refresh".into(),
857 },
858 )
859 .unwrap();
860
861 let capture = Arc::new(Capture(std::sync::Mutex::new(Vec::new())));
862 client.set_telemetry(capture.clone());
863
864 let _: ApiResponse = client.get("some/path").await.unwrap();
865 server.await.unwrap();
866
867 let events = capture.0.lock().unwrap();
868 assert_eq!(events.len(), 1, "exactly one http_request event");
869 let event = &events[0];
870 assert_eq!(event.operation, "http_request");
871 assert_eq!(event.outcome, Outcome::Success);
872 assert!(
873 event
874 .attributes
875 .iter()
876 .any(|(k, v)| *k == "method" && v == "GET")
877 );
878 assert!(
879 event
880 .attributes
881 .iter()
882 .any(|(k, v)| *k == "status" && v == "200")
883 );
884 }
885}