import asyncio
import json
import logging
from datetime import date, datetime, timedelta, timezone
from typing import Any

from apscheduler.schedulers.asyncio import AsyncIOScheduler
from sqlalchemy import text
from sqlalchemy.ext.asyncio import AsyncSession

from app.db import AsyncSessionLocal
from app.scrapers.base import random_delay
from app.scrapers.cire import scrape_cire_for_nif
from app.scrapers.cire_announcements import PAGE_CAP, CireAnnouncementSweeper
from app.scrapers.mj_irn import IrnVersionChanged, MjIrnClient
from app.scrapers.distribuicao import (
    discover_tribunal_codes,
    scrape_distribuicao_range,
    scrape_distribuicao_tribunal_sweep,
    short_search_term,
)
from app.services.nif_lookup import lookup_nif
from app.services.contracts_impic import sync_impic_contracts
from app.services.alvara_sync import sync_alvara_type
from app.services.dre import (
    fetch_and_archive_parte_e,
    fetch_dre_for_companies,
    retroactively_match_company,
    save_new_publications,
)
from app.services.cire_announcement_ingest import upsert_announcement
from app.services.cire_announcement_parser import parse_announcement_text, pdf_to_text
from app.services.court_filings import carregar_distintivos, match_filings, upsert_filing
from app.services.coverage import mark_checked
from app.services.ingest import upsert_process
from app.services.risk import refresh_all_risk
from app.services.insolvency_alerts import dispatch_pending
from app.services.insolvency_matching import match_announcement
from app.services.mj_irn_sync import sync_company as sync_mj_company
from app.services.registry_graph import (
    apply_identity_merges,
    parser_health,
    refresh_degrees,
    suggest_identity_links,
)
from app.services.registry_query import refresh_incident_summary
from app.services.registry_sweep import (
    enqueue_days,
    process_content,
    refresh_acts_count,
    refresh_priorities,
    sweep_days,
    sweep_progress,
)
from app.services.registry_sync import (
    enqueue_corporate_neighbours,
    process_queue,
    seed_all_companies,
    seed_monitored,
)
# app.services.mj_publications stays DORMANT. Google's anti-bot gate blocks
# headless Chromium (even with patchright stealth) from datacenter IPs before
# the audio captcha can be served to Whisper. Revive requires either a
# residential proxy (~€3-10/mês) or running the scraper from a real user
# browser — both out of scope for the €0 deployment.

logger = logging.getLogger(__name__)

scheduler = AsyncIOScheduler(timezone="Europe/Lisbon")


async def _log_start(session: AsyncSession, source: str, params: dict[str, Any]) -> str:
    row = (
        await session.execute(
            text(
                """
                INSERT INTO scraping_logs (source, started_at, status, params)
                VALUES (:s, now(), 'running', CAST(:p AS JSONB))
                RETURNING id::text
                """
            ),
            {"s": source, "p": json.dumps(params, default=str)},
        )
    ).first()
    await session.commit()
    return row[0]


async def _log_finish(
    session: AsyncSession,
    log_id: str,
    *,
    status: str,
    rows_seen: int = 0,
    rows_new: int = 0,
    error: str | None = None,
) -> None:
    await session.execute(
        text(
            """
            UPDATE scraping_logs
            SET finished_at = now(), status = :st, rows_seen = :rs,
                rows_new = :rn, error = :err
            WHERE id = :id
            """
        ),
        {
            "id": log_id,
            "st": status,
            "rs": rows_seen,
            "rn": rows_new,
            "err": (error or "")[:2000] or None,
        },
    )
    await session.commit()


async def _monitored_companies(
    session: AsyncSession,
    # Clientes entram nos scrapers judiciais: um cliente com execuções ou penhoras
    # é exactamente o sinal de risco de crédito que interessa, e sem isto a ficha
    # dele dizia "Processos (0)" para todos os 87. Na distribuição é praticamente
    # grátis — é um tribunal-sweep, só se acrescentam nomes ao cruzamento; no CIRE
    # são 87 consultas por NIF, uma vez, e depois só o que for novo.
    types: tuple[str, ...] = ("internal", "competitor", "client"),
    company_id: str | None = None,
    isolated: bool | None = False,
) -> list[dict[str, Any]]:
    """Companies eligible for scheduled scraping. Analysis-type companies are
    one-off and excluded by default. Internal first so partial failures still
    cover priority work.

    `company_id`: when set, return just that single company (used for on-create
    backfill). Still respects `monitored = TRUE` but ignores `isolated`.
    `isolated`: False excludes companies flagged `scrape_isolated` (the default
    for the main batch), True returns ONLY isolated, None means no filter."""
    if company_id is not None:
        # Explicit by-id path ignores `active`: on-create backfill must run
        # even for companies imported as inactive (one-shot historical scrape).
        rows = (
            await session.execute(
                text(
                    """
                    SELECT id::text, nif, legal_name, monitoring_type
                    FROM companies
                    WHERE id = :cid AND monitored = TRUE
                    """
                ),
                {"cid": company_id},
            )
        ).all()
    else:
        iso_clause = ""
        params_sql: dict[str, Any] = {"types": list(types)}
        if isolated is True:
            iso_clause = "AND scrape_isolated = TRUE"
        elif isolated is False:
            iso_clause = "AND scrape_isolated = FALSE"
        # isolated is None -> no filter (admin ad-hoc runs)
        rows = (
            await session.execute(
                text(
                    f"""
                    SELECT id::text, nif, legal_name, monitoring_type
                    FROM companies
                    WHERE monitored = TRUE
                      AND active = TRUE
                      AND monitoring_type = ANY(:types)
                      {iso_clause}
                    ORDER BY (monitoring_type = 'internal') DESC, legal_name
                    """
                ),
                params_sql,
            )
        ).all()
    return [
        {"id": r[0], "nif": r[1], "legal_name": r[2], "monitoring_type": r[3]} for r in rows
    ]


async def _last_successful_distribuicao_date(
    session: AsyncSession, isolated: bool = False
) -> date | None:
    """Latest date_to covered by a successful (or partial) distribuicao run.
    Isolated and non-isolated runs have separate catch-up histories so they
    don't clobber each other's progress."""
    row = (
        await session.execute(
            text(
                """
                SELECT (params->>'date_to')::date AS d
                FROM scraping_logs
                WHERE source = 'distribuicao'
                  AND status IN ('ok','partial')
                  AND params->>'date_to' IS NOT NULL
                  AND COALESCE(params->>'isolated','false') = :iso
                ORDER BY (params->>'date_to')::date DESC
                LIMIT 1
                """
            ),
            {"iso": "true" if isolated else "false"},
        )
    ).first()
    return row[0] if row else None


DISTRIBUICAO_MAX_LOOKBACK_DAYS = 180  # CITIUS only returns ~6 months server-side


