"""/changes (cursor feed) · /changes/daily ("Today in AI 2.0") · /changes/categories. API 1.1: feeds default to `is_backfill = false` and are keyed on `occurred_at` (= coalesce(effective_at, observed_at)); `include_backfill=1` restores the historical corpus, `date_field=observed` restores the v1 ordering. `/changes/daily` groups the events that share a `group_key` (one release across several documents) into one item with `sources: n` and `documents: [urls]`.""" from __future__ import annotations from datetime import UTC, datetime from typing import Any from fastapi import APIRouter, Query, Request from aiatlas.api.common import ( ENTITY_COLS, ENTITY_FROM, EVENT_COLS, EVENT_FROM, OPEN_CATEGORIES, ApiError, cached, change_event, csv, day_bounds, entity_summary, event_type_label, parse_date, parse_ts, resolve_id, ) from aiatlas.db import connection, fetch_all, fetch_val router = APIRouter(prefix="/api/v1/changes", tags=["changes"]) CATEGORY_LABELS = {"model": "Models", "price": "Pricing", "benchmark": "Benchmarks", "paper": "Research", "release": "Releases", "company": "Companies & labs", "provider": "Providers", "hardware": "Hardware", "framework": "Frameworks", "dataset": "Datasets", "regulation": "Regulation", "incident": "Incidents", "repository": "Repositories", "tool": "Tools", "source": "Sources & curation", "update": "Updates"} CATEGORY_ORDER = list(CATEGORY_LABELS) TOTAL_CAP = 10_000 # "Today in AI 2.0" sections — (key, label, SQL predicate on ev/e) SECTIONS: list[tuple[str, str, str]] = [ ("MAJOR_RELEASES", "Major releases", "ev.event_type in ('NEW_MODEL','RELEASE') and ev.importance >= 2 and e.entity_type = 'model'"), ("OPEN_WEIGHT_RELEASES", "Open-weight releases", "ev.event_type = 'NEW_MODEL' and e.entity_type = 'model' and e.attributes->>'openness' in (" + ", ".join(f"'{c}'" for c in OPEN_CATEGORIES) + ")"), ("PRICE_MOVES", "Price moves", "ev.category = 'price' or ev.event_type = 'PRICE_CHANGED'"), ("BENCHMARK_MOVES", "Benchmark moves", "ev.event_type in ('BENCHMARK_UPDATED','BENCHMARK_LEADER_CHANGED','NEW_BENCHMARK_LEADER')"), ("MODEL_CHANGES", "Model changes", "ev.event_type in ('CONTEXT_CHANGED','MAX_OUTPUT_CHANGED','STATUS_CHANGED','CAPABILITIES_CHANGED','PARAMETERS_CHANGED','KNOWLEDGE_CUTOFF_CHANGED','LICENSE_CHANGED','OPENNESS_CHANGED')"), ("RESEARCH", "Research", "ev.event_type = 'NEW_PAPER'"), ("DEPRECATIONS", "Deprecations & retirements", "ev.event_type in ('DEPRECATION_ANNOUNCED','RETIREMENT_ANNOUNCED') or (ev.event_type = 'STATUS_CHANGED' and ev.new_value::text ~* 'deprecated|retired')"), ("PROVIDER_CHANGES", "Provider listings", "ev.event_type in ('PROVIDER_LISTED','PROVIDER_DELISTED','NEW_PROVIDER')"), ("HARDWARE", "Hardware", "ev.category = 'hardware' or ev.event_type = 'NEW_HARDWARE'"), ] def _filters(*, category: str | None, types: list[str], entity_type: str | None, importance_min: int | None, since: datetime | None, until: datetime | None, q: str | None, include_documents: bool, entity_id: str | None = None, before: datetime | None = None, include_backfill: bool = False, date_field: str = "occurred") -> tuple[list[str], dict[str, Any]]: col = "ev.observed_at" if date_field == "observed" else "ev.occurred_at" where: list[str] = [] p: dict[str, Any] = {} if not include_backfill: where.append("ev.is_backfill = false") if category: where.append("ev.category = any(cast(:cats as text[]))") p["cats"] = csv(category) if types: where.append("ev.event_type = any(cast(:types as text[]))") p["types"] = [t.upper() for t in types] elif not include_documents: where.append("ev.event_type <> 'DOCUMENT_CHANGED'") if entity_type: where.append("e.entity_type = any(cast(:etypes as text[]))") p["etypes"] = csv(entity_type) if importance_min is not None: where.append("ev.importance >= :imp") p["imp"] = importance_min if since is not None: where.append(f"{col} >= :since") p["since"] = since if until is not None: where.append(f"{col} <= :until") p["until"] = until if before is not None: where.append(f"{col} < :before") p["before"] = before if q: where.append("(ev.summary ilike :qlike or e.canonical_name ilike :qlike)") p["qlike"] = f"%{q}%" if entity_id: where.append("ev.entity_id = :eid") p["eid"] = entity_id return where or ["true"], p def _event(r: dict[str, Any]) -> dict[str, Any]: ev = change_event(r) ev["occurred_at"] = r.get("occurred_at") ev["is_backfill"] = r.get("is_backfill") ev["group_key"] = r.get("group_key") return ev @router.get("") @cached(60) async def list_changes(request: Request, category: str | None = None, type: str | None = Query(None, alias="type"), entity_type: str | None = None, importance_min: int | None = Query(None, ge=0, le=3), since: str | None = None, until: str | None = None, q: str | None = Query(None, max_length=200), entity: str | None = None, limit: int = Query(50, ge=1, le=200), before: str | None = None, offset: int = Query(0, ge=0, le=10000), include_documents: int = Query(0, ge=0, le=1), include_backfill: int = Query(0, ge=0, le=1), date_field: str = Query("occurred", pattern="^(occurred|observed)$")) -> dict[str, Any]: before_ts, since_ts, until_ts = parse_ts(before, "before"), parse_ts(since, "since"), parse_ts(until, "until") col = "ev.observed_at" if date_field == "observed" else "ev.occurred_at" async with connection() as conn: eid = await resolve_id(conn, entity) if entity else None where, params = _filters(category=category, types=csv(type), entity_type=entity_type, importance_min=importance_min, since=since_ts, until=until_ts, q=q, include_documents=bool(include_documents), entity_id=eid, before=before_ts, include_backfill=bool(include_backfill), date_field=date_field) where_sql = " and ".join(where) rows = await fetch_all(conn, f"select {EVENT_COLS}, ev.occurred_at, ev.is_backfill, ev.group_key from {EVENT_FROM} where {where_sql} order by {col} desc, ev.id desc limit :lim offset :off", lim=limit, off=offset, **params) total = await fetch_val(conn, f"select count(*) from (select 1 from {EVENT_FROM} where {where_sql} limit {TOTAL_CAP}) t", **params) items = [_event(r) for r in rows] cursor_key = "observed_at" if date_field == "observed" else "occurred_at" return {"items": items, "total": int(total or 0), "limit": limit, "offset": offset, "next_before": items[-1][cursor_key] if len(items) == limit else None, "date_field": date_field, "include_backfill": bool(include_backfill)} def _group(items: list[dict[str, Any]]) -> list[dict[str, Any]]: """Fold events that share a `group_key` (one release across documents) into one item: the most important event + sources/documents.""" out: list[dict[str, Any]] = [] by_key: dict[str, dict[str, Any]] = {} for ev in items: k = ev.get("group_key") if not k: out.append({**ev, "sources": 1, "documents": [ev["source_url"]] if ev.get("source_url") else [], "grouped_events": 1}) continue g = by_key.get(k) if g is None: g = by_key[k] = {**ev, "sources": 0, "documents": [], "grouped_events": 0, "event_ids": []} out.append(g) g["grouped_events"] += 1 g["event_ids"].append(ev["id"]) if ev.get("source_url") and ev["source_url"] not in g["documents"]: g["documents"].append(ev["source_url"]) if (ev.get("importance") or 0) > (g.get("importance") or 0): for f in ("id", "event_type", "summary", "importance", "new_value", "old_value", "property", "entity"): g[f] = ev.get(f) for g in out: g["sources"] = max(1, len(g["documents"])) return out @router.get("/daily") @cached(120) async def changes_daily(request: Request, date: str | None = None, per_section: int = Query(30, ge=1, le=100), include_backfill: int = Query(0, ge=0, le=1), date_field: str = Query("occurred", pattern="^(occurred|observed)$")) -> dict[str, Any]: d = parse_date(date, "date") or datetime.now(UTC).date() start, end = day_bounds(d) col = "ev.observed_at" if date_field == "observed" else "ev.occurred_at" bf = "" if include_backfill else " and ev.is_backfill = false" base = f"{col} >= :s and {col} <= :e and ev.event_type <> 'DOCUMENT_CHANGED'{bf}" async with connection() as conn: rows = await fetch_all(conn, f"""select {EVENT_COLS}, ev.occurred_at, ev.is_backfill, ev.group_key from ( select ev.*, row_number() over (partition by ev.category order by ev.importance desc, ev.occurred_at desc) as rn from change_events ev where {base}) ev left join entities e on e.id = ev.entity_id left join entities eo on eo.id = e.organization_id where ev.rn <= :n order by ev.category, ev.rn""", s=start, e=end, n=per_section) counts = await fetch_all(conn, f"select ev.category, count(*) as n from change_events ev where {base} group by 1", s=start, e=end) backfill_excluded = 0 if include_backfill else int(await fetch_val(conn, f"select count(*) from change_events ev where {col} >= :s and {col} <= :e and ev.is_backfill and ev.event_type <> 'DOCUMENT_CHANGED'", s=start, e=end) or 0) new_models = await fetch_all(conn, f"select {ENTITY_COLS} from {ENTITY_FROM} where e.entity_type = 'model' and e.merged_into is null and e.first_seen_at >= :s and e.first_seen_at <= :e " f"order by e.first_seen_at desc limit 100", s=start, e=end) prev = await fetch_val(conn, f"select max(ev.occurred_at)::date from change_events ev where ev.occurred_at < :s and ev.event_type <> 'DOCUMENT_CHANGED'{bf}", s=start) nxt = await fetch_val(conn, f"select min(ev.occurred_at)::date from change_events ev where ev.occurred_at > :e and ev.event_type <> 'DOCUMENT_CHANGED'{bf}", e=end) sections2: list[dict[str, Any]] = [] for key, label, pred in SECTIONS: srows = await fetch_all(conn, f"""select {EVENT_COLS}, ev.occurred_at, ev.is_backfill, ev.group_key from {EVENT_FROM} where {base} and ({pred}) order by ev.importance desc, ev.occurred_at desc, ev.id limit :n""", s=start, e=end, n=per_section * 3) items = _group([_event(r) for r in srows])[:per_section] total_s = await fetch_val(conn, f"select count(*) from {EVENT_FROM} where {base} and ({pred})", s=start, e=end) if items: sections2.append({"key": key, "label": label, "items": items, "total": int(total_s or 0)}) by_cat: dict[str, list[dict[str, Any]]] = {} for r in rows: by_cat.setdefault(r["category"], []).append(_event(r)) order = {c: i for i, c in enumerate(CATEGORY_ORDER)} sections = [{"category": c, "label": CATEGORY_LABELS.get(c, c.title()), "items": items} for c, items in sorted(by_cat.items(), key=lambda kv: (order.get(kv[0], 99), kv[0]))] return {"date": d.isoformat(), "counts": {r["category"]: int(r["n"]) for r in counts}, "total": sum(int(r["n"]) for r in counts), "sections": sections, "today": sections2, "new_models": [entity_summary(r) for r in new_models], "labels": CATEGORY_LABELS, "backfill_excluded": backfill_excluded, "date_field": date_field, "previous_day": prev.isoformat() if prev else None, "next_day": nxt.isoformat() if nxt else None, "note": "Events that occurred on this UTC day (effective date when known, else observation date), excluding back-filled history and source-document " "changes. `today` groups events sharing a group_key (one release seen in several documents)."} @router.get("/categories") @cached(300) async def changes_categories(request: Request, days: int = Query(7, ge=1, le=365), include_backfill: int = Query(0, ge=0, le=1)) -> dict[str, Any]: bf = "" if include_backfill else " and is_backfill = false" async with connection() as conn: rows = await fetch_all(conn, f"""select category, event_type, count(*) as count from change_events where occurred_at > now() - make_interval(days => :d) and event_type <> 'DOCUMENT_CHANGED'{bf} group by 1, 2 order by 3 desc, 1, 2""", d=days) return {"days": days, "items": [{"category": r["category"], "label": CATEGORY_LABELS.get(r["category"], r["category"].title()), "event_type": r["event_type"], "event_label": event_type_label(r["event_type"]), "count": int(r["count"])} for r in rows]} __all__ = ["CATEGORY_LABELS", "SECTIONS", "ApiError", "router"]