import logging import time from datetime import datetime from ddtrace import tracer from ddtrace import config as dd_config from ddtrace.internal.utils.http import redact_url from stream.client import StreamClient as BaseStreamClient import requests from src import config assert config.STREAM_APP_ID, "Stream app id is required" assert config.STREAM_API_KEY, "Stream API key is required" assert config.STREAM_API_SECRET, "Stream API secret is required" assert config.STREAM_API_REGION, "Stream API region is required" # client = stream.connect(config.STREAM_API_KEY, config.STREAM_API_SECRET, location=config.STREAM_API_REGION) logger = logging.getLogger(__name__) saved_headers = {} class StreamClient(BaseStreamClient): def _make_request( self, method, relative_url, signature, service_name="api", params=None, data=None, ): if not tracer.enabled: return super()._make_request(method, relative_url, signature, service_name, params, data) resource_name = f'{method.__name__.upper()}: /{relative_url.strip("/")}' with tracer.trace('request', resource=resource_name, service='getstream'): return super()._make_request(method, relative_url, signature, service_name, params, data) def _parse_response(self, response: requests.Response | None): try: return super()._parse_response(response) finally: self.__finalize_request(response) @staticmethod def __finalize_request(response: requests.Response | None): if response is None: return # Save last call headers to check ratelimit global saved_headers saved_headers = response.headers if not tracer.enabled: return try: span = tracer.current_span() # Set http tags span.set_tag('http.status_code', response.status_code) request = getattr(response, 'request', None) if request is not None: url = redact_url(request.url, dd_config._obfuscation_query_string_pattern) span.set_tag_str('http.url', url) span.set_tag_str('http.method', request.method.upper()) # Set app tags ratelimit_limit = response.headers.get('X-Ratelimit-Limit') ratelimit_remaining = response.headers.get('X-Ratelimit-Remaining') ratelimit_reset = response.headers.get('X-Ratelimit-Reset') span.set_tag('stream.ratelimit_limit', int(ratelimit_limit) if ratelimit_limit else None) span.set_tag('stream.ratelimit_remaining', int(ratelimit_remaining) if ratelimit_remaining else None) span.set_tag('stream.ratelimit_reset', datetime.fromtimestamp(int(ratelimit_reset)) if ratelimit_reset else None) span.set_tag('stream.ratelimit_reset_in_seconds', (int(ratelimit_reset) - int(time.time())) if ratelimit_reset else None) except: logger.debug('datadog: error adding tags', exc_info=True) client = StreamClient(config.STREAM_API_KEY, config.STREAM_API_SECRET, None, location=config.STREAM_API_REGION)