async def run_distribuicao(
    date_from: date | None = None,
    date_to: date | None = None,
    types: tuple[str, ...] = ("internal", "competitor", "client"),
    company_id: str | None = None,
    isolated: bool | None = False,
) -> str:
    """Scrape Distribuição for a date range.

    Two modes:
      * **Sweep** (default) — one HTTP pass per tribunal without a party
        filter; each row matched locally against every monitored company's
        distintivo. ~50× fewer requests than the old per-company mode.
        Progress is tribunal-scoped.
      * **Per-company** (when `company_id` is passed) — legacy path, used
        for on-create backfill of a single new company (faster than warming
        up a full sweep just for one NIF).

    When both `date_from` and `date_to` are None (scheduled run), looks up
    the last successful scrape's `date_to` and scrapes from (last + 1)
    through yesterday. Capped at 30 days back. If no history: yesterday only.
    Manual windows hard-capped at 180 days (CITIUS serves only ~6 months).

    Commits per-tribunal (sweep) or per-company (legacy) — survives mid-run
    crashes either way.
    """
    today = date.today()
    yesterday = today - timedelta(days=1)
    auto_mode = date_from is None and date_to is None

    if auto_mode and company_id:
        date_from = yesterday - timedelta(days=DISTRIBUICAO_MAX_LOOKBACK_DAYS)
        date_to = yesterday
    elif auto_mode:
        async with AsyncSessionLocal() as session:
            last = await _last_successful_distribuicao_date(session, isolated=False)
        if last and last < yesterday:
            date_from = max(last + timedelta(days=1), yesterday - timedelta(days=30))
            date_to = yesterday
        else:
            date_from = date_to = yesterday
    elif date_from is None:
        date_from = date_to
    elif date_to is None:
        date_to = date_from
    assert date_from is not None and date_to is not None
    if date_from > date_to:
        date_from, date_to = date_to, date_from

    capped_from = date_to - timedelta(days=DISTRIBUICAO_MAX_LOOKBACK_DAYS)
    if date_from < capped_from:
        logger.info(
            "distribuicao: capping date_from %s -> %s (CITIUS 180d limit)",
            date_from, capped_from,
        )
        date_from = capped_from

    async with AsyncSessionLocal() as session:
        params_log: dict[str, Any] = {
            "date_from": date_from.isoformat(),
            "date_to": date_to.isoformat(),
            "types": list(types),
            "auto": auto_mode,
            "isolated": False,
            "mode": "per_company" if company_id else "sweep",
        }
        if company_id:
            params_log["company_id"] = company_id
            params_log["backfill_new"] = True
        log_id = await _log_start(session, "distribuicao", params_log)
        companies = await _monitored_companies(
            session, types=types, company_id=company_id, isolated=None,
        )

    try:
        if company_id:
            if not companies:
                async with AsyncSessionLocal() as session:
                    await _log_finish(session, log_id, status="ok", rows_seen=0, rows_new=0)
                return log_id
            return await _run_distribuicao_per_company(
                log_id, companies, date_from, date_to,
            )
        # O varrimento nacional não depende de haver empresas monitorizadas: o que
        # ele guarda é o arquivo, e o cruzamento vem depois. Ter isto à entrada
        # significava perder o dia inteiro por uma razão que não era sobre o dia.
        return await _run_distribuicao_sweep(log_id, date_from, date_to)
    except Exception as e:
        logger.exception("run_distribuicao failed: %r", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=repr(e))
        raise


async def _run_distribuicao_sweep(
    log_id: str,
    date_from: date,
    date_to: date,
) -> str:
    """Varrimento: uma passagem por tribunal, tudo guardado, cruzado a seguir.

    O `companies` já não é preciso aqui — o cruzamento deixou de viver dentro do
    scraper e passou para o `services/court_filings.py`, para poder voltar a
    correr sobre o arquivo. Uma empresa sem distintivo utilizável já não impede
    a recolha: antes, se nenhuma desse um distintivo seguro, a corrida terminava
    sem ir buscar nada, e a informação do dia perdia-se.

    **Um dia de cada vez, mesmo quando a janela é maior.** O apanhar-atrasos pode
    pedir até 30 dias de uma vez; numa consulta só, o Porto sozinho passava dos
    mil processos e aproximava-se do tecto de páginas — que trunca e regista `ok`
    à mesma. É a armadilha que o varrimento dos anúncios já pagou uma vez, e não
    se paga duas.
    """
    tribunals = await discover_tribunal_codes()
    dias = [
        date_from + timedelta(days=i) for i in range((date_to - date_from).days + 1)
    ]

    async with AsyncSessionLocal() as session:
        await session.execute(
            text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
            {"t": len(tribunals) * len(dias), "id": log_id},
        )
        await session.commit()

    guardadas = acertos = falhas = 0
    feitos_antes = 0

    async def progresso(feitos: int, g: int, a: int) -> None:
        async with AsyncSessionLocal() as session:
            await session.execute(
                text(
                    "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn,"
                    " progress_current = :pc WHERE id = :id"
                ),
                {
                    "rs": guardadas + g,
                    "rn": acertos + a,
                    "pc": feitos_antes + feitos,
                    "id": log_id,
                },
            )
            await session.commit()

    for dia in dias:
        g, a, f = await _varrer_e_guardar(tribunals, dia, dia, on_progress=progresso)
        guardadas += g
        acertos += a
        falhas += f
        feitos_antes += len(tribunals)

    status = "partial" if falhas else "ok"
    err_msg = f"{falhas} tribunal(s) failed" if falhas else None
    async with AsyncSessionLocal() as session:
        await _log_finish(
            session, log_id, status=status, rows_seen=guardadas,
            rows_new=acertos, error=err_msg,
        )
    return log_id


async def _varrer_e_guardar(
    tribunals: list[dict[str, str]],
    date_from: date,
    date_to: date,
    on_progress: Any = None,
) -> tuple[int, int, int]:
    """Varre os tribunais dados, guarda tudo e cruza. (guardadas, acertos, falhas).

    Partilhado entre a corrida diária e o backfill, para não haver duas versões
    da mesma sequência — guardar primeiro, cruzar depois — a divergir em silêncio.
    """
    from app.config import settings as _settings
    sem = asyncio.Semaphore(_settings.SCRAPE_CONCURRENCY)

    # Uma vez por corrida, não uma por tribunal: construí-lo em cada lote enchia
    # o log de avisos repetidos sobre os mesmos nomes ambíguos.
    async with AsyncSessionLocal() as session:
        distintivos = await carregar_distintivos(session)

    guardadas = 0
    acertos = 0
    falhas = 0
    feitos = 0

    async def worker(t: dict[str, str]) -> list[dict[str, Any]]:
        nonlocal falhas
        async with sem:
            try:
                return await scrape_distribuicao_tribunal_sweep(
                    t["value"], t["label"], date_from, date_to,
                )
            except Exception as e:
                falhas += 1
                logger.warning("distribuicao sweep tribunal=%s failed: %r", t["label"], e)
                return []

    # Cada tribunal é gravado assim que chega — o commit por tribunal já existia
    # e é o que faz uma falha a meio custar só os que estavam em voo.
    tasks = [asyncio.create_task(worker(t)) for t in tribunals]
    for coro in asyncio.as_completed(tasks):
        linhas = await coro
        feitos += 1
        if linhas:
            async with AsyncSessionLocal() as session:
                ids = []
                for row in linhas:
                    filing_id, _novo = await upsert_filing(session, row)
                    ids.append(filing_id)
                await session.commit()
                # O cruzamento vem a seguir e sobre o que já ficou guardado: se
                # falhar, a informação não se perde — repete-se depois.
                resultado = await match_filings(session, ids, distintivos)
                await session.commit()
            guardadas += len(linhas)
            acertos += resultado["hits_nif"] + resultado["hits_nome"]
        if on_progress:
            await on_progress(feitos, guardadas, acertos)
    return guardadas, acertos, falhas


