"""Adapters from the archived Jev-style sources to kev Examples. Each adapter yields one Example per question. `source.dataset` + `family` identify the leakage unit (a state and all its questions/variants); `source.panel` names the recipe part so metrics can be reported per panel. """ import json from collections.abc import Iterator from pathlib import Path from typing import cast from kev.types import Example, JSONValue, Label, Target ADAPTERS = ("kev-suite", "open-jev", "laya", "tasksource-jev") # Parquet adapters with `source` and `group_id` columns can be pre-sampled without loading states. GROUPED_PARQUET = ("open-jev", "tasksource-jev") def _hard(target: list[float]) -> int: return max(range(len(target)), key=target.__getitem__) def _is_one_hot(target: list[float]) -> bool: return max(target) > 1 - 1e-9 def kev_suite(path: Path, panel: str) -> Iterator[Example]: """kev-suites JSONL: one state with several typed questions and gold labels.""" with path.open(encoding="utf-8") as stream: for line in stream: if not line.strip(): continue raw = json.loads(line) meta = raw["_meta"] for key, question in raw["questions"].items(): kind = question["type"] converted: dict[str, JSONValue] = {"type": kind, "instructions": question.get("instructions")} if kind != "noul" or question.get("criteria"): converted["criteria"] = question["criteria"] label = cast(Label, question["label"]) yield cast(Example, { "id": f"{meta['id']}#{key}", "suite": str(meta.get("source") or question.get("src")), "family": str(meta.get("group_id") or meta["id"]), "label": label, "target": label, "state": raw["state"], "question": converted, "source": {"adapter": "kev-suite", "dataset": f"kev:{meta.get('source')}", "panel": panel, "file": str(path), "question_key": key, "row_sha256": meta.get("row_sha256"), "split": meta.get("split"), "variant": meta.get("variant")}, }) def open_jev(path: Path, panel: str, groups: set[str] | None = None) -> Iterator[Example]: """Open-Jev parquet: one question per record with an explicit option list and target distribution.""" import pyarrow.parquet as pq columns = ["id", "group_id", "split", "source", "kind", "question", "options", "target", "state_json"] for batch in pq.ParquetFile(path).iter_batches(batch_size=8192, columns=columns): for raw in batch.to_pylist(): if groups is not None and raw["group_id"] not in groups: continue kind, options, distribution = raw["kind"], list(raw["options"]), [float(p) for p in raw["target"]] question: dict[str, JSONValue] label: Label target: Target if kind == "noul": if options != ["no", "yes"]: raise ValueError(f"Unexpected noul options {options}: {raw['id']}") label, target = distribution[1] >= 0.5, distribution[1] question = {"type": "noul", "instructions": raw["question"]} elif kind == "score": label = _hard(distribution) target = label if _is_one_hot(distribution) else distribution question = {"type": "score", "instructions": raw["question"], "criteria": cast(JSONValue, options)} else: label = options[_hard(distribution)] target = label if _is_one_hot(distribution) else distribution question = {"type": "choice", "instructions": raw["question"], "criteria": {option: None for option in options}} yield cast(Example, { "id": raw["id"], "suite": raw["source"], "family": raw["group_id"], "label": label, "target": target, "state": json.loads(raw["state_json"]), "question": question, "source": {"adapter": "open-jev", "dataset": f"open-jev:{raw['source'].split('/')[0]}", "panel": panel, "file": str(path), "split": raw["split"]}, }) def tasksource_jev(path: Path, panel: str, groups: set[str] | None = None) -> Iterator[Example]: """tasksource-jev-typed-decisions parquet: one question per record; noul stores only P(yes).""" import pyarrow.parquet as pq columns = ["id", "group_id", "source", "kind", "question", "options", "target", "state", "variant"] for batch in pq.ParquetFile(path).iter_batches(batch_size=8192, columns=columns): for raw in batch.to_pylist(): if groups is not None and raw["group_id"] not in groups: continue kind, options, distribution = raw["kind"], list(raw["options"]), [float(p) for p in raw["target"]] question: dict[str, JSONValue] label: Label target: Target if kind == "noul": if len(distribution) != 1: raise ValueError(f"Expected one noul probability: {raw['id']}") label, target = distribution[0] >= 0.5, distribution[0] question = {"type": "noul", "instructions": raw["question"]} elif kind == "score": label = _hard(distribution) target = label if _is_one_hot(distribution) else distribution question = {"type": "score", "instructions": raw["question"], "criteria": cast(JSONValue, options)} else: label = options[_hard(distribution)] target = label if _is_one_hot(distribution) else distribution question = {"type": "choice", "instructions": raw["question"], "criteria": {option: None for option in options}} yield cast(Example, { "id": raw["id"], "suite": raw["source"], "family": raw["group_id"], "label": label, "target": target, "state": raw["state"], "question": question, "source": {"adapter": "tasksource-jev", "dataset": f"tasksource:{raw['source']}", "panel": panel, "file": str(path), "variant": raw["variant"]}, }) def laya(path: Path, panel: str) -> Iterator[Example]: """Laya typed-decisions parquet: five questions per case, gold as annotator distributions.""" import pyarrow.parquet as pq for raw in pq.read_table(path).to_pylist(): questions = json.loads(raw["questions"]) gold = json.loads(raw["gold"]) state = json.loads(raw["state"]) if isinstance(raw["state"], str) else raw["state"] for key, question in questions.items(): answer = gold[key] # Gold is stored to 6 decimals; renormalize so the distribution sums to 1 exactly. total = sum(float(value) for value in answer["probabilities"].values()) answer["probabilities"] = {name: float(value) / total for name, value in answer["probabilities"].items()} kind = question["type"] label: Label target: Target if kind == "noul": label, target = answer["label"] == "true", float(answer["probabilities"]["true"]) elif kind == "score": label = int(answer["label"]) target = [float(answer["probabilities"][str(i)]) for i in range(len(question["criteria"]))] else: label = answer["label"] target = [float(answer["probabilities"][option]) for option in question["criteria"]] yield cast(Example, { "id": f"laya:{raw['id']}#{key}", "suite": f"laya/{raw['workflow']}", "family": raw["id"], "label": label, "target": target, "state": state, "question": question, "source": {"adapter": "laya", "dataset": f"laya:{raw['workflow']}", "panel": panel, "file": str(path), "split": raw["split"]}, }) def group_index(path: Path) -> list[tuple[str, str]]: """(source, group_id) per record, read without loading states.""" import pyarrow.parquet as pq table = pq.read_table(path, columns=["source", "group_id"]) return list(zip(table.column("source").to_pylist(), table.column("group_id").to_pylist(), strict=True)) def load(adapter: str, path: Path, panel: str, groups: set[str] | None = None) -> Iterator[Example]: if groups is not None and adapter not in GROUPED_PARQUET: raise ValueError(f"Adapter {adapter} cannot pre-select groups") if adapter == "kev-suite": return kev_suite(path, panel) if adapter == "open-jev": return open_jev(path, panel, groups) if adapter == "tasksource-jev": return tasksource_jev(path, panel, groups) if adapter == "laya": return laya(path, panel) raise ValueError(f"Unknown adapter {adapter!r}; choose from {ADAPTERS}")