feat(harness): value grounding -- multi-column LSH + value_grounded (D14a)
Ports search/__init__.py (combined_search/RRF/aggregate) renamed psdwp3->nsp. New aggregate_lsh_multi (the D14a deviation): groups LSH hits by table keeping EVERY column where a value appears -- NOT collapsed to a single best column. The old _aggregate_lsh hid alternative groundings (e.g. 'ablazione' matching both a boolean flag and a free-text patologia field). aggregate_lsh_multi exposes all columns so the value-grounding widget lets the reviewer choose the anchor(s). Within one (table, column) the best-scored value is kept; columns ordered by score. value_grounded added to DecisionType (records the reviewer's anchor choice). L1: test_value_grounding (6 tests) -- multi-column exposure, grouping, within-column best-value, ordering, empty, and the value_grounded decision-type existence. Deferred: lshindex/ (needs vendor/thoth_lsh) and the L0 test_rrf.py land with the nsp lsh build command + index-building path; not needed for the pure L1 core here.
This commit is contained in:
@@ -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",
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,131 @@
|
||||
from pydantic import BaseModel
|
||||
|
||||
from nsp.vectorstore.store import VectorStore
|
||||
|
||||
|
||||
class SearchResult(BaseModel):
|
||||
key: str # column:<t>.<c> | table:<t> | evidence:<id>
|
||||
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]
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user