async def run_distribuicao_backfill(
    max_days: int = 400,
    dias_vazios_para_parar: int = 3,
    comecar_em: date | None = None,
) -> str:
    """Enche o arquivo para trás, um dia de cada vez, até a fonte secar.

    **Um dia de cada vez, nunca uma janela larga.** O paginador do CITIUS pára
    aos 100 ecrãs e a consulta devolve `ok` na mesma; foi assim que o varrimento
    dos anúncios deu janelas por cobertas sem o estarem. Um dia inteiro do país
    são ~1.000 processos, longe de qualquer tecto — e o scraper compara agora o
    que guardou com o "N Processos encontrados" da fonte, portanto uma truncagem
    passa a gritar em vez de passar despercebida.

    **Pára quando a fonte deixar de servir**, não numa contagem de dias
    codificada: são três dias úteis seguidos a zero. Medido a 03/08/2026, o
    limite caía a ~04/02/2026, mas isso desliza todos os dias e não é nosso.
    Fins-de-semana nem se consultam — não há distribuição e só gastavam pedidos.

    O `comecar_em` serve para retomar: um deploy mata a corrida e sem isto
    repetiam-se horas de pedidos já feitos, só para chegar ao mesmo sítio. A
    gravação é idempotente de qualquer forma, portanto repetir não estraga —
    apenas custa.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "distribuicao_backfill", {"max_days": max_days},
        )
        await session.execute(
            text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
            {"t": max_days, "id": log_id},
        )
        await session.commit()

    try:
        tribunals = await discover_tribunal_codes()
        dia = comecar_em or (date.today() - timedelta(days=1))
        vazios = 0
        total_guardadas = 0
        total_acertos = 0
        dias_feitos = 0

        while dias_feitos < max_days and vazios < dias_vazios_para_parar:
            if dia.weekday() >= 5:
                dia -= timedelta(days=1)
                continue
            guardadas, acertos, falhas = await _varrer_e_guardar(tribunals, dia, dia)
            total_guardadas += guardadas
            total_acertos += acertos
            dias_feitos += 1
            vazios = vazios + 1 if guardadas == 0 else 0
            logger.info(
                "backfill distribuicao %s: %d linhas, %d acertos%s",
                dia, guardadas, acertos, f", {falhas} tribunais falharam" if falhas else "",
            )
            async with AsyncSessionLocal() as session:
                await session.execute(
                    text(
                        "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn,"
                        " progress_current = :pc WHERE id = :id"
                    ),
                    {"rs": total_guardadas, "rn": total_acertos,
                     "pc": dias_feitos, "id": log_id},
                )
                await session.commit()
            dia -= timedelta(days=1)

        alcance = f"até {dia + timedelta(days=1)}"
        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status="ok",
                rows_seen=total_guardadas, rows_new=total_acertos,
                error=f"{dias_feitos} dias, {alcance}",
            )
        logger.info(
            "backfill distribuicao terminado: %d dias, %d linhas, %s",
            dias_feitos, total_guardadas, alcance,
        )
        return log_id
    except Exception as e:
        logger.exception("run_distribuicao_backfill failed: %r", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=repr(e))
        raise


async def _run_distribuicao_per_company(
    log_id: str,
    companies: list[dict[str, Any]],
    date_from: date,
    date_to: date,
) -> str:
    """Legacy per-company mode, retained for on-create backfill (single
    empresa → no point in warming up a full tribunal sweep).

    Kept aligned with the archived v1 code under scrapers/_archive/ — so a
    future revive of the old mode for bulk runs is a one-commit restore."""
    tribunals = await discover_tribunal_codes()
    total_seen = 0
    total_new = 0
    error_count = 0

    async with AsyncSessionLocal() as session:
        await session.execute(
            text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
            {"t": len(companies), "id": log_id},
        )
        await session.commit()

    from app.config import settings as _settings
    sem = asyncio.Semaphore(_settings.SCRAPE_CONCURRENCY)

    async def worker(
        company: dict[str, Any], t: dict[str, str], search: str
    ) -> list[dict[str, Any]]:
        nonlocal error_count
        async with sem:
            try:
                return await scrape_distribuicao_range(
                    t["value"], t["label"], date_from, date_to,
                    [company], party_filter=search,
                )
            except Exception as e:
                error_count += 1
                logger.warning(
                    "distribuicao company=%s tribunal=%s failed: %r",
                    company["legal_name"], t["label"], e,
                )
                return []

    for idx, company in enumerate(companies, 1):
        search = short_search_term(company["legal_name"])
        if not search:
            logger.info(
                "distribuicao [%d/%d] %s: SKIPPED (ambiguous/short name)",
                idx, len(companies), company["legal_name"][:40],
            )
            continue
        results = await asyncio.gather(
            *(worker(company, t, search) for t in tribunals),
            return_exceptions=False,
        )
        matches: list[dict[str, Any]] = [row for r in results for row in r]
        if matches:
            company_new = 0
            async with AsyncSessionLocal() as session:
                for row in matches:
                    _, is_new = await upsert_process(session, row)
                    if is_new:
                        company_new += 1
                await session.commit()
            total_seen += len(matches)
            total_new += company_new
        async with AsyncSessionLocal() as session:
            await session.execute(
                text(
                    "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn, progress_current = :pc WHERE id = :id"
                ),
                {"rs": total_seen, "rn": total_new, "pc": idx, "id": log_id},
            )
            await session.commit()

    status = "partial" if error_count else "ok"
    err_msg = f"{error_count} (company,tribunal) pair(s) failed" if error_count else None
    async with AsyncSessionLocal() as session:
        await _log_finish(
            session, log_id, status=status, rows_seen=total_seen,
            rows_new=total_new, error=err_msg,
        )
    return log_id


async def run_cire(
    days: str = "todos",
    types: tuple[str, ...] = ("internal", "competitor", "client"),
    company_id: str | None = None,
    isolated: bool | None = False,
) -> str:
    """Scrape CIRE for every monitored company. `days` is '15', '30', or 'todos'.
    `types` filters which monitoring_type categories to include.
    `company_id`: restrict to a single company (used for on-create backfill).
    `isolated`: see `run_distribuicao`."""
    async with AsyncSessionLocal() as session:
        params_log: dict[str, Any] = {
            "days": days,
            "types": list(types),
            "isolated": bool(isolated) if isolated is not None else False,
        }
        if company_id:
            params_log["company_id"] = company_id
            params_log["backfill_new"] = True
        log_id = await _log_start(session, "cire", params_log)
        companies = await _monitored_companies(
            session, types=types, company_id=company_id, isolated=isolated,
        )
    try:
        async with AsyncSessionLocal() as session:
            await session.execute(
                text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
                {"t": len(companies), "id": log_id},
            )
            await session.commit()
        total_seen = 0
        total_new = 0
        status = "ok"
        for idx, c in enumerate(companies, 1):
            try:
                matches = await scrape_cire_for_nif(c["nif"], c["id"], days=days)
                total_seen += len(matches)
                if matches:
                    async with AsyncSessionLocal() as session:
                        for row in matches:
                            _, is_new = await upsert_process(session, row)
                            if is_new:
                                total_new += 1
                        await session.commit()
            except Exception as e:
                logger.exception("cire nif=%s failed: %s", c["nif"], e)
                status = "partial"
            async with AsyncSessionLocal() as session:
                await session.execute(
                    text(
                        "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn, progress_current = :pc WHERE id = :id"
                    ),
                    {"rs": total_seen, "rn": total_new, "pc": idx, "id": log_id},
                )
                await session.commit()
        async with AsyncSessionLocal() as session:
            # Passámos por todas, com ou sem processos. É o que distingue "não
            # tem nada" de "nunca foi procurada" na ficha.
            await mark_checked(session, [c["id"] for c in companies])
            await refresh_all_risk(session, [c["id"] for c in companies])
            await _log_finish(
                session, log_id, status=status, rows_seen=total_seen, rows_new=total_new
            )
            await session.commit()
        return log_id
    except Exception as e:
        logger.exception("run_cire failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


CIRE_ANN_MAX_WINDOW_DAYS = 90  # tecto por corrida manual


async def _last_successful_announcement_date(session: AsyncSession) -> date | None:
    row = (
        await session.execute(
            text(
                """
                SELECT (params->>'date_to')::date AS d
                FROM scraping_logs
                WHERE source = 'cire_announcements'
                  AND status IN ('ok','partial')
                  AND params->>'date_to' IS NOT NULL
                ORDER BY (params->>'date_to')::date DESC
                LIMIT 1
                """
            )
        )
    ).first()
    return row[0] if row else None


async def run_cire_announcements(
    date_from: date | None = None,
    date_to: date | None = None,
) -> str:
    """Sweep nacional dos anúncios CIRE do CITIUS.

    Sem filtro de entidade: varre tudo o que foi publicado no intervalo e guarda,
    dê ou não match hoje. É o que permite detectar a insolvência de um devedor
    que ninguém tinha adicionado, e o que dá histórico ao cruzamento retroactivo
    quando entram clientes novos na watchlist.

    Sem argumentos (corrida agendada) apanha desde o dia a seguir ao último
    `date_to` bem-sucedido até ontem; sem histórico, só ontem. Janelas manuais
    são limitadas a 90 dias por corrida.

    **Varre dia a dia**, não a janela toda de uma vez. Uma janela de 90 dias tem
    ~18 mil anúncios e batia no tecto de páginas do sweeper, que truncava em
    silêncio e mesmo assim registava 'ok' — pior do que falhar, porque o
    catch-up dava a janela por coberta e os dias em falta nunca mais eram
    revisitados. Um dia são ~11 páginas, longe de qualquer tecto.
    """
    today = date.today()
    if date_from is None and date_to is None:
        last = await _run_in_session(_last_successful_announcement_date)
        date_to = today - timedelta(days=1)
        date_from = (last + timedelta(days=1)) if last else date_to
        if date_from > date_to:
            date_from = date_to
    date_from = date_from or (today - timedelta(days=1))
    date_to = date_to or (today - timedelta(days=1))
    if (date_to - date_from).days > CIRE_ANN_MAX_WINDOW_DAYS:
        date_from = date_to - timedelta(days=CIRE_ANN_MAX_WINDOW_DAYS)

    dias = [date_from + timedelta(days=i) for i in range((date_to - date_from).days + 1)]

    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session,
            "cire_announcements",
            {"date_from": str(date_from), "date_to": str(date_to), "dias": len(dias)},
        )
        await session.execute(
            text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
            {"t": len(dias), "id": log_id},
        )
        await session.commit()

    total_seen = 0
    total_new = 0
    total_hits = 0
    status = "ok"
    try:
        async with CireAnnouncementSweeper() as sweeper:
            for dia_idx, dia in enumerate(dias, 1):
                paginas = 0
                try:
                    async for page in sweeper.pages(dia, dia):
                        paginas += 1
                        for row in page["rows"]:
                            # PDF só para os actos que fixam prazo de reclamação:
                            # ~20/dia dos ~110, e os únicos com prazo a extrair.
                            if not row.get("wants_pdf"):
                                continue
                            try:
                                pdf = await sweeper.fetch_pdf_bytes(row.get("pdf_token"))
                                text_pdf = pdf_to_text(pdf)
                                row["pdf_text"] = text_pdf
                                row["detail"] = parse_announcement_text(text_pdf)
                            except Exception as e:  # noqa: BLE001
                                logger.warning(
                                    "anúncio %s: PDF falhou: %s", row.get("referencia"), e
                                )

                        async with AsyncSessionLocal() as session:
                            for row in page["rows"]:
                                try:
                                    ann_id, is_new = await upsert_announcement(session, row)
                                except Exception as e:  # noqa: BLE001
                                    logger.exception("upsert de anúncio falhou: %s", e)
                                    status = "partial"
                                    continue
                                total_seen += 1
                                if is_new:
                                    total_new += 1
                                    total_hits += await match_announcement(session, ann_id)
                            await session.commit()
                except Exception as e:  # noqa: BLE001
                    # Um dia que falhe não pode levar os restantes atrás.
                    logger.exception("cire_announcements %s falhou: %s", dia, e)
                    status = "partial"

                # Bater no tecto num único dia significaria >4000 anúncios nesse
                # dia. Improvável, mas se acontecer tem de ficar 'partial'.
                if paginas >= PAGE_CAP:
                    logger.error("cire_announcements %s: tecto de páginas atingido", dia)
                    status = "partial"

                async with AsyncSessionLocal() as session:
                    await session.execute(
                        text(
                            "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn,"
                            " progress_current = :pc WHERE id = :id"
                        ),
                        {"rs": total_seen, "rn": total_new, "pc": dia_idx, "id": log_id},
                    )
                    await session.commit()

        logger.info(
            "cire_announcements %s..%s (%d dias) -> vistos=%d novos=%d hits=%d [%s]",
            date_from, date_to, len(dias), total_seen, total_new, total_hits, status,
        )
        # Envio depois de tudo persistido: um SMTP em baixo não pode fazer
        # perder anúncios, e os hits ficam pendentes para a corrida seguinte.
        if total_hits:
            try:
                async with AsyncSessionLocal() as session:
                    await dispatch_pending(session)
            except Exception as e:  # noqa: BLE001
                logger.exception("envio de avisos de insolvência falhou: %s", e)

        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status=status, rows_seen=total_seen, rows_new=total_new
            )
        return log_id
    except Exception as e:
        logger.exception("run_cire_announcements failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def _run_in_session(fn):
    async with AsyncSessionLocal() as session:
        return await fn(session)


async def run_mj_irn(
    company_id: str | None = None,
    types: tuple[str, ...] = ("internal", "competitor", "client"),
    only_stale_days: int | None = 7,
) -> str:
    """Publicações do registo comercial, por NIF, via portal IRN.

    Substitui o `mj_publications.py` dormente e a extensão Chrome que obrigava
    alguém a navegar o portal à mão. Ver `app/scrapers/mj_irn.py` para o
    protocolo.

    `only_stale_days`: salta empresas sincronizadas há menos de N dias, para a
    corrida diária ser barata depois do backfill inicial. None força todas.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "mj_irn",
            {"types": list(types), "company_id": company_id, "stale_days": only_stale_days},
        )
        clause = "AND c.id = :cid" if company_id else ""
        stale = (
            "AND (c.registry_fetched_at IS NULL OR c.registry_fetched_at < now() - "
            f"make_interval(days => {int(only_stale_days)}))"
            if only_stale_days is not None and not company_id
            else ""
        )
        companies = (
            await session.execute(
                text(
                    f"""
                    SELECT c.id::text, c.nif, c.legal_name
                    FROM companies c
                    WHERE c.monitored = TRUE AND c.monitoring_type = ANY(:types)
                      {clause} {stale}
                    ORDER BY c.registry_fetched_at NULLS FIRST, c.legal_name
                    """
                ),
                {"types": list(types), **({"cid": company_id} if company_id else {})},
            )
        ).all()

    total_seen = total_new = total_people = 0
    status = "ok"
    try:
        async with AsyncSessionLocal() as session:
            await session.execute(
                text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
                {"t": len(companies), "id": log_id},
            )
            await session.commit()

        async with MjIrnClient() as client:
            for idx, (cid, nif, name) in enumerate(companies, 1):
                try:
                    async with AsyncSessionLocal() as session:
                        s = await sync_mj_company(session, client, cid, nif)
                        await session.commit()
                    total_seen += s["seen"]
                    total_new += s["new"]
                    total_people += s["people"]
                except IrnVersionChanged as e:
                    # O IRN fez deploy: os apiVersion codificados no scraper
                    # ficaram obsoletos e nada mais vai funcionar nesta corrida.
                    # Falhar alto é melhor do que gravar meia sincronização.
                    logger.error("mj_irn: %s", e)
                    async with AsyncSessionLocal() as session:
                        await _log_finish(session, log_id, status="error", error=str(e))
                    raise
                except Exception as e:  # noqa: BLE001
                    logger.exception("mj_irn nif=%s (%s) falhou: %s", nif, name, e)
                    status = "partial"
                async with AsyncSessionLocal() as session:
                    await session.execute(
                        text(
                            "UPDATE scraping_logs SET rows_seen = :rs, rows_new = :rn,"
                            " progress_current = :pc WHERE id = :id"
                        ),
                        {"rs": total_seen, "rn": total_new, "pc": idx, "id": log_id},
                    )
                    await session.commit()
                await random_delay()

        logger.info(
            "mj_irn: %d empresas, %d publicações vistas, %d novas, %d pessoas",
            len(companies), total_seen, total_new, total_people,
        )
        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status=status, rows_seen=total_seen, rows_new=total_new
            )
        return log_id
    except IrnVersionChanged:
        raise
    except Exception as e:
        logger.exception("run_mj_irn failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_registry_nacional(
    max_items: int = 200000, max_seconds: float = 40000.0, seed: bool = False
) -> str:
    """Fase A: listar o histórico por NIPC de todas as entidades colectivas.

    **Só lista.** Uma chamada por empresa, sem abrir corpos — medido, 102
    empresas/min, portanto as ~507 mil que não têm histórico próprio são cerca de
    83 horas. Os actos que aparecerem ficam na fila de conteúdo e são abertos
    depois, pela escada de prioridades que já põe a cap table mais recente de
    cada entidade à frente do resto do histórico.

    Porquê: a listagem diária do IRN **não é um índice completo**. Numa amostra
    de 30 empresas, 28 tinham actos que nunca tínhamos visto, e só temos ~1.200
    prestações de contas por ano num país onde centenas de milhares as entregam.

    A fila vive na base de dados: isto pára ao fim do orçamento e retoma na noite
    seguinte onde ficou, sem repetir nada.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "registry_nacional",
            {"max_items": max_items, "max_seconds": max_seconds, "seed": seed},
        )
        if seed:
            postos = await seed_all_companies(session)
            await session.commit()
            logger.info("registry_nacional: %d entidades semeadas", postos)
        pendentes = (
            await session.execute(
                text(
                    "SELECT count(*) FROM registry_entity_fetch"
                    " WHERE status IN ('pending','error') AND attempts < 3"
                )
            )
        ).scalar() or 0
        await session.execute(
            text("UPDATE scraping_logs SET progress_total = :t WHERE id = :id"),
            {"t": pendentes, "id": log_id},
        )
        await session.commit()

    try:
        totals = await process_queue(
            max_items=max_items, max_seconds=max_seconds, fetch_content=False
        )
        async with AsyncSessionLocal() as session:
            restam = (
                await session.execute(
                    text(
                        "SELECT count(*) FROM registry_entity_fetch"
                        " WHERE status IN ('pending','error') AND attempts < 3"
                    )
                )
            ).scalar() or 0
            await _log_finish(
                session, log_id,
                status="partial" if restam else "ok",
                rows_seen=totals["entidades"], rows_new=totals["novos"],
                error=f"{restam} entidades por listar" if restam else None,
            )
        logger.info("registry_nacional: %s | %d por listar", totals, restam)
        return log_id
    except IrnVersionChanged as e:
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise
    except Exception as e:
        logger.exception("run_registry_nacional failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_registry_expand(
    max_items: int = 60, max_seconds: float = 1800.0, seed: bool = True
) -> str:
    """Registo comercial por NIPC: as monitorizadas e, em BFS, os sócios-empresa.

    É a via que dá **histórico completo** de uma entidade, porque a pesquisa por
    NIPC do IRN não tem limite de datas. Complementa o varrimento nacional, que
    cobre em largura mas só a janela recolhida.

    Cobre os cinco tipos de monitorização, ao contrário do `run_mj_irn`: as 7 de
    análise estavam sem qualquer informação de registo e as 17 `related` nunca
    eram sincronizadas, o que travava o grafo no primeiro hop.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "registry_expand", {"max_items": max_items, "max_seconds": max_seconds},
        )
        if seed:
            seeded = await seed_monitored(session)
            grown = await enqueue_corporate_neighbours(session)
            await session.commit()
            logger.info("registry_expand: %d monitorizadas, %d vizinhas na fila", seeded, grown)

    try:
        totals = await process_queue(max_items=max_items, max_seconds=max_seconds)
        async with AsyncSessionLocal() as session:
            # A vizinhança cresce à medida que se descobrem sócios-empresa, por
            # isso volta-se a alargar a fila no fim de cada corrida.
            await enqueue_corporate_neighbours(session)
            await refresh_degrees(session)
            await suggest_identity_links(session)
            pending = (
                await session.execute(
                    text(
                        "SELECT count(*) FROM registry_entity_fetch"
                        " WHERE status IN ('pending','error') AND attempts < 3"
                    )
                )
            ).scalar() or 0
            await session.commit()

        # `partial` quando ficou fila por fazer: o orçamento acabou antes do
        # trabalho. Dar isto como `ok` faria parecer coberto o que não está.
        status = "partial" if pending else "ok"
        health = None
        async with AsyncSessionLocal() as session:
            health = await parser_health(session)
            await _log_finish(
                session, log_id, status=status,
                rows_seen=totals["listados"], rows_new=totals["novos"],
                error=(
                    f"{pending} entidades ainda na fila" if pending else None
                ),
            )
        if health and health.get("alarme"):
            logger.error("registo: extracção suspeita — %s", health)
        logger.info("registry_expand: %s (fila: %d)", totals, pending)
        return log_id
    except IrnVersionChanged as e:
        logger.error("registry_expand: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise
    except Exception as e:
        logger.exception("run_registry_expand failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_registry_sweep(
    date_from: date | None = None,
    date_to: date | None = None,
    max_days: int = 5,
    max_seconds: float = 1800.0,
) -> str:
    """Varrimento nacional dos actos de registo, dia a dia.

    É a via que responde a "em que outras empresas é que esta pessoa é sócia": o
    IRN só pesquisa por NIPC de entidade, e com o NIF de um indivíduo devolve
    zero. Sem índice nacional, essa pergunta não tem resposta.

    Um dia de cada vez porque é o único filtro fiável — intervalos explícitos
    devolvem contagens absurdas e sem filtro devolve zero.
    """
    async with AsyncSessionLocal() as session:
        if date_from and date_to:
            queued = await enqueue_days(session, date_from, date_to)
        else:
            # rotina diária: ontem, e o que tenha ficado para trás
            ate = date.today() - timedelta(days=1)
            queued = await enqueue_days(session, ate - timedelta(days=7), ate)
        await session.commit()
        log_id = await _log_start(
            session, "registry_sweep",
            {"max_days": max_days, "dias_novos_na_fila": queued,
             "date_from": str(date_from) if date_from else None,
             "date_to": str(date_to) if date_to else None},
        )
    try:
        totals = await sweep_days(max_days=max_days, max_seconds=max_seconds)
        async with AsyncSessionLocal() as session:
            await refresh_priorities(session)
            # O contador de actos por entidade acompanha a listagem, não o parse:
            # é aqui que os actos entram, é aqui que o número tem de ser acertado.
            await refresh_acts_count(session)
            progress = await sweep_progress(session)
            await session.commit()
        # `partial` enquanto houver dias por varrer: dar isto como `ok` faria
        # parecer coberto um arquivo que ainda vai a meio.
        status = "partial" if progress.get("dias_por_fazer") else "ok"
        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status=status,
                rows_seen=totals["listados"], rows_new=totals["novos"],
                error=(
                    f"{progress['dias_por_fazer']} dias por varrer"
                    if progress.get("dias_por_fazer") else None
                ),
            )
        logger.info("registry_sweep: %s | %s", totals, progress)
        return log_id
    except IrnVersionChanged as e:
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise
    except Exception as e:
        logger.exception("run_registry_sweep failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_registry_content(max_items: int = 2000, max_seconds: float = 3600.0) -> str:
    """Abre e parseia o conteúdo dos actos na fila, por prioridade.

    A listagem é barata (uma chamada por página) e o conteúdo são três chamadas
    por acto, portanto só se abre o que traz pessoas ou quotas. A escada de
    prioridades faz o útil aterrar primeiro: o nosso universo, a vizinhança do
    grafo, e a cap table mais recente de cada entidade antes do resto do
    histórico.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "registry_content", {"max_items": max_items, "max_seconds": max_seconds},
        )
    try:
        totals = await process_content(max_items=max_items, max_seconds=max_seconds)
        async with AsyncSessionLocal() as session:
            await refresh_degrees(session)
            await suggest_identity_links(session)
            # As sugestões de nome exacto na mesma empresa entram já
            # confirmadas; é aqui que passam a valer no grafo.
            await apply_identity_merges(session)
            progress = await sweep_progress(session)
            health = await parser_health(session)
            await session.commit()
        status = "partial" if progress.get("conteudo_por_abrir") else "ok"
        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status=status,
                rows_seen=totals["actos"], rows_new=totals["quotas"] + totals["cargos"],
                error=(
                    f"{progress['conteudo_por_abrir']} actos por abrir"
                    if progress.get("conteudo_por_abrir") else None
                ),
            )
        if health.get("alarme"):
            logger.error("registo: extracção suspeita — %s", health)
        logger.info("registry_content: %s | %s", totals, progress)
        return log_id
    except IrnVersionChanged as e:
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise
    except Exception as e:
        logger.exception("run_registry_content failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_irn_canary() -> str:
    """Um handshake e uma listagem, de hora a hora.

    Quando o IRN faz deploy, os `apiVersion` codificados no scraper ficam
    obsoletos e tudo passa a devolver 403. Sem isto sabia-se na manhã seguinte,
    com três jobs falhados; assim sabe-se dentro de uma hora.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, "irn_canary", {})
    try:
        async with MjIrnClient() as client:
            rows = await client.publications_for_day(date.today() - timedelta(days=1), max_pages=1)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="ok", rows_seen=len(rows))
        return log_id
    except Exception as e:
        logger.error("irn_canary: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_entity_incidents() -> str:
    """Incidentes por NIF, para cada nó do grafo poder ser pintado.

    O cruzamento é com o `cire_announcement_parties`, que é **nacional** — 240 296
    partes e 55 321 NIFs distintos — portanto cobre pessoas e entidades que não
    monitorizamos. É isso que permite ver, a dois hops de distância, que um sócio
    de um sócio tem uma insolvência a correr.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, "entity_incidents", {})
    try:
        async with AsyncSessionLocal() as session:
            rows = await refresh_incident_summary(session)
            await session.commit()
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="ok", rows_seen=rows, rows_new=rows)
        logger.info("entity_incidents: %d NIFs", rows)
        return log_id
    except Exception as e:
        logger.exception("run_entity_incidents failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_insolvency_digest() -> str:
    """Rede de segurança diária do envio de avisos.

    O sweep já dispara os avisos dos hits que cria. Isto apanha o que tenha
    ficado pendente por o SMTP estar em baixo, ou hits criados fora do sweep
    (por exemplo, quando um NIF novo entra na watchlist e cruza com o arquivo).
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, "insolvency_digest", {})
    try:
        async with AsyncSessionLocal() as session:
            sent = await dispatch_pending(session, limit=200)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="ok", rows_seen=sent, rows_new=sent)
        return log_id
    except Exception as e:
        logger.exception("run_insolvency_digest failed: %s", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=str(e))
        raise


