SPB Git forge

spb/ai-atlas

Public
41commits 1branches 0releases
4.6 MBsize
maindefault branch
12 days agolast push
HTML 77.2% TypeScript 10.5% Python 9.6% JavaScript 2.5%
38.1 KB · 532 lines python
Raw Blame History
1"""FactWriter — turns `Facts` into versioned database state.23Temporal rules (per entity × property):4  * no current claim          → insert current claim, set attribute, NEW_* event for new entities5  * same value                → confirm (observed_at bumped), nothing else6  * same assertion, older encoding (taxonomy) → re-encode the current claim in place (canonical value + value_raw), never an event7  * different, source ≥ tier  → supersede (valid_to = observed), insert new current, update attribute, CHANGE event for material props8  * different, same source AND same extractor → supersede (a source correcting itself); an LLM claim never supersedes a deterministic9    claim from the same URL — it is stored as conflicting10  * different, worse source   → store as `conflicting`, flag confidence, review-queue item — never overwrite11Prices and benchmark results have their own append-only tables with the same close/open semantics; results additionally keep ONE12current row per (model, benchmark, metric, config_key): a newer run group closes the older rows so leaderboards show the latest run.13Every fact carries `run_id` (batch rollback) and every event carries `recorded_at`, `is_backfill` and `group_key`.14"""15from __future__ import annotations1617import hashlib18import json19import logging20from datetime import UTC, datetime21from typing import Any2223from sqlalchemy.ext.asyncio import AsyncConnection2425from aiatlas.db import execute, fetch_one, jsonb26from aiatlas.ids import new_id, normalize_alias27from aiatlas.ontology import benchmarks as bench_ontology28from aiatlas.ontology.anomalies import check_price_movement, check_result29from aiatlas.ontology.models import effort_config30from aiatlas.ontology.taxonomy import TAXONOMY_PROPERTIES, normalize_property31from aiatlas.sdk.facts import (32    EVENT_CATEGORY_BY_TYPE,33    MATERIAL_PROPERTIES,34    NOISY_PREFIXES,35    EntityRef,36    Facts,37    PriceObs,38    ResultObs,39)40from aiatlas.sdk.resolution import Resolver41from aiatlas.services.events import classify_backfill, group_key_for, importance_for4243log = logging.getLogger(__name__)4445TIER_CONFIDENCE = {1: "high", 2: "medium", 3: "low", 4: "low"}46SOFT_PROPERTIES = {"description", "summary", "tagline", "availability_note", "abstract", "notes", "training_data_notes", "safety_notes", "hardware_requirements"}47# derived / secondary encodings of another property: never their own change event48NO_EVENT_PROPERTIES = {"license_key"}49NEW_IMPORTANCE = {"model": 3, "company": 2, "provider": 2, "paper": 1, "dataset": 1, "benchmark": 2, "framework": 1, "hardware": 2,50                  "tool": 1, "repository": 0, "release": 2, "regulation": 2, "incident": 2, "organization": 1, "researcher": 0,51                  "artifact": 0, "model_family": 1, "license": 0}52SILENT_NEW_TYPES = {"researcher", "country", "license", "quantization"}535455def _norm_value(v: Any) -> Any:56    if isinstance(v, datetime):57        return v.astimezone(UTC).isoformat(timespec="seconds")58    if isinstance(v, dict):59        return {k: _norm_value(x) for k, x in v.items()}60    if isinstance(v, (set, tuple, list)):61        items = [_norm_value(x) for x in v]62        return sorted(items) if all(isinstance(x, (str, int, float)) and not isinstance(x, bool) for x in items) else items63    if isinstance(v, float) and v.is_integer() and abs(v) < 1e15:64        return int(v)65    return v666768def _same(a: Any, b: Any) -> bool:69    return json.dumps(_norm_value(a), sort_keys=True, default=str) == json.dumps(_norm_value(b), sort_keys=True, default=str)707172def _cfg_hash(config: dict[str, Any] | None) -> str:73    return hashlib.sha1(json.dumps(config or {}, sort_keys=True, default=str).encode()).hexdigest()[:12]747576def result_dedupe_key(model_id: str, benchmark_id: str, config: dict[str, Any] | None, metric: str | None) -> str:77    return f"{model_id}:{benchmark_id}:{_cfg_hash(config)}:{metric or ''}"787980class WriteStats:81    def __init__(self) -> None:82        self.entities_created = 083        self.entities_updated = 084        self.claims = 085        self.relations = 086        self.events = 087        self.prices = 088        self.results = 089        self.conflicts = 090        self.folded_variants = 09192    def as_dict(self) -> dict[str, int]:93        return dict(self.__dict__)949596class FactWriter:97    def __init__(self, conn: AsyncConnection, *, source_id: str | None, snapshot_id: str | None, source_url: str | None,98                 tier: int = 2, connector_name: str | None = None, extractor: str = "deterministic", extractor_version: str = "1",99                 observed_at: datetime | None = None, run_id: str | None = None, source_key: str | None = None, is_first_run: bool = False,100                 variant_index: dict[str, str] | None = None):101        self.conn = conn102        self.source_id = source_id103        self.snapshot_id = snapshot_id104        self.source_url = source_url105        self.tier = tier106        self.connector_name = connector_name107        self.extractor = extractor108        self.extractor_version = extractor_version109        self.observed_at = observed_at or datetime.now(UTC)110        self.run_id = run_id111        self._source_key = source_key112        self.is_first_run = is_first_run           # connector's first successful run → every event is backfill (initial corpus)113        self.resolver = Resolver(conn, snapshot_id=snapshot_id, source_tier=tier, variant_index=variant_index)114        self.stats = WriteStats()115116    @property117    def derived(self) -> bool:118        """Derived writers (canonicalization) store disagreements as conflicting claims but never open review items."""119        return self.extractor == "derived"120121    async def source_key(self) -> str | None:122        if self._source_key is None and self.source_id:123            row = await fetch_one(self.conn, "select key from sources where id = :id", id=self.source_id)124            self._source_key = row["key"] if row else ""125        return self._source_key or None126127    # ---------------------------------------------------------------------------------------------- entry point128    async def write(self, facts: Facts) -> WriteStats:129        new_before = len(self.resolver.created)130        for ref in facts.entities:131            await self.resolver.resolve(ref)132            await self._apply_hierarchy(ref)133            for prop, value in ref.attributes.items():134                await self.write_claim(ref, prop, value)135        for c in facts.claims:136            await self.write_claim(c.entity, c.property, c.value, unit=c.unit, confidence=c.confidence, observed_at=c.observed_at,137                                   effective_at=c.effective_at, source_url=c.source_url)138        for r in facts.relations:139            await self.write_relation(r.subject, r.predicate, r.object, r.attributes, confidence=r.confidence, source_url=r.source_url)140        for p in facts.prices:141            await self.write_price(p)142        for res in facts.results:143            await self.write_result(res)144        for e in facts.events:145            eid = await self.resolver.resolve(e.entity) if e.entity else None146            await self.emit_event(e.event_type, e.category, e.summary, entity_id=eid, old_value=e.old_value, new_value=e.new_value,147                                  importance=e.importance, effective_at=e.effective_at, dedupe_key=e.dedupe_key, source_url=e.source_url, meta=e.meta)148        # NEW_* events for every entity created in this write149        for eid in self.resolver.created[new_before:]:150            await self._new_entity_event(eid)151        self.stats.entities_created = len(self.resolver.created)152        self.stats.entities_updated = len(self.resolver.updated - set(self.resolver.created))153        self.stats.folded_variants = len(self.resolver.folded)154        return self.stats155156    # ---------------------------------------------------------------------------------------------- hierarchy hints157    async def _apply_hierarchy(self, ref: EntityRef) -> None:158        """Materialise `family` / `canonical` hints: entities.family_id + member_of_family, entities.canonical_id + artifact_of, identity_confidence."""159        if not ref.id:160            return161        if ref.family is not None:162            if ref.family.entity_type != "model_family":163                ref.family.entity_type = "model_family"164            fid = await self.resolver.resolve(ref.family)165            if fid and fid != ref.id:166                await execute(self.conn, "update entities set family_id = :f where id = :e and family_id is distinct from :f", f=fid, e=ref.id)167                await self.write_relation(ref, "member_of_family", ref.family)168        if ref.canonical is not None:169            cid = await self.resolver.resolve(ref.canonical)170            if cid and cid != ref.id:171                await execute(self.conn, """update entities set canonical_id = :c, artifact_kind = coalesce(cast(:k as text), artifact_kind)172                                            where id = :e and (canonical_id is distinct from :c or (cast(:k as text) is not null and artifact_kind is null))""",173                              c=cid, k=ref.artifact_kind, e=ref.id)174                await self.write_relation(ref, "artifact_of", ref.canonical, {"artifact_kind": ref.artifact_kind} if ref.artifact_kind else None)175        if ref.identity_confidence:176            await execute(self.conn, "update entities set identity_confidence = :ic where id = :e and identity_confidence <> :ic", ic=ref.identity_confidence, e=ref.id)177178    # ---------------------------------------------------------------------------------------------- claims179    async def write_claim(self, ref: EntityRef, prop: str, value: Any, *, unit: str | None = None, confidence: str | None = None,180                          observed_at: datetime | None = None, effective_at: datetime | None = None, source_url: str | None = None) -> None:181        if value is None or value == "" or value == [] or value == {}:182            return183        eid = await self.resolver.resolve(ref)184        assert eid185        value = _norm_value(value)186        value_raw: str | None = None187        if prop in TAXONOMY_PROPERTIES:188            value, value_raw, mappings = normalize_property(ref.entity_type, prop, value)189            for domain, raw, canon in mappings:190                await self._record_mapping(domain, raw, canon)191            value = _norm_value(value)192        observed = observed_at or self.observed_at193        conf = confidence or TIER_CONFIDENCE.get(self.tier, "medium")194        url = source_url or self.source_url195        noisy = prop.startswith(NOISY_PREFIXES)196        current = await fetch_one(self.conn, """select id, value, value_raw, tier, source_url, observed_at, extractor from claims197                                                where entity_id = :e and property = :p and status = 'current' order by valid_from desc limit 1""",198                                  e=eid, p=prop)199        is_new_entity = eid in self.resolver.created200        await self._apply_claim(ref, eid, prop, value, value_raw, current, unit=unit, conf=conf, observed=observed, effective_at=effective_at, url=url,201                                noisy=noisy, is_new_entity=is_new_entity)202        if prop == "license" and isinstance(value, str):203            from aiatlas.ontology.licenses import LICENSES204205            if value in LICENSES:   # canonical licence key as its own property (filterable), never its own event206                await self.write_claim(ref, "license_key", value, confidence=confidence, observed_at=observed_at, effective_at=effective_at, source_url=source_url)207208    async def _apply_claim(self, ref: EntityRef, eid: str, prop: str, value: Any, value_raw: str | None, current: dict[str, Any] | None, *, unit: str | None,209                           conf: str, observed: datetime, effective_at: datetime | None, url: str | None, noisy: bool, is_new_entity: bool) -> None:210        if current is None:211            await self._insert_claim(eid, prop, value, unit, conf, "current", observed, effective_at, url, value_raw)212            await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw)213            return214        if _same(current["value"], value):215            await execute(self.conn, "update claims set observed_at = greatest(observed_at, :o) where id = :id", o=observed, id=current["id"])216            if noisy:217                await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw)218            return219        if prop in TAXONOMY_PROPERTIES and _same(normalize_property(ref.entity_type, prop, current["value"])[0], value):220            # same assertion in an older encoding ("apache-2.0" → "Apache-2.0"): re-encode in place, keep the source label, no event221            raw_keep = current["value_raw"] or (current["value"] if isinstance(current["value"], str) else json.dumps(current["value"], ensure_ascii=False))222            raw_keep = raw_keep if raw_keep != value else None223            await execute(self.conn, """update claims set value = cast(:v as jsonb), value_text = :vt, value_raw = :raw, observed_at = greatest(observed_at, :o) where id = :id""",224                          v=jsonb(value), vt=value[:2000] if isinstance(value, str) else None, raw=raw_keep, o=observed, id=current["id"])225            await self._set_attribute(eid, prop, value, unit, conf, url, observed, raw_keep, keep_provenance=True)226            return227        same_source = bool(url and current["source_url"] == url)228        same_extractor = (current["extractor"] or "deterministic") == self.extractor229        if prop in SOFT_PROPERTIES and not same_source:230            # soft text (descriptions, notes): first statement wins until *its own* source changes; never an event, never a conflict231            return232        if self.tier <= (current["tier"] or 2) or (same_source and same_extractor):233            # supersede234            await execute(self.conn, "update claims set status = 'superseded', valid_to = :o where id = :id", o=observed, id=current["id"])235            await self._insert_claim(eid, prop, value, unit, conf, "current", observed, effective_at, url, value_raw)236            await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw)237            if not noisy and not is_new_entity and prop not in SOFT_PROPERTIES and prop not in NO_EVENT_PROPERTIES:238                await self._property_change_event(eid, prop, current["value"], value, observed, effective_at, url)239            return240        # the same disagreement from the same source/extractor is recorded once (re-observed, not re-inserted)241        dup = await fetch_one(self.conn, """select id from claims where entity_id = :e and property = :p and status = 'conflicting' and value = cast(:v as jsonb)242                                            and extractor = :ex and coalesce(source_url, '') = :url limit 1""", e=eid, p=prop, v=jsonb(value), ex=self.extractor, url=url or "")243        if dup:244            await execute(self.conn, "update claims set observed_at = greatest(observed_at, :o) where id = :id", o=observed, id=dup["id"])245            return246        await self._insert_claim(eid, prop, value, unit, conf, "conflicting", observed, effective_at, url, value_raw)247        await execute(self.conn, "update claims set confidence = 'conflicted' where id = :id", id=current["id"])248        await execute(self.conn, """update entities set quality = quality || jsonb_build_object('conflicts', coalesce((quality->>'conflicts')::int, 0) + 1)249                                    where id = :id""", id=eid)250        self.stats.conflicts += 1251        if not self.derived:252            dedupe = f"conflict:{eid}:{prop}:{hashlib.sha1(json.dumps(value, sort_keys=True, default=str).encode()).hexdigest()[:10]}"253            await execute(self.conn, """insert into review_queue (id, kind, entity_ids, payload, reason, dedupe_key)254                                        values (:id, 'conflict', :ids, cast(:p as jsonb), :r, :d) on conflict (dedupe_key) do nothing""",255                          id=new_id("review"), ids=[eid], p=jsonb({"property": prop, "current": current["value"], "current_source": current["source_url"],256                                                                  "claimed": value, "claimed_source": url, "tier": self.tier, "extractor": self.extractor}),257                          r=f"source disagrees on {prop}", d=dedupe)258259    async def _record_mapping(self, domain: str, raw: str, canonical: str | None) -> None:260        if canonical is not None and raw == canonical:261            return   # identity mappings carry no information262        await execute(self.conn, """insert into taxonomy_mappings (domain, raw, canonical) values (:d, :r, :c)263                                    on conflict (domain, raw) do update set count = taxonomy_mappings.count + 1, last_seen_at = now(),264                                        canonical = coalesce(excluded.canonical, taxonomy_mappings.canonical)""", d=domain, r=raw[:300], c=canonical)265266    async def _insert_claim(self, eid: str, prop: str, value: Any, unit: str | None, conf: str, status: str, observed: datetime,267                            effective_at: datetime | None, url: str | None, value_raw: str | None = None) -> None:268        text_val = value if isinstance(value, str) else None269        num_val = float(value) if isinstance(value, (int, float)) and not isinstance(value, bool) else None270        await execute(self.conn, """insert into claims (id, entity_id, property, value, value_text, value_num, unit, source_id, snapshot_id, source_url, tier,271                                    confidence, status, extractor, extractor_version, observed_at, effective_at, valid_from, run_id, value_raw)272                                    values (:id, :e, :p, cast(:v as jsonb), :vt, :vn, :u, :src, :snap, :url, :tier, :conf, :status, :ex, :exv, :o, :eff, :vf, :run, :raw)""",273                      id=new_id("claim"), e=eid, p=prop, v=jsonb(value), vt=text_val[:2000] if text_val else None, vn=num_val, u=unit, src=self.source_id,274                      snap=self.snapshot_id, url=url, tier=self.tier, conf=conf, status=status, ex=self.extractor, exv=self.extractor_version,275                      o=observed, eff=effective_at, vf=effective_at or observed, run=self.run_id, raw=value_raw[:500] if value_raw else None)276        self.stats.claims += 1277278    async def _set_attribute(self, eid: str, prop: str, value: Any, unit: str | None, conf: str, url: str | None, observed: datetime,279                             value_raw: str | None = None, *, keep_provenance: bool = False) -> None:280        prov = {"source_id": self.source_id, "snapshot_id": self.snapshot_id, "url": url, "observed_at": observed.isoformat(timespec="seconds"),281                "tier": self.tier, "confidence": conf, "extractor": self.extractor}282        if unit:283            prov["unit"] = unit284        attrs = {prop: value}285        if value_raw is not None:286            attrs[f"{prop}_raw"] = value_raw287        if keep_provenance:288            await execute(self.conn, """update entities set attributes = (attributes - cast(:praw as text)) || cast(:a as jsonb), last_seen_at = greatest(last_seen_at, :o) where id = :id""",289                          praw=f"{prop}_raw", a=jsonb(attrs), o=observed, id=eid)290        else:291            await execute(self.conn, """update entities set attributes = (attributes - cast(:praw as text)) || cast(:a as jsonb),292                                        provenance = provenance || jsonb_build_object(cast(:p as text), cast(:prov as jsonb)), last_seen_at = greatest(last_seen_at, :o)293                                        where id = :id""", praw=f"{prop}_raw", a=jsonb(attrs), p=prop, prov=jsonb(prov), o=observed, id=eid)294        if prop == "status" and isinstance(value, str):295            await execute(self.conn, "update entities set status = :s where id = :id", s=value[:40], id=eid)296        if prop == "description" and isinstance(value, str):297            await execute(self.conn, "update entities set description = :d where id = :id", d=value[:4000], id=eid)298299    async def _property_change_event(self, eid: str, prop: str, old: Any, new: Any, observed: datetime, effective_at: datetime | None,300                                     url: str | None) -> None:301        event_type, importance = MATERIAL_PROPERTIES.get(prop, ("PROPERTY_CHANGED", 0))302        row = await fetch_one(self.conn, "select canonical_name, entity_type from entities where id = :id", id=eid)303        name = row["canonical_name"] if row else eid304        etype = row["entity_type"] if row else ""305        category = EVENT_CATEGORY_BY_TYPE.get(etype, "update")306        importance = importance_for(event_type, etype, old, new, tier=self.tier, default=importance)307        summary = f"{name}: {prop.replace('_', ' ')} changed from {_short(old)} to {_short(new)}"308        dedupe = f"{event_type}:{eid}:{prop}:{hashlib.sha1(json.dumps([old, new], sort_keys=True, default=str).encode()).hexdigest()[:12]}"309        await self.emit_event(event_type, category, summary, entity_id=eid, old_value=old, new_value=new, importance=importance,310                              effective_at=effective_at, dedupe_key=dedupe, source_url=url, meta={"property": prop}, observed_at=observed)311312    # ---------------------------------------------------------------------------------------------- relations313    async def write_relation(self, subject: EntityRef, predicate: str, obj: EntityRef, attributes: dict[str, Any] | None = None, *,314                             confidence: str | None = None, source_url: str | None = None) -> None:315        sid = await self.resolver.resolve(subject)316        oid = await self.resolver.resolve(obj)317        if not sid or not oid or sid == oid:318            return319        if predicate == "develops":320            # an organisation develops a *model*; a hub repository that is an artifact (quantisation, conversion) is merely published by its org321            otype = await fetch_one(self.conn, "select entity_type from entities where id = :id", id=oid)322            if otype and otype["entity_type"] == "artifact":323                predicate = "published_by"324                sid, oid = oid, sid325        conf = confidence or TIER_CONFIDENCE.get(self.tier, "medium")326        existing = await fetch_one(self.conn, "select id, attributes from relations where subject_id = :s and predicate = :p and object_id = :o and valid_to is null",327                                   s=sid, p=predicate, o=oid)328        if existing:329            if attributes and not _same(existing["attributes"], {**existing["attributes"], **attributes}):330                await execute(self.conn, "update relations set attributes = attributes || cast(:a as jsonb), observed_at = :o where id = :id",331                              a=jsonb(attributes), o=self.observed_at, id=existing["id"])332            else:333                await execute(self.conn, "update relations set observed_at = :o where id = :id", o=self.observed_at, id=existing["id"])334            return335        await execute(self.conn, """insert into relations (id, subject_id, predicate, object_id, attributes, source_id, snapshot_id, source_url, tier, confidence, observed_at, valid_from, run_id)336                                    values (:id, :s, :p, :o, cast(:a as jsonb), :src, :snap, :url, :tier, :conf, :obs, :obs, :run)""",337                      id=new_id("relation"), s=sid, p=predicate, o=oid, a=jsonb(attributes or {}), src=self.source_id, snap=self.snapshot_id,338                      url=source_url or self.source_url, tier=self.tier, conf=conf, obs=self.observed_at, run=self.run_id)339        self.stats.relations += 1340341    # ---------------------------------------------------------------------------------------------- prices342    async def write_price(self, p: PriceObs) -> None:343        mid = await self.resolver.resolve(p.model)344        pid = await self.resolver.resolve(p.provider)345        if not mid or not pid:346            return347        if all(v is None for v in p.price_tuple()[:-1]):348            return349        url = p.source_url or self.source_url350        current = await fetch_one(self.conn, """select * from prices where model_id = :m and provider_id = :p and coalesce(provider_model_id,'') = :pm and valid_to is null""",351                                  m=mid, p=pid, pm=p.provider_model_id or "")352        new_tuple = p.price_tuple()353        if current:354            cur_tuple = (current["input_per_mtok"], current["output_per_mtok"], current["cached_input_per_mtok"], current["cache_write_per_mtok"],355                         current["batch_input_per_mtok"], current["batch_output_per_mtok"], current["per_image"], current["per_request"], current["currency"])356            if _same(cur_tuple, new_tuple) and (p.context_length in (None, current["context_length"])):357                await execute(self.conn, "update prices set observed_at = :o where id = :id", o=self.observed_at, id=current["id"])358                return359            await execute(self.conn, "update prices set valid_to = :o where id = :id", o=self.observed_at, id=current["id"])360        price_id = new_id("price")361        await execute(self.conn, """insert into prices (id, model_id, provider_id, provider_model_id, input_per_mtok, output_per_mtok, cached_input_per_mtok, cache_write_per_mtok,362                                    batch_input_per_mtok, batch_output_per_mtok, per_image, per_request, currency, context_length, max_output_tokens, features,363                                    observed_at, valid_from, source_id, snapshot_id, source_url, tier, meta, run_id)364                                    values (:id, :m, :p, :pm, :i, :o, :ci, :cw, :bi, :bo, :img, :req, :cur, :ctx, :mo, cast(:f as jsonb), :obs, :obs, :src, :snap, :url, :tier, cast(:meta as jsonb), :run)""",365                      id=price_id, m=mid, p=pid, pm=p.provider_model_id, i=p.input_per_mtok, o=p.output_per_mtok, ci=p.cached_input_per_mtok,366                      cw=p.cache_write_per_mtok, bi=p.batch_input_per_mtok, bo=p.batch_output_per_mtok, img=p.per_image, req=p.per_request, cur=p.currency,367                      ctx=p.context_length, mo=p.max_output_tokens, f=jsonb(p.features), obs=self.observed_at, src=self.source_id, snap=self.snapshot_id,368                      url=url, tier=self.tier, meta=jsonb(p.meta), run=self.run_id)369        self.stats.prices += 1370        model = await fetch_one(self.conn, "select canonical_name, entity_type from entities where id = :id", id=mid)371        provider = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=pid)372        mname = model["canonical_name"] if model else mid373        pname = provider["canonical_name"] if provider else pid374        if current:375            old = {"input_per_mtok": current["input_per_mtok"], "output_per_mtok": current["output_per_mtok"]}376            new = {"input_per_mtok": p.input_per_mtok, "output_per_mtok": p.output_per_mtok}377            summary = f"{pname} changed pricing for {mname}: {_fmt_price(old)} → {_fmt_price(new)}"378            dedupe = f"PRICE_CHANGED:{mid}:{pid}:{p.provider_model_id or ''}:{hashlib.sha1(json.dumps([old, new], sort_keys=True, default=str).encode()).hexdigest()[:12]}"379            importance = importance_for("PRICE_CHANGED", model["entity_type"] if model else None, old, new, tier=self.tier)380            await self.emit_event("PRICE_CHANGED", "price", summary, entity_id=mid, old_value=old, new_value=new, importance=importance, dedupe_key=dedupe,381                                  source_url=url, meta={"provider_id": pid, "provider": pname})382            for anomaly in check_price_movement({**old, "model_id": mid, "provider_id": pid}, {**new, "model_id": mid, "provider_id": pid}):383                from aiatlas.services.anomalies import record384385                anomaly.detail["price_id"] = price_id386                await record(self.conn, anomaly)387        else:388            new = {"input_per_mtok": p.input_per_mtok, "output_per_mtok": p.output_per_mtok}389            summary = f"{pname} lists {mname} at {_fmt_price(new)}"390            dedupe = f"PROVIDER_LISTED:{mid}:{pid}:{p.provider_model_id or ''}"391            await self.emit_event("PROVIDER_LISTED", "provider", summary, entity_id=mid, new_value=new, importance=1, dedupe_key=dedupe, source_url=url,392                                  meta={"provider_id": pid, "provider": pname})393        await self.write_relation(p.model, "available_through", p.provider, {"provider_model_id": p.provider_model_id}, source_url=url)394395    # ---------------------------------------------------------------------------------------------- benchmark results396    async def write_result(self, r: ResultObs) -> None:397        mid = await self.resolver.resolve(r.model)398        bid = await self.resolver.resolve(r.benchmark)399        if not mid or not bid:400            return401        url = r.source_url or self.source_url402        # evaluation-effort variants ("gpt-5-4-mini-medium") are a *configuration* of the canonical model: the effort dict lands in the403        # config so folded results stay distinguishable and comparable (reasoning_effort is a condition key, not a task key)404        config = effort_config(r.model.name, r.config)405        metric = r.metric406        dedupe = result_dedupe_key(mid, bid, config, metric)407        config_key = bench_ontology.config_key(config, metric)408        trust = r.trust_level or bench_ontology.trust_level(await self.source_key(), config, extractor=self.extractor)409        variant = r.variant or bench_ontology.variant_from_config(config)410        run_group = r.run_group or bench_ontology.run_group_from_config(config)411        conf = r.confidence or TIER_CONFIDENCE.get(self.tier, "medium")412        lo, hi = bench_ontology.metric_bounds(metric, r.unit)413        out_of_range = (hi is not None and r.score > hi + 1e-9) or (lo is not None and r.score < lo - 1e-9)414        if out_of_range:415            conf = "low"416        existing = await fetch_one(self.conn, "select id, score from benchmark_results where dedupe_key = :d", d=dedupe)417        names: dict[str, Any] | None = None418        if existing:419            if abs((existing["score"] or 0) - r.score) > 1e-9:420                await execute(self.conn, "update benchmark_results set valid_to = :o, is_current = false, dedupe_key = dedupe_key || ':' || :suffix where id = :id",421                              o=self.observed_at, suffix=new_id("result")[-10:], id=existing["id"])422                names = await self._names(mid, bid)423                await self.emit_event("BENCHMARK_UPDATED", "benchmark", f"{names['model']} on {names['bench']}: {existing['score']:g} → {r.score:g}",424                                      entity_id=mid, old_value=existing["score"], new_value=r.score, importance=1,425                                      dedupe_key=f"BENCHMARK_UPDATED:{dedupe}:{r.score:g}", source_url=url, meta={"benchmark_id": bid})426            else:427                await execute(self.conn, "update benchmark_results set observed_at = :o where id = :id", o=self.observed_at, id=existing["id"])428                return429        rid = new_id("result")430        await execute(self.conn, """insert into benchmark_results (id, model_id, benchmark_id, score, metric, unit, higher_is_better, config, evaluated_at, observed_at,431                                    source_id, snapshot_id, source_url, tier, confidence, dedupe_key, config_key, trust_level, variant, run_group, is_current, extractor, run_id)432                                    values (:id, :m, :b, :s, :metric, :unit, :hib, cast(:cfg as jsonb), :ev, :obs, :src, :snap, :url, :tier, :conf, :d, :ck, :trust, :variant, :rg, true, :ex, :run)""",433                      id=rid, m=mid, b=bid, s=r.score, metric=metric, unit=r.unit, hib=r.higher_is_better, cfg=jsonb(config), ev=r.evaluated_at, obs=self.observed_at,434                      src=self.source_id, snap=self.snapshot_id, url=url, tier=self.tier, conf=conf, d=dedupe, ck=config_key, trust=trust, variant=variant, rg=run_group,435                      ex=self.extractor, run=self.run_id)436        self.stats.results += 1437        # one current row per (model, benchmark, metric, config_key): a different (older) run group is closed so the leaderboard shows the latest run438        await execute(self.conn, """update benchmark_results set is_current = false, valid_to = coalesce(valid_to, :obs)439                                    where model_id = :m and benchmark_id = :b and coalesce(metric, '') = :metric and config_key = :ck and id <> :id and is_current440                                      and coalesce(run_group, '') <> :rg and observed_at <= :obs""",441                      obs=self.observed_at, m=mid, b=bid, metric=metric or "", ck=config_key, id=rid, rg=run_group or "")442        if out_of_range:443            from aiatlas.services.anomalies import record444445            names = names or await self._names(mid, bid)446            for anomaly in check_result({"id": rid, "model_id": mid, "benchmark_id": bid, "score": r.score, "metric": metric, "unit": r.unit,447                                         "model_name": names["model"], "benchmark_name": names["bench"]}):448                await record(self.conn, anomaly)449        if not existing:450            names = names or await self._names(mid, bid)451            await self.emit_event("BENCHMARK_RESULT", "benchmark", f"{names['model']} scores {r.score:g}{r.unit or ''} on {names['bench']}",452                                  entity_id=mid, new_value=r.score, importance=1, dedupe_key=f"BENCHMARK_RESULT:{dedupe}", source_url=url,453                                  meta={"benchmark_id": bid, "metric": metric, "config_key": config_key})454        await self.write_relation(r.model, "evaluated_on", r.benchmark, source_url=url)455456    async def _names(self, mid: str, bid: str) -> dict[str, Any]:457        model = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=mid)458        bench = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=bid)459        return {"model": model["canonical_name"] if model else mid, "bench": bench["canonical_name"] if bench else bid}460461    # ---------------------------------------------------------------------------------------------- events462    async def emit_event(self, event_type: str, category: str, summary: str, *, entity_id: str | None = None, old_value: Any = None,463                         new_value: Any = None, importance: int = 2, effective_at: datetime | None = None, dedupe_key: str | None = None,464                         source_url: str | None = None, meta: dict[str, Any] | None = None, observed_at: datetime | None = None,465                         entity_first_seen: datetime | None = None, is_backfill: bool | None = None) -> None:466        dedupe = dedupe_key or f"{event_type}:{entity_id or ''}:{normalize_alias(summary)[:120]}"467        observed = observed_at or self.observed_at468        if self.derived:469            # a derived writer re-encodes what we already know: its events are bookkeeping, never news470            is_backfill, importance = True, min(importance, 1)471        backfill = is_backfill if is_backfill is not None else classify_backfill(event_type, effective_at, observed, is_first_run=self.is_first_run,472                                                                                 entity_first_seen=entity_first_seen)473        group_key = group_key_for(event_type, entity_id, effective_at, observed)474        await execute(self.conn, """insert into change_events (id, entity_id, event_type, category, property, old_value, new_value, summary, importance, observed_at,475                                    effective_at, source_id, snapshot_id, source_url, connector_name, dedupe_key, meta, run_id, recorded_at, is_backfill, group_key)476                                    values (:id, :e, :t, :c, :p, cast(:o as jsonb), cast(:n as jsonb), :s, :imp, :obs, :eff, :src, :snap, :url, :conn, :d, cast(:meta as jsonb),477                                            :run, now(), :bf, :gk)478                                    on conflict (dedupe_key) do nothing""",479                      id=new_id("change_event"), e=entity_id, t=event_type, c=category, p=(meta or {}).get("property"), o=jsonb(old_value) if old_value is not None else None,480                      n=jsonb(new_value) if new_value is not None else None, s=summary[:500], imp=max(0, min(3, importance)), obs=observed,481                      eff=effective_at, src=self.source_id, snap=self.snapshot_id, url=source_url or self.source_url, conn=self.connector_name, d=dedupe[:400],482                      meta=jsonb(meta or {}), run=self.run_id, bf=backfill, gk=group_key)483        self.stats.events += 1484485    async def _new_entity_event(self, eid: str) -> None:486        row = await fetch_one(self.conn, "select canonical_name, entity_type, attributes, organization_id from entities where id = :id", id=eid)487        if not row:488            return489        etype = row["entity_type"]490        if etype in SILENT_NEW_TYPES:491            return492        org = None493        org_models = 0494        if row["organization_id"]:495            o = await fetch_one(self.conn, """select canonical_name, (select count(*) from entities m where m.organization_id = e.id and m.entity_type = 'model'496                                              and m.merged_into is null) as models from entities e where e.id = :id""", id=row["organization_id"])497            org = o["canonical_name"] if o else None498            org_models = int(o["models"] or 0) if o else 0499        attrs = row["attributes"] or {}500        if etype == "model":501            importance = importance_for("NEW_MODEL", etype, tier=self.tier, org_model_count=org_models, openness=attrs.get("openness"))502        else:503            importance = importance_for(f"NEW_{etype.upper()}", etype, tier=self.tier, default=NEW_IMPORTANCE.get(etype, 1))504        label = etype.replace("_", " ")505        summary = f"New {label}: {row['canonical_name']}" + (f" ({org})" if org else "")506        effective = None507        rel = attrs.get("release_date") or attrs.get("published_at")508        if isinstance(rel, str):509            from aiatlas.sdk.extract.dates import parse_datetime510511            effective = parse_datetime(rel)512        await self.emit_event(f"NEW_{etype.upper()}", EVENT_CATEGORY_BY_TYPE.get(etype, "update"), summary, entity_id=eid, importance=importance,513                              dedupe_key=f"NEW_{etype.upper()}:{eid}", effective_at=effective, entity_first_seen=effective)514515516def _short(v: Any) -> str:517    s = json.dumps(v, default=str, ensure_ascii=False) if not isinstance(v, str) else v518    return s if len(s) <= 60 else s[:57] + "…"519520521def _fmt_price(p: dict[str, Any]) -> str:522    i, o = p.get("input_per_mtok"), p.get("output_per_mtok")523    parts = []524    if i is not None:525        parts.append(f"${i:g} in")526    if o is not None:527        parts.append(f"${o:g} out")528    return " / ".join(parts) + " per 1M tokens" if parts else "n/a"529530531__all__ = ["FactWriter", "WriteStats", "result_dedupe_key"]532