1use crate::routers::RouterError;
2use crate::routers::common::message_size_check;
3use crate::routers::fcm::error::FcmError;
4use crate::routers::fcm::settings::{FcmServerCredential, FcmSettings};
5use actix_web::http::header::HttpDate;
6use reqwest::StatusCode;
7use serde::Deserialize;
8use std::collections::HashMap;
9use std::path::Path;
10use std::time::{Duration, SystemTime};
11use url::Url;
12use yup_oauth2::authenticator::DefaultAuthenticator;
13use yup_oauth2::{ServiceAccountAuthenticator, ServiceAccountKey};
14
15const OAUTH_SCOPES: &[&str] = &["https://www.googleapis.com/auth/firebase.messaging"];
16
17const MIN_RETRY_AFTER_SECS: u64 = 10;
19const MAX_RETRY_AFTER_SECS: u64 = 3600;
21
22fn parse_retry_after(headers: &reqwest::header::HeaderMap) -> Option<u64> {
29 let raw = headers
30 .get(reqwest::header::RETRY_AFTER)?
31 .to_str()
32 .ok()?
33 .trim();
34
35 let secs = match raw.parse::<u64>() {
36 Ok(secs) => secs,
37 Err(_) => {
38 let when: SystemTime = raw.parse::<HttpDate>().ok()?.into();
39 when.duration_since(SystemTime::now())
40 .map_or(0, |delay| delay.as_secs())
41 }
42 };
43
44 Some(secs.clamp(MIN_RETRY_AFTER_SECS, MAX_RETRY_AFTER_SECS))
45}
46
47pub struct FcmClient {
50 endpoint: Url,
51 timeout: Duration,
52 max_data: usize,
53 authenticator: Option<DefaultAuthenticator>,
54 http_client: reqwest::Client,
55}
56
57impl FcmClient {
58 pub async fn new(
60 settings: &FcmSettings,
61 server_credential: FcmServerCredential,
62 http: reqwest::Client,
63 ) -> std::io::Result<Self> {
64 let auth = if server_credential.server_access_token.contains('{') {
69 trace!(
70 "Reading credential for {} from string...",
71 &server_credential.project_id
72 );
73 let key_data =
74 serde_json::from_str::<ServiceAccountKey>(&server_credential.server_access_token)?;
75 Some(
76 ServiceAccountAuthenticator::builder(key_data)
77 .build()
78 .await?,
79 )
80 } else {
81 if Path::new(&server_credential.server_access_token).exists() {
83 warn!(
84 "Reading credential for {} from file...",
85 &server_credential.project_id
86 );
87 let content = std::fs::read_to_string(&server_credential.server_access_token)?;
88 let key_data = serde_json::from_str::<ServiceAccountKey>(&content)?;
89 Some(
90 ServiceAccountAuthenticator::builder(key_data)
91 .build()
92 .await?,
93 )
94 } else {
95 trace!("Presuming {} is GCM", &server_credential.project_id);
96 None
97 }
98 };
99 Ok(FcmClient {
100 endpoint: settings
101 .base_url
102 .join(&format!(
103 "v1/projects/{}/messages:send",
104 server_credential.project_id
105 ))
106 .expect("Project ID is not URL-safe"),
107 timeout: Duration::from_secs(settings.timeout as u64),
108 max_data: settings.max_data,
109 authenticator: auth,
110 http_client: http,
111 })
112 }
113
114 pub async fn send(
116 &self,
117 data: HashMap<&'static str, String>,
118 routing_token: String,
119 ttl: u64,
120 ) -> Result<(), RouterError> {
121 let data_json = serde_json::to_string(&data).unwrap();
124 message_size_check(data_json.as_bytes(), self.max_data)?;
125
126 let message = serde_json::json!({
128 "message": {
129 "token": routing_token,
130 "android": {
131 "ttl": format!("{ttl}s"),
132 "data": data
133 }
134 }
135 });
136
137 let server_access_token = self
138 .authenticator
139 .as_ref()
140 .unwrap()
141 .token(OAUTH_SCOPES)
142 .await
143 .map_err(FcmError::OAuthToken)?;
144 let token = server_access_token.token().ok_or(FcmError::NoOAuthToken)?;
145
146 let response = self
148 .http_client
149 .post(self.endpoint.clone())
150 .header("Authorization", format!("Bearer {token}"))
151 .header("Content-Type", "application/json")
152 .json(&message)
153 .timeout(self.timeout)
154 .send()
155 .await
156 .map_err(|e| {
157 if e.is_timeout() {
158 RouterError::RequestTimeout
159 } else {
160 RouterError::Connect(e)
161 }
162 })?;
163
164 let status = response.status();
166 if status.is_client_error() || status.is_server_error() {
167 let retry_after = parse_retry_after(response.headers());
168 let raw_data = response
169 .bytes()
170 .await
171 .map_err(FcmError::DeserializeResponse)?;
172 if raw_data.is_empty() {
173 warn!("Empty FCM response [{status}]");
174 return Err(FcmError::EmptyResponse(status).into());
175 }
176 let data: FcmResponse = serde_json::from_slice(&raw_data).map_err(|e| {
177 let s = String::from_utf8(raw_data.to_vec()).unwrap_or_else(|e| e.to_string());
178 warn!("Invalid FCM response [{status}] \"{s}\"");
179 FcmError::InvalidResponse(e, s, status)
180 })?;
181
182 return Err(match (status, data.error) {
184 (StatusCode::UNAUTHORIZED, _) => RouterError::Authentication,
185 (StatusCode::NOT_FOUND, _) => RouterError::NotFound,
186 (_, Some(error)) => {
187 info!("🌉Bridge Error: {:?}, {:?}", error.message, &self.endpoint);
188 FcmError::Upstream {
189 error_code: error.status, message: error.message,
191 retry_after,
192 }
193 .into()
194 }
195 (_, None) => {
199 warn!(
200 "🌉Unknown Bridge Error: {:?}, <{:?}>, [{:?}]",
201 status.to_string(),
202 &self.endpoint,
203 raw_data,
204 );
205 FcmError::Upstream {
206 error_code: "UNKNOWN".to_string(),
207 message: format!("Unknown reason: {:?}", status.to_string()),
208 retry_after,
209 }
210 }
211 .into(),
212 });
213 }
214
215 Ok(())
216 }
217}
218
219#[derive(Deserialize)]
220struct FcmResponse {
221 error: Option<FcmErrorResponse>,
222}
223
224#[derive(Deserialize)]
226struct FcmErrorResponse {
227 status: String,
229 message: String,
230}
231
232#[cfg(test)]
233pub mod tests {
234 use crate::routers::RouterError;
235 use crate::routers::fcm::client::{FcmClient, MIN_RETRY_AFTER_SECS};
236 use crate::routers::fcm::error::FcmError;
237 use crate::routers::fcm::settings::{FcmServerCredential, FcmSettings};
238 use actix_web::http::header::HttpDate;
239 use std::collections::HashMap;
240 use std::time::{Duration, SystemTime};
241 use url::Url;
242
243 pub const PROJECT_ID: &str = "yup-test-243420";
244 const ACCESS_TOKEN: &str = "ya29.c.ElouBywiys0LyNaZoLPJcp1Fdi2KjFMxzvYKLXkTdvM-rDfqKlvEq6PiMhGoGHx97t5FAvz3eb_ahdwlBjSStxHtDVQB4ZPRJQ_EOi-iS7PnayahU2S9Jp8S6rk";
245 pub const GCM_PROJECT_ID: &str = "valid_gcm_access_token";
246
247 pub fn make_service_key(server: &mockito::ServerGuard) -> String {
249 serde_json::json!({
251 "type": "service_account",
252 "project_id": PROJECT_ID,
253 "private_key_id": "26de294916614a5ebdf7a065307ed3ea9941902b",
254 "private_key": "-----BEGIN PRIVATE KEY-----\nMIIEvwIBADANBgkqhkiG9w0BAQEFAASCBKkwggSlAgEAAoIBAQDemmylrvp1KcOn\n9yTAVVKPpnpYznvBvcAU8Qjwr2fSKylpn7FQI54wCk5VJVom0jHpAmhxDmNiP8yv\nHaqsef+87Oc0n1yZ71/IbeRcHZc2OBB33/LCFqf272kThyJo3qspEqhuAw0e8neg\nLQb4jpm9PsqR8IjOoAtXQSu3j0zkXemMYFy93PWHjVpPEUX16NGfsWH7oxspBHOk\n9JPGJL8VJdbiAoDSDgF0y9RjJY5I52UeHNhMsAkTYs6mIG4kKXt2+T9tAyHw8aho\nwmuytQAfydTflTfTG8abRtliF3nil2taAc5VB07dP1b4dVYy/9r6M8Z0z4XM7aP+\nNdn2TKm3AgMBAAECggEAWi54nqTlXcr2M5l535uRb5Xz0f+Q/pv3ceR2iT+ekXQf\n+mUSShOr9e1u76rKu5iDVNE/a7H3DGopa7ZamzZvp2PYhSacttZV2RbAIZtxU6th\n7JajPAM+t9klGh6wj4jKEcE30B3XVnbHhPJI9TCcUyFZoscuPXt0LLy/z8Uz0v4B\nd5JARwyxDMb53VXwukQ8nNY2jP7WtUig6zwE5lWBPFMbi8GwGkeGZOruAK5sPPwY\nGBAlfofKANI7xKx9UXhRwisB4+/XI1L0Q6xJySv9P+IAhDUI6z6kxR+WkyT/YpG3\nX9gSZJc7qEaxTIuDjtep9GTaoEqiGntjaFBRKoe+VQKBgQDzM1+Ii+REQqrGlUJo\nx7KiVNAIY/zggu866VyziU6h5wjpsoW+2Npv6Dv7nWvsvFodrwe50Y3IzKtquIal\nVd8aa50E72JNImtK/o5Nx6xK0VySjHX6cyKENxHRDnBmNfbALRM+vbD9zMD0lz2q\nmns/RwRGq3/98EqxP+nHgHSr9QKBgQDqUYsFAAfvfT4I75Glc9svRv8IsaemOm07\nW1LCwPnj1MWOhsTxpNF23YmCBupZGZPSBFQobgmHVjQ3AIo6I2ioV6A+G2Xq/JCF\nmzfbvZfqtbbd+nVgF9Jr1Ic5T4thQhAvDHGUN77BpjEqZCQLAnUWJx9x7e2xvuBl\n1A6XDwH/ewKBgQDv4hVyNyIR3nxaYjFd7tQZYHTOQenVffEAd9wzTtVbxuo4sRlR\nNM7JIRXBSvaATQzKSLHjLHqgvJi8LITLIlds1QbNLl4U3UVddJbiy3f7WGTqPFfG\nkLhUF4mgXpCpkMLxrcRU14Bz5vnQiDmQRM4ajS7/kfwue00BZpxuZxst3QKBgQCI\nRI3FhaQXyc0m4zPfdYYVc4NjqfVmfXoC1/REYHey4I1XetbT9Nb/+ow6ew0UbgSC\nUZQjwwJ1m1NYXU8FyovVwsfk9ogJ5YGiwYb1msfbbnv/keVq0c/Ed9+AG9th30qM\nIf93hAfClITpMz2mzXIMRQpLdmQSR4A2l+E4RjkSOwKBgQCB78AyIdIHSkDAnCxz\nupJjhxEhtQ88uoADxRoEga7H/2OFmmPsqfytU4+TWIdal4K+nBCBWRvAX1cU47vH\nJOlSOZI0gRKe0O4bRBQc8GXJn/ubhYSxI02IgkdGrIKpOb5GG10m85ZvqsXw3bKn\nRVHMD0ObF5iORjZUqD0yRitAdg==\n-----END PRIVATE KEY-----\n",
255 "client_email": "yup-test-sa-1@yup-test-243420.iam.gserviceaccount.com",
256 "client_id": "102851967901799660408",
257 "auth_uri": "https://accounts.google.com/o/oauth2/auth",
258 "token_uri": server.url() + "/token",
259 "auth_provider_x509_cert_url": "https://www.googleapis.com/oauth2/v1/certs",
260 "client_x509_cert_url": "https://www.googleapis.com/robot/v1/metadata/x509/yup-test-sa-1%40yup-test-243420.iam.gserviceaccount.com"
261 }).to_string()
262 }
263
264 pub async fn mock_token_endpoint(server: &mut mockito::ServerGuard) -> mockito::Mock {
266 server
267 .mock("POST", "/token")
268 .with_body(
269 serde_json::json!({
270 "access_token": ACCESS_TOKEN,
271 "expires_in": 3600,
272 "token_type": "Bearer"
273 })
274 .to_string(),
275 )
276 .create_async()
277 .await
278 }
279
280 pub fn mock_fcm_endpoint_builder(server: &mut mockito::ServerGuard, id: &str) -> mockito::Mock {
282 server.mock("POST", format!("/v1/projects/{id}/messages:send").as_str())
283 }
284
285 async fn make_client(
287 server: &mockito::ServerGuard,
288 credential: FcmServerCredential,
289 ) -> FcmClient {
290 FcmClient::new(
291 &FcmSettings {
292 base_url: Url::parse(&server.url()).unwrap(),
293 server_credentials: serde_json::json!(credential).to_string(),
294 ..Default::default()
295 },
296 credential,
297 reqwest::Client::new(),
298 )
299 .await
300 .unwrap()
301 }
302
303 #[tokio::test]
306 async fn sends_correct_fcm_request() {
307 let mut server = mockito::Server::new_async().await;
308
309 let client = make_client(
310 &server,
311 FcmServerCredential {
312 project_id: PROJECT_ID.to_owned(),
313 is_gcm: None,
314 server_access_token: make_service_key(&server),
315 },
316 )
317 .await;
318 let _token_mock = mock_token_endpoint(&mut server).await;
319 let fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
320 .match_header("Authorization", format!("Bearer {ACCESS_TOKEN}").as_str())
321 .match_header("Content-Type", "application/json")
322 .match_body(r#"{"message":{"android":{"data":{"is_test":"true"},"ttl":"42s"},"token":"test-token"}}"#)
323 .create();
324
325 let mut data = HashMap::new();
326 data.insert("is_test", "true".to_string());
327
328 let result = client.send(data, "test-token".to_string(), 42).await;
329 assert!(result.is_ok(), "result = {result:?}");
330 fcm_mock.assert();
331 }
332
333 #[tokio::test]
335 async fn unauthorized() {
336 let mut server = mockito::Server::new_async().await;
337
338 let client = make_client(
339 &server,
340 FcmServerCredential {
341 project_id: PROJECT_ID.to_owned(),
342 is_gcm: None,
343 server_access_token: make_service_key(&server),
344 },
345 )
346 .await;
347 let _token_mock = mock_token_endpoint(&mut server).await;
348 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
349 .with_status(401)
350 .with_body(r#"{"error":{"status":"UNAUTHENTICATED","message":"test-message"}}"#)
351 .create_async()
352 .await;
353
354 let result = client
355 .send(HashMap::new(), "test-token".to_string(), 42)
356 .await;
357 assert!(result.is_err());
358 assert!(
359 matches!(result.as_ref().unwrap_err(), RouterError::Authentication),
360 "result = {result:?}"
361 );
362 }
363
364 #[tokio::test]
366 async fn not_found() {
367 let mut server = mockito::Server::new_async().await;
368
369 let client = make_client(
370 &server,
371 FcmServerCredential {
372 project_id: PROJECT_ID.to_owned(),
373 is_gcm: None,
374 server_access_token: make_service_key(&server),
375 },
376 )
377 .await;
378 let _token_mock = mock_token_endpoint(&mut server).await;
379 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
380 .with_status(404)
381 .with_body(r#"{"error":{"status":"NOT_FOUND","message":"test-message"}}"#)
382 .create_async()
383 .await;
384
385 let result = client
386 .send(HashMap::new(), "test-token".to_string(), 42)
387 .await;
388 assert!(result.is_err());
389 assert!(
390 matches!(result.as_ref().unwrap_err(), RouterError::NotFound),
391 "result = {result:?}"
392 );
393 }
394
395 #[tokio::test]
397 async fn resource_exhausted_is_throttled() {
398 let mut server = mockito::Server::new_async().await;
399
400 let client = make_client(
401 &server,
402 FcmServerCredential {
403 project_id: PROJECT_ID.to_owned(),
404 is_gcm: Some(false),
405 server_access_token: make_service_key(&server),
406 },
407 )
408 .await;
409 let _token_mock = mock_token_endpoint(&mut server).await;
410 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
411 .with_status(429)
412 .with_header("Retry-After", "30")
413 .with_body(r#"{"error":{"status":"RESOURCE_EXHAUSTED","message":"quota"}}"#)
414 .create_async()
415 .await;
416
417 let result = client
418 .send(HashMap::new(), "test-token".to_string(), 42)
419 .await;
420 let err = result.unwrap_err();
421
422 assert_eq!(err.status().as_u16(), 429);
423 assert_eq!(err.errno(), Some(201));
424 assert_eq!(err.retry_after(), Some(30));
425 }
426
427 #[tokio::test]
429 async fn retry_after_accepts_http_date() {
430 let mut server = mockito::Server::new_async().await;
431
432 let client = make_client(
433 &server,
434 FcmServerCredential {
435 project_id: PROJECT_ID.to_owned(),
436 is_gcm: Some(false),
437 server_access_token: make_service_key(&server),
438 },
439 )
440 .await;
441 let _token_mock = mock_token_endpoint(&mut server).await;
442 let when = HttpDate::from(SystemTime::now() + Duration::from_secs(120)).to_string();
444 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
445 .with_status(429)
446 .with_header("Retry-After", &when)
447 .with_body(r#"{"error":{"status":"RESOURCE_EXHAUSTED","message":"quota"}}"#)
448 .create_async()
449 .await;
450
451 let result = client
452 .send(HashMap::new(), "test-token".to_string(), 42)
453 .await;
454 let retry_after = result.unwrap_err().retry_after().expect("a delay");
455 assert!(
458 (110..=120).contains(&retry_after),
459 "retry_after = {retry_after}"
460 );
461 }
462
463 #[tokio::test]
465 async fn retry_after_ignores_unparseable() {
466 let mut server = mockito::Server::new_async().await;
467
468 let client = make_client(
469 &server,
470 FcmServerCredential {
471 project_id: PROJECT_ID.to_owned(),
472 is_gcm: Some(false),
473 server_access_token: make_service_key(&server),
474 },
475 )
476 .await;
477 let _token_mock = mock_token_endpoint(&mut server).await;
478 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
479 .with_status(429)
480 .with_header("Retry-After", "soon")
481 .with_body(r#"{"error":{"status":"RESOURCE_EXHAUSTED","message":"quota"}}"#)
482 .create_async()
483 .await;
484
485 let result = client
486 .send(HashMap::new(), "test-token".to_string(), 42)
487 .await;
488 assert_eq!(result.unwrap_err().retry_after(), None);
489 }
490
491 #[tokio::test]
494 async fn retry_after_floors_past_http_date() {
495 let mut server = mockito::Server::new_async().await;
496
497 let client = make_client(
498 &server,
499 FcmServerCredential {
500 project_id: PROJECT_ID.to_owned(),
501 is_gcm: Some(false),
502 server_access_token: make_service_key(&server),
503 },
504 )
505 .await;
506 let _token_mock = mock_token_endpoint(&mut server).await;
507 let when = HttpDate::from(SystemTime::now() - Duration::from_secs(60)).to_string();
508 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
509 .with_status(429)
510 .with_header("Retry-After", &when)
511 .with_body(r#"{"error":{"status":"RESOURCE_EXHAUSTED","message":"quota"}}"#)
512 .create_async()
513 .await;
514
515 let result = client
516 .send(HashMap::new(), "test-token".to_string(), 42)
517 .await;
518 assert_eq!(
519 result.unwrap_err().retry_after(),
520 Some(MIN_RETRY_AFTER_SECS)
521 );
522 }
523
524 #[tokio::test]
526 async fn other_fcm_error() {
527 let mut server = mockito::Server::new_async().await;
528
529 let client = make_client(
530 &server,
531 FcmServerCredential {
532 project_id: PROJECT_ID.to_owned(),
533 is_gcm: Some(false),
534 server_access_token: make_service_key(&server),
535 },
536 )
537 .await;
538 let _token_mock = mock_token_endpoint(&mut server).await;
539 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
540 .with_status(400)
541 .with_body(r#"{"error":{"status":"TEST_ERROR","message":"test-message"}}"#)
542 .create_async()
543 .await;
544
545 let result = client
546 .send(HashMap::new(), "test-token".to_string(), 42)
547 .await;
548 assert!(result.is_err());
549 assert!(
550 matches!(
551 result.as_ref().unwrap_err(),
552 RouterError::Fcm(FcmError::Upstream{ error_code, message, .. })
553 if error_code == "TEST_ERROR" && message == "test-message"
554 ),
555 "result = {result:?}"
556 );
557 }
558
559 #[tokio::test]
561 async fn unknown_fcm_error() {
562 let mut server = mockito::Server::new_async().await;
563
564 let client = make_client(
565 &server,
566 FcmServerCredential {
567 project_id: PROJECT_ID.to_owned(),
568 is_gcm: Some(true),
569 server_access_token: make_service_key(&server),
570 },
571 )
572 .await;
573 let _token_mock = mock_token_endpoint(&mut server).await;
574 let _fcm_mock = mock_fcm_endpoint_builder(&mut server, PROJECT_ID)
575 .with_status(400)
576 .with_body("{}")
577 .create_async()
578 .await;
579
580 let result = client
581 .send(HashMap::new(), "test-token".to_string(), 42)
582 .await;
583 assert!(result.is_err());
584 assert!(
585 matches!(
586 result.as_ref().unwrap_err(),
587 RouterError::Fcm(FcmError::Upstream { error_code, message, .. })
588 if error_code == "UNKNOWN" && message.starts_with("Unknown reason")
589 ),
590 "result = {result:?}"
591 );
592 }
593}