Skip to content

Commit 38927a1

Browse files
committed
feat(fcm): Migrate topic management to FCM v1 API
1 parent 62e559b commit 38927a1

2 files changed

Lines changed: 648 additions & 21 deletions

File tree

‎firebase_admin/messaging.py‎

Lines changed: 316 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -15,14 +15,17 @@
1515
"""Firebase Cloud Messaging module."""
1616

1717
from __future__ import annotations
18-
from typing import Any, Callable, Dict, List, Optional, cast
18+
import asyncio
1919
import concurrent.futures
2020
import json
21-
import asyncio
2221
import logging
22+
import re
23+
from typing import Any, Callable, Dict, List, Optional, cast
24+
import urllib.parse
2325
import warnings
24-
import requests
26+
2527
import httpx
28+
import requests
2629

2730
import firebase_admin
2831
from firebase_admin import (
@@ -73,7 +76,11 @@
7376
'send_each_for_multicast',
7477
'send_each_for_multicast_async',
7578
'subscribe_to_topic',
79+
'subscribe_to_topic_async',
80+
'subscribe_to_topic_legacy',
7681
'unsubscribe_from_topic',
82+
'unsubscribe_from_topic_async',
83+
'unsubscribe_from_topic_legacy',
7784
]
7885

7986

@@ -255,6 +262,44 @@ def send_each_for_multicast(multicast_message, dry_run=False, app=None):
255262
def subscribe_to_topic(tokens, topic, app=None):
256263
"""Subscribes a list of registration tokens to an FCM topic.
257264
265+
Args:
266+
tokens: A non-empty list of device registration tokens. List may not have more than 1000
267+
elements.
268+
topic: Name of the topic to subscribe to. May contain the ``/topics/`` prefix.
269+
app: An App instance (optional).
270+
271+
Returns:
272+
TopicManagementResponse: A ``TopicManagementResponse`` instance.
273+
274+
Raises:
275+
FirebaseError: If an error occurs while communicating with the FCM service.
276+
ValueError: If the input arguments are invalid.
277+
"""
278+
return _get_messaging_service(app).subscribe_to_topic(tokens, topic)
279+
280+
async def subscribe_to_topic_async(tokens, topic, app=None):
281+
"""Subscribes a list of registration tokens to an FCM topic asynchronously.
282+
283+
Args:
284+
tokens: A non-empty list of device registration tokens. List may not have more than 1000
285+
elements.
286+
topic: Name of the topic to subscribe to. May contain the ``/topics/`` prefix.
287+
app: An App instance (optional).
288+
289+
Returns:
290+
TopicManagementResponse: A ``TopicManagementResponse`` instance.
291+
292+
Raises:
293+
FirebaseError: If an error occurs while communicating with the FCM service.
294+
ValueError: If the input arguments are invalid.
295+
"""
296+
return await _get_messaging_service(app).subscribe_to_topic_async(tokens, topic)
297+
298+
def subscribe_to_topic_legacy(tokens, topic, app=None):
299+
"""Subscribes a list of registration tokens to an FCM topic using the legacy Instance ID API.
300+
301+
subscribe_to_topic_legacy is deprecated. Use subscribe_to_topic instead.
302+
258303
Args:
259304
tokens: A non-empty list of device registration tokens. List may not have more than 1000
260305
elements.
@@ -268,12 +313,55 @@ def subscribe_to_topic(tokens, topic, app=None):
268313
FirebaseError: If an error occurs while communicating with instance ID service.
269314
ValueError: If the input arguments are invalid.
270315
"""
316+
warnings.warn(
317+
'subscribe_to_topic_legacy is deprecated. Use subscribe_to_topic instead.',
318+
DeprecationWarning,
319+
stacklevel=2)
271320
return _get_messaging_service(app).make_topic_management_request(
272321
tokens, topic, 'iid/v1:batchAdd')
273322

