Repository navigation
Expand file tree
/
Copy pathsolve.py
More file actions
3272 lines (2961 loc) · 130 KB
/
Copy pathsolve.py
File metadata and controls
3272 lines (2961 loc) · 130 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""solve.py — drive coding agents to solve NatureBench task packages.
Architecture overview:
1. Start the Evaluation Service (HTTP) on the host
2. Launch a solver container for each task
3. Container mounts: /task/problem/ (read-only), /workspace (read-write)
4. The evaluation/ directory is not mounted into the container (the agent cannot see the answers directly)
5. Start the agent CLI inside the container (claude/codex)
6. The agent calls the Evaluation Service over HTTP to get scores and iterate
7. After the timeout the container is force-stopped, and the best score becomes the final score
Usage:
python solve.py --config config.yaml
python solve.py --task-set ./task-set/cpu.txt --data-dir ./data/tasks --out-dir ./results --agent claude --model <model>
"""
from __future__ import annotations
import argparse
import json
import logging
import os
import platform
import queue
import re
import secrets
import shlex
import shutil
import subprocess
import sys
import threading
import time
import urllib.request
import urllib.error
import urllib.parse
import uuid as _uuid_mod
from datetime import datetime, timezone
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Dict, List, Optional, Set, Tuple
from concurrent.futures import ThreadPoolExecutor, as_completed
from config_loader import load_yaml_config, merge_args_with_config
from eval_service import (
CONTROL_TOKEN_HEADER,
DEFAULT_CONTROL_TOKEN_PATH,
ScoreTracker,
load_or_create_control_token,
start_server_background,
)
from tqdm import tqdm
# Canonical in-container command builders + resume prompts. These live in
# agent/cli_commands.py so solve.py and the agent adapter layer share one
# source of truth for what gets run inside the container. The leading-underscore
# aliases preserve the names used throughout this module.
from agent.cli_commands import (
_CODEX_HOME,
build_claude_cmd as _build_claude_cmd,
build_gemini_cmd as _build_gemini_cmd,
build_codex_cmd as _build_codex_cmd,
codex_exec_cmd as _codex_exec_cmd,
RESUME_PROMPT_CLAUDE as _RESUME_PROMPT_CLAUDE,
RESUME_PROMPT_CODEX as _RESUME_PROMPT_CODEX,
RESUME_PROMPT_GEMINI as _RESUME_PROMPT_GEMINI,
)
# Agent registry. Importing cli_adapters registers the built-in CLI adapters
# (claude/codex/gemini); custom agents register themselves the same way.
from agent.adapter import REGISTRY, AgentRunContext
import agent.cli_adapters # noqa: F401 (registers built-in adapters on import)
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [solve] %(levelname)s %(message)s",
)
logger = logging.getLogger("cns_bench.solve")
# ---------------------------------------------------------------------------
# Task-set reader (shared with eval.py)
# ---------------------------------------------------------------------------
def _read_task_set(task_file: Path) -> List[str]:
"""Read task names from a text file (one per line, or JSON lines)."""
if not task_file.exists():
raise FileNotFoundError(f"Task list file not found: {task_file}")
tasks: List[str] = []
with task_file.open("r", encoding="utf-8") as handle:
for raw_line in handle:
line = raw_line.strip()
if not line or line.startswith("#"):
continue
try:
spec = json.loads(line)
task_name = (
spec.get("task_name")
or spec.get("name")
or spec.get("task_file")
)
if not task_name:
raise ValueError
tasks.append(task_name)
continue
except (json.JSONDecodeError, ValueError):
pass
tasks.append(line)
if not tasks:
raise ValueError(f"No tasks found inside {task_file}")
return tasks
# ---------------------------------------------------------------------------
# Docker helpers
# ---------------------------------------------------------------------------
@dataclass
class DockerfileSpec:
"""Parsed Dockerfile instructions for skip-build mode."""
run_commands: List[str] = field(default_factory=list) # Shell commands from RUN
env_vars: Dict[str, str] = field(default_factory=dict) # ENV key=value
copy_srcs: List[Tuple[str, str]] = field(default_factory=list) # (src, dest) from COPY/ADD
def _parse_dockerfile(dockerfile: Path) -> DockerfileSpec:
"""Parse a task Dockerfile into structured instructions.
Handles: RUN (pip, apt-get, wget, git+, etc.), COPY/ADD, ENV.
Skips: FROM, comments, blank lines.
"""
spec = DockerfileSpec()
if not dockerfile.exists():
return spec
text = dockerfile.read_text(encoding="utf-8")
# Join backslash-continued lines
text = text.replace("\\\n", " ")
for line in text.splitlines():
line = line.strip()
if not line or line.startswith("#"):
continue
upper = line.upper()
if upper.startswith("FROM"):
continue
if upper.startswith("RUN "):
# Preserve the full shell command (apt-get, pip, wget, etc.)
cmd = line[4:].strip()
if cmd:
spec.run_commands.append(cmd)
elif upper.startswith("ENV "):
# ENV KEY=VALUE or ENV KEY VALUE
rest = line[4:].strip()
if "=" in rest:
key, _, val = rest.partition("=")
spec.env_vars[key.strip()] = val.strip().strip('"')
else:
parts = rest.split(None, 1)
if len(parts) == 2:
spec.env_vars[parts[0]] = parts[1].strip('"')
elif upper.startswith("COPY ") or upper.startswith("ADD "):
prefix_len = 5 if upper.startswith("COPY ") else 4
rest = line[prefix_len:].strip()
parts = rest.split()
if len(parts) >= 2:
spec.copy_srcs.append((parts[0], parts[-1]))
return spec
def _ensure_clash_proxy_container(
container_name: str,
network_name: str,
) -> None:
"""Validate that a shared Clash proxy container exists on the target network."""
inspect = subprocess.run(
["docker", "inspect", container_name],
capture_output=True,
text=True,
)
if inspect.returncode != 0:
raise RuntimeError(
f"Clash proxy container {container_name!r} does not exist. "
"Create it first before enabling Codex proxying."
)
try:
payload = json.loads(inspect.stdout)[0]
except (json.JSONDecodeError, IndexError, KeyError) as exc:
raise RuntimeError(
f"Failed to inspect Clash proxy container {container_name!r}: {exc}"
) from exc
running = bool(payload.get("State", {}).get("Running"))
networks = payload.get("NetworkSettings", {}).get("Networks", {}) or {}
if not running:
raise RuntimeError(
f"Clash proxy container {container_name!r} is not running. "
"Start it first before enabling Codex proxying."
)
if network_name not in networks:
joined = ", ".join(sorted(networks)) or "<none>"
raise RuntimeError(
f"Clash proxy container {container_name!r} is not attached to Docker "
f"network {network_name!r} (current: {joined})."
)
def _host_proxy_env_args() -> List[str]:
"""Pass through standard host proxy environment variables when present."""
args: List[str] = []
for key in (
"HTTP_PROXY", "HTTPS_PROXY", "ALL_PROXY", "NO_PROXY",
"http_proxy", "https_proxy", "all_proxy", "no_proxy",
):
val = os.environ.get(key)
if val:
args.extend(["-e", f"{key}={val}"])
return args
def _codex_auth_mount_args(codex_state_dir: Path) -> List[str]:
"""Mount per-task Codex state (auth + sessions) into the container's HOME."""
return [
"-v", f"{str(codex_state_dir)}:{_CODEX_HOME}/.codex",
]
def _prepare_task_codex_state_dir(
task_out_dir: Path,
source_auth_dir: Path,
*,
reseed: bool,
) -> Path:
"""Seed an isolated Codex state directory (auth + sessions) per task.
Fresh runs (reseed=True) wipe any prior state and copy **only** the auth
files the CLI needs (`auth.json`, `config.toml`). Session history stays
empty so that ``_capture_codex_session_id`` cannot pick a stale sid from
the host.
Resume runs (reseed=False) require the prior state dir to still exist;
never silently re-bootstrap it.
"""
state_dir = task_out_dir / ".codex_state"
if not reseed:
if not state_dir.is_dir():
raise RuntimeError(
f"Codex resume requested but {state_dir} is missing"
)
return state_dir
if state_dir.exists():
shutil.rmtree(state_dir)
state_dir.mkdir(parents=True)
for name in ("auth.json", "config.toml"):
src = source_auth_dir / name
if src.is_file():
shutil.copy2(src, state_dir / name)
(state_dir / "sessions").mkdir(exist_ok=True)
return state_dir
def _proxy_env_args(
proxy_host: str,
http_port: int,
socks_port: int,
) -> List[str]:
"""Build proxy env args for a Codex container.
``proxy_host`` is the address visible from inside the task container —
use ``127.0.0.1`` for embedded mode (clash runs in the same network
namespace) or the sidecar container name for sidecar mode.
"""
http_proxy = f"http://{proxy_host}:{http_port}"
socks_proxy = f"socks5://{proxy_host}:{socks_port}"
no_proxy_value = "127.0.0.1,localhost,::1,host.docker.internal"
merged_no_proxy = []
for key in ("NO_PROXY", "no_proxy"):
val = os.environ.get(key)
if val:
merged_no_proxy.extend([x.strip() for x in val.split(",") if x.strip()])
if merged_no_proxy:
merged_no_proxy = list(dict.fromkeys(no_proxy_value.split(",") + merged_no_proxy))
no_proxy_value = ",".join(merged_no_proxy)
return [
"-e", f"HTTP_PROXY={http_proxy}",
"-e", f"HTTPS_PROXY={http_proxy}",
"-e", f"ALL_PROXY={socks_proxy}",
"-e", f"NO_PROXY={no_proxy_value}",
"-e", f"http_proxy={http_proxy}",
"-e", f"https_proxy={http_proxy}",
"-e", f"all_proxy={socks_proxy}",
"-e", f"no_proxy={no_proxy_value}",
]
def _embedded_clash_setup_cmds() -> List[str]:
"""Setup commands that bring up the in-container Clash proxy.
The bundle is mounted read-only at /clash-bundle; clash needs a writable
config dir for its cache.db, so we copy into /tmp/clash-rw first.
Each returned string is joined with ``&&`` by the caller, so avoid bare
``&`` backgrounding mid-chain (that would create ``& &&`` which is a
bash parse error). The backgrounded launch + readiness loop is wrapped
into a single subshell entry.
"""
launch_and_wait = (
"{ "
"nohup /clash-bundle/clash -d /tmp/clash-rw > /tmp/clash.log 2>&1 & "
"for i in $(seq 1 60); do "
" curl -sf -m 1 http://127.0.0.1:9090/version >/dev/null 2>&1 && break; "
" sleep 0.5; "
"done; "
"curl -sf -m 1 http://127.0.0.1:9090/version >/dev/null "
"|| { echo '[clash] failed to start' >&2; tail -50 /tmp/clash.log >&2; exit 1; }; "
"}"
)
return [
"mkdir -p /tmp/clash-rw",
"cp -r /clash-bundle/config/. /tmp/clash-rw/",
launch_and_wait,
]
def _probe_codex_login(
image_tag: str,
codex_auth_dir: Path,
proxy_mode: str = "none",
proxy_container: Optional[str] = None,
proxy_network: Optional[str] = None,
proxy_bundle: Optional[Path] = None,
proxy_http_port: int = 7890,
proxy_socks_port: int = 7891,
) -> None:
"""Verify that Codex login state is usable in the same image/env as the task.
``proxy_mode`` selects how the probe reaches the network:
* ``host`` — pass host proxy environment variables
* ``none`` — no proxy, direct egress
* ``sidecar`` — attach to ``proxy_network`` and use the
``proxy_container`` container name as proxy host
* ``embedded`` — mount ``proxy_bundle`` and start clash inside
the probe container before running ``codex login status``
"""
probe_cmd: List[str] = ["docker", "run", "--rm"]
if proxy_mode == "sidecar":
if not proxy_container or not proxy_network:
raise ValueError(
"sidecar proxy mode requires proxy_container and proxy_network"
)
probe_cmd.extend(["--network", proxy_network])
if platform.system() == "Linux":
probe_cmd.extend(["--add-host=host.docker.internal:host-gateway"])
if proxy_mode == "sidecar":
probe_cmd.extend(
_proxy_env_args(
proxy_container,
proxy_http_port,
proxy_socks_port,
)
)
elif proxy_mode == "embedded":
if not proxy_bundle or not proxy_bundle.is_dir():
raise ValueError(
f"embedded proxy mode requires a clash bundle dir; got {proxy_bundle}"
)
probe_cmd.extend([
"-v", f"{str(proxy_bundle.resolve())}:/clash-bundle:ro",
])
probe_cmd.extend(
_proxy_env_args(
"127.0.0.1",
proxy_http_port,
proxy_socks_port,
)
)
elif proxy_mode == "host":
probe_cmd.extend(_host_proxy_env_args())
probe_cmd.extend(_codex_auth_mount_args(codex_auth_dir))
if proxy_mode == "embedded":
clash_setup = " && ".join(_embedded_clash_setup_cmds())
inner = f"{clash_setup} && HOME={_CODEX_HOME} codex login status"
else:
inner = f"HOME={_CODEX_HOME} codex login status"
probe_cmd.extend([
"--entrypoint", "/bin/bash",
image_tag,
"-lc", inner,
])
probe = subprocess.run(probe_cmd, capture_output=True, text=True)
if probe.returncode != 0:
output = (probe.stdout + "\n" + probe.stderr).strip()
if "not logged in" in output.lower() or "login" in output.lower():
hint = "Codex auth dir is not logged in"
else:
hint = "Codex probe failed"
raise RuntimeError(
f"{hint}. Run `codex login --device-auth` once on the host and ensure {codex_auth_dir} "
f"contains the login state; status output: {probe.stdout[-200:]} {probe.stderr[-200:]}"
)
def _ensure_task_image(task_name: str, data_dir: Path, base_image: str = "naturebench-base:v3", dockerfile_name: str = "Dockerfile.v3") -> str:
"""Ensure naturebench-task-<task>:base image exists. Build from Dockerfile if needed."""
image_tag = f"naturebench-task-{task_name}:base"
# Check if image already exists
result = subprocess.run(
["docker", "images", "-q", image_tag],
capture_output=True, text=True,
)
if result.stdout.strip():
logger.info("[%s] Image %s already exists", task_name, image_tag)
return image_tag
# Try to build from task Dockerfile
dockerfile = data_dir / task_name / "environment" / dockerfile_name
if dockerfile.exists():
logger.info("[%s] Building image %s from %s", task_name, image_tag, dockerfile)
build_cmd = [
"docker", "build",
"-t", image_tag,
"-f", str(dockerfile.resolve()),
str(dockerfile.parent.resolve()),
]
proc = subprocess.run(build_cmd, capture_output=True, text=True)
if proc.returncode != 0:
logger.error("[%s] Docker build failed:\n%s", task_name, proc.stderr[-2000:])
raise RuntimeError(f"Failed to build image {image_tag}")
logger.info("[%s] Image built: %s", task_name, image_tag)
return image_tag
# Fall back to base image
logger.warning("[%s] No Dockerfile found, using %s", task_name, base_image)
return base_image
def _remove_container(container_name: str) -> None:
"""Force remove a container if it exists."""
subprocess.run(
["docker", "rm", "-f", container_name],
capture_output=True, text=True,
)
# ---------------------------------------------------------------------------
# Single-task solver
# ---------------------------------------------------------------------------
def _host_url(eval_service_url: str) -> str:
"""Convert the container-visible host.docker.internal URL to the host localhost URL."""
return eval_service_url.replace("host.docker.internal", "localhost")
# ---------------------------------------------------------------------------
# Resume support helpers
# ---------------------------------------------------------------------------
# Default Claude Code mounts ~/.claude → /root/.claude inside the container.
# We pin the workspace cwd to /workspace, so the CLI stores its session jsonl at
# `~/.claude/projects/-workspace/<sid>.jsonl`.
# The resume prompts themselves are defined in agent/cli_commands.py and
# imported above as _RESUME_PROMPT_CLAUDE / _RESUME_PROMPT_CODEX /
# _RESUME_PROMPT_GEMINI.
# Backward-compat alias (in case anything still imports the old name).
_RESUME_PROMPT = _RESUME_PROMPT_CLAUDE
def _setup_session_id(
task_out_dir: Path,
task_name: str,
is_resume: bool,
) -> str:
"""Return the Claude session UUID for this task.
Behavior:
* is_resume=True : require existing claude_session_id.txt; raise if missing.
* is_resume=False : if a sid file already exists, archive it with a timestamp
suffix (.bak.<epoch>) before writing a fresh UUID, so we
never silently clobber a recoverable session.
"""
sid_file = task_out_dir / "claude_session_id.txt"
if is_resume:
if not sid_file.exists():
raise RuntimeError(
f"[{task_name}] resume requested but {sid_file} not found"
)
return sid_file.read_text(encoding="utf-8").strip()
if sid_file.exists():
bak = sid_file.with_suffix(f".bak.{int(time.time())}")
sid_file.rename(bak)
logger.warning(
"[%s] archived previous session id to %s (fresh run requested)",
task_name, bak.name,
)
new_sid = str(_uuid_mod.uuid4())
sid_file.write_text(new_sid, encoding="utf-8")
return new_sid
def _capture_codex_session_id(
state_dir: Path,
*,
since_mtime: float = 0.0,
) -> Optional[str]:
"""Extract this run's Codex session UUID from the per-task state dir.
Only filenames are inspected and files older than ``since_mtime`` are
skipped, to avoid grabbing a stale sid from a seeded session file.
"""
sessions = state_dir / "sessions"
if not sessions.is_dir():
return None
sid_re = re.compile(r"[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}")
candidates: List[Tuple[float, Path]] = []
for f in sessions.rglob("*"):
if not f.is_file():
continue
try:
mtime = f.stat().st_mtime
except OSError:
continue
if mtime < since_mtime:
continue
candidates.append((mtime, f))
candidates.sort(reverse=True)
for _, f in candidates:
m = sid_re.search(f.name)
if m:
return m.group(0)
return None
def _capture_gemini_session_id(jsonl_path: Path) -> Optional[str]:
"""Parse the stream-json init event from gemini.jsonl and return its UUID.
Gemini CLI's first stdout line in --output-format=stream-json mode is an
``init`` event that contains ``session_id``. We scan early lines (defensive
against minor format variations) and return the first uuid we find.
"""
if not jsonl_path.exists():
return None
sid_re = re.compile(
r"[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}"
)
try:
with jsonl_path.open("r", encoding="utf-8", errors="replace") as f:
for _ in range(64):
line = f.readline()
if not line:
break
try:
obj = json.loads(line)
except Exception:
continue
if obj.get("type") == "init":
sid = obj.get("session_id")
if isinstance(sid, str) and sid_re.fullmatch(sid):
return sid
# Fallback: any top-level "session_id" that looks like a uuid
sid = obj.get("session_id") if isinstance(obj, dict) else None
if isinstance(sid, str) and sid_re.fullmatch(sid):
return sid
except OSError:
return None
return None
def _has_prior_state(task_out_dir: Path) -> bool:
"""Return True if the task directory already contains state from a prior run."""
for name in (
"result.json",
"submissions.jsonl",
"claude_session_id.txt",
".claude_state",
"codex_session_id.txt",
".codex_state",
"claude.jsonl",
"claude.err",
"codex.jsonl",
"codex.err",
"judge_verdict.json",
):
if (task_out_dir / name).exists():
return True
return False
def _load_or_create_eval_token(
task_out_dir: Path,
control_token: Optional[str] = None,
) -> str:
"""Return the stable opaque token used by this task's agent session."""
task_out_dir.mkdir(parents=True, exist_ok=True)
token_path = task_out_dir / "eval_token.txt"
if token_path.exists():
token = token_path.read_text(encoding="utf-8").strip()
if token and token != control_token:
return token
return _replace_eval_token(task_out_dir, control_token)
def _replace_eval_token(
task_out_dir: Path,
control_token: Optional[str] = None,
) -> str:
"""Generate and atomically persist a new opaque evaluation token."""
task_out_dir.mkdir(parents=True, exist_ok=True)
token_path = task_out_dir / "eval_token.txt"
token = secrets.token_urlsafe(32)
while token == control_token:
token = secrets.token_urlsafe(32)
tmp_path = task_out_dir / ".eval_token.txt.tmp"
tmp_path.write_text(token + "\n", encoding="utf-8")
tmp_path.replace(token_path)
return token
def _load_resume_eval_token(task_out_dir: Path, control_token: str) -> str:
"""Read the existing task token for resume without creating or replacing it."""
token_path = task_out_dir / "eval_token.txt"
if not token_path.is_file():
raise RuntimeError(f"resume eval token file is missing: {token_path}")
try:
token = token_path.read_text(encoding="utf-8").strip()
except OSError as e:
raise RuntimeError(f"could not read resume eval token file {token_path}: {e}") from e
if not token:
raise RuntimeError(f"resume eval token file is empty: {token_path}")
if secrets.compare_digest(token, control_token):
raise RuntimeError("resume eval token conflicts with the control token")
return token
def _build_resume_eval_notice(
eval_service_url: str,
eval_token: str,
*,
include_submit: bool = False,
) -> str:
"""Provide resumed agents with the current opaque evaluation protocol."""
submit = ""
if include_submit:
submit = (
f'\nSubmit the most recent evaluation with:\n'
f'curl -s -X POST {eval_service_url}/submit '
f'-H "Content-Type: application/json" '
f'-d \'{{"eval_token":"{eval_token}"}}\'\n'
)
return (
"\n\n[EVALUATION ACCESS]\n"
f'Check service health: curl -s "{eval_service_url}/health"\n'
f'Evaluate with: curl -s -X POST {eval_service_url}/evaluate '
f'-H "Content-Type: application/json" '
f'-d \'{{"eval_token":"{eval_token}"}}\'\n'
f'Check best score: curl -s "{eval_service_url}/best_score?eval_token={eval_token}"\n'
f'Check remaining time: curl -s "{eval_service_url}/time_remaining?eval_token={eval_token}"\n'
f"{submit}"
)
def _archive_prior_state_for_force_fresh(task_out_dir: Path, task_name: str) -> None:
"""For --force-fresh: move ALL prior artifacts in the task output directory
out of the way (no delete) so the next fresh run starts from clean.
The task output directory (``out_dir/<task>``) holds only prior-run output —
the task's input data lives in the task package (``data_dir/<task>``) — and
this runs before any fresh artifact is created, so archiving the whole
directory is safe and agent-agnostic (it covers every built-in CLI and any
custom agent without a hardcoded file list). We rename rather than rm -rf to
keep history recoverable, and leave any earlier archive directories in place
rather than nesting them.
"""
if not task_out_dir.is_dir():
return
ts = int(time.time())
bak_root = task_out_dir / f"_force_fresh_archive_{ts}"
# Snapshot entries before creating bak_root so it is never moved into itself.
entries = [p for p in sorted(task_out_dir.iterdir())
if not p.name.startswith("_force_fresh_archive_")]
moved = []
for src in entries:
bak_root.mkdir(parents=True, exist_ok=True)
src.rename(bak_root / src.name)
moved.append(src.name)
if moved:
logger.warning(
"[%s] --force-fresh: archived %d item(s) into %s/",
task_name, len(moved), bak_root.name,
)
def _setup_gemini_session_id(
task_out_dir: Path,
task_name: str,
is_resume: bool,
) -> Optional[str]:
"""Return the Gemini session UUID for this task, or None for a fresh run.
Gemini CLI does NOT support pinning a fresh session to a caller-supplied
UUID (no --session-id flag). On fresh runs we therefore let the CLI
generate its own session id; ``_capture_gemini_session_id`` is responsible
for writing the real id to gemini_session_id.txt after the run completes.
On resume we require gemini_session_id.txt to exist and pass its uuid via
``--resume <uuid>``.
"""
sid_file = task_out_dir / "gemini_session_id.txt"
if is_resume:
if not sid_file.exists():
raise RuntimeError(
f"[{task_name}] gemini resume requested but {sid_file} not found"
)
return sid_file.read_text(encoding="utf-8").strip()
# Fresh run — archive old sid file (if any) so a previous run's id is not
# mistaken for the current one. The CLI itself will allocate a new UUID;
# we capture it from the stream-json ``init`` event after the run finishes.
if sid_file.exists():
ts = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%S")
sid_file.rename(task_out_dir / f"gemini_session_id_{ts}.txt.bak")
return None
def _load_resume_history(task_out_dir: Path) -> List[Dict[str, Any]]:
"""Read the cumulative resume_history from the existing result.json (if any)."""
res_path = task_out_dir / "result.json"
if not res_path.exists():
return []
try:
old = json.loads(res_path.read_text(encoding="utf-8"))
except Exception:
return []
history = old.get("resume_history") or []
if not isinstance(history, list):
return []
return history
def _resume_eligible(task_out_dir: Path, agent_name: str) -> Tuple[bool, str]:
"""Decide whether a task is eligible for resume.
Returns (ok, reason). Caller is expected to use this for soft validation
only — final guard is `_setup_session_id` requiring sid file.
"""
res_path = task_out_dir / "result.json"
if not res_path.exists():
return False, "no prior result.json"
if agent_name == "claude":
sid_path = task_out_dir / "claude_session_id.txt"
state_path = task_out_dir / ".claude_state"
if not sid_path.exists():
return False, "no claude_session_id.txt"
if not state_path.exists():
return False, "no .claude_state directory"
return True, "ok"
if agent_name == "codex":
sid_path = task_out_dir / "codex_session_id.txt"
state_path = task_out_dir / ".codex_state"
if not sid_path.exists():
return False, "no codex_session_id.txt"
if not state_path.exists():
return False, "no .codex_state directory"
return True, "ok"
if agent_name == "gemini":
sid_path = task_out_dir / "gemini_session_id.txt"
state_path = task_out_dir / ".gemini_state"
if not sid_path.exists():
return False, "no gemini_session_id.txt"
if not state_path.exists():
return False, "no .gemini_state directory"
return True, "ok"
return False, f"resume not supported for agent {agent_name}"
def _query_time_remaining(
eval_service_url: str,
eval_token: str,
task_name: str,
) -> Optional[float]:
"""GET /time_remaining; return remaining seconds or None on error."""
host_url = _host_url(eval_service_url)
url = f"{host_url}/time_remaining?eval_token={urllib.parse.quote(eval_token)}"
try:
with urllib.request.urlopen(url, timeout=10) as resp:
data = json.loads(resp.read().decode())
return data.get("remaining_seconds")
except Exception as e:
logger.warning("[%s] time_remaining query failed: %s", task_name, e)
return None
def _control_headers(control_token: str) -> Dict[str, str]:
return {
"Content-Type": "application/json",
CONTROL_TOKEN_HEADER: control_token,
}
def _notify_timer_resume(
eval_service_url: str,
task_name: str,
control_token: str,
batch_name: Optional[str] = None,
) -> None:
"""POST /resume_timer; ignore failures (eval_service may not support it)."""
host_url = _host_url(eval_service_url)
url = f"{host_url}/resume_timer"
body = {"task_name": task_name}
if batch_name:
body["batch_name"] = batch_name
try:
req = urllib.request.Request(
url,
data=json.dumps(body).encode(),
headers=_control_headers(control_token),
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as _resp:
_ = _resp.read()
except Exception as e:
logger.debug("[%s] /resume_timer not honored: %s", task_name, e)
def _notify_timer_pause(
eval_service_url: str,
task_name: str,
control_token: str,
batch_name: Optional[str] = None,
) -> None:
"""POST /pause_timer when the agent container is stopping but the task is
not "done" (i.e. could resume later). Ignore HTTP errors so legacy
eval_service builds without /pause_timer don't break solve.py.
"""
host_url = _host_url(eval_service_url)
url = f"{host_url}/pause_timer"
body = {"task_name": task_name}
if batch_name:
body["batch_name"] = batch_name
try:
req = urllib.request.Request(
url,
data=json.dumps(body).encode(),
headers=_control_headers(control_token),
method="POST",
)
with urllib.request.urlopen(req, timeout=10) as _resp:
_ = _resp.read()
except Exception as e:
logger.debug("[%s] /pause_timer not honored: %s", task_name, e)
def _wait_eval_drain(eval_service_url: str, task_name: str,
eval_token: str,
poll_interval: float = 2.0,
max_unreachable_polls: int = 30) -> bool:
"""Wait until the remote eval_service has no in-flight evaluator for task.
In external-eval mode, solve.py's local tracker does not know the real
active_evals count, so we must poll /time_remaining and wait for
is_paused=False before calling /pause_timer. Otherwise /pause_timer is
rejected with 409 while evaluator threads are still running, and the task
timer keeps advancing after the agent container exits.
There is intentionally NO fixed time cap: an evaluator invocation may run
up to its own 3600s subprocess cap, so a short deadline would abandon a
legitimately-running evaluation and re-introduce the timer leak this
function exists to prevent. The wait terminates when:
- is_paused flips to False (evaluation drained) -> return True
- the task is unknown (404) -> return True
- the eval_service is unreachable for
``max_unreachable_polls`` consecutive polls -> return False
(the service is dead; there is no live timer left to protect).
While the wait is in progress the task timer is already paused
(active_evals > 0), so this wait does not consume the agent's solve budget.
"""
host_url = _host_url(eval_service_url)
url = f"{host_url}/time_remaining?eval_token={urllib.parse.quote(eval_token)}"
unreachable = 0
waited = 0.0
while True:
try:
with urllib.request.urlopen(url, timeout=10) as resp:
data = json.loads(resp.read().decode())
unreachable = 0
if not data.get("is_paused", False):
return True
except urllib.error.HTTPError as e:
if e.code == 404:
return True
# Service is alive (it answered), just an HTTP-level error.
unreachable = 0
logger.debug("[%s] drain poll HTTP error: %s", task_name, e)
except Exception as e:
unreachable += 1
logger.debug("[%s] drain poll failed (%d/%d): %s",
task_name, unreachable, max_unreachable_polls, e)
if unreachable >= max_unreachable_polls:
logger.warning(
"[%s] eval_service unreachable for %d consecutive polls; "
"abandoning drain wait (/pause_timer may be skipped)",
task_name, unreachable,
)
return False
time.sleep(poll_interval)
waited += poll_interval
if waited % 60 < poll_interval:
logger.info(
"[%s] still waiting for in-flight evaluation to drain (%.0fs)...",
task_name, waited,
)
# ---------------------------------------------------------------------------
# GPU pool (in-memory or cross-process file-backed)
# ---------------------------------------------------------------------------
@dataclass
class _GpuLease:
"""A single task's GPU allocation.
`kind="normal"` means `gpu_id` came from the existing exclusive GPU pool.
`kind="shared"` means `slot_id` came from the shared slot pool, while Docker
still receives the full physical `gpu_id` via `--gpus device=<gpu_id>`.
"""
gpu_id: int
kind: str
task_name: str
slot_id: Optional[int] = None
class _InMemoryGpuPool:
"""Original single-process queue-based pool. `get(block=True)` blocks
until a GPU is freed; `get(block=False)` raises queue.Empty."""
def __init__(self, gpu_ids: List[int]) -> None:
self._q: "queue.Queue[int]" = queue.Queue()
for gid in gpu_ids:
self._q.put(gid)
def get(self, block: bool = True, timeout: Optional[float] = None) -> int:
return self._q.get(block=block, timeout=timeout)
def put(self, gid: int) -> None:
self._q.put(gid)
class _FileGpuPool:
"""Cross-process GPU pool. Coordinates multiple solve.py instances via a
JSON file under fcntl.flock. Acquire blocks until a GPU is free; release
marks it free. Stale holders (dead pid) are auto-reclaimed.
State file schema:
{"gpus": {"<id>": {"holder_pid": int|null, "task": str|null,
"acquired_at": float|null}}}
"""
POLL_INTERVAL = 1.0 # seconds between retries when blocked
_PROBE_CACHE_TTL = 5.0 # seconds — reuse nvidia-smi result within this window
def __init__(self, gpu_ids: List[int], pool_path: Path,
skip_busy_mb: int = 0, skip_busy_util: int = 100) -> None:
import fcntl as _fcntl
self._fcntl = _fcntl
self._pool_path = pool_path
self._lock_path = pool_path.with_suffix(pool_path.suffix + ".lock")
self._registered: List[int] = list(gpu_ids)
self._held: List[int] = [] # GPUs this process currently holds
# External-busy detection thresholds. skip_busy_mb<=0 disables.
self._skip_busy_mb = int(skip_busy_mb)
self._skip_busy_util = int(skip_busy_util)
self._probe_cache: Dict[int, tuple] = {} # gid -> (mem_mb, util_pct, ts)
# initialize / merge own gpu_ids into the file pool
self._with_lock(lambda state: self._merge_ids(state, gpu_ids))
def _probe_external_busy(self, gid: int) -> bool:
"""Return True if GPU `gid` is currently used by an *external* process
(memory or utilization above thresholds). Cached for _PROBE_CACHE_TTL
to limit nvidia-smi calls. Falls back to "not busy" on probe error.
"""