diff --git a/harness/nsp/decisions.py b/harness/nsp/decisions.py index 00072e03..7c7de0a7 100644 --- a/harness/nsp/decisions.py +++ b/harness/nsp/decisions.py @@ -35,6 +35,10 @@ DecisionType = Literal[ # D15: marker di ritrazione. subject = "phase:N", retracts = decision_seq ritirata. # Resta nel log di audit (append-only); effective_decisions() la esclude dalla vista. "decision_retracted", + # D14a: valore citato nella domanda ancorato a una o piu' colonne. subject = + # "phase:4", detail = il valore (es. "ablazione"), rationale = la/e colonna/e scelta/e + # dal reviewer (aggregate_lsh_multi le espone tutte senza collassare al miglior match). + "value_grounded", ] diff --git a/harness/nsp/search/__init__.py b/harness/nsp/search/__init__.py new file mode 100644 index 00000000..d2c47a30 --- /dev/null +++ b/harness/nsp/search/__init__.py @@ -0,0 +1,131 @@ +from pydantic import BaseModel + +from nsp.vectorstore.store import VectorStore + + +class SearchResult(BaseModel): + key: str # column:. | table: | evidence: + label: str + kind: str # schema_column | schema_table | evidence | values + signals: dict # {"lsh": {"rank","score","value"?}, "vector": {"rank","score"}} + rrf: float + status: str = "" # status evidence, se applicabile + content: str = "" # testo matchato (per --explain) + + +def rrf_fuse(rankings: dict[str, list[tuple[str, float]]], k: int) -> dict[str, dict]: + """rankings: nome_segnale -> lista (key, raw_score) gia' ordinata per rilevanza. + Ritorna key -> {"rrf": float, "signals": {segnale: {"rank", "score"}}}.""" + fused: dict[str, dict] = {} + for signal, ranked in rankings.items(): + for rank, (key, score) in enumerate(ranked, start=1): + entry = fused.setdefault(key, {"rrf": 0.0, "signals": {}}) + entry["rrf"] += 1.0 / (k + rank) + entry["signals"][signal] = {"rank": rank, "score": round(score, 4)} + return fused + + +def _aggregate_lsh(lsh_hits: list[tuple[str, str, str, float]]) -> list[tuple[str, float, str]]: + """Aggrega i match LSH per tabella.colonna tenendo il migliore: (key, score, value).""" + best: dict[str, tuple[float, str]] = {} + for table, column, value, score in lsh_hits: + key = f"column:{table}.{column}" + if key not in best or score > best[key][0]: + best[key] = (score, value) + ordered = sorted(best.items(), key=lambda kv: kv[1][0], reverse=True) + return [(key, score, value) for key, (score, value) in ordered] + + +def aggregate_lsh_multi(hits: list[dict]) -> dict[str, list[dict]]: + """Value grounding (spec D14a): group LSH hits by table, keeping EVERY column + where the value appears -- NOT collapsed to a single best column. + + The old _aggregate_lsh collapsed matches to one column per table.column key, + hiding alternative groundings (e.g. 'ablazione' matching both a boolean flag + and a free-text patologia field). This function exposes all of them so the + value-grounding widget can let the reviewer choose which column(s) anchor a + cited value. + + hits: list of {table, column, value, score}. + Returns: {table -> [{column, value, score}, ...]}, each table's columns ordered + by score desc; within one (table, column) the best-scored value is kept. + """ + best: dict[tuple[str, str], dict] = {} + for h in hits: + key = (h["table"], h["column"]) + cur = best.get(key) + if cur is None or h["score"] > cur["score"]: + best[key] = {"column": h["column"], "value": h["value"], "score": h["score"]} + grouped: dict[str, list[dict]] = {} + for (table, _), row in best.items(): + grouped.setdefault(table, []).append(row) + for table in grouped: + grouped[table].sort(key=lambda r: r["score"], reverse=True) + return grouped + + +def _vector_key(hit) -> str: + if hit.kind == "schema_column": + return f"column:{hit.ref}" + if hit.kind == "schema_table": + return f"table:{hit.ref}" + return f"evidence:{hit.id}" + + +def schema_tables(results: list["SearchResult"], top_tables: int) -> list[tuple[str, float]]: + """Aggrega i risultati schema a livello di tabella per lo schema-linking: ogni chunk + (tabella o colonna) contribuisce alla sua tabella tenendo il miglior RRF. Ritorna le + prime top_tables tabelle, ordinate per RRF desc (poi nome). I chunk non-schema sono + ignorati.""" + best: dict[str, float] = {} + for r in results: + if r.kind not in ("schema_table", "schema_column"): + continue + table = r.key.split(":", 1)[1].split(".", 1)[0] + if table not in best or r.rrf > best[table]: + best[table] = r.rrf + ordered = sorted(best.items(), key=lambda kv: (-kv[1], kv[0])) + return ordered[:top_tables] + + +def combined_search( + keyword: str, + *, + lsh_hits: list[tuple[str, str, str, float]] | None, + store: VectorStore, + embedder, + top: int, + rrf_k: int, + kinds: list[str] | None, +) -> list[SearchResult]: + """Fonde LSH (valori di campo) e pgvector con Reciprocal Rank Fusion.""" + rankings: dict[str, list[tuple[str, float]]] = {} + lsh_values: dict[str, str] = {} + if lsh_hits: + aggregated = _aggregate_lsh(lsh_hits) + rankings["lsh"] = [(key, score) for key, score, _ in aggregated] + lsh_values = {key: value for key, _, value in aggregated} + + vector_hits = store.search(embedder.embed_query(keyword), top_n=top * 2, kinds=kinds) + rankings["vector"] = [(_vector_key(h), h.similarity) for h in vector_hits] + by_key = {_vector_key(h): h for h in vector_hits} + + fused = rrf_fuse(rankings, k=rrf_k) + results: list[SearchResult] = [] + for key, data in fused.items(): + hit = by_key.get(key) + if "lsh" in data["signals"] and key in lsh_values: + data["signals"]["lsh"]["value"] = lsh_values[key] + results.append( + SearchResult( + key=key, + label=hit.title if hit else key.removeprefix("column:"), + kind=hit.kind if hit else "values", + signals=data["signals"], + rrf=data["rrf"], + status=(hit.metadata.get("status", "") if hit else ""), + content=(hit.content if hit else lsh_values.get(key, "")), + ) + ) + results.sort(key=lambda r: r.rrf, reverse=True) + return results[:top] diff --git a/harness/tests/test_value_grounding.py b/harness/tests/test_value_grounding.py new file mode 100644 index 00000000..df3e17a8 --- /dev/null +++ b/harness/tests/test_value_grounding.py @@ -0,0 +1,70 @@ +"""L1: value grounding -- multi-column LSH exposure (spec D14a). + +The D14a deviation: a value cited in the question (e.g. 'ablazione') may match +MULTIPLE columns (a boolean flag, a free-text patologia field). The old behavior +collapsed matches to a single best column, hiding the alternative grounding. +aggregate_lsh_multi exposes every column where the value appears, grouped by table, +so the value-grounding widget can let the reviewer choose which column(s) anchor +the value. +""" +from nsp.search import aggregate_lsh_multi + + +def test_value_in_multiple_columns_returns_all(): + hits = [ + {"table": "t", "column": "c1", "value": "ablazione", "score": 0.9}, + {"table": "t", "column": "c2", "value": "ablazione", "score": 0.7}, + ] + result = aggregate_lsh_multi(hits) + # NOT collapsed to single best -- both columns exposed + cols = {h["column"] for h in result["t"]} + assert cols == {"c1", "c2"} + + +def test_values_grouped_by_table(): + hits = [ + {"table": "pazienti", "column": "flag_abl", "value": "ablazione", "score": 0.9}, + {"table": "ricoveri", "column": "procedura", "value": "ablazione", "score": 0.6}, + ] + result = aggregate_lsh_multi(hits) + assert set(result.keys()) == {"pazienti", "ricoveri"} + assert result["pazienti"][0]["column"] == "flag_abl" + assert result["ricoveri"][0]["column"] == "procedura" + + +def test_within_column_keeps_best_value(): + # two hits on the SAME column: keep the best-scored value (no duplicate rows + # for one column), but the column still appears once. + hits = [ + {"table": "t", "column": "c", "value": "ablazione", "score": 0.9}, + {"table": "t", "column": "c", "value": "ablaz", "score": 0.5}, + ] + result = aggregate_lsh_multi(hits) + rows = result["t"] + assert len(rows) == 1 + assert rows[0]["value"] == "ablazione" # best score kept + assert rows[0]["score"] == 0.9 + + +def test_columns_ordered_by_score_desc_within_table(): + hits = [ + {"table": "t", "column": "low", "value": "x", "score": 0.3}, + {"table": "t", "column": "high", "value": "x", "score": 0.95}, + {"table": "t", "column": "mid", "value": "x", "score": 0.6}, + ] + result = aggregate_lsh_multi(hits) + cols = [h["column"] for h in result["t"]] + assert cols == ["high", "mid", "low"] + + +def test_empty_hits_returns_empty(): + assert aggregate_lsh_multi([]) == {} + + +def test_value_grounded_decision_type_exists(): + # D14a adds the value_grounded decision type so the gate can record the + # reviewer's choice of which column(s) anchor a cited value. + from nsp.decisions import DecisionType + import typing + args = typing.get_args(DecisionType) + assert "value_grounded" in args