import logging import os import queue import re import threading from datetime import datetime from logging.config import dictConfig from logging.handlers import QueueHandler, QueueListener from typing import Optional import pytz import structlog from cachetools import cached from ddtrace.helpers import get_correlation_ids from flask import Response, request from flask.logging import wsgi_errors_stream from rapidjson import dumps, loads import delphi_api from delphi_api.auth import ProxyResourceProtector from delphi_api.const import CLIENTS_LIST_SECRET_NAME, ENVIRONMENT, LOG_LEVEL from delphi_api.core.caches import TracingCache from delphi_api.core.secrets_manager import get_secret from delphi_api.errors import AuthError LOG = structlog.get_logger(__name__) def thread_info_injection(logger, method_name, event_dict): thread = threading.current_thread() event_dict['thread_no'] = thread.ident event_dict['thread_name'] = thread.name event_dict['process_pid'] = os.getpid() return event_dict def tracer_injection(logger, log_method, event_dict) -> dict: """Add DataDog tracing fields to our structured log entries via structlog processor interface """ # get correlation ids from current tracer context trace_id, span_id = get_correlation_ids() # add ids to structlog event dictionary # if no trace present, set ids to 0 event_dict['dd.trace_id'] = trace_id or 0 event_dict['dd.span_id'] = span_id or 0 stack = event_dict.pop('stack', None) if stack: event_dict['error'] = event_dict.get('error', {}) event_dict['error']['stack'] = stack return event_dict def configure_logging(cache_logger: bool): """Configure our (global) logging format, handlers, and processors """ # Logging config before app instantiation dictConfig({ 'version': 1, 'disable_existing_loggers': False, 'formatters': { 'default': { 'format': '%(message)s', } }, 'root': { 'level': LOG_LEVEL, 'handlers': [] }, }) # instantiate queue & attach it to handler for async/non-blocking logging log_queue = queue.Queue(-1) # no limit on size queue_handler = QueueHandler(log_queue) handler = logging.StreamHandler(stream=wsgi_errors_stream) listener = QueueListener(log_queue, handler) root = logging.getLogger() root.addHandler(queue_handler) listener.start() structlog.configure( processors=[ # This performs the initial filtering, so we don't # evaluate e.g. DEBUG when unnecessary structlog.stdlib.filter_by_level, # Adds logger=module_name (e.g __main__) structlog.stdlib.add_logger_name, # Adds level=info, debug, etc. structlog.stdlib.add_log_level, # Performs the % string interpolation as expected structlog.stdlib.PositionalArgumentsFormatter(), # Include the stack when stack_info=True structlog.processors.StackInfoRenderer(), # Include the exception when exc_info=True # e.g log.exception() or log.warning(exc_info=True)'s behavior structlog.processors.format_exc_info, # Decodes the unicode values in any kv pairs structlog.processors.UnicodeDecoder(), # Creates the necessary args, kwargs for log() structlog.stdlib.render_to_log_kwargs, thread_info_injection, tracer_injection, structlog.processors.JSONRenderer(dumps, sort_keys=True), ], context_class=structlog.threadlocal.wrap_dict(dict), logger_factory=structlog.stdlib.LoggerFactory(), cache_logger_on_first_use=cache_logger, ) @cached(TracingCache) def get_client_name(client_id: str) -> Optional[str]: if not client_id: return secret = get_secret(CLIENTS_LIST_SECRET_NAME) if not secret: return secret = loads(secret) return secret.get(client_id, '') @cached(TracingCache) def get_client_id(token: str) -> Optional[str]: if not token: return try: decoded_token = ProxyResourceProtector.decode_token(token) except AuthError: # don't raise an error for log tracing return return decoded_token.get('azp', '') def _log_request(response: Response) -> Response: """Flask hook for requests to log detailed, structured information by default. Call the public :func:`log_request` function instead of this method. """ # Do not log the following requests if they include: if request.path in {'/favicon.ico', '/health'}: return response elif any([ bool(re.match(r'/(v[0-9]*)/(openapi|ui)', request.path, flags=re.IGNORECASE)), request.path.startswith('/static'), request.path.startswith('/ui'), ]): return response ip = request.headers.get('X-Forwarded-For', request.remote_addr) host = request.host.split(':', 1)[0] args = dict(request.args) # note: avoid adding any DD reserved attributes below # see: https://docs.datadoghq.com/logs/log_collection/?tab=tcpussite#reserved-attributes params = { 'app_env': ENVIRONMENT, 'hostname': host, 'ip': ip, 'method': request.method, 'params': args, 'path': request.path, 'pkg_version': delphi_api.__version__, 'resp_status': response.status, 'resp_status_code': response.status_code, 'timestamp': datetime.now(tz=pytz.utc).timestamp(), } client_name = None client_id = None if request.headers.get('Authorization'): parts = request.headers['Authorization'].split() if len(parts) == 2: token = parts[1] client_id = get_client_id(token) client_name = get_client_name(client_id) if client_id: params.update({'client_id': client_id}) if client_name: params.update({'client_name': client_name}) message = f'{response.status_code} {request.method} {request.path}' LOG.info(message, **params) return response def log_request(response: Response) -> Response: """Flask hook for requests to log detailed, structured information by default.""" try: return _log_request(response) except Exception as e: LOG.exception('Unable to log request details: %s' % str(e)) return response