Files
2026-05-28 14:19:27 +04:00

100 lines
3.5 KiB
Python

import asyncio
from contextlib import asynccontextmanager
from fastapi import FastAPI
from fastapi.responses import Response
from . import metrics
from .api import agents, checks, events, health, heartbeat, incidents, outbox, system_summary
from .db import dispose_engine, get_engine, get_sessionmaker
from .errors import install_error_handlers
from .logging_config import configure_logging, get_logger
from .middleware.auth import BearerAuthMiddleware
from .middleware.body_limit import BodyLimitMiddleware
from .middleware.request_id import RequestIdMiddleware
from .services.detector import run_detector
from .services.notifier_worker import run_notifier_worker
from .services.partitioning import ensure_event_partitions, run_partition_maintenance
from .settings import get_settings
log = get_logger("monlet.main")
async def _ensure_partitions_with_retry(engine, attempts: int = 6, base_delay: float = 1.0) -> None:
last_exc: Exception | None = None
for i in range(attempts):
try:
await ensure_event_partitions(engine)
return
except Exception as exc:
last_exc = exc
log.warning("ensure_event_partitions_retry", attempt=i + 1, error=str(exc))
await asyncio.sleep(min(base_delay * (2**i), 10.0))
assert last_exc is not None
raise last_exc
@asynccontextmanager
async def lifespan(app: FastAPI):
s = get_settings()
configure_logging(s.log_level)
engine = get_engine()
sm = get_sessionmaker()
# Mandatory: current month partition must exist before accepting writes.
# Retry briefly so a transient DB unavailability at startup (compose race)
# doesn't crash-loop the server.
await _ensure_partitions_with_retry(engine)
stop_event = asyncio.Event()
tasks: list[asyncio.Task] = []
# Partition maintenance is independent of detector — must run even when detector is off,
# otherwise events_future_partitions runs out and ingest fails.
tasks.append(
asyncio.create_task(run_partition_maintenance(sm, stop_event), name="partition_maintenance")
)
if s.enable_detector:
tasks.append(asyncio.create_task(run_detector(sm, stop_event), name="detector"))
if s.enable_notifier_worker:
tasks.append(
asyncio.create_task(run_notifier_worker(sm, stop_event), name="notifier_worker")
)
log.info("server_started", host=s.host, port=s.port)
try:
yield
finally:
stop_event.set()
for t in tasks:
try:
await asyncio.wait_for(t, timeout=5.0)
except TimeoutError:
t.cancel()
await dispose_engine()
def create_app() -> FastAPI:
app = FastAPI(title="Monlet Server", version="1.0.0", lifespan=lifespan)
app.add_middleware(BearerAuthMiddleware)
app.add_middleware(BodyLimitMiddleware)
app.add_middleware(RequestIdMiddleware)
install_error_handlers(app)
app.include_router(health.router, prefix="/api/v1")
app.include_router(heartbeat.router, prefix="/api/v1")
app.include_router(events.router, prefix="/api/v1")
app.include_router(agents.router, prefix="/api/v1")
app.include_router(checks.router, prefix="/api/v1")
app.include_router(incidents.router, prefix="/api/v1")
app.include_router(outbox.router, prefix="/api/v1")
app.include_router(system_summary.router, prefix="/api/v1")
@app.get("/metrics")
async def metrics_endpoint() -> Response:
body, ct = metrics.render()
return Response(content=body, media_type=ct)
return app
app = create_app()