- Remove duplicate function definitions with hardcoded "AI-оператор" strings (ai_voice_runtime, ai_orchestrator, voice_name_config, voice.py) - Remove unreachable dead code after return in ai_voice_runtime - Add SQL LIMIT to 17 unbounded queries across 12 services to prevent OOM - Move Python-side filtering to SQL WHERE in reporting_service - Downgrade 19 logger.warning to logger.info for normal-flow events in media_runtime Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
96 lines
2.9 KiB
Python
96 lines
2.9 KiB
Python
from __future__ import annotations
|
|
|
|
import threading
|
|
|
|
from fastapi import Depends, FastAPI, HTTPException
|
|
from sqlalchemy import select
|
|
|
|
from services.shared.core import Role
|
|
from services.shared.db import get_session
|
|
from services.shared.event_bus import (
|
|
dispatch_outbox_batch,
|
|
event_bus_enabled,
|
|
outbox_item_from_row,
|
|
poll_forever,
|
|
retry_outbox_event,
|
|
)
|
|
from services.shared.models import EventOutboxItem, HealthResponse
|
|
from services.shared.security import require_roles
|
|
from services.shared.sql_init import init_sql_schema
|
|
from services.shared.sql_models import EventOutboxRow
|
|
|
|
app = FastAPI(title="event-bus-service", version="1.0.0")
|
|
|
|
init_sql_schema()
|
|
|
|
|
|
@app.get("/health", response_model=HealthResponse)
|
|
def health() -> HealthResponse:
|
|
suffix = "enabled" if event_bus_enabled() else "disabled"
|
|
return HealthResponse(status="ok", service=f"event-bus-service ({suffix})")
|
|
|
|
|
|
@app.get("/bus/outbox", response_model=list[EventOutboxItem])
|
|
def list_outbox(
|
|
status: str | None = None,
|
|
limit: int = 100,
|
|
_: dict = Depends(require_roles(Role.ADMIN)),
|
|
) -> list[EventOutboxItem]:
|
|
session = get_session()
|
|
try:
|
|
stmt = select(EventOutboxRow).order_by(EventOutboxRow.id.desc())
|
|
if status:
|
|
stmt = stmt.where(EventOutboxRow.status == status)
|
|
rows = session.execute(stmt.limit(max(limit, 1))).scalars().all()
|
|
return [outbox_item_from_row(row) for row in rows]
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
@app.get("/bus/outbox/{event_id}", response_model=EventOutboxItem)
|
|
def get_outbox_event(
|
|
event_id: str,
|
|
_: dict = Depends(require_roles(Role.ADMIN)),
|
|
) -> EventOutboxItem:
|
|
session = get_session()
|
|
try:
|
|
row = session.execute(select(EventOutboxRow).where(EventOutboxRow.event_id == event_id)).scalar_one_or_none()
|
|
if not row:
|
|
raise HTTPException(status_code=404, detail="Outbox event not found")
|
|
return outbox_item_from_row(row)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
@app.post("/bus/outbox/{event_id}/retry", response_model=EventOutboxItem)
|
|
def retry_outbox(
|
|
event_id: str,
|
|
_: dict = Depends(require_roles(Role.ADMIN)),
|
|
) -> EventOutboxItem:
|
|
session = get_session()
|
|
try:
|
|
row = session.execute(select(EventOutboxRow).where(EventOutboxRow.event_id == event_id)).scalar_one_or_none()
|
|
if not row:
|
|
raise HTTPException(status_code=404, detail="Outbox event not found")
|
|
retry_outbox_event(session, row)
|
|
session.commit()
|
|
session.refresh(row)
|
|
return outbox_item_from_row(row)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
def _dispatch_once() -> None:
|
|
session = get_session()
|
|
try:
|
|
dispatch_outbox_batch(session)
|
|
finally:
|
|
session.close()
|
|
|
|
|
|
@app.on_event("startup")
|
|
def _startup() -> None:
|
|
if not event_bus_enabled():
|
|
return
|
|
threading.Thread(target=lambda: poll_forever(_dispatch_once), daemon=True).start()
|