Skip to content
216 changes: 19 additions & 197 deletions nemo_skills/pipeline/eval.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,177 +41,6 @@ class SingleNodeMode(str, enum.Enum):
parallel = "parallel"


def _create_comet_judge_tasks(
exp,
expname,
benchmark,
judge_pipeline_args,
rerun_done,
log_dir,
server_parameters,
cluster_config,
judge_server_gpus,
judge_server_nodes,
partition,
run_after,
reuse_code_exp,
reuse_code,
dependent_tasks,
all_tasks,
_task_dependencies,
installation_command,
skip_hf_home_check,
sbatch_kwargs,
):
"""Create tasks for Comet judge evaluation."""
from nemo_skills.pipeline.utils.generation import get_remaining_jobs

output_dir_path = judge_pipeline_args.get("output_dir")
input_file = judge_pipeline_args.get("input_file")
comet_model_path = judge_pipeline_args.get("judge_model")

# Determine seeds to check
if input_file is None:
num_seeds = judge_pipeline_args.get("num_random_seeds", 1)
random_seeds = list(range(num_seeds))
else:
random_seeds = [None]

remaining_jobs = get_remaining_jobs(
cluster_config=cluster_config,
output_dir=output_dir_path,
random_seeds=random_seeds,
chunk_ids=[None], # No chunking for judge task
rerun_done=rerun_done,
)

if not remaining_jobs or all(not chunks for chunks in remaining_jobs.values()):
LOG.info(f"Skipping Comet judge for {benchmark} - all output files and .done markers exist")
return []

# Build command to run xCOMET-XXL judge script
script_args = [f"--output-dir {output_dir_path} --comet-model-path {comet_model_path}"]

if input_file is None:
input_dir = judge_pipeline_args.get("input_dir")
script_args.append(f"--input-dir {input_dir}")
script_args.append(f"--num-seeds {num_seeds}")
else:
script_args.append(f"--input-file {input_file}")

run_cmd = f"pip install unbabel-comet && python3 -I /nemo_run/code/nemo_skills/evaluation/evaluator/comet.py {' '.join(script_args)}"

# Create task with GPU support for Comet
judge_task = pipeline_utils.add_task(
exp,
cmd=run_cmd,
task_name=f"{expname}-{benchmark}-comet-judge",
log_dir=log_dir + "/judge",
container=cluster_config["containers"]["vllm"],
cluster_config=cluster_config,
num_gpus=judge_server_gpus or 1,
num_nodes=judge_server_nodes or 1,
partition=partition,
run_after=run_after,
reuse_code_exp=reuse_code_exp,
reuse_code=reuse_code,
task_dependencies=(
dependent_tasks if cluster_config["executor"] == "slurm" else all_tasks + _task_dependencies
),
installation_command=installation_command,
skip_hf_home_check=skip_hf_home_check,
sbatch_kwargs=sbatch_kwargs,
)
return [judge_task]


def _create_nvembed_judge_tasks(
exp,
expname,
benchmark,
judge_pipeline_args,
rerun_done,
log_dir,
server_parameters,
cluster_config,
judge_server_gpus,
judge_server_nodes,
partition,
run_after,
reuse_code_exp,
reuse_code,
dependent_tasks,
all_tasks,
_task_dependencies,
installation_command,
skip_hf_home_check,
sbatch_kwargs,
):
"""Create tasks for NVEmbed judge evaluation."""
from nemo_skills.pipeline.utils.generation import get_remaining_jobs

output_dir_path = judge_pipeline_args.get("output_dir")
input_file = judge_pipeline_args.get("input_file")

# Determine seeds to check
if input_file is None:
num_seeds = judge_pipeline_args.get("num_random_seeds", 1)
random_seeds = list(range(num_seeds))
else:
random_seeds = [None]

remaining_jobs = get_remaining_jobs(
cluster_config=cluster_config,
output_dir=output_dir_path,
random_seeds=random_seeds,
chunk_ids=[None], # No chunking for judge task
rerun_done=rerun_done,
)

if not remaining_jobs or all(not chunks for chunks in remaining_jobs.values()):
LOG.info(f"Skipping NVEmbed judge for {benchmark} - all output files and .done markers exist")
return []

# Build command to run NVEmbed judge script
script_args = [f"--output-dir {output_dir_path}"]

if input_file is None:
input_dir = judge_pipeline_args.get("input_dir")
script_args.append(f"--input-dir {input_dir}")
script_args.append(f"--num-seeds {num_seeds}")
else:
script_args.append(f"--input-file {input_file}")

# Add skip-existing flag unless rerun_done is set
if not rerun_done:
script_args.append("--skip-existing")

run_cmd = f"python3 -I /nemo_run/code/nemo_skills/evaluation/evaluator/nvembed_judge.py {' '.join(script_args)}"

