File size: 2,429 Bytes
732b14f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
c893230
732b14f
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
"""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
            report.generation_section_total = None
            report.error_message = (
                f"Generation timed out after {timeout_s}s. "
                "The background worker may have stopped or Temporal may be unreachable."
            )[: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)