async def run_nif_refresh() -> str:
    """Enriquece empresas sem CAE / nome / morada, via VIES + SICAE.

    Substitui o `run_ptdata_refresh`, que estava a ser um no-op silencioso: só
    gravava quando `data["source"] == "ptdata"`, e como o VIES responde primeiro
    a condição nunca era verdadeira. As últimas corridas liam 95 empresas e
    actualizavam 0 — mas ficavam registadas como 'ok'.

    Passa a aceitar qualquer fonte e a gravar só o que veio preenchido.
    """
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, "nif_refresh", {})
        stale = (
            await session.execute(
                text(
                    """
                    SELECT id::text, nif FROM companies
                    WHERE cae IS NULL
                       OR ptdata_fetched_at IS NULL
                       OR ptdata_fetched_at < now() - interval '30 days'
                    ORDER BY (cae IS NULL) DESC, ptdata_fetched_at NULLS FIRST
                    LIMIT 400
                    """
                )
            )
        ).all()

    refreshed = 0
    for cid, nif in stale:
        try:
            data = await lookup_nif(nif)
            if data.get("source") == "local":
                continue
            async with AsyncSessionLocal() as session:
                res = await session.execute(
                    text(
                        """
                        UPDATE companies SET
                            ptdata_payload = CAST(:pl AS JSONB),
                            ptdata_fetched_at = now(),
                            legal_name = COALESCE(NULLIF(:ln, ''), legal_name),
                            cae = COALESCE(:cae, cae),
                            address = COALESCE(NULLIF(:addr, ''), address),
                            updated_at = now()
                        WHERE id = :id
                          AND (cae IS DISTINCT FROM COALESCE(:cae, cae)
                               OR ptdata_fetched_at IS NULL
                               OR address IS NULL)
                        """
                    ),
                    {
                        "pl": json.dumps(data.get("raw")),
                        "ln": data.get("legal_name") or "",
                        "cae": data.get("cae"),
                        "addr": data.get("address") or "",
                        "id": cid,
                    },
                )
                await session.commit()
                if res.rowcount:
                    refreshed += 1
        except Exception as e:  # noqa: BLE001
            logger.warning("nif_refresh nif=%s falhou: %s", nif, e)
        await random_delay()

    async with AsyncSessionLocal() as session:
        await _log_finish(session, log_id, status="ok", rows_seen=len(stale), rows_new=refreshed)
    logger.info("nif_refresh: %d lidas, %d actualizadas", len(stale), refreshed)
    return log_id


