"""Wrapper for sync getstream client class.""" import logging from typing import TYPE_CHECKING, Collection, Type from ddtrace import tracer from stream.client import StreamClient as BaseStreamClient from .. import config from ..feed import SyncFeed from ..utils import datadog, retry from .base import Activity, BaseExtendedClient if TYPE_CHECKING: from requests import Response from stream.feed.base import BaseFeed __all__ = ["SyncStreamClient"] logger = logging.getLogger(__name__) class SyncStreamClient(BaseExtendedClient, BaseStreamClient): """Extended Stream (getstream.io) Client. It does several things: 1. Save ratelimit info from response headers 2. Add DataDog traces to GetStream API calls 3. Add trace context to stream messages 4. Add correlation id to stream messages 5. Add throttling mechanism to prevent API rate limiting from being reached or exceeded. 6. Add retry mechanism to API calls in case of API rate limit been reached. """ def _get_feed_cls(self) -> Type["BaseFeed"]: return SyncFeed @tracer.wrap("stream.add_to_many", service=config.SERVICE_NAME) def add_to_many(self, activity: Activity, feeds: Collection[str]): """Override to add trace context to activity.""" self.inject_trace_data(activity) return super().add_to_many(activity, feeds) @retry.on_ratelimit_reached def _make_request( self, method, relative_url, signature, service_name="api", params=None, data=None, ): """Override to add retry and throttling functionality.""" self._throttler.delay(self) return super()._make_request(method, relative_url, signature, service_name, params, data) def _parse_response(self, response: "Response"): # Save last call ratelimit values self._ratelimit_info.set_from_headers(response.headers) request = getattr(response, "request", None) url = request.url if request else None method = request.method if request else None datadog.set_tags(response.status_code, url, method, self.ratelimit_info.to_dict()) return super()._parse_response(response)