Загрузить файлы в «alert-processor/app»
This commit is contained in:
@@ -0,0 +1,216 @@
|
||||
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)
|
||||
Reference in New Issue
Block a user