"""FactWriter — turns `Facts` into versioned database state. Temporal rules (per entity × property): * no current claim → insert current claim, set attribute, NEW_* event for new entities * same value → confirm (observed_at bumped), nothing else * same assertion, older encoding (taxonomy) → re-encode the current claim in place (canonical value + value_raw), never an event * different, source ≥ tier → supersede (valid_to = observed), insert new current, update attribute, CHANGE event for material props * different, same source AND same extractor → supersede (a source correcting itself); an LLM claim never supersedes a deterministic claim from the same URL — it is stored as conflicting * different, worse source → store as `conflicting`, flag confidence, review-queue item — never overwrite Prices and benchmark results have their own append-only tables with the same close/open semantics; results additionally keep ONE current row per (model, benchmark, metric, config_key): a newer run group closes the older rows so leaderboards show the latest run. Every fact carries `run_id` (batch rollback) and every event carries `recorded_at`, `is_backfill` and `group_key`. """ from __future__ import annotations import hashlib import json import logging from datetime import UTC, datetime from typing import Any from sqlalchemy.ext.asyncio import AsyncConnection from aiatlas.db import execute, fetch_one, jsonb from aiatlas.ids import new_id, normalize_alias from aiatlas.ontology import benchmarks as bench_ontology from aiatlas.ontology.anomalies import check_price_movement, check_result from aiatlas.ontology.models import effort_config from aiatlas.ontology.taxonomy import TAXONOMY_PROPERTIES, normalize_property from aiatlas.sdk.facts import ( EVENT_CATEGORY_BY_TYPE, MATERIAL_PROPERTIES, NOISY_PREFIXES, EntityRef, Facts, PriceObs, ResultObs, ) from aiatlas.sdk.resolution import Resolver from aiatlas.services.events import classify_backfill, group_key_for, importance_for log = logging.getLogger(__name__) TIER_CONFIDENCE = {1: "high", 2: "medium", 3: "low", 4: "low"} SOFT_PROPERTIES = {"description", "summary", "tagline", "availability_note", "abstract", "notes", "training_data_notes", "safety_notes", "hardware_requirements"} # derived / secondary encodings of another property: never their own change event NO_EVENT_PROPERTIES = {"license_key"} NEW_IMPORTANCE = {"model": 3, "company": 2, "provider": 2, "paper": 1, "dataset": 1, "benchmark": 2, "framework": 1, "hardware": 2, "tool": 1, "repository": 0, "release": 2, "regulation": 2, "incident": 2, "organization": 1, "researcher": 0, "artifact": 0, "model_family": 1, "license": 0} SILENT_NEW_TYPES = {"researcher", "country", "license", "quantization"} def _norm_value(v: Any) -> Any: if isinstance(v, datetime): return v.astimezone(UTC).isoformat(timespec="seconds") if isinstance(v, dict): return {k: _norm_value(x) for k, x in v.items()} if isinstance(v, (set, tuple, list)): items = [_norm_value(x) for x in v] return sorted(items) if all(isinstance(x, (str, int, float)) and not isinstance(x, bool) for x in items) else items if isinstance(v, float) and v.is_integer() and abs(v) < 1e15: return int(v) return v def _same(a: Any, b: Any) -> bool: return json.dumps(_norm_value(a), sort_keys=True, default=str) == json.dumps(_norm_value(b), sort_keys=True, default=str) def _cfg_hash(config: dict[str, Any] | None) -> str: return hashlib.sha1(json.dumps(config or {}, sort_keys=True, default=str).encode()).hexdigest()[:12] def result_dedupe_key(model_id: str, benchmark_id: str, config: dict[str, Any] | None, metric: str | None) -> str: return f"{model_id}:{benchmark_id}:{_cfg_hash(config)}:{metric or ''}" class WriteStats: def __init__(self) -> None: self.entities_created = 0 self.entities_updated = 0 self.claims = 0 self.relations = 0 self.events = 0 self.prices = 0 self.results = 0 self.conflicts = 0 self.folded_variants = 0 def as_dict(self) -> dict[str, int]: return dict(self.__dict__) class FactWriter: def __init__(self, conn: AsyncConnection, *, source_id: str | None, snapshot_id: str | None, source_url: str | None, tier: int = 2, connector_name: str | None = None, extractor: str = "deterministic", extractor_version: str = "1", observed_at: datetime | None = None, run_id: str | None = None, source_key: str | None = None, is_first_run: bool = False, variant_index: dict[str, str] | None = None): self.conn = conn self.source_id = source_id self.snapshot_id = snapshot_id self.source_url = source_url self.tier = tier self.connector_name = connector_name self.extractor = extractor self.extractor_version = extractor_version self.observed_at = observed_at or datetime.now(UTC) self.run_id = run_id self._source_key = source_key self.is_first_run = is_first_run # connector's first successful run → every event is backfill (initial corpus) self.resolver = Resolver(conn, snapshot_id=snapshot_id, source_tier=tier, variant_index=variant_index) self.stats = WriteStats() @property def derived(self) -> bool: """Derived writers (canonicalization) store disagreements as conflicting claims but never open review items.""" return self.extractor == "derived" async def source_key(self) -> str | None: if self._source_key is None and self.source_id: row = await fetch_one(self.conn, "select key from sources where id = :id", id=self.source_id) self._source_key = row["key"] if row else "" return self._source_key or None # ---------------------------------------------------------------------------------------------- entry point async def write(self, facts: Facts) -> WriteStats: new_before = len(self.resolver.created) for ref in facts.entities: await self.resolver.resolve(ref) await self._apply_hierarchy(ref) for prop, value in ref.attributes.items(): await self.write_claim(ref, prop, value) for c in facts.claims: await self.write_claim(c.entity, c.property, c.value, unit=c.unit, confidence=c.confidence, observed_at=c.observed_at, effective_at=c.effective_at, source_url=c.source_url) for r in facts.relations: await self.write_relation(r.subject, r.predicate, r.object, r.attributes, confidence=r.confidence, source_url=r.source_url) for p in facts.prices: await self.write_price(p) for res in facts.results: await self.write_result(res) for e in facts.events: eid = await self.resolver.resolve(e.entity) if e.entity else None await self.emit_event(e.event_type, e.category, e.summary, entity_id=eid, old_value=e.old_value, new_value=e.new_value, importance=e.importance, effective_at=e.effective_at, dedupe_key=e.dedupe_key, source_url=e.source_url, meta=e.meta) # NEW_* events for every entity created in this write for eid in self.resolver.created[new_before:]: await self._new_entity_event(eid) self.stats.entities_created = len(self.resolver.created) self.stats.entities_updated = len(self.resolver.updated - set(self.resolver.created)) self.stats.folded_variants = len(self.resolver.folded) return self.stats # ---------------------------------------------------------------------------------------------- hierarchy hints async def _apply_hierarchy(self, ref: EntityRef) -> None: """Materialise `family` / `canonical` hints: entities.family_id + member_of_family, entities.canonical_id + artifact_of, identity_confidence.""" if not ref.id: return if ref.family is not None: if ref.family.entity_type != "model_family": ref.family.entity_type = "model_family" fid = await self.resolver.resolve(ref.family) if fid and fid != ref.id: 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) await self.write_relation(ref, "member_of_family", ref.family) if ref.canonical is not None: cid = await self.resolver.resolve(ref.canonical) if cid and cid != ref.id: await execute(self.conn, """update entities set canonical_id = :c, artifact_kind = coalesce(cast(:k as text), artifact_kind) where id = :e and (canonical_id is distinct from :c or (cast(:k as text) is not null and artifact_kind is null))""", c=cid, k=ref.artifact_kind, e=ref.id) await self.write_relation(ref, "artifact_of", ref.canonical, {"artifact_kind": ref.artifact_kind} if ref.artifact_kind else None) if ref.identity_confidence: await execute(self.conn, "update entities set identity_confidence = :ic where id = :e and identity_confidence <> :ic", ic=ref.identity_confidence, e=ref.id) # ---------------------------------------------------------------------------------------------- claims async def write_claim(self, ref: EntityRef, prop: str, value: Any, *, unit: str | None = None, confidence: str | None = None, observed_at: datetime | None = None, effective_at: datetime | None = None, source_url: str | None = None) -> None: if value is None or value == "" or value == [] or value == {}: return eid = await self.resolver.resolve(ref) assert eid value = _norm_value(value) value_raw: str | None = None if prop in TAXONOMY_PROPERTIES: value, value_raw, mappings = normalize_property(ref.entity_type, prop, value) for domain, raw, canon in mappings: await self._record_mapping(domain, raw, canon) value = _norm_value(value) observed = observed_at or self.observed_at conf = confidence or TIER_CONFIDENCE.get(self.tier, "medium") url = source_url or self.source_url noisy = prop.startswith(NOISY_PREFIXES) current = await fetch_one(self.conn, """select id, value, value_raw, tier, source_url, observed_at, extractor from claims where entity_id = :e and property = :p and status = 'current' order by valid_from desc limit 1""", e=eid, p=prop) is_new_entity = eid in self.resolver.created await self._apply_claim(ref, eid, prop, value, value_raw, current, unit=unit, conf=conf, observed=observed, effective_at=effective_at, url=url, noisy=noisy, is_new_entity=is_new_entity) if prop == "license" and isinstance(value, str): from aiatlas.ontology.licenses import LICENSES if value in LICENSES: # canonical licence key as its own property (filterable), never its own event await self.write_claim(ref, "license_key", value, confidence=confidence, observed_at=observed_at, effective_at=effective_at, source_url=source_url) 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, conf: str, observed: datetime, effective_at: datetime | None, url: str | None, noisy: bool, is_new_entity: bool) -> None: if current is None: await self._insert_claim(eid, prop, value, unit, conf, "current", observed, effective_at, url, value_raw) await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw) return if _same(current["value"], value): await execute(self.conn, "update claims set observed_at = greatest(observed_at, :o) where id = :id", o=observed, id=current["id"]) if noisy: await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw) return if prop in TAXONOMY_PROPERTIES and _same(normalize_property(ref.entity_type, prop, current["value"])[0], value): # same assertion in an older encoding ("apache-2.0" → "Apache-2.0"): re-encode in place, keep the source label, no event raw_keep = current["value_raw"] or (current["value"] if isinstance(current["value"], str) else json.dumps(current["value"], ensure_ascii=False)) raw_keep = raw_keep if raw_keep != value else None 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""", v=jsonb(value), vt=value[:2000] if isinstance(value, str) else None, raw=raw_keep, o=observed, id=current["id"]) await self._set_attribute(eid, prop, value, unit, conf, url, observed, raw_keep, keep_provenance=True) return same_source = bool(url and current["source_url"] == url) same_extractor = (current["extractor"] or "deterministic") == self.extractor if prop in SOFT_PROPERTIES and not same_source: # soft text (descriptions, notes): first statement wins until *its own* source changes; never an event, never a conflict return if self.tier <= (current["tier"] or 2) or (same_source and same_extractor): # supersede await execute(self.conn, "update claims set status = 'superseded', valid_to = :o where id = :id", o=observed, id=current["id"]) await self._insert_claim(eid, prop, value, unit, conf, "current", observed, effective_at, url, value_raw) await self._set_attribute(eid, prop, value, unit, conf, url, observed, value_raw) if not noisy and not is_new_entity and prop not in SOFT_PROPERTIES and prop not in NO_EVENT_PROPERTIES: await self._property_change_event(eid, prop, current["value"], value, observed, effective_at, url) return # the same disagreement from the same source/extractor is recorded once (re-observed, not re-inserted) 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) and extractor = :ex and coalesce(source_url, '') = :url limit 1""", e=eid, p=prop, v=jsonb(value), ex=self.extractor, url=url or "") if dup: await execute(self.conn, "update claims set observed_at = greatest(observed_at, :o) where id = :id", o=observed, id=dup["id"]) return await self._insert_claim(eid, prop, value, unit, conf, "conflicting", observed, effective_at, url, value_raw) await execute(self.conn, "update claims set confidence = 'conflicted' where id = :id", id=current["id"]) await execute(self.conn, """update entities set quality = quality || jsonb_build_object('conflicts', coalesce((quality->>'conflicts')::int, 0) + 1) where id = :id""", id=eid) self.stats.conflicts += 1 if not self.derived: dedupe = f"conflict:{eid}:{prop}:{hashlib.sha1(json.dumps(value, sort_keys=True, default=str).encode()).hexdigest()[:10]}" await execute(self.conn, """insert into review_queue (id, kind, entity_ids, payload, reason, dedupe_key) values (:id, 'conflict', :ids, cast(:p as jsonb), :r, :d) on conflict (dedupe_key) do nothing""", id=new_id("review"), ids=[eid], p=jsonb({"property": prop, "current": current["value"], "current_source": current["source_url"], "claimed": value, "claimed_source": url, "tier": self.tier, "extractor": self.extractor}), r=f"source disagrees on {prop}", d=dedupe) async def _record_mapping(self, domain: str, raw: str, canonical: str | None) -> None: if canonical is not None and raw == canonical: return # identity mappings carry no information await execute(self.conn, """insert into taxonomy_mappings (domain, raw, canonical) values (:d, :r, :c) on conflict (domain, raw) do update set count = taxonomy_mappings.count + 1, last_seen_at = now(), canonical = coalesce(excluded.canonical, taxonomy_mappings.canonical)""", d=domain, r=raw[:300], c=canonical) async def _insert_claim(self, eid: str, prop: str, value: Any, unit: str | None, conf: str, status: str, observed: datetime, effective_at: datetime | None, url: str | None, value_raw: str | None = None) -> None: text_val = value if isinstance(value, str) else None num_val = float(value) if isinstance(value, (int, float)) and not isinstance(value, bool) else None await execute(self.conn, """insert into claims (id, entity_id, property, value, value_text, value_num, unit, source_id, snapshot_id, source_url, tier, confidence, status, extractor, extractor_version, observed_at, effective_at, valid_from, run_id, value_raw) values (:id, :e, :p, cast(:v as jsonb), :vt, :vn, :u, :src, :snap, :url, :tier, :conf, :status, :ex, :exv, :o, :eff, :vf, :run, :raw)""", 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, snap=self.snapshot_id, url=url, tier=self.tier, conf=conf, status=status, ex=self.extractor, exv=self.extractor_version, o=observed, eff=effective_at, vf=effective_at or observed, run=self.run_id, raw=value_raw[:500] if value_raw else None) self.stats.claims += 1 async def _set_attribute(self, eid: str, prop: str, value: Any, unit: str | None, conf: str, url: str | None, observed: datetime, value_raw: str | None = None, *, keep_provenance: bool = False) -> None: prov = {"source_id": self.source_id, "snapshot_id": self.snapshot_id, "url": url, "observed_at": observed.isoformat(timespec="seconds"), "tier": self.tier, "confidence": conf, "extractor": self.extractor} if unit: prov["unit"] = unit attrs = {prop: value} if value_raw is not None: attrs[f"{prop}_raw"] = value_raw if keep_provenance: 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""", praw=f"{prop}_raw", a=jsonb(attrs), o=observed, id=eid) else: await execute(self.conn, """update entities set attributes = (attributes - cast(:praw as text)) || cast(:a as jsonb), provenance = provenance || jsonb_build_object(cast(:p as text), cast(:prov as jsonb)), last_seen_at = greatest(last_seen_at, :o) where id = :id""", praw=f"{prop}_raw", a=jsonb(attrs), p=prop, prov=jsonb(prov), o=observed, id=eid) if prop == "status" and isinstance(value, str): await execute(self.conn, "update entities set status = :s where id = :id", s=value[:40], id=eid) if prop == "description" and isinstance(value, str): await execute(self.conn, "update entities set description = :d where id = :id", d=value[:4000], id=eid) async def _property_change_event(self, eid: str, prop: str, old: Any, new: Any, observed: datetime, effective_at: datetime | None, url: str | None) -> None: event_type, importance = MATERIAL_PROPERTIES.get(prop, ("PROPERTY_CHANGED", 0)) row = await fetch_one(self.conn, "select canonical_name, entity_type from entities where id = :id", id=eid) name = row["canonical_name"] if row else eid etype = row["entity_type"] if row else "" category = EVENT_CATEGORY_BY_TYPE.get(etype, "update") importance = importance_for(event_type, etype, old, new, tier=self.tier, default=importance) summary = f"{name}: {prop.replace('_', ' ')} changed from {_short(old)} to {_short(new)}" dedupe = f"{event_type}:{eid}:{prop}:{hashlib.sha1(json.dumps([old, new], sort_keys=True, default=str).encode()).hexdigest()[:12]}" await self.emit_event(event_type, category, summary, entity_id=eid, old_value=old, new_value=new, importance=importance, effective_at=effective_at, dedupe_key=dedupe, source_url=url, meta={"property": prop}, observed_at=observed) # ---------------------------------------------------------------------------------------------- relations async def write_relation(self, subject: EntityRef, predicate: str, obj: EntityRef, attributes: dict[str, Any] | None = None, *, confidence: str | None = None, source_url: str | None = None) -> None: sid = await self.resolver.resolve(subject) oid = await self.resolver.resolve(obj) if not sid or not oid or sid == oid: return if predicate == "develops": # an organisation develops a *model*; a hub repository that is an artifact (quantisation, conversion) is merely published by its org otype = await fetch_one(self.conn, "select entity_type from entities where id = :id", id=oid) if otype and otype["entity_type"] == "artifact": predicate = "published_by" sid, oid = oid, sid conf = confidence or TIER_CONFIDENCE.get(self.tier, "medium") 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", s=sid, p=predicate, o=oid) if existing: if attributes and not _same(existing["attributes"], {**existing["attributes"], **attributes}): await execute(self.conn, "update relations set attributes = attributes || cast(:a as jsonb), observed_at = :o where id = :id", a=jsonb(attributes), o=self.observed_at, id=existing["id"]) else: await execute(self.conn, "update relations set observed_at = :o where id = :id", o=self.observed_at, id=existing["id"]) return 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) values (:id, :s, :p, :o, cast(:a as jsonb), :src, :snap, :url, :tier, :conf, :obs, :obs, :run)""", id=new_id("relation"), s=sid, p=predicate, o=oid, a=jsonb(attributes or {}), src=self.source_id, snap=self.snapshot_id, url=source_url or self.source_url, tier=self.tier, conf=conf, obs=self.observed_at, run=self.run_id) self.stats.relations += 1 # ---------------------------------------------------------------------------------------------- prices async def write_price(self, p: PriceObs) -> None: mid = await self.resolver.resolve(p.model) pid = await self.resolver.resolve(p.provider) if not mid or not pid: return if all(v is None for v in p.price_tuple()[:-1]): return url = p.source_url or self.source_url 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""", m=mid, p=pid, pm=p.provider_model_id or "") new_tuple = p.price_tuple() if current: cur_tuple = (current["input_per_mtok"], current["output_per_mtok"], current["cached_input_per_mtok"], current["cache_write_per_mtok"], current["batch_input_per_mtok"], current["batch_output_per_mtok"], current["per_image"], current["per_request"], current["currency"]) if _same(cur_tuple, new_tuple) and (p.context_length in (None, current["context_length"])): await execute(self.conn, "update prices set observed_at = :o where id = :id", o=self.observed_at, id=current["id"]) return await execute(self.conn, "update prices set valid_to = :o where id = :id", o=self.observed_at, id=current["id"]) price_id = new_id("price") 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, batch_input_per_mtok, batch_output_per_mtok, per_image, per_request, currency, context_length, max_output_tokens, features, observed_at, valid_from, source_id, snapshot_id, source_url, tier, meta, run_id) 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)""", 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, 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, 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, url=url, tier=self.tier, meta=jsonb(p.meta), run=self.run_id) self.stats.prices += 1 model = await fetch_one(self.conn, "select canonical_name, entity_type from entities where id = :id", id=mid) provider = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=pid) mname = model["canonical_name"] if model else mid pname = provider["canonical_name"] if provider else pid if current: old = {"input_per_mtok": current["input_per_mtok"], "output_per_mtok": current["output_per_mtok"]} new = {"input_per_mtok": p.input_per_mtok, "output_per_mtok": p.output_per_mtok} summary = f"{pname} changed pricing for {mname}: {_fmt_price(old)} → {_fmt_price(new)}" 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]}" importance = importance_for("PRICE_CHANGED", model["entity_type"] if model else None, old, new, tier=self.tier) await self.emit_event("PRICE_CHANGED", "price", summary, entity_id=mid, old_value=old, new_value=new, importance=importance, dedupe_key=dedupe, source_url=url, meta={"provider_id": pid, "provider": pname}) for anomaly in check_price_movement({**old, "model_id": mid, "provider_id": pid}, {**new, "model_id": mid, "provider_id": pid}): from aiatlas.services.anomalies import record anomaly.detail["price_id"] = price_id await record(self.conn, anomaly) else: new = {"input_per_mtok": p.input_per_mtok, "output_per_mtok": p.output_per_mtok} summary = f"{pname} lists {mname} at {_fmt_price(new)}" dedupe = f"PROVIDER_LISTED:{mid}:{pid}:{p.provider_model_id or ''}" await self.emit_event("PROVIDER_LISTED", "provider", summary, entity_id=mid, new_value=new, importance=1, dedupe_key=dedupe, source_url=url, meta={"provider_id": pid, "provider": pname}) await self.write_relation(p.model, "available_through", p.provider, {"provider_model_id": p.provider_model_id}, source_url=url) # ---------------------------------------------------------------------------------------------- benchmark results async def write_result(self, r: ResultObs) -> None: mid = await self.resolver.resolve(r.model) bid = await self.resolver.resolve(r.benchmark) if not mid or not bid: return url = r.source_url or self.source_url # evaluation-effort variants ("gpt-5-4-mini-medium") are a *configuration* of the canonical model: the effort dict lands in the # config so folded results stay distinguishable and comparable (reasoning_effort is a condition key, not a task key) config = effort_config(r.model.name, r.config) metric = r.metric dedupe = result_dedupe_key(mid, bid, config, metric) config_key = bench_ontology.config_key(config, metric) trust = r.trust_level or bench_ontology.trust_level(await self.source_key(), config, extractor=self.extractor) variant = r.variant or bench_ontology.variant_from_config(config) run_group = r.run_group or bench_ontology.run_group_from_config(config) conf = r.confidence or TIER_CONFIDENCE.get(self.tier, "medium") lo, hi = bench_ontology.metric_bounds(metric, r.unit) out_of_range = (hi is not None and r.score > hi + 1e-9) or (lo is not None and r.score < lo - 1e-9) if out_of_range: conf = "low" existing = await fetch_one(self.conn, "select id, score from benchmark_results where dedupe_key = :d", d=dedupe) names: dict[str, Any] | None = None if existing: if abs((existing["score"] or 0) - r.score) > 1e-9: await execute(self.conn, "update benchmark_results set valid_to = :o, is_current = false, dedupe_key = dedupe_key || ':' || :suffix where id = :id", o=self.observed_at, suffix=new_id("result")[-10:], id=existing["id"]) names = await self._names(mid, bid) await self.emit_event("BENCHMARK_UPDATED", "benchmark", f"{names['model']} on {names['bench']}: {existing['score']:g} → {r.score:g}", entity_id=mid, old_value=existing["score"], new_value=r.score, importance=1, dedupe_key=f"BENCHMARK_UPDATED:{dedupe}:{r.score:g}", source_url=url, meta={"benchmark_id": bid}) else: await execute(self.conn, "update benchmark_results set observed_at = :o where id = :id", o=self.observed_at, id=existing["id"]) return rid = new_id("result") await execute(self.conn, """insert into benchmark_results (id, model_id, benchmark_id, score, metric, unit, higher_is_better, config, evaluated_at, observed_at, source_id, snapshot_id, source_url, tier, confidence, dedupe_key, config_key, trust_level, variant, run_group, is_current, extractor, run_id) 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)""", 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, 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, ex=self.extractor, run=self.run_id) self.stats.results += 1 # one current row per (model, benchmark, metric, config_key): a different (older) run group is closed so the leaderboard shows the latest run await execute(self.conn, """update benchmark_results set is_current = false, valid_to = coalesce(valid_to, :obs) where model_id = :m and benchmark_id = :b and coalesce(metric, '') = :metric and config_key = :ck and id <> :id and is_current and coalesce(run_group, '') <> :rg and observed_at <= :obs""", obs=self.observed_at, m=mid, b=bid, metric=metric or "", ck=config_key, id=rid, rg=run_group or "") if out_of_range: from aiatlas.services.anomalies import record names = names or await self._names(mid, bid) for anomaly in check_result({"id": rid, "model_id": mid, "benchmark_id": bid, "score": r.score, "metric": metric, "unit": r.unit, "model_name": names["model"], "benchmark_name": names["bench"]}): await record(self.conn, anomaly) if not existing: names = names or await self._names(mid, bid) await self.emit_event("BENCHMARK_RESULT", "benchmark", f"{names['model']} scores {r.score:g}{r.unit or ''} on {names['bench']}", entity_id=mid, new_value=r.score, importance=1, dedupe_key=f"BENCHMARK_RESULT:{dedupe}", source_url=url, meta={"benchmark_id": bid, "metric": metric, "config_key": config_key}) await self.write_relation(r.model, "evaluated_on", r.benchmark, source_url=url) async def _names(self, mid: str, bid: str) -> dict[str, Any]: model = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=mid) bench = await fetch_one(self.conn, "select canonical_name from entities where id = :id", id=bid) return {"model": model["canonical_name"] if model else mid, "bench": bench["canonical_name"] if bench else bid} # ---------------------------------------------------------------------------------------------- events async def emit_event(self, event_type: str, category: str, summary: str, *, entity_id: str | None = None, old_value: Any = None, new_value: Any = None, importance: int = 2, effective_at: datetime | None = None, dedupe_key: str | None = None, source_url: str | None = None, meta: dict[str, Any] | None = None, observed_at: datetime | None = None, entity_first_seen: datetime | None = None, is_backfill: bool | None = None) -> None: dedupe = dedupe_key or f"{event_type}:{entity_id or ''}:{normalize_alias(summary)[:120]}" observed = observed_at or self.observed_at if self.derived: # a derived writer re-encodes what we already know: its events are bookkeeping, never news is_backfill, importance = True, min(importance, 1) backfill = is_backfill if is_backfill is not None else classify_backfill(event_type, effective_at, observed, is_first_run=self.is_first_run, entity_first_seen=entity_first_seen) group_key = group_key_for(event_type, entity_id, effective_at, observed) await execute(self.conn, """insert into change_events (id, entity_id, event_type, category, property, old_value, new_value, summary, importance, observed_at, effective_at, source_id, snapshot_id, source_url, connector_name, dedupe_key, meta, run_id, recorded_at, is_backfill, group_key) 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), :run, now(), :bf, :gk) on conflict (dedupe_key) do nothing""", 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, n=jsonb(new_value) if new_value is not None else None, s=summary[:500], imp=max(0, min(3, importance)), obs=observed, 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], meta=jsonb(meta or {}), run=self.run_id, bf=backfill, gk=group_key) self.stats.events += 1 async def _new_entity_event(self, eid: str) -> None: row = await fetch_one(self.conn, "select canonical_name, entity_type, attributes, organization_id from entities where id = :id", id=eid) if not row: return etype = row["entity_type"] if etype in SILENT_NEW_TYPES: return org = None org_models = 0 if row["organization_id"]: 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' and m.merged_into is null) as models from entities e where e.id = :id""", id=row["organization_id"]) org = o["canonical_name"] if o else None org_models = int(o["models"] or 0) if o else 0 attrs = row["attributes"] or {} if etype == "model": importance = importance_for("NEW_MODEL", etype, tier=self.tier, org_model_count=org_models, openness=attrs.get("openness")) else: importance = importance_for(f"NEW_{etype.upper()}", etype, tier=self.tier, default=NEW_IMPORTANCE.get(etype, 1)) label = etype.replace("_", " ") summary = f"New {label}: {row['canonical_name']}" + (f" ({org})" if org else "") effective = None rel = attrs.get("release_date") or attrs.get("published_at") if isinstance(rel, str): from aiatlas.sdk.extract.dates import parse_datetime effective = parse_datetime(rel) await self.emit_event(f"NEW_{etype.upper()}", EVENT_CATEGORY_BY_TYPE.get(etype, "update"), summary, entity_id=eid, importance=importance, dedupe_key=f"NEW_{etype.upper()}:{eid}", effective_at=effective, entity_first_seen=effective) def _short(v: Any) -> str: s = json.dumps(v, default=str, ensure_ascii=False) if not isinstance(v, str) else v return s if len(s) <= 60 else s[:57] + "…" def _fmt_price(p: dict[str, Any]) -> str: i, o = p.get("input_per_mtok"), p.get("output_per_mtok") parts = [] if i is not None: parts.append(f"${i:g} in") if o is not None: parts.append(f"${o:g} out") return " / ".join(parts) + " per 1M tokens" if parts else "n/a" __all__ = ["FactWriter", "WriteStats", "result_dedupe_key"]