|
"""LAB-251 — live recall validation of the curated-seeding path (67 pins / 56 pages). |
|
|
|
Per technique: acquire the organic pool ONCE (live, keyless Stage-1 providers), then run the |
|
production selection path twice from that identical pool — BEFORE (no seeding = the 6cc50bdc |
|
baseline behaviour on current spine) and AFTER (curated seeding, authoritative pins). Selection |
|
calls the SAME node functions in the SAME order as run_staged_acquisition; the only node step |
|
skipped is the live Haiku design pre-pass (GEO_ABSTRACT_DESIGN_CLASSIFY — needs an LLM key; |
|
same omission as the original geo315 audit, caveat 3). |
|
|
|
Run: uv run --project <repo> python lab251_validate.py [--tech TECH_023] [--all] |
|
Outputs: out/{tech_id}.json (resume-safe: existing files are skipped), out/_failures.json |
|
""" |
|
|
|
from __future__ import annotations |
|
|
|
import argparse |
|
import json |
|
import logging |
|
import os |
|
import time |
|
import traceback |
|
from pathlib import Path |
|
|
|
HERE = Path(__file__).resolve().parent |
|
OUT = HERE / "out" |
|
OUT.mkdir(exist_ok=True) |
|
|
|
logging.basicConfig(level=logging.WARNING, format="%(levelname)s %(name)s: %(message)s") |
|
|
|
from ambient_intelligence.workflows.evidence_taxonomy.insight_labs_web_first.curated_refs import ( # noqa: E402 |
|
apply_curated_references, |
|
pin_record_key, |
|
) |
|
from ambient_intelligence.workflows.evidence_taxonomy.insight_labs_web_first.nodes.evidence_acquisition import ( # noqa: E402 |
|
_STAGED_MAX_SELECTED_REFERENCES, |
|
_abstract_design_enabled, |
|
_classify_designs_via_llm, |
|
_diversity_select_enabled, |
|
_extract_outcome_dimension, |
|
_extract_outcome_dimensions, |
|
_merged_evidence_to_selected_reference, |
|
_packet_scope_aliases, |
|
_partition_admissible_records, |
|
_partition_topically_relevant_records, |
|
_quality_first_reorder, |
|
_select_with_quota, |
|
_source_context_terms, |
|
_topical_relevance_gate_enabled, |
|
load_input_packet, |
|
) |
|
from ambient_intelligence.workflows.evidence_taxonomy.nodes.evidence_acquisition.dedupe import ( # noqa: E402 |
|
apply_cochrane_design_override, |
|
) |
|
from ambient_intelligence.workflows.evidence_taxonomy.nodes.evidence_acquisition.identity_norm import ( # noqa: E402 |
|
normalize_doi, |
|
normalize_pmid, |
|
) |
|
from ambient_intelligence.workflows.evidence_taxonomy.nodes.evidence_acquisition.staged_orchestrator import ( # noqa: E402 |
|
acquire_evidence_staged_with_taxonomy, |
|
) |
|
|
|
|
|
def registry_key(raw: str) -> str | None: |
|
"""Identity key for a missing-literature.csv doi cell: normalized DOI, else pmid:{digits}. |
|
|
|
Mirrors curated_refs._key_for_value ordering (PMID first) but tolerates the CSV's |
|
annotated cells ("PMID 30936803 / PMCID PMC6438091", "10.1136/bmj.a884 (PMID ...)"). |
|
""" |
|
raw = (raw or "").strip() |
|
if not raw: |
|
return None |
|
if raw.upper().startswith("PMID"): |
|
pmid = normalize_pmid(raw.split("/")[0]) |
|
return f"pmid:{pmid}" if pmid else None |
|
return normalize_doi(raw.split(" ")[0] if raw.startswith("10.") else raw) |
|
|
|
|
|
def record_keys(rec) -> set[str]: |
|
keys = set() |
|
if nd := normalize_doi(getattr(rec, "doi", None)): |
|
keys.add(nd) |
|
if pmid := normalize_pmid(getattr(rec, "pmid", None)): |
|
keys.add(f"pmid:{pmid}") |
|
return keys |
|
|
|
|
|
def run_lane(pool, curated, *, tech_id, technique, scope_ctx, registry_keys_for_tech): |
|
prov: dict = {} |
|
records = apply_curated_references(list(pool), curated, live=True, tech_id=tech_id, provenance=prov) |
|
pins = frozenset(prov.get("injected") or ()) | frozenset(prov.get("matched_present") or ()) |
|
# Lever 3 (node parity, seed → LLM classify → cochrane): abstract-aware Haiku design pass on the |
|
# coarse residual. Skips any record already stamped (curated-pin assertions protected) and is |
|
# edge-cached, so the second lane re-reads identical verdicts; in-place stamps on shared pool |
|
# objects keep both lanes consistent even cache-cold. Soft-fails internally; requires live keys. |
|
if _abstract_design_enabled(): |
|
_classify_designs_via_llm(records) |
|
apply_cochrane_design_override(records) |
|
reordered = _quality_first_reorder(records, scope_ctx=scope_ctx, pinned_keys=pins) |
|
admitted, gate_rejected = _partition_admissible_records(reordered) |
|
relevance_rejected: list = [] |
|
if _topical_relevance_gate_enabled(): |
|
admitted, relevance_rejected = _partition_topically_relevant_records(admitted, scope_ctx=scope_ctx, pinned_keys=pins) |
|
if _diversity_select_enabled(): |
|
capped, _ = _select_with_quota(admitted) |
|
else: |
|
capped = admitted[: _STAGED_MAX_SELECTED_REFERENCES] |
|
|
|
selected = [ |
|
_merged_evidence_to_selected_reference( |
|
rec, |
|
rank=i, |
|
technique=technique, |
|
aliases=scope_ctx["aliases"], |
|
modality=scope_ctx["modality"], |
|
mechanism_terms=scope_ctx["mechanism_terms"], |
|
category=scope_ctx["category"], |
|
is_pinned=pin_record_key(rec) in pins, |
|
) |
|
for i, rec in enumerate(capped, start=1) |
|
] |
|
capped_keys = {k for r in capped for k in record_keys(r)} |
|
scope_counts: dict[str, int] = {} |
|
for ref in selected: |
|
scope = ref.get("evidence_scope") or "indirect" |
|
scope_counts[scope] = scope_counts.get(scope, 0) + 1 |
|
pool_keys = {k for r in records for k in record_keys(r)} |
|
return { |
|
"n_selected": len(selected), |
|
"scope_counts": scope_counts, |
|
"registry_on_page": sorted(k for k in registry_keys_for_tech if k in capped_keys), |
|
"registry_in_pool": sorted(k for k in registry_keys_for_tech if k in pool_keys), |
|
"capped_out": sorted(pins - {k for r in capped if (k := pin_record_key(r))}), |
|
"gate_rejected": len(gate_rejected), |
|
"relevance_rejected": len(relevance_rejected), |
|
"seed_prov": { |
|
k: prov.get(k) |
|
for k in ( |
|
"injected", |
|
"matched_present", |
|
"fetch_failed", |
|
"retracted_refused", |
|
"citation_mismatch", |
|
"drop_conflicts", |
|
"invalid", |
|
"asserted", |
|
"pins", |
|
) |
|
if prov.get(k) |
|
}, |
|
"selected_pinned_ranks": [ |
|
{"rank": i, "key": pin_record_key(rec), "scope": selected[i - 1].get("evidence_scope")} |
|
for i, rec in enumerate(capped, start=1) |
|
if pin_record_key(rec) in pins |
|
], |
|
} |
|
|
|
|
|
def run_tech(tech_id: str, registry_keys_for_tech: set[str]) -> dict: |
|
t0 = time.time() |
|
outdir = OUT / "runs" / tech_id |
|
outdir.mkdir(parents=True, exist_ok=True) |
|
state: dict = {"evidence_mode": "fresh", "tech_id": tech_id, "live": True, "output_dir": str(outdir)} |
|
state = load_input_packet(state) |
|
packet = state["input_packet"] |
|
technique = packet.get("technique") or packet.get("display_label") or tech_id |
|
outcome = _extract_outcome_dimension(state) |
|
extras: list[str] = [] |
|
if os.environ.get("GEO_MULTI_OUTCOME_ACQUISITION", "1") != "0": |
|
extras = _extract_outcome_dimensions(state)[1:] |
|
|
|
# All 7 providers enabled — the audit's default config. The 3 keyless ones (consensus/openai/ |
|
# perplexity) contribute 0 records in this runtime, so the organic pool is production-minus-those- |
|
# keys; keeping them enabled ALSO paces the SS/OpenAlex index calls (consensus's transient retry |
|
# delay spaces the burst) which, without an S2 API key, is what keeps them under the anonymous |
|
# 429 rate limit. Disabling them sped each tech up but throttled the real index providers into |
|
# 120s timeouts — net slower and a degraded pool. Slow-and-paced is both faithful and reliable. |
|
result = acquire_evidence_staged_with_taxonomy(tech_id, outcome, extra_outcome_dimensions=extras) |
|
pool = result.records |
|
|
|
scope_ctx = { |
|
"technique": technique, |
|
"aliases": _packet_scope_aliases(packet), |
|
"modality": str(packet.get("modality") or ""), |
|
"mechanism_terms": _source_context_terms(packet).get("mechanism", []), |
|
"category": str(packet.get("category") or ""), |
|
} |
|
|
|
provider_status = { |
|
name: {"status": ps.status, "records": ps.record_count} for name, ps in (result.provider_status or {}).items() |
|
} |
|
|
|
before = run_lane( |
|
pool, |
|
{"pinned": [], "dropped": []}, |
|
tech_id=tech_id, |
|
technique=technique, |
|
scope_ctx=scope_ctx, |
|
registry_keys_for_tech=registry_keys_for_tech, |
|
) |
|
after = run_lane( |
|
pool, |
|
packet.get("curated_references") or {}, |
|
tech_id=tech_id, |
|
technique=technique, |
|
scope_ctx=scope_ctx, |
|
registry_keys_for_tech=registry_keys_for_tech, |
|
) |
|
return { |
|
"tech_id": tech_id, |
|
"technique": technique, |
|
"outcome_dimension": outcome, |
|
"extra_outcome_dimensions": extras, |
|
"pool_size": len(pool), |
|
"provider_status": provider_status, |
|
"registry_keys": sorted(registry_keys_for_tech), |
|
"before": before, |
|
"after": after, |
|
"elapsed_s": round(time.time() - t0, 1), |
|
} |
|
|
|
|
|
def load_registry() -> dict[str, set[str]]: |
|
import csv |
|
|
|
repo_data = ( |
|
Path(__file__).resolve().parents[1] |
|
/ "ambient-intelligence/ambient_intelligence/workflows/evidence_taxonomy/insight_labs_web_first/data/missing-literature.csv" |
|
) |
|
reg: dict[str, set[str]] = {} |
|
for row in csv.DictReader(open(repo_data, encoding="utf-8")): |
|
key = registry_key(row["doi"]) |
|
if key: |
|
reg.setdefault(row["tech_id"], set()).add(key) |
|
return reg |
|
|
|
|
|
def main() -> None: |
|
ap = argparse.ArgumentParser() |
|
ap.add_argument("--tech", action="append", default=None) |
|
ap.add_argument("--all", action="store_true") |
|
ap.add_argument("--shard", default=None, help="i/n — process techs where index %% n == i (parallel workers)") |
|
args = ap.parse_args() |
|
|
|
registry = load_registry() |
|
techs = args.tech or (sorted(registry) if args.all else []) |
|
if not techs: |
|
ap.error("pass --tech TECH_xxx or --all") |
|
if args.shard: |
|
i, n = (int(x) for x in args.shard.split("/")) |
|
techs = [t for idx, t in enumerate(techs) if idx % n == i] |
|
|
|
failures = {} |
|
fail_path = OUT / "_failures.json" |
|
if fail_path.is_file(): |
|
failures = json.loads(fail_path.read_text()) |
|
for tech_id in techs: |
|
dest = OUT / f"{tech_id}.json" |
|
if dest.is_file(): |
|
print(f"{tech_id}: exists, skipped", flush=True) |
|
continue |
|
try: |
|
res = run_tech(tech_id, registry.get(tech_id, set())) |
|
except Exception as exc: # record and continue — a one-tech failure must not kill the sweep |
|
failures[tech_id] = {"error": f"{type(exc).__name__}: {exc}", "trace": traceback.format_exc()[-2000:]} |
|
fail_path.write_text(json.dumps(failures, indent=2)) |
|
print(f"{tech_id}: FAILED {type(exc).__name__}: {exc}", flush=True) |
|
continue |
|
failures.pop(tech_id, None) |
|
fail_path.write_text(json.dumps(failures, indent=2)) |
|
dest.write_text(json.dumps(res, indent=2)) |
|
b, a = res["before"], res["after"] |
|
print( |
|
f"{tech_id}: pool={res['pool_size']} on_page {len(b['registry_on_page'])}->{len(a['registry_on_page'])}" |
|
f" of {len(res['registry_keys'])} | injected={len(a['seed_prov'].get('injected') or [])}" |
|
f" matched={len(a['seed_prov'].get('matched_present') or [])}" |
|
f" fetch_failed={len(a['seed_prov'].get('fetch_failed') or [])}" |
|
f" capped_out={len(a['capped_out'])} | {res['elapsed_s']}s", |
|
flush=True, |
|
) |
|
|
|
|
|
if __name__ == "__main__": |
|
main() |