HTML 77.2%
TypeScript 10.5%
Python 9.6%
JavaScript 2.5%
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