async def run_contracts_impic_sync(force: bool = False) -> str:
    """Bulk sync public contracts straight from IMPIC dumps on dados.gov.pt.

    The replacement for the old ptdata-based `run_contracts_sync`. Runs weekly
    (IMPIC refreshes the dumps biweekly). Downloads the per-year ZIPs whose
    `last_modified` moved forward, stream-parses the JSON, keeps only rows
    where a monitored NIF appears as adjudicatário or adjudicante, upserts
    `public_contracts`, and refreshes the summary cache on `companies`.

    `force=True` re-downloads every year regardless of sync state — used for
    the initial full-history backfill."""
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(
            session, "contracts", {"source": "impic_dados.gov.pt", "force": force},
        )
    try:
        async with AsyncSessionLocal() as session:
            stats = await sync_impic_contracts(session, force=force)
        async with AsyncSessionLocal() as session:
            await _log_finish(
                session, log_id, status="ok",
                rows_seen=sum((stats.get("rows_kept_by_year") or {}).values()),
                rows_new=stats.get("rows_inserted", 0),
            )
        return log_id
    except Exception as e:
        logger.exception("run_contracts_impic_sync failed: %r", e)
        async with AsyncSessionLocal() as session:
            await _log_finish(session, log_id, status="error", error=repr(e))
        raise



