"""Digest data (spec §142): weekly company digest and industry / country market digests as JSON. No e-mail sending here — the API or a future mailer renders these. Everything comes from measured tables (events, metrics, jobs, signals); empty sections stay empty. """ from __future__ import annotations from datetime import UTC, datetime, timedelta from typing import Any from companyatlas.db import fetch_all, fetch_one, transaction EVENT_COLS = "id, event_type, event_subtype, importance, confidence, confidence_label, title, summary, detected_at, source_url, surface, origin" async def company_digest(slug_or_id: str, *, days: int = 7, now: datetime | None = None) -> dict[str, Any] | None: now = now or datetime.now(UTC) since = now - timedelta(days=days) async with transaction() as conn: co = await fetch_one(conn, "select id, slug, display_name, canonical_domain, country, industries, last_event_at from companies where slug = :s or id = :s", s=slug_or_id) if co is None: return None events = await fetch_all(conn, f"select {EVENT_COLS} from events where company_id = :c and detected_at >= :since and status = 'active' order by importance desc, detected_at desc limit 50", c=co["id"], since=since) by_type: dict[str, int] = {} for e in events: by_type[e["event_type"]] = by_type.get(e["event_type"], 0) + 1 metrics = {r["metric"]: {"value": r["value"], "confidence": r["confidence"], "computed_at": r["computed_at"]} for r in await fetch_all(conn, "select metric, value, confidence, computed_at from metrics_current where company_id = :c", c=co["id"])} series = await fetch_all(conn, "select metric, day, value from metric_series where company_id = :c and day >= :d and metric in ('activity_score', 'open_jobs') order by day", c=co["id"], d=(now - timedelta(days=days * 2)).date()) jobs = await fetch_one(conn, """select count(*) filter (where status = 'open') as open, count(*) filter (where first_seen_at >= :since) as new, count(*) filter (where removed_at >= :since) as no_longer_listed, count(*) filter (where status = 'open' and is_ai) as ai_open from jobs where company_id = :c""", c=co["id"], since=since) signals = await fetch_all(conn, "select kind, strength, confidence, title, explanation, detected_at from signals where company_id = :c and status = 'active' order by strength desc", c=co["id"]) sensors = await fetch_one(conn, "select count(*) filter (where status = 'active') as active, count(*) as total, max(last_success_at) as last_checked from sensors where company_id = :c", c=co["id"]) return {"company": co, "window": {"days": days, "since": since.isoformat(), "until": now.isoformat()}, "highlights": events[:8], "events": events, "events_by_type": by_type, "metrics": metrics, "series": {m: [{"day": r["day"].isoformat(), "value": r["value"]} for r in series if r["metric"] == m] for m in ("activity_score", "open_jobs")}, "jobs": jobs, "signals": signals, "coverage": sensors, "generated_at": now.isoformat()} async def scope_digest(scope: str, key: str, *, days: int = 7, now: datetime | None = None, limit: int = 10) -> dict[str, Any]: """scope ∈ industry | country. Movers = highest activity, top events = most important, hiring = aggregate momentum.""" now = now or datetime.now(UTC) since = now - timedelta(days=days) if scope == "country": where = "co.country = :key" elif scope == "industry": where = "(:key = any(co.industries) or co.industry_primary = :key)" else: raise ValueError("scope must be industry or country") async with transaction() as conn: companies_n = await fetch_val_int(conn, f"select count(*) from companies co where {where}", key=key) events = await fetch_all(conn, f"""select e.{EVENT_COLS.replace(', ', ', e.')}, co.slug, co.display_name from events e join companies co on co.id = e.company_id where {where} and e.detected_at >= :since and e.status = 'active' order by e.importance desc, e.detected_at desc limit :limit""", key=key, since=since, limit=limit * 2) by_type = await fetch_all(conn, f"""select e.event_type, count(*) as n from events e join companies co on co.id = e.company_id where {where} and e.detected_at >= :since and e.status = 'active' group by 1 order by n desc""", key=key, since=since) movers = await fetch_all(conn, f"""select co.slug, co.display_name, m.value as activity_score from metrics_current m join companies co on co.id = m.company_id where {where} and m.metric = 'activity_score' order by m.value desc limit :limit""", key=key, limit=limit) hiring = await fetch_one(conn, f"""select count(*) filter (where j.status = 'open') as open, count(*) filter (where j.first_seen_at >= :since) as new, count(*) filter (where j.removed_at >= :since) as no_longer_listed from jobs j join companies co on co.id = j.company_id where {where}""", key=key, since=since) momentum = await fetch_one(conn, f"""select avg(m.value) as avg_momentum_30d, count(*) as companies from metrics_current m join companies co on co.id = m.company_id where {where} and m.metric = 'hiring_momentum_30d'""", key=key) signals = await fetch_all(conn, "select kind, strength, title, explanation, evidence, detected_at from signals where scope = :s and scope_key = :k and status = 'active' order by strength desc", s=scope, k=key) return {"scope": scope, "key": key, "window": {"days": days, "since": since.isoformat(), "until": now.isoformat()}, "companies": companies_n, "top_events": events, "events_by_type": {r["event_type"]: r["n"] for r in by_type}, "movers": movers, "hiring": {**(hiring or {}), **(momentum or {})}, "signals": signals, "generated_at": now.isoformat()} async def fetch_val_int(conn, sql: str, **params: Any) -> int: # type: ignore[no-untyped-def] row = await fetch_one(conn, sql, **params) return int(next(iter(row.values()))) if row else 0 __all__ = ["company_digest", "scope_digest"]