274323
def unsubscribe_from_topic(tokens, topic, app=None):
275324
"""Unsubscribes a list of registration tokens from an FCM topic.
276325
326+
Args:
327+
tokens: A non-empty list of device registration tokens. List may not have more than 1000
328+
elements.
329+
topic: Name of the topic to unsubscribe from. May contain the ``/topics/`` prefix.
330+
app: An App instance (optional).
331+
332+
Returns:
333+
TopicManagementResponse: A ``TopicManagementResponse`` instance.
334+
335+
Raises:
336+
FirebaseError: If an error occurs while communicating with the FCM service.
337+
ValueError: If the input arguments are invalid.
338+
"""
339+
return _get_messaging_service(app).unsubscribe_from_topic(tokens, topic)
340+
341+
async def unsubscribe_from_topic_async(tokens, topic, app=None):
342+
"""Unsubscribes a list of registration tokens from an FCM topic asynchronously.
343+
344+
Args:
345+
tokens: A non-empty list of device registration tokens. List may not have more than 1000
346+
elements.
347+
topic: Name of the topic to unsubscribe from. May contain the ``/topics/`` prefix.
348+
app: An App instance (optional).
349+
350+
Returns:
351+
TopicManagementResponse: A ``TopicManagementResponse`` instance.
352+
353+
Raises:
354+
FirebaseError: If an error occurs while communicating with the FCM service.
355+
ValueError: If the input arguments are invalid.
356+
"""
357+
return await _get_messaging_service(app).unsubscribe_from_topic_async(tokens, topic)
358+
359+
def unsubscribe_from_topic_legacy(tokens, topic, app=None):
360+
"""Unsubscribes a list of registration tokens from an FCM topic using the legacy
361+
Instance ID API.
362+
363+
unsubscribe_from_topic_legacy is deprecated. Use unsubscribe_from_topic instead.
364+
277365
Args:
278366
tokens: A non-empty list of device registration tokens. List may not have more than 1000
279367
elements.
@@ -287,6 +375,10 @@ def unsubscribe_from_topic(tokens, topic, app=None):
287375
FirebaseError: If an error occurs while communicating with instance ID service.
288376
ValueError: If the input arguments are invalid.
289377
"""
378+
warnings.warn(
379+
'unsubscribe_from_topic_legacy is deprecated. Use unsubscribe_from_topic instead.',
380+
DeprecationWarning,
381+
stacklevel=2)
290382
return _get_messaging_service(app).make_topic_management_request(
291383
tokens, topic, 'iid/v1:batchRemove')
292384

@@ -410,7 +502,9 @@ def __init__(self, app: App) -> None:
410502
'Project ID is required to access Cloud Messaging service. Either set the '
411503
'projectId option, or use service account credentials. Alternatively, set the '
412504
'GOOGLE_CLOUD_PROJECT environment variable.')
505+
self._project_id = project_id
413506
self._fcm_url = _MessagingService.FCM_URL.format(project_id)
507+
self._fcm_topic_url = f'https://fcm.googleapis.com/v1/projects/{project_id}/registrations'
414508
self._fcm_headers = {
415509
'X-GOOG-API-FORMAT-VERSION': '2',
416510
'X-FIREBASE-CLIENT': f'fire-admin-python/{firebase_admin.__version__}',
@@ -499,6 +593,225 @@ async def send_data(data):
499593
message=f'Unknown error while making remote service calls: {error}',
500594
cause=error)
501595

