diff --git a/evalharness/agent/envs/tau2_official.py b/evalharness/agent/envs/tau2_official.py index 6b12ddd..bb0d6b5 100644 --- a/evalharness/agent/envs/tau2_official.py +++ b/evalharness/agent/envs/tau2_official.py @@ -123,8 +123,17 @@ class Tau2Environment(Environment): task = Task.model_validate(task_json if not isinstance(task_json, str) else json.loads(task_json)) domain = (sample.metadata or {}).get('domain') or 'airline' - res = run_task(domain=domain, task=task, agent='llm_agent', - user='user_simulator', max_steps=max_turns) + # the official engine is SYNCHRONOUS and calls our adapter back via a + # private event loop in a worker thread; run it off the main loop so it + # never blocks the runner's other benches + import concurrent.futures + + loop = asyncio.get_running_loop() + with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool: + res = await loop.run_in_executor( + pool, lambda: run_task(domain=domain, task=task, + agent='llm_agent', user='user_simulator', + max_steps=max_turns)) rewards = {} try: info = res.reward_info diff --git a/evalharness/model/adapter.py b/evalharness/model/adapter.py index b4f8925..9e24153 100644 --- a/evalharness/model/adapter.py +++ b/evalharness/model/adapter.py @@ -25,6 +25,10 @@ from .output import ModelOutput, ToolCall, Usage ADAPTER_REGISTRY = EvalRegistry('model adapter') +_ADAPTER_CACHE = {} # spec -> shared instance; keeps pool round-robin state + # GLOBAL across benches (else each pool restarts at the + # first backend and starves the rest) + def register_adapter(name: str): def decorator(cls): @@ -87,9 +91,13 @@ def resolve_adapter(spec: str, deploy_fn=None) -> ModelAdapter: endpoint = deploy_fn(m.group('deployer'), m.group('model')) spec = f"openai/{endpoint['api_base']}?{endpoint['model']}" parsed = parse_model_spec(spec) + if spec in _ADAPTER_CACHE: + return _ADAPTER_CACHE[spec] cls = ADAPTER_REGISTRY.get(parsed['adapter']) key = parsed.get('api_base') and _key_for(parsed['api_base']) - return cls(model=parsed['model'], api_base=parsed['api_base'], api_key=key) + inst = cls(model=parsed['model'], api_base=parsed['api_base'], api_key=key) + _ADAPTER_CACHE[spec] = inst + return inst def _key_for(api_base: str) -> str: @@ -219,6 +227,10 @@ class OpenAICompatible(ModelAdapter): args = fn.get('arguments') or '{}' try: args_dict = json.loads(args) + if isinstance(args_dict, str): # double-encoded JSON string + args_dict = json.loads(args_dict) + if not isinstance(args_dict, dict): + args_dict = {'raw': args_dict} except (ValueError, TypeError): args_dict = {} calls.append(ToolCall(id=c.get('id', ''), name=fn.get('name', ''), diff --git a/evalharness/model/runner.py b/evalharness/model/runner.py index 064b4e7..051f7db 100644 --- a/evalharness/model/runner.py +++ b/evalharness/model/runner.py @@ -382,7 +382,16 @@ def _make_adapter(spec: str) -> ModelAdapter: - 'openai-pool/?model' with {port} placeholder: e.g. 'openai-pool/http://127.0.0.1:{8123..8130}/v1?Qwen3-8B' -> N ports - else resolve_adapter(spec) single endpoint + + Pooled specs are CACHED per spec: all benches share one pool so the + round-robin counter stays global (independent pools would each restart + at the first backend and starve the rest). """ + from .adapter import _ADAPTER_CACHE as _CACHE + + cache_key = spec + if cache_key in _CACHE: + return _CACHE[cache_key] opts = {} while True: for f in ('!nothink', '!textools'): @@ -417,6 +426,7 @@ def _make_adapter(spec: str) -> ModelAdapter: a.extra['no_think'] = True if opts.get('!textools'): a.extra['tools_mode'] = 'text' + _CACHE[cache_key] = adapter return adapter