Skip to content
Draft
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
44 changes: 44 additions & 0 deletions docs/source_en/Usage Guide/Agentic-Evaluator.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
# Agentic Evaluator

Install evaluation support separately:

```bash
pip install 'twinkle-kit[eval]'
```

`twinkle_agentic.evaluator.Evaluator` is a single-use facade for EvalScope's native runner. EvalScope owns datasets, agent loops, tools, judges, caches, metrics, and reports; Twinkle only adapts the candidate model boundary.

## Protocol API

```python
from twinkle_agentic.evaluator import Evaluator
from twinkle_agentic.protocol.openai import OpenAI

reports = Evaluator(
api=OpenAI(model='qwen-plus', api_key='...', base_url='https://example.com/v1'),
datasets=['gsm8k', 'bfcl_v4'],
task_config={'generation_config': {'temperature': 0.0}},
).run()
```

## Sampler

```python
from twinkle_agentic.evaluator import Evaluator

reports = Evaluator(
sampler=sampler,
datasets=['gsm8k'],
template=template,
sampler_kwargs={'adapter_uri': 'twinkle://my-adapter'},
task_config={'limit': 100, 'generation_config': {'max_tokens': 2048}},
).run()
```

The sampler path micro-batches compatible EvalScope requests. `template.parse_tool_call()` is required for tool benchmarks unless the sampler returns structured assistant tool calls. The HTTP sampler accepts either Twinkle `SamplingParams` or a mapping.

## Capability boundary

Text/chat, function-calling, and EvalScope native AgentLoop tasks are supported. Image, audio, video, streaming, non-native EvalScope backends, and nested Twinkle rollout loops are not. Explicit generation options are mapped exactly or rejected before evaluation begins; Twinkle does not approximate unsupported parameters.

`run()` returns EvalScope's `dict[str, Report]` unchanged. `resolved_task_config` exposes the constructed EvalScope configuration, and `output_dir` becomes available after EvalScope resolves the run directory. See EvalScope `TaskConfig` for pass-through fields and install its benchmark-specific extras separately when needed.
1 change: 1 addition & 0 deletions docs/source_en/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ Twinkle DOCUMENTATION
Usage Guide/NPU-Support.md
Usage Guide/Train-as-a-Service.md
Usage Guide/Agentic-RL-Deployment-and-Training.md
Usage Guide/Agentic-Evaluator.md
Usage Guide/Introduction-with-Qwen3.5.md
Usage Guide/Embedding-Training.md

Expand Down
1 change: 1 addition & 0 deletions docs/source_zh/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ Twinkle DOCUMENTATION
使用指引/NPU的支持.md
使用指引/训练服务.md
使用指引/Agentic RL部署与训练.md
使用指引/Agentic评测.md
使用指引/Qwen3.5最佳实践.md
使用指引/Embedding训练.md

Expand Down
42 changes: 42 additions & 0 deletions docs/source_zh/使用指引/Agentic评测.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
# Agentic 评测

评测能力是可选依赖:

```bash
pip install 'twinkle-kit[eval]'
```

`twinkle_agentic.evaluator.Evaluator` 是 EvalScope Native runner 的一次性轻量封装。数据集、AgentLoop、工具、Judge、缓存、指标和报告仍由 EvalScope 负责,Twinkle 只适配被评测模型。

## Protocol API

```python
from twinkle_agentic.evaluator import Evaluator
from twinkle_agentic.protocol.openai import OpenAI

reports = Evaluator(
api=OpenAI(model='qwen-plus', api_key='...', base_url='https://example.com/v1'),
datasets=['gsm8k', 'bfcl_v4'],
task_config={'generation_config': {'temperature': 0.0}},
).run()
```

## Sampler

```python
reports = Evaluator(
sampler=sampler,
datasets=['gsm8k'],
template=template,
sampler_kwargs={'adapter_uri': 'twinkle://my-adapter'},
task_config={'limit': 100, 'generation_config': {'max_tokens': 2048}},
).run()
```