async def run_psp_alvara(alvara_type: str = "A") -> str:
    """Daily scrape of the PSP sigesponline roster for the given
    alvará type. Upserts alvara_companies and auto-registers newly
    seen NIPCs as monitoring_type='related'."""
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, f"psp_alvara_{alvara_type.lower()}", {})
        try:
            result = await sync_alvara_type(session, alvara_type)
            await _log_finish(
                session, log_id, status="ok",
                rows_seen=result.get("count") or 0,
                rows_new=len(result.get("new") or []),
            )
            return log_id
        except Exception as e:
            logger.exception("run_psp_alvara failed: %r", e)
            await _log_finish(session, log_id, status="error", error=repr(e))
            raise


async def run_dre_fetch() -> str:
    """Fetch official DRE RSS feeds (Série I + II) once for per-company
    matching, AND archive Série II Parte E into `dre_raw` for retroactive
    match when a new company is added later."""
    async with AsyncSessionLocal() as session:
        log_id = await _log_start(session, "dre", {})
        companies = await _monitored_companies(session)
    total_seen = 0
    total_new = 0
    archived_new = 0
    status = "ok"
    try:
        # 1. Archive Parte E raw stream (independent of our company set —
        #    feeds the retroactive match on new company creation).
        try:
            archived_new = await fetch_and_archive_parte_e()
        except Exception as e:
            logger.warning("dre raw archive failed: %s", e)
        # 2. Per-company match on current feeds (ongoing flow, unchanged).
        by_company = await fetch_dre_for_companies(companies)
        for cid, items in by_company.items():
            total_seen += len(items)
            async with AsyncSessionLocal() as session:
                total_new += await save_new_publications(session, cid, items)
    except Exception as e:
        logger.exception("dre run failed: %s", e)
        status = "error"
    async with AsyncSessionLocal() as session:
        await _log_finish(
            session, log_id, status=status,
            rows_seen=total_seen + archived_new,
            rows_new=total_new + archived_new,
        )
    return log_id


