Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
103 changes: 19 additions & 84 deletions judgearena/benchmarks/pairwise/runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,26 +12,14 @@
from judgearena.artifacts import prepare_run_directory, write_run_metadata_safely
from judgearena.benchmarks.execution import build_generation_kwargs, build_judge
from judgearena.benchmarks.pairwise.baselines import resolve_baseline_plan
from judgearena.benchmarks.pairwise.scoring import (
DEFAULT_PAIRWISE_SCORER,
resolve_pairwise_scorer,
)
from judgearena.datasets import load_instructions
from judgearena.datasets.fluency import is_fluency_task as task_is_fluency
from judgearena.datasets.fluency import load_fluency_contexts
from judgearena.datasets.pairwise import PairwiseTaskData, load_pairwise_task_data
from judgearena.benchmarks.pairwise.scoring import resolve_pairwise_scorer
from judgearena.datasets.pairwise import load_pairwise_task_data
from judgearena.evaluate import judge_and_parse_prefs, resolve_run_judge_prompt
from judgearena.generate import generate_base, generate_instructions
from judgearena.log import get_logger
from judgearena.tasks.registry import get_packaged_task
from judgearena.tasks.schema import ResolvedTaskSpec
from judgearena.utils import (
cache_function_dataframe,
data_root,
download_hf,
generation_cache_token,
read_df,
)
from judgearena.utils import cache_function_dataframe, generation_cache_token
from judgearena.utils.eval import BattleReport

if TYPE_CHECKING:
Expand All @@ -40,39 +28,6 @@
logger = get_logger(__name__)


def _try_load_legacy_dataset_completions(
dataset: str, model: str, n_instructions: int | None
) -> pd.DataFrame | None:
"""Try loading pre-existing completions for an unregistered legacy task.

Registered tasks load outputs through ``PairwiseTaskData`` instead.
"""
local_path_tables = data_root / "tables"
download_hf(name=dataset, local_path=local_path_tables)
output_path = local_path_tables / "model_outputs" / f"{dataset}.csv.zip"
if not output_path.exists():
return None
df_outputs = read_df(output_path)
df_outputs.loc[:, "output"] = df_outputs.loc[:, "output"].fillna("")
df_outputs = df_outputs.pivot_table(
index="instruction_index", columns="model", values="output", aggfunc="last"
).sort_index()
if model not in df_outputs.columns:
return None
logger.info(
"Found pre-existing completions for '%s' in dataset '%s'.", model, dataset
)
completions = df_outputs.loc[:, model]
if n_instructions is not None:
completions = completions.head(n_instructions)
return pd.DataFrame(
{
"completion": completions.values,
"instruction_index": completions.index.tolist(),
}
)


