From 8008752375e284312848d2608b600b819f82ac1f Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Fri, 18 Sep 2026 23:59:59 +0000 Subject: [PATCH] V2 P0: Evidence Graph writer, authority policy, and explain API (#69). Add EvidenceGraphWriter (SUPPORTS/CONTRADICTS/ASSERTS), configurable authority scores, EvidenceExplainer, and GET /evidence/explain/{claim_id}. Co-authored-by: Abhinaysai Kamineni --- .env.example | 2 + api/evidence.py | 48 +++++++ api/main.py | 2 + evidence/__init__.py | 5 +- evidence/authority.py | 54 ++++++++ evidence/explain.py | 81 ++++++++++++ evidence/graph.py | 190 +++++++++++++++++++++++++++ tests/api/test_evidence.py | 42 ++++++ tests/evidence/test_graph_explain.py | 107 +++++++++++++++ 9 files changed, 530 insertions(+), 1 deletion(-) create mode 100644 api/evidence.py create mode 100644 evidence/authority.py create mode 100644 evidence/explain.py create mode 100644 evidence/graph.py create mode 100644 tests/api/test_evidence.py create mode 100644 tests/evidence/test_graph_explain.py diff --git a/.env.example b/.env.example index 2c779e4..8f07f5c 100644 --- a/.env.example +++ b/.env.example @@ -125,6 +125,8 @@ CORTEX_CONTRADICTION_KAFKA=true CORTEX_SEMANTIC_ENABLED=false # V2 Evidence Graph writers (Claim/Evidence). Schema V010 applies regardless; writes gated here. CORTEX_EVIDENCE_GRAPH=false +# Optional JSON overrides for authority class scores (merged onto defaults). +# CORTEX_AUTHORITY_SCORES_JSON={"casual_chat":0.4,"merged_PR":0.85} # Decay worker interval when Compose profile `api` is up (seconds). Default 24h. CORTEX_DECAY_INTERVAL_SECONDS=86400 diff --git a/api/evidence.py b/api/evidence.py new file mode 100644 index 0000000..dd8aba3 --- /dev/null +++ b/api/evidence.py @@ -0,0 +1,48 @@ +"""Evidence Graph HTTP routes — explain Claim confidence.""" + +from __future__ import annotations + +from typing import Any + +from fastapi import APIRouter, HTTPException, status +from pydantic import BaseModel, Field + +from api.deps import RolesDep +from evidence.explain import EvidenceExplainer + +router = APIRouter(prefix="/evidence", tags=["evidence"]) + + +class ExplainResponse(BaseModel): + claim: str + claim_id: str + confidence: float | None = None + status: str | None = None + assertion_type: str | None = None + temporal: dict[str, Any] = Field(default_factory=dict) + supporting_evidence: list[dict[str, Any]] = Field(default_factory=list) + conflicting_evidence: list[dict[str, Any]] = Field(default_factory=list) + + +@router.get( + "/explain/{claim_id}", + response_model=ExplainResponse, + summary="Explain a claim via supporting and conflicting evidence", +) +def explain_claim( + claim_id: str, + workspace_id: str, + _roles: RolesDep, +) -> ExplainResponse: + """Return evidence-backed explanation for a Claim node.""" + explainer = EvidenceExplainer() + try: + payload = explainer.explain(claim_id=claim_id, workspace_id=workspace_id) + finally: + explainer.close() + if payload is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail=f"Claim {claim_id} not found in workspace {workspace_id}", + ) + return ExplainResponse(**payload) diff --git a/api/main.py b/api/main.py index 87c7358..dabdc8d 100644 --- a/api/main.py +++ b/api/main.py @@ -26,6 +26,7 @@ from slowapi.errors import RateLimitExceeded from api.contradictions import router as contradictions_router +from api.evidence import router as evidence_router from api.gdpr import router as gdpr_router from api.metrics import record_http_request, record_query, render_metrics from api.decisions import router as decisions_router @@ -105,6 +106,7 @@ async def lifespan(app: FastAPI) -> AsyncIterator[None]: app.include_router(gdpr_router) app.include_router(decisions_router) app.include_router(remember_router) +app.include_router(evidence_router) @app.middleware("http") diff --git a/evidence/__init__.py b/evidence/__init__.py index 235ceef..7121c9f 100644 --- a/evidence/__init__.py +++ b/evidence/__init__.py @@ -1,4 +1,4 @@ -"""Evidence graph package — Claim/Evidence adapters and (later) explain APIs.""" +"""Evidence graph package — adapters, authority, writer, explain.""" from evidence.adapters import ( claim_evidence_bundle, @@ -6,10 +6,13 @@ decision_to_evidence, evidence_graph_enabled, ) +from evidence.authority import compare_authority, score_for_class __all__ = [ "claim_evidence_bundle", + "compare_authority", "decision_to_claim", "decision_to_evidence", "evidence_graph_enabled", + "score_for_class", ] diff --git a/evidence/authority.py b/evidence/authority.py new file mode 100644 index 0000000..38edca7 --- /dev/null +++ b/evidence/authority.py @@ -0,0 +1,54 @@ +"""Configurable source-authority policy for Evidence Graph.""" + +from __future__ import annotations + +import json +import os +from typing import Any + +from shared.v2_models import AuthorityClass + +# Default ordering: higher = more authoritative (do not hardcode without tests). +DEFAULT_AUTHORITY_SCORES: dict[str, float] = { + "formal_policy": 0.98, + "approved_ADR": 0.95, + "production_state": 0.90, + "merged_PR": 0.82, + "approved_ticket": 0.75, + "incident_record": 0.78, + "verified_human_statement": 0.65, + "casual_chat": 0.45, + "agent_inference": 0.35, + "external_unverified": 0.25, +} + + +def load_authority_scores() -> dict[str, float]: + """Load authority scores from CORTEX_AUTHORITY_SCORES_JSON or defaults.""" + raw = os.environ.get("CORTEX_AUTHORITY_SCORES_JSON", "").strip() + if not raw: + return dict(DEFAULT_AUTHORITY_SCORES) + try: + parsed: dict[str, Any] = json.loads(raw) + out = dict(DEFAULT_AUTHORITY_SCORES) + for key, value in parsed.items(): + out[str(key)] = float(value) + return out + except (json.JSONDecodeError, TypeError, ValueError): + return dict(DEFAULT_AUTHORITY_SCORES) + + +def score_for_class(authority_class: AuthorityClass | str) -> float: + """Return configured authority score for a class.""" + scores = load_authority_scores() + return float(scores.get(str(authority_class), 0.30)) + + +def compare_authority(a: str, b: str) -> int: + """Return 1 if a > b, -1 if a < b, 0 if equal.""" + sa, sb = score_for_class(a), score_for_class(b) + if sa > sb: + return 1 + if sa < sb: + return -1 + return 0 diff --git a/evidence/explain.py b/evidence/explain.py new file mode 100644 index 0000000..580c982 --- /dev/null +++ b/evidence/explain.py @@ -0,0 +1,81 @@ +"""Explain Claim confidence via supporting / conflicting evidence.""" + +from __future__ import annotations + +import os +from typing import Any + +import structlog +from neo4j import Driver, GraphDatabase + +log = structlog.get_logger(__name__) + +_EXPLAIN_CYPHER = """ +MATCH (c:Claim {id: $claim_id}) +WHERE c.workspace_id = $workspace_id +OPTIONAL MATCH (s:Evidence)-[:SUPPORTS]->(c) +OPTIONAL MATCH (x:Evidence)-[:CONTRADICTS]->(c) +RETURN c { + .id, .workspace_id, .subject, .predicate, .object, .content, + .confidence, .status, .assertion_type, .decision_id, + .valid_from, .valid_to, .observed_at +} AS claim, +collect(DISTINCT s { + .id, .source_type, .source_id, .authority_class, .authority_score, + .snippet, .integrity_state +}) AS supporting, +collect(DISTINCT x { + .id, .source_type, .source_id, .authority_class, .authority_score, + .snippet, .integrity_state +}) AS conflicting +""" + + +class EvidenceExplainer: + """Assemble explainability payload for a Claim.""" + + def __init__(self, driver: Driver | None = None) -> None: + self._owns_driver = driver is None + if driver is None: + uri = os.environ.get("NEO4J_URI", "bolt://localhost:7687") + user = os.environ.get("NEO4J_USER", "neo4j") + password = os.environ.get("NEO4J_PASSWORD", "cortex_local") + self._driver = GraphDatabase.driver(uri, auth=(user, password)) + else: + self._driver = driver + + def explain(self, *, claim_id: str, workspace_id: str) -> dict[str, Any] | None: + """Return claim + evidence breakdown or None if missing.""" + with self._driver.session() as session: + record = session.run( + _EXPLAIN_CYPHER, + claim_id=claim_id, + workspace_id=workspace_id, + ).single() + if record is None or record["claim"] is None: + return None + supporting = [e for e in (record["supporting"] or []) if e and e.get("id")] + conflicting = [e for e in (record["conflicting"] or []) if e and e.get("id")] + claim = dict(record["claim"]) + return { + "claim": claim.get("content") or self._claim_text(claim), + "claim_id": claim.get("id"), + "confidence": claim.get("confidence"), + "status": claim.get("status"), + "assertion_type": claim.get("assertion_type"), + "temporal": { + "valid_from": claim.get("valid_from"), + "valid_to": claim.get("valid_to"), + "observed_at": claim.get("observed_at"), + }, + "supporting_evidence": supporting, + "conflicting_evidence": conflicting, + } + + @staticmethod + def _claim_text(claim: dict[str, Any]) -> str: + return f"{claim.get('subject')} {claim.get('predicate')} {claim.get('object')}" + + def close(self) -> None: + if self._owns_driver: + self._driver.close() diff --git a/evidence/graph.py b/evidence/graph.py new file mode 100644 index 0000000..068acc4 --- /dev/null +++ b/evidence/graph.py @@ -0,0 +1,190 @@ +"""Neo4j writer for Claim/Evidence nodes (gated by CORTEX_EVIDENCE_GRAPH).""" + +from __future__ import annotations + +import os +from typing import Any, Literal + +import structlog +from neo4j import Driver, GraphDatabase + +from evidence.adapters import evidence_graph_enabled +from evidence.authority import score_for_class +from shared.v2_models import Claim, Evidence + +log = structlog.get_logger(__name__) + +RelationKind = Literal["SUPPORTS", "CONTRADICTS"] + +_UPSERT_CLAIM = """ +MERGE (c:Claim {id: $id}) +ON CREATE SET + c.workspace_id = $workspace_id, + c.subject = $subject, + c.predicate = $predicate, + c.object = $object, + c.claim_type = $claim_type, + c.assertion_type = $assertion_type, + c.confidence = $confidence, + c.status = $status, + c.valid_from = $valid_from, + c.valid_to = $valid_to, + c.observed_at = $observed_at, + c.invalidated_at = $invalidated_at, + c.superseded_by = $superseded_by, + c.decision_id = $decision_id, + c.content = $content, + c.access_policy = $access_policy, + c.created_at = datetime(), + c.updated_at = datetime() +ON MATCH SET + c.confidence = $confidence, + c.status = $status, + c.updated_at = datetime(), + c.content = $content +RETURN c.id AS id +""" + +_UPSERT_EVIDENCE = """ +MERGE (e:Evidence {id: $id}) +ON CREATE SET + e.workspace_id = $workspace_id, + e.source_type = $source_type, + e.source_id = $source_id, + e.source_uri = $source_uri, + e.author = $author, + e.content_hash = $content_hash, + e.authority_class = $authority_class, + e.authority_score = $authority_score, + e.captured_at = $captured_at, + e.observed_at = $observed_at, + e.integrity_state = $integrity_state, + e.snippet = $snippet, + e.access_policy = $access_policy +ON MATCH SET + e.authority_score = $authority_score, + e.integrity_state = $integrity_state, + e.snippet = $snippet +RETURN e.id AS id +""" + +_LINK = """ +MATCH (e:Evidence {id: $evidence_id}) +MATCH (c:Claim {id: $claim_id}) +MERGE (e)-[r:%s]->(c) +ON CREATE SET r.created_at = datetime() +RETURN type(r) AS rel +""" + +_ASSERTS = """ +MATCH (d:Decision {id: $decision_id}) +MATCH (c:Claim {id: $claim_id}) +MERGE (d)-[r:ASSERTS]->(c) +ON CREATE SET r.created_at = datetime() +RETURN type(r) AS rel +""" + + +class EvidenceGraphWriter: + """Persists Claim/Evidence when the evidence-graph feature flag is on.""" + + def __init__( + self, + uri: str | None = None, + user: str | None = None, + password: str | None = None, + ) -> None: + self._uri = uri or os.environ.get("NEO4J_URI", "bolt://localhost:7687") + self._user = user or os.environ.get("NEO4J_USER", "neo4j") + self._password = password or os.environ.get("NEO4J_PASSWORD", "cortex_local") + self._driver: Driver = GraphDatabase.driver( + self._uri, auth=(self._user, self._password) + ) + + def write_claim_with_evidence( + self, + claim: Claim, + evidence_items: list[Evidence], + *, + relation: RelationKind = "SUPPORTS", + force: bool = False, + ) -> str | None: + """Write claim + evidence edges. Returns claim id or None if flag off.""" + if not force and not evidence_graph_enabled(): + log.debug("evidence_graph.skipped", reason="flag_off") + return None + + with self._driver.session() as session: + session.execute_write(self._tx_write, claim, evidence_items, relation) + + log.info( + "evidence_graph.written", + claim_id=claim.id, + evidence_count=len(evidence_items), + relation=relation, + ) + return claim.id + + @staticmethod + def _tx_write( + tx: Any, + claim: Claim, + evidence_items: list[Evidence], + relation: RelationKind, + ) -> None: + tx.run( + _UPSERT_CLAIM, + id=claim.id, + workspace_id=claim.workspace_id, + subject=claim.subject, + predicate=claim.predicate, + object=claim.object, + claim_type=claim.claim_type, + assertion_type=claim.assertion_type, + confidence=claim.confidence, + status=claim.status, + valid_from=claim.valid_from.isoformat() if claim.valid_from else None, + valid_to=claim.valid_to.isoformat() if claim.valid_to else None, + observed_at=claim.observed_at.isoformat() if claim.observed_at else None, + invalidated_at=( + claim.invalidated_at.isoformat() if claim.invalidated_at else None + ), + superseded_by=claim.superseded_by, + decision_id=claim.decision_id, + content=claim.content, + access_policy=str(claim.access_policy or {}), + ) + link_cypher = _LINK % relation + for ev in evidence_items: + score = ev.authority_score or score_for_class(ev.authority_class) + tx.run( + _UPSERT_EVIDENCE, + id=ev.id, + workspace_id=ev.workspace_id, + source_type=ev.source_type, + source_id=ev.source_id, + source_uri=ev.source_uri, + author=ev.author, + content_hash=ev.content_hash, + authority_class=ev.authority_class, + authority_score=score, + captured_at=ev.captured_at.isoformat(), + observed_at=ev.observed_at.isoformat() if ev.observed_at else None, + integrity_state=ev.integrity_state, + snippet=ev.snippet, + access_policy=str(ev.access_policy or {}), + ) + tx.run( + link_cypher, + evidence_id=ev.id, + claim_id=claim.id, + ) + if claim.decision_id: + tx.run( + _ASSERTS, + decision_id=claim.decision_id, + claim_id=claim.id, + ) + + def close(self) -> None: + self._driver.close() diff --git a/tests/api/test_evidence.py b/tests/api/test_evidence.py new file mode 100644 index 0000000..d61e295 --- /dev/null +++ b/tests/api/test_evidence.py @@ -0,0 +1,42 @@ +"""API tests for GET /evidence/explain/{claim_id}.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +from fastapi.testclient import TestClient + +from api.main import app + +client = TestClient(app) + + +def test_explain_404_when_missing() -> None: + with patch("api.evidence.EvidenceExplainer") as cls: + inst = MagicMock() + inst.explain.return_value = None + cls.return_value = inst + resp = client.get("/evidence/explain/missing?workspace_id=ws") + assert resp.status_code == 404 + + +def test_explain_200() -> None: + payload = { + "claim": "payments uses CockroachDB", + "claim_id": "c1", + "confidence": 0.94, + "status": "ACTIVE", + "assertion_type": "ASSERTED", + "temporal": {}, + "supporting_evidence": [{"id": "e1", "type": "ADR"}], + "conflicting_evidence": [], + } + with patch("api.evidence.EvidenceExplainer") as cls: + inst = MagicMock() + inst.explain.return_value = payload + cls.return_value = inst + resp = client.get("/evidence/explain/c1?workspace_id=ws") + assert resp.status_code == 200 + body = resp.json() + assert body["claim_id"] == "c1" + assert body["confidence"] == 0.94 diff --git a/tests/evidence/test_graph_explain.py b/tests/evidence/test_graph_explain.py new file mode 100644 index 0000000..2a05678 --- /dev/null +++ b/tests/evidence/test_graph_explain.py @@ -0,0 +1,107 @@ +"""Tests for evidence authority policy and explainer assembly.""" + +from __future__ import annotations + +from unittest.mock import MagicMock, patch + +import pytest + +from evidence.authority import compare_authority, load_authority_scores, score_for_class +from evidence.explain import EvidenceExplainer +from evidence.graph import EvidenceGraphWriter +from shared.v2_models import Claim, Evidence + + +def test_default_authority_ordering() -> None: + assert score_for_class("formal_policy") > score_for_class("casual_chat") + assert compare_authority("merged_PR", "casual_chat") == 1 + assert compare_authority("agent_inference", "approved_ADR") == -1 + + +def test_authority_json_override(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv( + "CORTEX_AUTHORITY_SCORES_JSON", + '{"casual_chat": 0.9, "merged_PR": 0.1}', + ) + scores = load_authority_scores() + assert scores["casual_chat"] == 0.9 + assert scores["merged_PR"] == 0.1 + + +def test_writer_skips_when_flag_off(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("CORTEX_EVIDENCE_GRAPH", "false") + with patch("evidence.graph.GraphDatabase") as gdb: + writer = EvidenceGraphWriter(uri="bolt://x", user="u", password="p") + claim = Claim( + workspace_id="ws", + subject="s", + predicate="p", + object="o", + content="c", + ) + assert writer.write_claim_with_evidence(claim, []) is None + gdb.driver.return_value.session.assert_not_called() + writer.close() + + +def test_writer_runs_when_forced(monkeypatch: pytest.MonkeyPatch) -> None: + monkeypatch.setenv("CORTEX_EVIDENCE_GRAPH", "false") + mock_session = MagicMock() + mock_session.__enter__ = MagicMock(return_value=mock_session) + mock_session.__exit__ = MagicMock(return_value=False) + mock_driver = MagicMock() + mock_driver.session.return_value = mock_session + with patch("evidence.graph.GraphDatabase") as gdb: + gdb.driver.return_value = mock_driver + writer = EvidenceGraphWriter(uri="bolt://x", user="u", password="p") + writer._driver = mock_driver + claim = Claim( + workspace_id="ws", + subject="payments", + predicate="uses", + object="CockroachDB", + content="payments uses CockroachDB", + decision_id="dec-1", + ) + ev = Evidence( + workspace_id="ws", + source_type="github", + source_id="pr-1", + authority_class="merged_PR", + ) + result = writer.write_claim_with_evidence(claim, [ev], force=True) + assert result == claim.id + mock_session.execute_write.assert_called_once() + writer.close() + + +def test_explainer_filters_empty_evidence() -> None: + mock_record = { + "claim": { + "id": "c1", + "content": "payments uses CockroachDB", + "confidence": 0.9, + "status": "ACTIVE", + "assertion_type": "ASSERTED", + "valid_from": None, + "valid_to": None, + "observed_at": None, + }, + "supporting": [{"id": "e1", "source_type": "github"}, None, {}], + "conflicting": [None], + } + mock_result = MagicMock() + mock_result.single.return_value = mock_record + mock_session = MagicMock() + mock_session.__enter__ = MagicMock(return_value=mock_session) + mock_session.__exit__ = MagicMock(return_value=False) + mock_session.run.return_value = mock_result + mock_driver = MagicMock() + mock_driver.session.return_value = mock_session + + explainer = EvidenceExplainer(driver=mock_driver) + out = explainer.explain(claim_id="c1", workspace_id="ws") + assert out is not None + assert out["claim"] == "payments uses CockroachDB" + assert len(out["supporting_evidence"]) == 1 + assert out["conflicting_evidence"] == []