async def vacuum_logs() -> None:
    """Maintenance: drop old logs, expire sessions, recover stale 'running' jobs
    (killed by container restart)."""
    async with AsyncSessionLocal() as session:
        await session.execute(
            text("DELETE FROM scraping_logs WHERE started_at < now() - interval '180 days'")
        )
        await session.execute(text("DELETE FROM user_sessions WHERE expires_at < now()"))
        # Stale 'running' rows: a job that started >2h ago without finishing is
        # almost certainly orphaned (backend restart). Mark as failed so the UI
        # doesn't show a permanently-spinning row.
        await session.execute(
            text(
                """
                UPDATE scraping_logs
                SET status = 'error',
                    finished_at = now(),
                    error = COALESCE(error, 'killed by restart or timeout')
                WHERE status = 'running'
                  AND started_at < now() - interval '2 hours'
                """
            )
        )
        await session.commit()


async def reap_stale_running_on_startup() -> None:
    """Called once on backend lifespan startup. Any 'running' row from a previous
    process is by definition orphaned — mark them failed immediately."""
    async with AsyncSessionLocal() as session:
        await session.execute(
            text(
                """
                UPDATE scraping_logs
                SET status = 'error',
                    finished_at = now(),
                    error = COALESCE(error, 'killed by backend restart')
                WHERE status = 'running'
                """
            )
        )
        await session.commit()


async def _runner_for(source: str, fn, *args: Any, **kwargs: Any) -> None:
    try:
        await fn(*args, **kwargs)
    except Exception as e:
        logger.exception("background scrape %s failed: %s", source, e)


def schedule_background(source: str, fn, *args: Any, **kwargs: Any) -> asyncio.Task:
    loop = asyncio.get_running_loop()
    return loop.create_task(_runner_for(source, fn, *args, **kwargs))


async def _run_dre_retro_match(company_id: str) -> None:
    """Scheduler wrapper around `retroactively_match_company` — resolves the
    company's legal_name from DB and runs the retro match."""
    async with AsyncSessionLocal() as session:
        row = (
            await session.execute(
                text("SELECT legal_name FROM companies WHERE id = :id"),
                {"id": company_id},
            )
        ).first()
        if not row:
            return
        await retroactively_match_company(session, company_id, row[0])


def schedule_backfill_for_company(company_id: str, nif: str | None = None) -> None:
    """Fire distribuicao + cire + contracts history for a newly-created
    company, off the request thread.

    All three paths use the same canonical sources as the daily/weekly crons
    (no third-party aggregators):
      * distribuicao — CITIUS per-company sweep, 180-day window
      * cire         — CITIUS per-NIF, days=todos
      * contracts    — IMPIC dumps on dados.gov.pt, force=True so old ZIPs
        are re-processed with the new NIF included in the monitored set

    The contracts path re-downloads ~500 MB of ZIPs (~2 min); accepted as a
    one-off on-create cost in exchange for using the same authoritative
    pipeline as the weekly cron (no ptdata 504 mystery, no stale cache).

    Also fires the publicacoes.mj.pt historical scrape (Playwright + Whisper
    audio bypass of reCAPTCHA v2) — that's the full atos societários archive."""
    schedule_background(
        "distribuicao",
        run_distribuicao,
        None, None, ("internal", "competitor", "analysis"), company_id,
    )
    schedule_background(
        "cire",
        run_cire,
        "todos", ("internal", "competitor", "analysis"), company_id,
    )
    schedule_background(
        "contracts",
        run_contracts_impic_sync,
        True,  # force=True so every year file is re-processed with the new NIF
    )
    # Retroactively match the last 12 months of `dre_raw` (archived Série II
    # Parte E) against this company's distintivo, so the DRE timeline is
    # populated on day 1 instead of starting empty.
    schedule_background(
        "dre_retro",
        _run_dre_retro_match,
        company_id,
    )
    _ = nif  # reserved for publicacoes.mj.pt revive (see mj_publications docstring)


_VALID_TYPES = {"internal", "competitor", "analysis"}


def _coerce_types(raw: Any) -> tuple[str, ...]:
    """Default to (internal, competitor) when caller omits a filter."""
    if not raw:
        return ("internal", "competitor")
    if isinstance(raw, str):
        raw = [raw]
    out = tuple(t for t in raw if t in _VALID_TYPES)
    return out or ("internal", "competitor")


async def trigger_manual_scrape(
    source: str, options: dict[str, Any] | None = None
) -> str:
    """Enqueue a manual scrape. Returns a log_id placeholder (the real log is
    inserted by the job itself)."""
    options = options or {}
    types = _coerce_types(options.get("types"))
    if source == "distribuicao":
        date_from = options.get("date_from")
        date_to = options.get("date_to")
        if isinstance(date_from, str):
            date_from = datetime.strptime(date_from, "%Y-%m-%d").date()
        if isinstance(date_to, str):
            date_to = datetime.strptime(date_to, "%Y-%m-%d").date()
        task = schedule_background(
            "distribuicao", run_distribuicao, date_from, date_to, types
        )
    elif source == "registry_nacional":
        task = schedule_background(
            "registry_nacional", run_registry_nacional,
            int(options.get("max_items") or 200000),
            float(options.get("max_seconds") or 39000.0),
            bool(options.get("seed")),
        )
    elif source == "distribuicao_backfill":
        # Enche o arquivo nacional para trás até a fonte secar. Corre em segundo
        # plano porque são horas: ~180 dias úteis a 197 tribunais.
        task = schedule_background(
            "distribuicao_backfill",
            run_distribuicao_backfill,
            int(options.get("max_days") or 400),
        )
    elif source == "cire":
        days = options.get("days", "todos")
        task = schedule_background("cire", run_cire, days, types)
    elif source in ("ptdata", "nif"):
        # "ptdata" mantido: e o valor que o botao do /admin/scraping ja envia
        task = schedule_background("nif_refresh", run_nif_refresh)
    elif source == "dre":
        task = schedule_background("dre", run_dre_fetch)
    elif source in ("registry_sweep", "registry_expand", "registry_content"):
        # O backfill histórico corre por aqui, com orçamento grande. A fila vive
        # na base de dados, portanto um restart do contentor não perde trabalho —
        # ao contrário do backfill do CIRE, que foi lançado com `docker exec -d` e
        # que um rebuild teria morto a meio.
        date_from = options.get("date_from")
        date_to = options.get("date_to")
        if isinstance(date_from, str):
            date_from = datetime.strptime(date_from, "%Y-%m-%d").date()
        if isinstance(date_to, str):
            date_to = datetime.strptime(date_to, "%Y-%m-%d").date()
        budget = float(options.get("max_seconds") or 3600.0)
        items = int(options.get("max_items") or 4000)
        if source == "registry_sweep":
            task = schedule_background(
                "registry_sweep", run_registry_sweep, date_from, date_to,
                int(options.get("max_days") or 400), budget,
            )
        elif source == "registry_expand":
            task = schedule_background("registry_expand", run_registry_expand, items, budget)
        else:
            task = schedule_background("registry_content", run_registry_content, items, budget)
    elif source == "contracts":
        # Admin UI button runs the IMPIC dataset sync. For a per-company
        # backfill use schedule_backfill_for_company (on POST /companies).
        force = bool(options.get("force"))
        task = schedule_background("contracts", run_contracts_impic_sync, force)
    else:
        raise ValueError(f"unknown source: {source}")
    # Return task repr as a correlation id; the real log row is created by the job
    return f"task-{id(task)}"