def run_pairwise(cfg: "RunConfig", resolved_task: ResolvedTaskSpec | None = None):
"""
1) take as input:
Expand All @@ -92,28 +47,15 @@ def run_pairwise(cfg: "RunConfig", resolved_task: ResolvedTaskSpec | None = None
# set_langchain_cache()
ignore_cache = cfg.run.ignore_cache

# Currrently, we run context evaluation
is_fluency_task = task_is_fluency(cfg.task)
resolved_task = resolved_task or get_packaged_task(cfg.task)
task_data: PairwiseTaskData | None = None
if resolved_task is not None:
task_data = load_pairwise_task_data(
resolved_task,
n_instructions=cfg.generation.n_instructions,
)
instructions_df = task_data.instructions
instructions = instructions_df.loc[:, "instruction"]
elif is_fluency_task:
# if cfg.task = "fluency-french", we map to the "French" config of
# https://huggingface.co/datasets/geoalgo/multilingual-fluency
instructions = load_fluency_contexts(data_root, cfg.task)
instructions_df = pd.DataFrame({"instruction": instructions.values})
instructions_df.index = instructions.index
else:
instructions_df = load_instructions(
dataset=cfg.task, n_instructions=cfg.generation.n_instructions
)
instructions = instructions_df.loc[:, "instruction"]
if resolved_task is None:
raise ValueError(f"Unknown task {cfg.task!r}.")
task_data = load_pairwise_task_data(
resolved_task,
n_instructions=cfg.generation.n_instructions,
)
instructions_df = task_data.instructions
instructions = instructions_df.loc[:, "instruction"]

n_instructions = (
cfg.generation.n_instructions
Expand Down Expand Up @@ -155,7 +97,11 @@ def run_pairwise(cfg: "RunConfig", resolved_task: ResolvedTaskSpec | None = None
)

# TODO currently we just support base models for fluency, we could also support instruction-tuned models
generation_function = generate_base if is_fluency_task else generate_instructions
generation_function = (
generate_base
if resolved_task.spec.protocol.generation.mode == "base_completion"
else generate_instructions
)

def _run_generation(
model_spec: str, *, generation_kwargs: dict[str, object]
Expand All @@ -172,16 +118,9 @@ def _align_completion_series(df: pd.DataFrame) -> pd.Series:
return df.set_index("instruction_index").loc[instructions.index, "completion"]

def _load_or_generate_completions(model_spec: str, *, role: str) -> pd.Series:
if task_data is not None:
preloaded = task_data.completions_for(model_spec)
if preloaded is not None:
return preloaded.loc[instructions.index]
else:
preloaded = _try_load_legacy_dataset_completions(
cfg.task, model_spec, n_instructions
)
preloaded = task_data.completions_for(model_spec)
if preloaded is not None:
return _align_completion_series(preloaded)
return preloaded.loc[instructions.index]
# Fold the resolved generation kwargs into the cache key so that changing
# any sampling param (temperature, seed, top_p/k, max_tokens, ...) busts
# the cached completions instead of silently reusing a stale run.
Expand Down Expand Up @@ -266,11 +205,7 @@ def _load_or_generate_completions(model_spec: str, *, role: str) -> pd.Series:

df.to_csv(res_folder / f"{name}-annotations.csv", index=False)

scorer = resolve_pairwise_scorer(
resolved_task.spec.protocol.scoring.adapter
if resolved_task is not None
else DEFAULT_PAIRWISE_SCORER
)
scorer = resolve_pairwise_scorer(resolved_task.spec.protocol.scoring.adapter)
summary = scorer.summarize(prefs)

report = BattleReport(
Expand Down
20 changes: 4 additions & 16 deletions judgearena/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,6 @@
)

from judgearena.benchmarks.pairwise.baselines import native_pairwise_baseline
from judgearena.datasets.fluency import is_fluency_task
from judgearena.tasks.registry import get_packaged_task
from judgearena.tasks.schema import EloProtocol

Expand Down Expand Up @@ -390,21 +389,10 @@ class RunConfig(BaseSettings):
def _validate(self) -> RunConfig:
resolved_task = get_packaged_task(self.task)
if resolved_task is None:
if not is_fluency_task(self.task):
raise ValueError(
f"Unknown task {self.task!r}; use 'judgearena tasks list' to "
"inspect packaged tasks."
)
# Fluency tasks are not packaged yet and run through the legacy path.
if self.elo is not None:
raise ValueError("elo config is only valid for ELO tasks.")
if self.model.name is None:
raise ValueError("model.name is required.")
if self.model.baseline is None:
raise ValueError(
f"model.baseline is required for task {self.task!r}."
)
return self
raise ValueError(
f"Unknown task {self.task!r}; use 'judgearena tasks list' to "
"inspect packaged tasks."
)

protocol = resolved_task.spec.protocol
task_judge = protocol.judge
Expand Down
46 changes: 0 additions & 46 deletions judgearena/dataset_revisions.py

This file was deleted.

118 changes: 62 additions & 56 deletions judgearena/datasets/fluency.py
Original file line number Diff line number Diff line change
@@ -1,14 +1,10 @@
"""Fluency benchmark: multilingual sentence-completion contexts for base models.

Loads per-language contexts from ``geoalgo/multilingual-fluency`` (the
successor to ``geoalgo/multilingual-contexts-to-be-completed``), which is
generated by ``scripts/fluency/generate_fluency.py``. Each language is a
separate HF dataset config/folder, so a task like ``fluency-mandarin-chinese``
resolves to the ``Mandarin Chinese`` config.

