Spaces:
Sleeping
Sleeping
| """Recover reports stuck in ``generating`` after worker/API crashes.""" | |
| from __future__ import annotations | |
| import logging | |
| from datetime import UTC, datetime, timedelta | |
| from sqlalchemy import and_, or_, select | |
| from app.config import settings | |
| from app.db.database import get_session_factory | |
| from app.db.models import Report, ReportStatus | |
| logger = logging.getLogger(__name__) | |
| async def sweep_stale_generating_reports() -> int: | |
| """Mark long-running ``generating`` reports as failed so the UI can retry.""" | |
| timeout_s = int(getattr(settings, "generation_timeout_seconds", 3600)) | |
| cutoff = datetime.now(UTC) - timedelta(seconds=timeout_s) | |
| factory = get_session_factory() | |
| async with factory() as db: | |
| result = await db.execute( | |
| select(Report).where( | |
| Report.status == ReportStatus.generating, | |
| or_( | |
| Report.generation_started_at < cutoff, | |
| and_( | |
| Report.generation_started_at.is_(None), | |
| Report.updated_at < cutoff, | |
| ), | |
| ), | |
| ) | |
| ) | |
| stale = list(result.scalars().all()) | |
| if not stale: | |
| return 0 | |
| for report in stale: | |
| report.status = ReportStatus.failed | |
| report.generation_started_at = None | |
| sla_s = int(getattr(settings, "generation_sla_seconds", 600)) | |
| report.error_message = ( | |
| f"Generation timed out after {timeout_s}s " | |
| f"(target SLA {sla_s}s / {sla_s // 60} min). " | |
| "Try fewer sections, medium interference, or re-generate failed sections." | |
| )[:2000] | |
| await db.commit() | |
| for report in stale: | |
| logger.warning( | |
| "Marked stale generating report failed id=%s tenant=%s", | |
| report.id, | |
| report.tenant_id, | |
| ) | |
| return len(stale) | |
| async def generation_stale_sweeper_loop() -> None: | |
| """Background loop; no-op when interval is zero.""" | |
| import asyncio | |
| interval = int(getattr(settings, "generation_stale_sweep_seconds", 120)) | |
| if interval <= 0: | |
| return | |
| while True: | |
| try: | |
| n = await sweep_stale_generating_reports() | |
| if n: | |
| logger.info("generation_stale_sweep recovered %d report(s)", n) | |
| except Exception: # noqa: BLE001 | |
| logger.exception("generation_stale_sweep failed") | |
| await asyncio.sleep(interval) | |