Skip to content
Merged
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
20 changes: 18 additions & 2 deletions docs/api/reward.rst
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,9 @@ Last updated: |today| (API docstrings are auto-generated).
VeRL-Omni reward pipelines support both rule-based scoring (e.g. JPEG
compressibility) and model-based generative reward models (e.g. OCR via a
vision-language model served behind an OpenAI-compatible router). Reward
computation is dispatched per sample by the
:class:`~verl_omni.reward_loop.reward_manager.VisualRewardManager`, which
computation is dispatched per sample by modality-specific reward managers,
including :class:`~verl_omni.reward_loop.reward_manager.VisualRewardManager`
and :class:`~verl_omni.reward_loop.reward_manager.AudioRewardManager`, which
plugs into :class:`~verl_omni.reward_loop.reward_loop.OmniRewardLoopManager` —
verl's :class:`~verl.experimental.reward_loop.RewardLoopManager` extended with
profiler control over the reward-model rollout servers.
Expand All @@ -17,8 +18,10 @@ profiler control over the reward-model rollout servers.

verl_omni.reward_loop.reward_loop.OmniRewardLoopManager
verl_omni.reward_loop.reward_manager.VisualRewardManager
verl_omni.reward_loop.reward_manager.AudioRewardManager
verl_omni.utils.reward_score.default_compute_score_image
verl_omni.utils.reward_score.http_scorer_client.compute_score
verl_omni.utils.reward_score.audio_http_scorer_client.compute_score
verl_omni.utils.reward_score.unified_reward.compute_score_unified_reward

Reward Loop Manager
Expand All @@ -33,6 +36,13 @@ Reward Manager
.. autoclass:: verl_omni.reward_loop.reward_manager.VisualRewardManager
:members: __init__, run_single

.. autoclass:: verl_omni.reward_loop.reward_manager.AudioRewardManager
:members: __init__, run_single

``AudioRewardManager`` reads ``audio`` and ``audio_sample_rate`` from rollout
``extra_info``, validates a finite CPU float waveform, and calls a synchronous
or asynchronous custom scorer with ``solution_audio=(waveform, sample_rate)``.

Default Score Dispatcher
~~~~~~~~~~~~~~~~~~~~~~~~~

Expand Down Expand Up @@ -60,6 +70,12 @@ HTTP Scorer Client
.. automodule:: verl_omni.utils.reward_score.http_scorer_client
:members: compute_score

Audio HTTP Scorer Client
^^^^^^^^^^^^^^^^^^^^^^^^

.. automodule:: verl_omni.utils.reward_score.audio_http_scorer_client
:members: compute_score

UnifiedReward Scorer
^^^^^^^^^^^^^^^^^^^^^

Expand Down
22 changes: 21 additions & 1 deletion docs/start/http_scorer.md
Original file line number Diff line number Diff line change
@@ -1,10 +1,30 @@
(http_scorer)=
# Using an External HTTP Scorer Service

Last updated: 08/09/2026
Last updated: 09/04/2026

VeRL-Omni ships a generic HTTP reward client (`verl_omni.utils.reward_score.http_scorer_client`) that sends generated images to an external scorer service over HTTP and returns the score. This is useful when your reward model is too large to co-locate with training, needs a different runtime (e.g., a separate GPU pool), or is shared across multiple experiments.

Audio rollouts use `verl_omni.utils.reward_score.audio_http_scorer_client`.
That client sends JSON containing a base64-encoded float32 waveform,
`sample_rate`, target `prompt`, and scalar metadata. The service returns a JSON
object with a finite `score` and optional diagnostics:

```json
{
"protocol_version": "1",
"waveform_f32_base64": "...",
"num_samples": 24000,
"sample_rate": 24000,
"prompt": "Text to synthesize",
"metadata": {"id": "stable-id"}
}
```

The endpoint must return `{"score": 1.25}` and may include additional scalar
diagnostics. The client retries only transient network, timeout, HTTP 408/429,
and 5xx failures; malformed or non-finite responses fail closed.

## How it works