596+
def _validate_topic_management_args(self, tokens, topic):
597+
"""Validates and formats topic management arguments."""
598+
if isinstance(tokens, str):
599+
tokens = [tokens]
600+
if not isinstance(tokens, list) or not tokens:
601+
raise ValueError('Tokens must be a string or a non-empty list of strings.')
602+
invalid_str = [t for t in tokens if not isinstance(t, str) or not t]
603+
if invalid_str:
604+
raise ValueError('Tokens must be non-empty strings.')
605+
if len(tokens) > 1000:
606+
raise ValueError('tokens must not contain more than 1000 elements.')
607+
608+
if not isinstance(topic, str) or not topic:
609+
raise ValueError('Topic must be a non-empty string.')
610+
topic_name = topic
611+
if topic_name.startswith('/topics/'):
612+
topic_name = topic_name[len('/topics/'):]
613+
if not topic_name or not re.match(r'^[a-zA-Z0-9-_\.~%]+$', topic_name):
614+
raise ValueError('Malformed topic name.')
615+
616+
return tokens, topic_name
617+
618+
def subscribe_to_topic(self, tokens, topic) -> TopicManagementResponse:
619+
"""Subscribes a list of registration tokens to an FCM topic via the FCM v1 API."""
620+
return self._make_topic_management_request_v1(tokens, topic, is_subscribe=True)
621+
622+
def unsubscribe_from_topic(self, tokens, topic) -> TopicManagementResponse:
623+
"""Unsubscribes a list of registration tokens from an FCM topic via the FCM v1 API."""
624+
return self._make_topic_management_request_v1(tokens, topic, is_subscribe=False)
625+
626+
async def subscribe_to_topic_async(self, tokens, topic) -> TopicManagementResponse:
627+
"""Subscribes a list of registration tokens to an FCM topic asynchronously
628+
via the FCM v1 API."""
629+
return await self._make_topic_management_request_v1_async(
630+
tokens, topic, is_subscribe=True)
631+
632+
async def unsubscribe_from_topic_async(self, tokens, topic) -> TopicManagementResponse:
633+
"""Unsubscribes a list of registration tokens from an FCM topic asynchronously
634+
via the FCM v1 API."""
635+
return await self._make_topic_management_request_v1_async(
636+
tokens, topic, is_subscribe=False)
637+
638+
def _make_topic_management_request_v1(
639+
self, tokens, topic, is_subscribe: bool
640+
) -> TopicManagementResponse:
641+
"""Helper method that sends topic subscription requests via FCM v1 API."""
642+
tokens_list, topic_name = self._validate_topic_management_args(tokens, topic)
643+
644+
def send_request(token: str):
645+
encoded_token = urllib.parse.quote(token, safe='')
646+
encoded_topic = urllib.parse.quote(topic_name, safe='')
647+
base_url = f'{self._fcm_topic_url}/{encoded_token}/topicSubscriptions'
648+
if is_subscribe:
649+
url = f'{base_url}?topic_name={encoded_topic}'
650+
method = 'post'
651+
json_data = {}
652+
else:
653+
url = f'{base_url}/{encoded_topic}?allow_missing=true'
654+
method = 'delete'
655+
json_data = None
656+
657+
try:
658+
self._client.request(
659+
method,
660+
url=url,
661+
headers=self._fcm_headers,
662+
json=json_data,
663+
)
664+
return {'success': True}
665+
except requests.exceptions.RequestException as error:
666+
return self._build_topic_subscription_result_from_requests_error(
667+
error, is_subscribe)
668+
669+
try:
670+
with concurrent.futures.ThreadPoolExecutor(
671+
max_workers=min(len(tokens_list), 100)
672+
) as executor:
673+
results = list(executor.map(send_request, tokens_list))
674+
return self._parse_topic_management_results(results)
675+
except Exception as error:
676+
raise exceptions.UnknownError(
677+
message=f'Unknown error while making remote service calls: {error}',
678+
cause=error)
679+
680+
async def _make_topic_management_request_v1_async(
681+
self, tokens, topic, is_subscribe: bool
682+
) -> TopicManagementResponse:
683+
"""Helper method that sends topic subscription requests asynchronously via FCM v1 API."""
684+
tokens_list, topic_name = self._validate_topic_management_args(tokens, topic)
685+
semaphore = asyncio.Semaphore(100)
686+
687+
async def send_request_async(token: str):
688+
encoded_token = urllib.parse.quote(token, safe='')
689+
encoded_topic = urllib.parse.quote(topic_name, safe='')
690+
base_url = f'{self._fcm_topic_url}/{encoded_token}/topicSubscriptions'
691+
if is_subscribe:
692+
url = f'{base_url}?topic_name={encoded_topic}'
693+
method = 'post'
694+
json_data = {}
695+
else:
696+
url = f'{base_url}/{encoded_topic}?allow_missing=true'
697+
method = 'delete'
698+
json_data = None
699+
700+
async with semaphore:
701+
try:
702+
await self._async_client.request(
703+
method,
704+
url=url,
705+
headers=self._fcm_headers,
706+
json=json_data,
707+
)
708+
return {'success': True}
709+
except httpx.HTTPError as error:
710+
return self._build_topic_subscription_result_from_httpx_error(
711+
error, is_subscribe)
712+
except requests.exceptions.RequestException as error:
713+
return self._build_topic_subscription_result_from_requests_error(
714+
error, is_subscribe)
715+
716+
try:
717+
results = await asyncio.gather(*[send_request_async(token) for token in tokens_list])
718+
return self._parse_topic_management_results(results)
719+
except Exception as error:
720+
raise exceptions.UnknownError(
721+
message=f'Unknown error while making remote service calls: {error}',
722+
cause=error)
723+
724+
@classmethod
725+
def _get_topic_error_code(cls, error_dict: dict, status_code: int) -> str:
726+
"""Extracts the error code for a topic subscription error response."""
727+
error_data = error_dict.get('error')
728+
if isinstance(error_data, str) and error_data:
729+
return error_data
730+
if isinstance(error_data, dict):
731+
details = error_data.get('details')
732+
if isinstance(details, list):
733+
fcm_error_type = 'type.googleapis.com/google.firebase.fcm.v1.FcmError'
734+
for element in details:
735+
if isinstance(element, dict) and element.get('@type') == fcm_error_type:
736+
code = element.get('errorCode')
737+
if code:
738+
return code
739+
status = error_data.get('status')
740+
if status:
741+
return status
742+
message = error_data.get('message')
743+
if message:
744+
return message
745+
746+
status_map = {
747+
400: 'INVALID_ARGUMENT',
748+
401: 'PERMISSION_DENIED',
749+
403: 'PERMISSION_DENIED',
750+
404: 'NOT_FOUND',
751+
429: 'RESOURCE_EXHAUSTED',
752+
500: 'INTERNAL',
753+
503: 'DEADLINE_EXCEEDED',
754+
}
755+
return status_map.get(status_code, 'UNKNOWN_ERROR')
756+
757+
def _build_topic_subscription_result_from_requests_error(self, error, is_subscribe):
758+
"""Constructs a result dict from a requests error."""
759+
if error.response is not None:
760+
if is_subscribe and error.response.status_code == 409:
761+
return {'success': True}
762+
error_dict = {}
763+
try:
764+
parsed = error.response.json()
765+
if isinstance(parsed, dict):
766+
error_dict = parsed
767+
except ValueError:
768+
pass
769+
770+
error_data = error_dict.get('error')
771+
if is_subscribe and isinstance(error_data, dict) and (
772+
error_data.get('status') == 'ALREADY_EXISTS'
773+
):
774+
return {'success': True}
775+
776+
error_code = self._get_topic_error_code(error_dict, error.response.status_code)
777+
return {'success': False, 'error': error_code}
778+
779+
return {'success': False, 'error': 'UNKNOWN_ERROR'}
780+
781+
def _build_topic_subscription_result_from_httpx_error(self, error, is_subscribe):
782+
"""Constructs a result dict from an httpx error."""
783+
if isinstance(error, httpx.HTTPStatusError):
784+
if is_subscribe and error.response.status_code == 409:
785+
return {'success': True}
786+
error_dict = {}
787+
try:
788+
parsed = error.response.json()
789+
if isinstance(parsed, dict):
790+
error_dict = parsed
791+
except ValueError:
792+
pass
793+
794+
error_data = error_dict.get('error')
795+
if is_subscribe and isinstance(error_data, dict) and (
796+
error_data.get('status') == 'ALREADY_EXISTS'
797+
):
798+
return {'success': True}
799+
800+
error_code = self._get_topic_error_code(error_dict, error.response.status_code)
801+
return {'success': False, 'error': error_code}
802+
803+
return {'success': False, 'error': 'UNKNOWN_ERROR'}
804+
805+
def _parse_topic_management_results(self, results) -> TopicManagementResponse:
806+
"""Parses individual request results into a TopicManagementResponse."""
807+
formatted_results = []
808+
for result in results:
809+
if result.get('success'):
810+
formatted_results.append({})
811+
else:
812+
formatted_results.append({'error': result.get('error', 'UNKNOWN_ERROR')})
813+
return TopicManagementResponse({'results': formatted_results})
814+
502815
def make_topic_management_request(self, tokens, topic, operation):
503816
"""Invokes the IID service for topic management functionality."""
504817
if isinstance(tokens, str):

0 commit comments

Comments
 (0)