Sampler 路径会将兼容的 EvalScope 请求微批处理。工具评测要求 sampler 返回结构化 tool calls,或提供 `template.parse_tool_call()`。HTTP sampler 同时接受 Twinkle `SamplingParams` 和字典。

## 能力边界

支持文本/对话、函数调用和 EvalScope Native AgentLoop;不支持多模态、流式、EvalScope 非 Native backend,以及在 adapter 内嵌套 Twinkle rollout。显式 generation 参数必须精确映射,否则会在评测前报错,不会静默近似或忽略。

`run()` 原样返回 EvalScope 的 `dict[str, Report]`。`resolved_task_config` 可获取实际构建的配置,EvalScope 创建输出目录后可通过 `output_dir` 获取。其余透传字段请参考 EvalScope `TaskConfig`;benchmark 的额外依赖仍应按 EvalScope 文档单独安装。
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ twinkle-server = "twinkle.server.cli:main"
twinkle-auto = "twinkle_client.auto:main"

[project.optional-dependencies]
eval = ["evalscope>=1.11,<1.12"]
megatron = ["megatron-core>=0.12.0", "transformer-engine[pytorch]", "mcore_bridge"]
data = ["py-data-juicer"]
rl = [
Expand Down
12 changes: 12 additions & 0 deletions src/twinkle_agentic/evaluator/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
"""EvalScope-backed evaluation for Twinkle Agentic backends."""

from ._contracts import BackendContractError, EvaluatorConfigError, SamplerBatchError, UnsupportedCapabilityError
from .evaluator import Evaluator

__all__ = [
'BackendContractError',
'Evaluator',
'EvaluatorConfigError',
'SamplerBatchError',
'UnsupportedCapabilityError',
]
157 changes: 157 additions & 0 deletions src/twinkle_agentic/evaluator/_batcher.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
"""A single-worker micro-batcher for structurally compatible samplers."""

from collections import deque
from concurrent.futures import Future
from dataclasses import dataclass
from threading import Condition, Thread
from time import monotonic
from typing import Any, Hashable, Mapping

from twinkle.data_format import SamplingParams, Trajectory

from ._contracts import BackendContractError, SamplerBatchError


def _freeze(value: Any) -> Hashable:
if isinstance(value, Mapping):
return tuple(sorted((key, _freeze(item)) for key, item in value.items()))
if isinstance(value, (list, tuple)):
return tuple(_freeze(item) for item in value)
if isinstance(value, (str, int, float, bool, type(None))):
return value
try:
hash(value)
except TypeError as exc:
raise ValueError(f'Cannot safely batch unhashable value of type {type(value).__name__}') from exc
return value


def _params_key(params: SamplingParams) -> Hashable:
return _freeze(vars(params))


@dataclass
class _BatchRequest:
trajectory: Trajectory
sampling_params: SamplingParams
sampler_kwargs: Mapping[str, Any]
compatibility_key: Hashable
future: Future
request_id: int


class SamplerBatcher:
"""Serialize sampler calls while coalescing equal requests into batches."""

def __init__(self, sampler: Any, *, batch_size: int, batch_wait_ms: float, sampler_kwargs: Mapping[str, Any]):
self._sampler = sampler
self._batch_size = batch_size
self._batch_wait_seconds = batch_wait_ms / 1000
self._sampler_kwargs = dict(sampler_kwargs)
self._queue: deque[_BatchRequest] = deque()
self._condition = Condition()
self._closed = False
self._request_id = 0
self._worker = Thread(target=self._run, name='twinkle-evaluator-sampler-batcher', daemon=False)
self._worker.start()

def submit(self, trajectory: Trajectory, sampling_params: SamplingParams) -> Any:
future: Future = Future()
key = (_params_key(sampling_params), _freeze(self._sampler_kwargs))
with self._condition:
if self._closed:
raise RuntimeError('Sampler batcher is closed')
request = _BatchRequest(
trajectory=trajectory,
sampling_params=sampling_params,
sampler_kwargs=self._sampler_kwargs,
compatibility_key=key,
future=future,
request_id=self._request_id,
)
self._request_id += 1
self._queue.append(request)
self._condition.notify()
return future.result()

def _pop_first(self) -> _BatchRequest | None:
return self._queue.popleft() if self._queue else None

def _take_compatible(self, key: Hashable, capacity: int) -> list[_BatchRequest]:
selected: list[_BatchRequest] = []
retained: deque[_BatchRequest] = deque()
while self._queue:
request = self._queue.popleft()
if request.compatibility_key == key and len(selected) < capacity:
selected.append(request)
else:
retained.append(request)
self._queue = retained
return selected

def _minimum_physical_batch_size(self) -> int:
mesh = getattr(self._sampler, 'device_mesh', None)
for name in ('data_world_size', 'dp_world_size'):
value = getattr(mesh, name, None)
if isinstance(value, int) and value > 0:
return value
return 1

def _complete_batch(self, requests: list[_BatchRequest]) -> None:
inputs = [request.trajectory for request in requests]
physical_size = max(len(inputs), self._minimum_physical_batch_size())
physical_inputs = inputs + [inputs[-1]] * (physical_size - len(inputs))
try:
responses = list(self._sampler.sample(
physical_inputs,
sampling_params=requests[0].sampling_params,
**requests[0].sampler_kwargs,
))
if len(responses) != physical_size:
raise BackendContractError(
f'Sampler returned {len(responses)} responses for physical batch size {physical_size}')
except Exception as exc:
error = SamplerBatchError(f'Sampler batch for {len(requests)} request(s) failed: {exc}')
for request in requests:
if not request.future.done():
request.future.set_exception(error)
return
for request, response in zip(requests, responses):
if not request.future.done():
request.future.set_result(response)

def _run(self) -> None:
while True:
with self._condition:
while not self._queue and not self._closed:
self._condition.wait()
if self._closed and not self._queue:
return
first = self._pop_first()
assert first is not None
deadline = monotonic() + self._batch_wait_seconds
selected = [first]
selected.extend(self._take_compatible(first.compatibility_key, self._batch_size - len(selected)))
while len(selected) < self._batch_size and not self._closed:
remaining = deadline - monotonic()
if remaining <= 0:
break
self._condition.wait(remaining)
selected.extend(self._take_compatible(first.compatibility_key, self._batch_size - len(selected)))
if len(selected) > self._batch_size:
overflow = selected[self._batch_size:]
selected = selected[:self._batch_size]
self._queue.extendleft(reversed(overflow))
self._complete_batch(selected)

def close(self) -> None:
with self._condition:
if self._closed:
return
self._closed = True
while self._queue:
request = self._queue.popleft()
if not request.future.done():
request.future.set_exception(RuntimeError('Sampler batcher was closed'))
self._condition.notify_all()
self._worker.join()
38 changes: 38 additions & 0 deletions src/twinkle_agentic/evaluator/_contracts.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
"""Small public-facing contracts shared by evaluator internals."""

from typing import Any, Mapping, Protocol, Sequence, runtime_checkable

from twinkle.data_format import SamplingParams, Trajectory


class EvaluatorConfigError(ValueError):
"""The evaluator constructor or its owned configuration is invalid."""


class UnsupportedCapabilityError(ValueError):
"""The selected backend cannot represent an explicit request exactly."""


class BackendContractError(RuntimeError):
"""A caller-owned API or sampler returned an invalid value."""


class SamplerBatchError(RuntimeError):
"""A single physical sampler batch failed."""


@runtime_checkable
class SamplerLike(Protocol):
def sample(
self,
inputs: list[Trajectory],
sampling_params: SamplingParams | Mapping[str, Any],
**kwargs: Any,
) -> Sequence[Any]: ...


def read_value(value: Any, name: str, default: Any = None) -> Any:
"""Read an attribute or mapping key without treating falsy values as absent."""
if isinstance(value, Mapping):
return value.get(name, default)
return getattr(value, name, default)
Loading
Loading