from __future__ import annotations import logging from fastapi import FastAPI, Header, HTTPException, Query, status from redis.asyncio import Redis from app.audit_logger import AuditLogger from app.config import settings from app.models import IngestAck, ProcessorForwardEnvelope from app.queue_repo import RedisQueueRepository logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s %(message)s", ) logger = logging.getLogger(__name__) app = FastAPI( title="alert-processor-ingest", version="2.1.0", description="Ingress endpoint that queues normalized events for async processing.", ) redis_client: Redis | None = None queue_repo: RedisQueueRepository | None = None audit_logger: AuditLogger | None = None def _extract_token( x_internal_token: str | None, authorization: str | None, ) -> str | None: if x_internal_token: return x_internal_token.strip() if authorization: auth = authorization.strip() if auth.lower().startswith("bearer "): return auth[7:].strip() return auth return None def _validate_token( x_internal_token: str | None, authorization: str | None, ) -> None: if not settings.require_internal_api_token: return provided = _extract_token(x_internal_token, authorization) if not provided or provided != settings.internal_api_token: raise HTTPException( status_code=status.HTTP_401_UNAUTHORIZED, detail="Invalid or missing internal API token", ) @app.on_event("startup") async def startup_check() -> None: global redis_client, queue_repo, audit_logger if settings.require_internal_api_token and not settings.internal_api_token: raise RuntimeError( "INTERNAL_API_TOKEN is required, but not set. " "Set it in environment variables or in .env file." ) if not settings.redis_enabled: raise RuntimeError("REDIS_ENABLED must be true for async queue mode") if not settings.queue_enabled: raise RuntimeError("QUEUE_ENABLED must be true for async queue mode") redis_client = Redis.from_url( settings.redis_url, encoding="utf-8", decode_responses=False, ) queue_repo = RedisQueueRepository( client=redis_client, queue_name=settings.queue_name, processing_name=settings.queue_processing_name, deadletter_name=settings.queue_deadletter_name, dedup_ttl_seconds=settings.queue_dedup_ttl_seconds, ) await queue_repo.ping() if settings.audit_enabled: audit_logger = AuditLogger( client=redis_client, key_prefix=settings.audit_key_prefix, ttl_seconds=settings.audit_ttl_seconds, max_stage_records=settings.audit_max_stage_records, ) logger.info( "Async ingest initialized: redis_url=%s queue_name=%s processing_name=%s deadletter_name=%s audit_enabled=%s", settings.redis_url, settings.queue_name, settings.queue_processing_name, settings.queue_deadletter_name, settings.audit_enabled, ) @app.on_event("shutdown") async def shutdown_event() -> None: global redis_client if redis_client is not None: await redis_client.close() logger.info("Redis connection closed") @app.get("/health") async def health() -> dict[str, str | bool | int]: return { "status": "ok", "service": settings.app_name, "mode": "async_ingest", "redis_enabled": settings.redis_enabled, "redis_connected": redis_client is not None, "queue_enabled": settings.queue_enabled, "queue_configured": queue_repo is not None, "queue_name": settings.queue_name, "queue_processing_name": settings.queue_processing_name, "queue_deadletter_name": settings.queue_deadletter_name, "audit_enabled": settings.audit_enabled, "audit_configured": audit_logger is not None, } @app.post( "/internal/events", response_model=IngestAck, status_code=status.HTTP_202_ACCEPTED, ) async def receive_internal_event( envelope: ProcessorForwardEnvelope, x_internal_token: str | None = Header(default=None), authorization: str | None = Header(default=None), ) -> IngestAck: _validate_token(x_internal_token, authorization) if queue_repo is None: raise HTTPException( status_code=status.HTTP_503_SERVICE_UNAVAILABLE, detail="Queue repository is not initialized", ) event = envelope.event job_id = await queue_repo.enqueue_event(envelope) if audit_logger is not None: await audit_logger.log_ingest_queued( envelope=envelope, job_id=job_id, queue_name=settings.queue_name, ) logger.info( "Queued event for async processing: job_id=%s correlation_id=%s event_id=%s severity=%s host=%s trigger=%s", job_id, event.correlation_id, event.event_id, event.severity, event.host, event.trigger_name, ) return IngestAck( accepted=True, correlation_id=event.correlation_id, queued=True, job_id=job_id, message="Event queued for async processing", ) @app.get("/audit/events/{correlation_id}") async def get_audit_by_correlation(correlation_id: str) -> dict: if audit_logger is None: raise HTTPException(status_code=503, detail="Audit logger is not enabled") payload = await audit_logger.get_event_audit(correlation_id=correlation_id) if payload is None: raise HTTPException(status_code=404, detail="Audit record not found") return payload @app.get("/audit/by-event/{event_id}") async def get_audit_by_event_id(event_id: str) -> dict: if audit_logger is None: raise HTTPException(status_code=503, detail="Audit logger is not enabled") payload = await audit_logger.get_event_audit(event_id=event_id) if payload is None: raise HTTPException(status_code=404, detail="Audit record not found") return payload @app.get("/audit/recent") async def get_recent_audit( limit: int = Query(default=20, ge=1, le=100), ) -> list[dict]: if audit_logger is None: raise HTTPException(status_code=503, detail="Audit logger is not enabled") return await audit_logger.list_recent(limit=limit)