Mirrors ``judgearena/instruction_dataset/m_arenahard.py``: the whole dataset
repo is snapshot-downloaded (pinned via ``dataset_revisions``) and then
filtered locally to the requested language's parquet file(s).
Dataset adapter for the packaged ``fluency`` task family. Contexts come from
``geoalgo/multilingual-fluency`` (generated by
``scripts/fluency/generate_fluency.py``), where each language is a separate
HF dataset config/folder; a task like ``fluency-mandarin-chinese`` selects the
``Mandarin Chinese`` folder through the task's language variants.
"""

from __future__ import annotations
Expand All @@ -18,9 +14,7 @@
import pandas as pd
from huggingface_hub import snapshot_download

from judgearena.dataset_revisions import hf_revision

FLUENCY_TASK_PREFIX = "fluency-"
from judgearena.tasks.schema import HuggingFaceDatasetSource, ResolvedTaskSpec

FLUENCY_HF_REPO = "geoalgo/multilingual-fluency"

Expand Down Expand Up @@ -83,59 +77,71 @@ def fluency_language_slug(language: str) -> str:
}


def fluency_task_name(language: str) -> str:
"""Build the ``--task`` value for a given language, e.g. "French" -> "fluency-french"."""
return f"{FLUENCY_TASK_PREFIX}{fluency_language_slug(language)}"


def is_fluency_task(task: str) -> bool:
return task.startswith(FLUENCY_TASK_PREFIX)


def resolve_fluency_language(task: str) -> str:
"""Map a ``fluency-{slug}`` task name to its dataset config/folder name."""
if not is_fluency_task(task):
raise ValueError(f"Not a fluency task: {task!r}")
slug = task[len(FLUENCY_TASK_PREFIX) :]
if slug not in _SLUG_TO_LANGUAGE:
def _source(task: ResolvedTaskSpec) -> HuggingFaceDatasetSource:
source = task.spec.dataset.sources.get("examples")
if not isinstance(source, HuggingFaceDatasetSource):
raise ValueError(
f"Unknown fluency language {slug!r} from task {task!r}. "
f"Known languages: {sorted(_SLUG_TO_LANGUAGE)}"
f"Task {task.task!r} must define a Hugging Face dataset source "
"named 'examples'."
)
return _SLUG_TO_LANGUAGE[slug]
return source


def _fluency_local_root(local_path: Path) -> Path:
local_subdir = FLUENCY_HF_REPO.split("/", 1)[1]
return local_path / local_subdir
def _source_local_dir(source: HuggingFaceDatasetSource, root: Path) -> Path:
"""Keep raw source snapshots separate from normalized JudgeArena tables."""
return root / "_sources" / source.repo_id.replace("/", "--")


def download_fluency_dataset(local_path: Path) -> Path:
"""Download every language config of the fluency dataset; return the local root."""
root = _fluency_local_root(local_path)
def download_task_sources(task: ResolvedTaskSpec, local_tables_path: Path) -> None:
"""Download the pinned fluency contexts declared by the task."""
if task.spec.dataset.adapter != "fluency":
raise ValueError(f"Task {task.task!r} does not use the fluency adapter.")
source = _source(task)
snapshot_download(
repo_id=FLUENCY_HF_REPO,
repo_id=source.repo_id,
repo_type="dataset",
allow_patterns="*",
local_dir=root,
revision=source.revision,
allow_patterns=list(source.allow_patterns) or None,
local_dir=_source_local_dir(source, local_tables_path),
force_download=False,
revision=hf_revision(FLUENCY_HF_REPO),
)
return root


def load_fluency_contexts(local_path: Path, task: str) -> pd.Series:
"""Load the ``instruction`` contexts (sentence prefixes) for a fluency task."""
language = resolve_fluency_language(task)
root = download_fluency_dataset(local_path)
matches = sorted((root / language).glob("*.parquet"))
assert matches, f"No parquet found for language {language!r} under {root}."
df = pd.concat([pd.read_parquet(path) for path in matches], ignore_index=True)
return df.loc[:, "sentence"].rename("instruction")


if __name__ == "__main__":
from judgearena.paths import data_root

contexts = load_fluency_contexts(data_root, task="fluency-french")
print(contexts.head())
def _selected_languages(task: ResolvedTaskSpec) -> tuple[str, ...]:
variants = task.spec.variants
if variants is None or variants.selector != "language":
raise ValueError(f"Task {task.task!r} must define language suffix variants.")
slugs = task.selection.values if task.selection is not None else variants.values
return tuple(_SLUG_TO_LANGUAGE[slug] for slug in slugs)


def load_task_instructions(
task: ResolvedTaskSpec, local_tables_path: Path
) -> pd.DataFrame:
"""Load the sentence contexts for the languages selected by the invocation."""
download_task_sources(task, local_tables_path)
root = _source_local_dir(_source(task), local_tables_path)

frames: list[pd.DataFrame] = []
for language in _selected_languages(task):
matches = sorted((root / language).glob("*.parquet"))
if not matches:
raise FileNotFoundError(
f"No parquet found for language {language!r} under {root}."
)
frame = pd.concat([pd.read_parquet(path) for path in matches])
frame["instruction_index"] = [
f"{fluency_language_slug(language)}-{i}" for i in range(len(frame))
]
frames.append(frame)

df = pd.concat(frames, ignore_index=True)
fields = task.spec.dataset.fields
return df.rename(columns={fields.instruction: "instruction"})


def load_task_model_outputs(
task: ResolvedTaskSpec, local_tables_path: Path
) -> pd.DataFrame | None:
"""Fluency has no pre-generated model outputs; completions are generated."""
return None
Loading
Loading