```text
Expand Down
263 changes: 263 additions & 0 deletions tests/reward_loop/test_audio_reward_manager_on_cpu.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,263 @@
# Copyright 2026 Bytedance Ltd. and/or its affiliates
#
# 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.
"""CPU tests for generic waveform reward routing."""

import importlib.util
import threading
from functools import partial
from pathlib import Path
from unittest.mock import MagicMock

import numpy as np
import pytest
import torch
from omegaconf import OmegaConf
from verl import DataProto
from verl.utils.reward_score import default_compute_score


def _load_audio_reward_manager():
path = Path(__file__).parents[2] / "verl_omni/reward_loop/reward_manager/audio.py"
spec = importlib.util.spec_from_file_location("audio_reward_manager_under_test", path)
module = importlib.util.module_from_spec(spec)
assert spec.loader is not None
spec.loader.exec_module(module)
return module.AudioRewardManager


AudioRewardManager = _load_audio_reward_manager()


def _config():
return OmegaConf.create({"reward": {}})


def _manager(compute_score):
return AudioRewardManager(_config(), MagicMock(), compute_score=compute_score)


def _data(audio=None, sample_rate=24_000, *, layout="streamed"):
fields = {}
if audio is not None:
fields = {"audio": audio, "audio_sample_rate": sample_rate}
non_tensors = {
"data_source": ["tts_reward"],
"reward_model": [{"ground_truth": "ni3 hao3"}],
"extra_info": [{"id": "sample-0"}],
"tool_extra_fields": [fields if layout == "streamed" else {}],
}
if layout == "finalized" and audio is not None:
non_tensors.update({"audio": [audio], "audio_sample_rate": [sample_rate]})
return DataProto.from_dict(
tensors={"responses": torch.zeros(1, 4, dtype=torch.long)},
non_tensors=non_tensors,
)


def test_assemble_scores_preserves_one_reward_per_sample():
data = DataProto.from_dict(
tensors={
"prompts": torch.zeros(3, 2, dtype=torch.long),
"responses": torch.zeros(3, 4, dtype=torch.long),
"attention_mask": torch.tensor(
[
[1, 1, 1, 1, 1, 1],
[1, 1, 1, 1, 0, 0],
[1, 1, 1, 1, 1, 0],
]
),
}
)

scores = AudioRewardManager.assemble_rm_scores(data, [0.1, -1.0, 2.5])

assert scores.shape == (3, 4)
assert scores.dtype == torch.float32
torch.testing.assert_close(
scores,
torch.tensor(
[
[0.0, 0.0, 0.0, 0.1],
[0.0, -1.0, 0.0, 0.0],
[0.0, 0.0, 2.5, 0.0],
]
),
)


def test_run_single_rejects_multi_sample_batch():
data = DataProto.from_dict(
tensors={"responses": torch.zeros(2, 4, dtype=torch.long)},
non_tensors={
"data_source": ["tts_reward", "tts_reward"],
"reward_model": [{"ground_truth": "first"}, {"ground_truth": "second"}],
},
)
manager = _manager(lambda **kwargs: 0.0)

with pytest.raises(ValueError, match="batch size 2"):
manager.loop.run_until_complete(manager.run_single(data))


@pytest.mark.parametrize("compute_score", [default_compute_score, partial(default_compute_score)])
def test_default_text_reward_function_is_rejected(compute_score):
with pytest.raises(ValueError, match="custom_reward_function"):
_manager(compute_score)


def test_run_single_passes_waveform_and_returns_diagnostics():
def compute_score(data_source, solution_audio, ground_truth, extra_info):
waveform, sample_rate = solution_audio
assert data_source == "tts_reward"
assert waveform.dtype == np.float32
assert waveform.shape == (24_000,)
assert sample_rate == 24_000
assert ground_truth == "ni3 hao3"
assert extra_info["id"] == "sample-0"
assert "global_steps" not in extra_info
return {"score": 0.75, "pinyin_error_rate": 0.1}

manager = _manager(compute_score)
result = manager.loop.run_until_complete(manager.run_single(_data(np.ones(24_000, dtype=np.float32))))

assert result == {
"reward_score": 0.75,
"reward_extra_info": {"pinyin_error_rate": 0.1},
}


def test_run_single_forwards_reward_router_arguments():
expected_reward_model_tokenizer = MagicMock()

def compute_score(reward_router_address, reward_model_tokenizer, model_name, **kwargs):
assert reward_router_address == "reward-router:8000"
assert reward_model_tokenizer is expected_reward_model_tokenizer
assert model_name == "reward-model"
return 0.5

config = OmegaConf.create({"reward": {"reward_model": {"model_path": "reward-model"}}})
manager = AudioRewardManager(
config,
MagicMock(),
compute_score=compute_score,
reward_router_address="reward-router:8000",
reward_model_tokenizer=expected_reward_model_tokenizer,
)

result = manager.loop.run_until_complete(manager.run_single(_data(np.ones(8, dtype=np.float32))))

assert result["reward_score"] == 0.5


def test_run_single_reads_finalized_top_level_audio_layout():
def compute_score(solution_audio, extra_info, **kwargs):
waveform, sample_rate = solution_audio
np.testing.assert_array_equal(waveform, np.ones(8, dtype=np.float32))
assert sample_rate == 16_000
assert extra_info["id"] == "sample-0"
return 0.5

manager = _manager(compute_score)
result = manager.loop.run_until_complete(
manager.run_single(_data(np.ones(8, dtype=np.float32), 16_000, layout="finalized"))
)

assert result["reward_score"] == 0.5


@pytest.mark.parametrize(
("data", "message"),
[
(_data(), "requires extra_info\\['audio'\\]"),
(_data([], 24_000), "empty waveform"),
(_data([0.0, float("nan")], 24_000), "NaN or infinity"),
(_data([0.0], 0), "positive integer"),
(_data([0.0], 24_000.5), "positive integer"),
],
)
def test_invalid_audio_fails_closed(data, message):
manager = _manager(lambda **kwargs: 0.0)

with pytest.raises((KeyError, ValueError), match=message):
manager.loop.run_until_complete(manager.run_single(data))


def test_missing_sample_rate_fails_closed():
data = _data([0.0])
data.non_tensor_batch["tool_extra_fields"][0].pop("audio_sample_rate")
manager = _manager(lambda **kwargs: 0.0)

with pytest.raises(KeyError, match="audio_sample_rate"):
manager.loop.run_until_complete(manager.run_single(data))


def test_chunked_waveform_is_rejected_instead_of_guessed():
manager = _manager(lambda **kwargs: 0.5)

with pytest.raises(ValueError, match="could not convert"):
manager.loop.run_until_complete(
manager.run_single(_data([torch.tensor([0.1, 0.2]), torch.tensor([0.3])], 16_000))
)


def test_two_dimensional_waveform_is_rejected_instead_of_downmixed_on_the_wrong_axis():
manager = _manager(lambda **kwargs: 0.5)

with pytest.raises(ValueError, match="one mono waveform"):
manager.loop.run_until_complete(manager.run_single(_data(np.zeros((128, 2)), 16_000)))


@pytest.mark.asyncio
async def test_async_score_function_is_supported():
async def compute_score(solution_audio, **kwargs):
assert solution_audio[1] == 24_000
return {"score": -0.25, "judge_margin": 3.0}

result = await _manager(compute_score).run_single(_data(torch.ones(32)))

assert result == {"reward_score": -0.25, "reward_extra_info": {"judge_margin": 3.0}}


@pytest.mark.asyncio
async def test_waveform_extraction_runs_off_the_event_loop(monkeypatch):
event_loop_thread = threading.get_ident()
extraction_threads = []

def extract_audio(extra_info):
extraction_threads.append(threading.get_ident())
return np.zeros(8, dtype=np.float32), 24_000

async def compute_score(**kwargs):
return 0.5

monkeypatch.setattr(AudioRewardManager, "_extract_audio", staticmethod(extract_audio))
result = await _manager(compute_score).run_single(_data(np.ones(8, dtype=np.float32)))

assert result["reward_score"] == 0.5
assert extraction_threads and extraction_threads[0] != event_loop_thread


@pytest.mark.parametrize("score", [float("nan"), float("inf"), -float("inf")])
def test_non_finite_reward_fails_closed(score):
manager = _manager(lambda **kwargs: score)

with pytest.raises(ValueError, match="must be finite"):
manager.loop.run_until_complete(manager.run_single(_data([0.0])))


def test_reward_dictionary_requires_score():
manager = _manager(lambda **kwargs: {"metric": 1.0})

with pytest.raises(ValueError, match="missing 'score'"):
manager.loop.run_until_complete(manager.run_single(_data([0.0])))
Loading
Loading