def register_jobs() -> None:
    # Daily scrapes cover internal + competitor.
    # _monitored_companies orders internal first so a partial failure still
    # has internal results. Manual triggers via /api/admin/scrape can pass
    # types=['internal'] or types=['competitor'] to scope explicitly.
    # Main daily batch. With the tribunal-sweep scraper the CARLOS-SILVA-style
    # "noisy distintivo" problem no longer hurts: each tribunal is pulled once,
    # matched locally — the cost is constant regardless of any single company's
    # ambiguity. The old `scrape_isolated` / 20:00 late job is retired; column
    # kept in the DB for potential future revive (see migration 0005).
    scheduler.add_job(
        run_distribuicao, "cron", hour=3, minute=0,
        id="distribuicao_daily", replace_existing=True,
    )
    scheduler.add_job(
        run_cire, "cron", hour=3, minute=30,
        id="cire_daily", replace_existing=True,
        kwargs={"days": "30"},
    )

    # Sweep nacional dos anúncios de insolvência. Corre antes dos scrapes por
    # empresa para o prazo de reclamação estar disponível logo de manhã — o
    # relógio dos 30 dias já está a contar desde a publicação.
    scheduler.add_job(
        run_cire_announcements, "cron", hour=1, minute=30,
        id="cire_announcements_daily", replace_existing=True,
    )
    scheduler.add_job(
        run_insolvency_digest, "cron", hour=8, minute=0,
        id="insolvency_digest_daily", replace_existing=True,
    )

    # Publicações do registo comercial (portal IRN). Depois do backfill inicial
    # só toca nas empresas sincronizadas há mais de 7 dias, portanto a corrida
    # diária é curta.
    scheduler.add_job(
        run_mj_irn, "cron", hour=5, minute=45,
        id="mj_irn_daily", replace_existing=True,
    )

    # Registo comercial por NIPC + expansão pelos sócios-empresa. Corre depois do
    # mj_irn para as fichas já estarem enriquecidas, e com orçamento de 40 min
    # para não arrastar a noite quando a fila estiver longa.
    scheduler.add_job(
        run_registry_expand, "cron", hour=6, minute=20,
        kwargs={"max_items": 120, "max_seconds": 2400.0},
        id="registry_expand_daily", replace_existing=True,
    )

    # Fase A do índice nacional: listar o histórico por NIPC de todas as
    # entidades colectivas. Corre na mesma janela nocturna do resto da recolha e
    # cede a vez pelo `cpuunits` do CT. O orçamento de tempo é a noite inteira
    # menos uma folga; a fila vive na base, portanto o que não couber fica para a
    # noite seguinte sem se repetir.
    scheduler.add_job(
        run_registry_nacional, "cron", hour=20, minute=30,
        kwargs={"max_items": 200000, "max_seconds": 39000.0},
        id="registry_nacional_nightly", replace_existing=True,
    )

    # Depois da recolha, para os nós novos já aparecerem com os seus incidentes.
    scheduler.add_job(
        run_entity_incidents, "cron", hour=7, minute=10,
        id="entity_incidents_daily", replace_existing=True,
    )

    # Varrimento nacional: ontem e o que tenha ficado para trás. O backfill
    # histórico corre à parte, com orçamento maior, pelo /admin.
    scheduler.add_job(
        run_registry_sweep, "cron", hour=2, minute=15,
        kwargs={"max_days": 10, "max_seconds": 1200.0},
        id="registry_sweep_daily", replace_existing=True,
    )

    # Conteúdo dos actos, **fora do horário de trabalho**.
    #
    # O tecto de ritmo ao IRN já baixava a metade de dia, mas isso protege o
    # portal deles e não o CPU daqui: cada acto aberto é um parse e uma dúzia de
    # escritas, e o Postgres deste contentor media 1,26 cores em contínuo num
    # host de 4 que partilha com o Sabichão. O OCR do Sabichão é processador
    # puro e é uma pessoa à espera do ecrã — passou a arrastar-se e a estoirar
    # nos 180 segundos por não ter CPU disponível.
    #
    # Este backfill é um robô com dias de folga. Corre das 20h às 8h, e o
    # `cpuunits` do CT (20 contra os 100 do Sabichão) trata do resto: quando não
    # há disputa usa o que quiser, quando há cede a vez.
    #
    # **O orçamento que manda é o tempo, não a contagem.** Uma corrida de 90
    # minutos ao ritmo configurado abre ~20 mil actos; um `max_items` de 4000
    # cortava-a a um quinto e esticava o backfill de 2 dias e meio para quase
    # duas semanas, sem nada no log a explicar por quê. A contagem fica só como
    # travão de segurança, bem acima do que o tempo permite.
    scheduler.add_job(
        run_registry_content, "cron", hour="20-23,0-7", minute=40,
        kwargs={"max_items": 200000, "max_seconds": 5400.0},
        id="registry_content_loop", replace_existing=True,
    )

    # Canário do IRN: avisa dentro de uma hora quando um deploy do portal
    # invalidar os apiVersion, em vez de se descobrir com os jobs da noite.
    scheduler.add_job(
        run_irn_canary, "cron", minute=5,
        id="irn_canary_hourly", replace_existing=True,
    )

    # ptdata enrichment — once a day for any company missing CAE/etc.
    scheduler.add_job(
        run_nif_refresh, "cron", hour=5, minute=15,
        id="refresh_nif", replace_existing=True,
    )

    # PSP alvará A — daily scrape of the public security-company
    # roster. New NIPCs are auto-added as monitoring_type=related.
    scheduler.add_job(
        run_psp_alvara, "cron", hour=6, minute=30,
        id="psp_alvara_daily", replace_existing=True,
        kwargs={"alvara_type": "A"},
    )

    # DRE — fetched twice a day. Morning 07:00 catches the daily bulletin;
    # night 22:30 sweeps up anything added later in the day. Dedup hash in
    # dre_publications prevents duplicates.
    scheduler.add_job(
        run_dre_fetch, "cron", hour=7, minute=0,
        id="fetch_dre_morning", replace_existing=True,
    )
    scheduler.add_job(
        run_dre_fetch, "cron", hour=22, minute=30,
        id="fetch_dre_night", replace_existing=True,
    )

    # Public contracts (BASE/IMPIC via dados.gov.pt direct). Dataset refreshes
    # biweekly, so weekly is sufficient. Monday 04:30 avoids collision with
    # the 03:00 distribuicao + 03:30 cire jobs.
    scheduler.add_job(
        run_contracts_impic_sync, "cron",
        day_of_week="mon", hour=4, minute=30,
        id="contracts_weekly", replace_existing=True,
    )

    # Cleanup
    scheduler.add_job(
        vacuum_logs, "cron", day_of_week="sun", hour=2, minute=0,
        id="vacuum_logs", replace_existing=True,
    )


def list_scheduled_jobs() -> list[dict[str, Any]]:
    """Snapshot of all scheduled jobs and their next run times for the UI."""
    out: list[dict[str, Any]] = []
    for j in scheduler.get_jobs():
        out.append({
            "id": j.id,
            "name": j.func.__name__ if j.func else j.id,
            "trigger": str(j.trigger),
            "next_run_time": j.next_run_time.isoformat() if j.next_run_time else None,
            "kwargs": {k: list(v) if isinstance(v, tuple) else v for k, v in (j.kwargs or {}).items()},
        })
    return out


def start() -> None:
    register_jobs()
    scheduler.start()


def shutdown() -> None:
    scheduler.shutdown(wait=False)