# Create task with GPU support for NVEmbed
judge_task = pipeline_utils.add_task(
exp,
cmd=run_cmd,
task_name=f"{expname}-{benchmark}-nvembed-judge",
log_dir=log_dir + "/judge",
container=cluster_config["containers"]["vllm"],
cluster_config=cluster_config,
num_gpus=judge_server_gpus or 1,
num_nodes=judge_server_nodes or 1,
partition=partition,
run_after=run_after,
reuse_code_exp=reuse_code_exp,
reuse_code=reuse_code,
task_dependencies=(
dependent_tasks if cluster_config["executor"] == "slurm" else all_tasks + _task_dependencies
),
installation_command=installation_command,
skip_hf_home_check=skip_hf_home_check,
sbatch_kwargs=sbatch_kwargs,
)
return [judge_task]


def _create_llm_judge_tasks(
ctx,
expname,
Expand Down Expand Up @@ -325,7 +154,7 @@ def eval(
help="Path to the entrypoint of the server. "
"If not specified, will use the default entrypoint for the server type.",
),
judge_type: str = typer.Option("llm", help="Type of judge to use: 'llm' (default) or 'nvembed'"),
judge_type: str = typer.Option("llm", help="Type of judge to use: llm (default), nvembed, comet or custom"),
judge_model: str = typer.Option(None, help="Path to the model to be used as a judge (if applicable)"),
judge_server_address: str = typer.Option(None, help="Address of the server hosting the judge model"),
judge_server_type: pipeline_utils.SupportedServers = typer.Option(
Expand Down Expand Up @@ -647,39 +476,32 @@ def eval(
benchmark_judge_type = judge_pipeline_args.pop("judge_type", judge_type)

# Create judge tasks based on judge type
# Use locate() pattern for extensible judge loading
judge_creator_path = None

if benchmark_judge_type == "nvembed":
judge_tasks = _create_nvembed_judge_tasks(
exp=exp,
expname=expname,
benchmark=benchmark,
judge_pipeline_args=judge_pipeline_args,
rerun_done=rerun_done,
log_dir=log_dir,
server_parameters=server_parameters,
cluster_config=cluster_config,
judge_server_gpus=judge_server_gpus,
judge_server_nodes=judge_server_nodes,
partition=partition,
run_after=run_after,
reuse_code_exp=reuse_code_exp,
reuse_code=reuse_code,
dependent_tasks=dependent_tasks,
all_tasks=all_tasks,
_task_dependencies=_task_dependencies,
installation_command=installation_command,
skip_hf_home_check=skip_hf_home_check,
sbatch_kwargs=sbatch_kwargs,
)
judge_creator_path = "nemo_skills.pipeline.judges.nvembed_judge::create_judge_tasks"
elif benchmark_judge_type == "comet":
judge_pipeline_args["judge_model"] = judge_model
judge_tasks = _create_comet_judge_tasks(
judge_creator_path = "nemo_skills.pipeline.judges.comet_judge::create_judge_tasks"
elif benchmark_judge_type == "custom":
judge_creator_path = judge_pipeline_args.pop("judge_creator_fn")

if judge_creator_path:
# Use locate() to dynamically load judge creator function
from nemo_skills.dataset.utils import locate

judge_creator_fn = locate(judge_creator_path)

# Call with standardized parameters
judge_tasks = judge_creator_fn(
exp=exp,
expname=expname,
benchmark=benchmark,
benchmarks=[benchmark],

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

what's the reason to have multiple benchmarks here?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No specific reason. Reverted this.

judge_pipeline_args=judge_pipeline_args,
rerun_done=rerun_done,
log_dir=log_dir,
server_parameters=server_parameters,
output_dir=output_dir,
cluster_config=cluster_config,
judge_server_gpus=judge_server_gpus,
judge_server_nodes=judge_server_nodes,
Expand Down
15 changes: 15 additions & 0 deletions nemo_skills/pipeline/judges/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""Judge implementations for evaluation pipeline."""
133 changes: 133 additions & 0 deletions nemo_skills/pipeline/judges/comet_judge.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
# Copyright (c) 2025, NVIDIA CORPORATION. All rights reserved.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.

"""Comet judge implementation for translation quality evaluation."""

import logging

from nemo_skills.pipeline.utils import add_task
from nemo_skills.pipeline.utils.generation import get_remaining_jobs
from nemo_skills.utils import get_logger_name

LOG = logging.getLogger(get_logger_name(__file__))


def create_judge_tasks(
exp,
expname,
benchmarks,
judge_pipeline_args,
rerun_done,
log_dir,
output_dir,
cluster_config,
judge_server_gpus,
judge_server_nodes,
partition,
run_after,
reuse_code_exp,
reuse_code,
dependent_tasks,
all_tasks,
_task_dependencies,
installation_command,
skip_hf_home_check,
sbatch_kwargs,
):
"""Create tasks for Comet judge evaluation.

Args:
exp: NeMo-Run experiment object
expname: Name of the experiment
benchmarks: List of benchmarks to evaluate (typically single benchmark)
judge_pipeline_args: Configuration for judge pipeline
rerun_done: Whether to rerun already completed jobs
log_dir: Directory for logs
output_dir: Output directory (unused, kept for interface compatibility)
cluster_config: Cluster configuration dict
judge_server_gpus: Number of GPUs for judge
judge_server_nodes: Number of nodes for judge
partition: SLURM partition
run_after: Dependencies to run after
reuse_code_exp: Experiment to reuse code from
reuse_code: Whether to reuse code
dependent_tasks: List of dependent tasks
all_tasks: List of all tasks
_task_dependencies: Additional task dependencies
installation_command: Installation command
skip_hf_home_check: Whether to skip HF_HOME check
sbatch_kwargs: Additional sbatch kwargs

Returns:
List of judge tasks created
"""
benchmark = benchmarks[0] # Comet judge works on single benchmark

output_dir_path = judge_pipeline_args.get("output_dir")
input_file = judge_pipeline_args.get("input_file")
comet_model_path = judge_pipeline_args.get("judge_model")
Comment on lines +77 to +79

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🛠️ Refactor suggestion | 🟠 Major

Use direct dict access for required judge_pipeline_args keys.

output_dir, judge_model, and input_dir (inside the if input_file is None branch) are all expected to be present. Using .get() on them silently yields None, which then propagates into both get_remaining_jobs(output_dir=None, …) and the constructed command string (e.g., --output-dir None --comet-model-path None), producing confusing downstream failures instead of an immediate, clear KeyError. Per coding guidelines, use direct dict access dict[key] for required keys.

input_file correctly uses .get() since None is its legitimate "not-provided" sentinel.

♻️ Proposed fix
-    output_dir_path = judge_pipeline_args.get("output_dir")
-    input_file = judge_pipeline_args.get("input_file")
-    comet_model_path = judge_pipeline_args.get("judge_model")
+    output_dir_path = judge_pipeline_args["output_dir"]
+    input_file = judge_pipeline_args.get("input_file")   # intentional sentinel
+    comet_model_path = judge_pipeline_args["judge_model"]
     if input_file is None:
-        input_dir = judge_pipeline_args.get("input_dir")
+        input_dir = judge_pipeline_args["input_dir"]
         script_args.append(f"--input-dir {input_dir}")

As per coding guidelines, "Do not use .get() for accessing dictionary keys if the code expects them to be present; use direct dictionary access dict[key] instead to allow proper error handling and fail fast with clear errors."

Also applies to: 101-101, 104-105

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@nemo_skills/pipeline/judges/comet_judge.py` around lines 77 - 79, Replace
silent .get() calls for required judge_pipeline_args keys with direct dict
access so missing required values raise immediately: change output_dir_path =
judge_pipeline_args.get("output_dir") and comet_model_path =
judge_pipeline_args.get("judge_model") to output_dir_path =
judge_pipeline_args["output_dir"] and comet_model_path =
judge_pipeline_args["judge_model"], and likewise when you reference "input_dir"
in the branch that runs when input_file is None use
judge_pipeline_args["input_dir"] instead of .get(); leave input_file using
.get() since None is a valid sentinel, and ensure the values passed into
get_remaining_jobs(...) and the constructed command string come from these
direct accesses to fail fast with a clear KeyError.


# Determine seeds to check
if input_file is None:
num_seeds = judge_pipeline_args.get("num_random_seeds", 1)
random_seeds = list(range(num_seeds))
else:
random_seeds = [None]

remaining_jobs = get_remaining_jobs(
cluster_config=cluster_config,
output_dir=output_dir_path,
random_seeds=random_seeds,
chunk_ids=[None], # No chunking for judge task
rerun_done=rerun_done,
)

if not remaining_jobs or all(not chunks for chunks in remaining_jobs.values()):
LOG.info(f"Skipping Comet judge for {benchmark} - all output files and .done markers exist")
return []

# Build command to run xCOMET-XXL judge script
script_args = [f"--output-dir {output_dir_path} --comet-model-path {comet_model_path}"]

if input_file is None:
input_dir = judge_pipeline_args.get("input_dir")
script_args.append(f"--input-dir {input_dir}")
script_args.append(f"--num-seeds {num_seeds}")
else:
script_args.append(f"--input-file {input_file}")

run_cmd = f"pip install unbabel-comet && python3 -I /nemo_run/code/nemo_skills/evaluation/evaluator/comet.py {' '.join(script_args)}"

# Create task with GPU support for Comet
judge_task = add_task(
exp,
cmd=run_cmd,
task_name=f"{expname}-{benchmark}-comet-judge",
log_dir=log_dir + "/judge",
container=cluster_config["containers"]["vllm"],
cluster_config=cluster_config,
num_gpus=judge_server_gpus or 1,
num_nodes=judge_server_nodes or 1,
partition=partition,
run_after=run_after,
reuse_code_exp=reuse_code_exp,
reuse_code=reuse_code,
task_dependencies=(
dependent_tasks if cluster_config["executor"] == "slurm" else all_tasks + _task_dependencies
),
installation_command=installation_command,
skip_hf_home_check=skip_hf_home_check,
sbatch_kwargs=sbatch_kwargs,
)
return [judge_task]
Loading
Loading