Spaces:
Paused
Paused
| """ | |
| Structured Logging with Correlation ID — CrowData | |
| """ | |
| import logging | |
| import json | |
| import uuid | |
| import sys | |
| from datetime import datetime | |
| from contextvars import ContextVar | |
| from typing import Optional, Dict, Any | |
| from functools import wraps | |
| import asyncio | |
| # Context variable to store correlation ID across async calls | |
| correlation_id_var: ContextVar[str] = ContextVar('correlation_id', default='') | |
| request_id_var: ContextVar[str] = ContextVar('request_id', default='') | |
| class JSONFormatter(logging.Formatter): | |
| """JSON formatter for structured logging.""" | |
| def format(self, record: logging.LogRecord) -> str: | |
| log_data = { | |
| "timestamp": datetime.utcnow().isoformat() + "Z", | |
| "level": record.levelname, | |
| "logger": record.name, | |
| "message": record.getMessage(), | |
| "module": record.module, | |
| "function": record.funcName, | |
| "line": record.lineno, | |
| } | |
| # Add correlation IDs | |
| corr_id = correlation_id_var.get() | |
| if corr_id: | |
| log_data["correlation_id"] = corr_id | |
| req_id = request_id_var.get() | |
| if req_id: | |
| log_data["request_id"] = req_id | |
| # Add extra fields from record | |
| for key, value in record.__dict__.items(): | |
| if key not in ('name', 'msg', 'args', 'created', 'filename', 'funcName', | |
| 'levelname', 'levelno', 'lineno', 'module', 'msecs', | |
| 'message', 'msg', 'name', 'pathname', 'process', | |
| 'processName', 'relativeCreated', 'thread', | |
| 'threadName', 'exc_info', 'exc_text', 'stack_info'): | |
| log_data[key] = value | |
| # Handle exceptions | |
| if record.exc_info: | |
| log_data["exception"] = self.formatException(record.exc_info) | |
| return json.dumps(log_data, ensure_ascii=False) | |
| class StructuredLogger: | |
| """Wrapper for structured logging with correlation IDs.""" | |
| def __init__(self, name: str): | |
| self.logger = logging.getLogger(name) | |
| self._setup() | |
| def _setup(self): | |
| handler = logging.StreamHandler(sys.stdout) | |
| handler.setFormatter(JSONFormatter()) | |
| self.logger.addHandler(handler) | |
| self.logger.setLevel(logging.INFO) | |
| self.logger.propagate = False | |
| def _log(self, level: int, message: str, **kwargs): | |
| extra = {"extra_fields": kwargs} | |
| self.logger.log(level, message, extra=extra) | |
| def debug(self, message: str, **kwargs): | |
| self._log(logging.DEBUG, message, **kwargs) | |
| def info(self, message: str, **kwargs): | |
| self._log(logging.INFO, message, **kwargs) | |
| def warning(self, message: str, **kwargs): | |
| self._log(logging.WARNING, message, **kwargs) | |
| def error(self, message: str, **kwargs): | |
| self._log(logging.ERROR, message, **kwargs) | |
| def critical(self, message: str, **kwargs): | |
| self._log(logging.CRITICAL, message, **kwargs) | |
| def exception(self, message: str, **kwargs): | |
| kwargs["exc_info"] = True | |
| self._log(logging.ERROR, message, **kwargs) | |
| # Context management | |
| def set_correlation_id(correlation_id: str = None) -> str: | |
| """Set correlation ID for current context. Returns the ID.""" | |
| if not correlation_id: | |
| correlation_id = str(uuid.uuid4())[:8] | |
| correlation_id_var.set(correlation_id) | |
| return correlation_id | |
| def get_correlation_id() -> str: | |
| return correlation_id_var.get() | |
| def set_request_id(request_id: str): | |
| request_id_var.set(request_id) | |
| def get_request_id() -> str: | |
| return request_id_var.get() | |
| def with_correlation_id(func): | |
| """Decorator to inject correlation ID into async functions.""" | |
| async def wrapper(*args, **kwargs): | |
| corr_id = get_correlation_id() | |
| if not corr_id: | |
| corr_id = set_correlation_id() | |
| return await func(*args, **kwargs) | |
| return wrapper | |
| class CorrelationMiddleware: | |
| """Middleware to inject correlation ID into requests.""" | |
| def __init__(self, app): | |
| self.app = app | |
| async def __call__(self, scope, receive, send): | |
| if scope["type"] != "http": | |
| await self.app(scope, receive, send) | |
| return | |
| # Extract or generate correlation ID | |
| correlation_id = None | |
| headers = dict(scope.get("headers", [])) | |
| for k, v in headers: | |
| if k.decode() == "x-correlation-id": | |
| correlation_id = v.decode() | |
| break | |
| if not correlation_id: | |
| correlation_id = str(uuid.uuid4())[:8] | |
| set_correlation_id(correlation_id) | |
| async def send_wrapper(message): | |
| if message["type"] == "http.response.start": | |
| headers = list(message.get("headers", [])) | |
| headers.append((b"x-correlation-id", correlation_id.encode())) | |
| message["headers"] = headers | |
| await send(message) | |
| await self.app(scope, receive, send_wrapper) | |
| # Helper for creating structured log entries | |
| def log_structured(logger: logging.Logger, level: int, message: str, | |
| correlation_id: str = None, **fields): | |
| """Log structured message with optional correlation ID.""" | |
| extra = {"extra_fields": fields} | |
| if correlation_id: | |
| extra["correlation_id"] = correlation_id | |
| logger.log(level, message, extra=extra) |