Skip to content

Commit 0e5940f

Browse files
samsjacursoragent
andauthored
refactor logging on environment in orch (#1594)
* remove logging * add per env logging * log to one file when having multiple worker * intercept verifier log * intercept verifier log * intercept verifier log * Add worker tag to env worker logger for identification Co-authored-by: sami <sami@primeintellect.ai> * Update changelog for logging config changes Co-authored-by: sami <sami@primeintellect.ai> * Rename env worker log config flag Co-authored-by: sami <sami@primeintellect.ai> --------- Signed-off-by: samsja <55492238+samsja@users.noreply.github.com> Co-authored-by: sami jaghouar <sami@primeintellect.ai> Co-authored-by: Cursor Agent <cursoragent@cursor.com>
1 parent 969df59 commit 0e5940f

12 files changed

Lines changed: 77 additions & 70 deletions

File tree

CHANGELOG.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,4 +33,6 @@ Documenting changes which affect configuration usage patterns (added/moved/remov
3333
- **`model.lora.alpha`**: Changed default from 16.0 to 32.0 (2026-01-10)
3434
- **`orchestrator.env.log`**: Added logging configuration for environment workers. If set, enables logging with `level` (str, default: "warn") and `vf_level` (str, default: "warn") fields. If None (default), logging is disabled (#1561, 2026-01-13)
3535
- **`eval.watcher`**: Added flag (default `False`) to watch `weights_dir` for newly-created stable checkpoints and evaluate them as they appear (2026-01-14)
36+
- **`orchestrator.log.env_worker_logs`**: Added flag (default `True`) to write env worker logs to `logs/env_workers/{env_name}.log` (2026-01-15)
37+
- **`orchestrator.env.log`**: Removed. Use `orchestrator.log` for env worker logging instead (2026-01-15)
3638
- **`orchestrator.eval.retry.reraise`**: Changed default from `True` to `False`. When `False`, raises `tenacity.RetryError` after retries are exhausted instead of the original exception, allowing failed eval environments to be skipped with a warning (#1586, 2026-01-14)

docs/logging.md

Lines changed: 11 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -57,32 +57,28 @@ For RL training, logs are organized by component:
5757

5858
## Per-Environment Worker Logging
5959

60-
Environment workers run in **separate subprocesses** to isolate event loop lag. By default, they don't write to separate log files. To enable per-env logging, set the `log` field on the environment config:
60+
Environment workers run in **separate subprocesses** to isolate event loop lag. Worker logging is controlled at the orchestrator level via `orchestrator.log`:
6161

6262
```toml
63-
[[orchestrator.env]]
64-
id = "math-env"
65-
name = "math"
66-
log = { level = "debug", vf_level = "info" }
63+
[orchestrator.log]
64+
level = "debug" # Log level for prime-rl logger
65+
vf_level = "info" # Log level for verifiers library
66+
env_worker_logs = true # Enable file logging for env workers
6767
```
6868

69-
This produces:
69+
When `env_worker_logs = true`, logs are written to:
7070
```
7171
output_dir/
7272
└── logs/
7373
└── env_workers/
74-
└── math/
75-
├── worker_0.log
76-
├── worker_1.log
77-
└── ...
74+
├── {env_name_1}.log
75+
├── {env_name_2}.log
76+
└── ...
7877
```
7978

80-
### Configuration Options
79+
All workers for an environment share the same log file.
8180

82-
| Field | Description | Default |
83-
|-------|-------------|---------|
84-
| `level` | Log level for prime-rl logger (`debug`, `info`, `warn`, `error`) | `warn` |
85-
| `vf_level` | Log level for verifiers library logger | `warn` |
81+
Set `env_worker_logs = false` to disable worker file logging (workers inherit parent process logging).
8682

8783
## Torchrun
8884

src/prime_rl/eval/eval.py

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,5 @@
11
import asyncio
22

3-
import verifiers as vf
4-
53
from prime_rl.eval.config import OfflineEvalConfig
64
from prime_rl.eval.utils import run_evals
75
from prime_rl.orchestrator.utils import set_semaphore
@@ -14,7 +12,7 @@
1412
setup_evals_client,
1513
update_weights,
1614
)
17-
from prime_rl.utils.logger import setup_logger
15+
from prime_rl.utils.logger import intercept_verifiers_logging, setup_logger
1816
from prime_rl.utils.monitor import setup_monitor
1917
from prime_rl.utils.pydantic_config import parse_argv
2018
from prime_rl.utils.utils import clean_exit, get_env_ids_to_install, get_step_path, install_env
@@ -26,7 +24,7 @@ async def eval(config: OfflineEvalConfig):
2624
logger = setup_logger(
2725
config.log.level, log_file=config.output_dir / "logs" / "eval.log" if config.log.file else None
2826
)
29-
vf.setup_logging(level=config.log.vf_level.upper())
27+
intercept_verifiers_logging(level=config.log.vf_level)
3028

3129
env_names = [env.name or env.id for env in config.env]
3230
logger.info(f"Starting evals for {config.model.name} in environments {', '.join(env_names)}")

src/prime_rl/eval/utils.py

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,7 @@
2626
from prime_rl.synthesize.utils import merge_reasoning_content, save_result
2727
from prime_rl.utils.client import setup_clients, setup_evals_client
2828
from prime_rl.utils.config import ClientConfig
29-
from prime_rl.utils.logger import get_logger, reset_logger, setup_logger
29+
from prime_rl.utils.logger import get_logger, intercept_verifiers_logging, reset_logger, setup_logger
3030
from prime_rl.utils.monitor import get_monitor
3131
from prime_rl.utils.utils import capitalize, get_eval_dir, get_step_path
3232
from prime_rl.utils.vf import (
@@ -689,6 +689,7 @@ def _run_evals_in_subprocess(
689689
# Setup logger for subprocess (reset first since we inherit parent's global state when forked)
690690
reset_logger()
691691
logger = setup_logger("info")
692+
intercept_verifiers_logging(level="info")
692693
logger.info(f"Eval subprocess started for checkpoint step {ckpt_step}")
693694

694695
# Create fresh clients in subprocess

src/prime_rl/orchestrator/config.py

Lines changed: 0 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -297,30 +297,12 @@ class RetryConfig(BaseConfig):
297297
] = False
298298

299299

300-
class EnvLogConfig(BaseConfig):
301-
"""Configures logging for an environment worker."""
302-
303-
level: Annotated[
304-
str,
305-
Field(description="Log level for prime-rl logger in worker (debug, info, warn, error)."),
306-
] = "warn"
307-
308-
vf_level: Annotated[
309-
str,
310-
Field(description="Log level for verifiers logger in worker (debug, info, warn, error)."),
311-
] = "warn"
312-
313-
314300
class EnvConfig(BaseConfig):
315301
"""Configures an environment for training."""
316302

317303
id: Annotated[str, Field(description="ID of the environment to use.")] = "reverse-text"
318304
args: Annotated[dict, Field(description="Arguments to pass to the environment.")] = {}
319305
name: Annotated[str | None, Field(description="Name of the environment to use.")] = None
320-
log: Annotated[
321-
EnvLogConfig | None,
322-
Field(description="Logging config for this env's workers. If None, logging is disabled."),
323-
] = None
324306

325307

326308
class EvalEnvConfig(EnvConfig):

src/prime_rl/orchestrator/env_worker.py

Lines changed: 5 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@
55
"""
66

77
import asyncio
8-
import logging
98
import queue
109
import uuid
1110
from dataclasses import dataclass
@@ -18,7 +17,7 @@
1817

1918
from prime_rl.utils.client import setup_clients
2019
from prime_rl.utils.config import ClientConfig
21-
from prime_rl.utils.logger import reset_logger, setup_logger
20+
from prime_rl.utils.logger import intercept_verifiers_logging, reset_logger, setup_logger
2221

2322

2423
class WorkerDiedError(Exception):
@@ -188,19 +187,14 @@ def worker_main(
188187
log_level: str,
189188
vf_log_level: str,
190189
log_file: str | None,
190+
worker_name: str | None = None,
191191
):
192192
"""Main entry point for worker process."""
193193
# Reset logger inherited from parent process, then setup fresh logger for this worker
194194
if log_file:
195195
reset_logger()
196-
setup_logger(log_level, log_file=Path(log_file))
197-
vf.setup_logging(level=vf_log_level.upper())
198-
# Redirect verifiers to file instead of inherited stderr
199-
vf_logger = logging.getLogger("verifiers")
200-
vf_logger.handlers.clear()
201-
vf_handler = logging.FileHandler(log_file)
202-
vf_handler.setFormatter(logging.Formatter("%(asctime)s %(levelname)7s %(message)s", datefmt="%H:%M:%S"))
203-
vf_logger.addHandler(vf_handler)
196+
setup_logger(log_level, log_file=Path(log_file), append=True, tag=worker_name)
197+
intercept_verifiers_logging(level=vf_log_level)
204198

205199
# Load environment
206200
env = vf.load_environment(env_id, **env_args)
@@ -293,6 +287,7 @@ def start(self):
293287
self.log_level,
294288
self.vf_log_level,
295289
self.log_file,
290+
self.worker_name,
296291
),
297292
daemon=True,
298293
)

src/prime_rl/orchestrator/orchestrator.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@
4444
update_weights,
4545
)
4646
from prime_rl.utils.heartbeat import Heartbeat
47-
from prime_rl.utils.logger import setup_logger
47+
from prime_rl.utils.logger import intercept_verifiers_logging, setup_logger
4848
from prime_rl.utils.monitor import setup_monitor
4949
from prime_rl.utils.pydantic_config import parse_argv
5050
from prime_rl.utils.utils import (
@@ -66,7 +66,7 @@ async def orchestrate(config: OrchestratorConfig):
6666
logger = setup_logger(
6767
config.log.level, log_file=config.output_dir / "logs" / "orchestrator.log" if config.log.file else None
6868
)
69-
vf.setup_logging(level=config.log.vf_level.upper())
69+
intercept_verifiers_logging(level=config.log.vf_level)
7070
logger.info("Starting orchestrator")
7171

7272
event_loop_lag_monitor = EventLoopLagMonitor()

src/prime_rl/orchestrator/scheduler.py

Lines changed: 9 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@
1919
from prime_rl.utils.client import update_weights
2020
from prime_rl.utils.config import ClientConfig
2121
from prime_rl.utils.logger import get_logger
22-
from prime_rl.utils.pathing import get_env_worker_log_dir
22+
from prime_rl.utils.pathing import get_env_worker_log_file
2323
from prime_rl.utils.utils import (
2424
get_broadcast_dir,
2525
get_latest_ckpt_step,
@@ -92,12 +92,11 @@ def __init__(
9292
self.env_names.append(env_name)
9393
self.workers[env_name] = []
9494

95-
# Setup log directory if logging is enabled for this env
96-
env_log = env_config.log
97-
env_log_dir = None
98-
if env_log is not None and output_dir is not None:
99-
env_log_dir = get_env_worker_log_dir(output_dir, env_name)
100-
env_log_dir.mkdir(parents=True, exist_ok=True)
95+
# Setup log file if env worker file logging is enabled (all workers share one file)
96+
env_log_file = None
97+
if config.log.env_worker_logs and output_dir is not None:
98+
env_log_file = get_env_worker_log_file(output_dir, env_name)
99+
env_log_file.parent.mkdir(parents=True, exist_ok=True)
101100

102101
for worker_idx in range(self.workers_per_env):
103102
worker = EnvWorker(
@@ -111,9 +110,9 @@ def __init__(
111110
example_lookup=self.example_lookups[env_name],
112111
sampling_args=self.sampling_args,
113112
worker_name=f"{env_name}_{worker_idx}",
114-
log_level=env_log.level if env_log else "warn",
115-
vf_log_level=env_log.vf_level if env_log else "warn",
116-
log_file=str(env_log_dir / f"worker_{worker_idx}.log") if env_log_dir else None,
113+
log_level=config.log.level,
114+
vf_log_level=config.log.vf_level,
115+
log_file=str(env_log_file) if env_log_file else None,
117116
)
118117
self.workers[env_name].append(worker)
119118

src/prime_rl/synthesize/synthesize.py

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,5 @@
11
import asyncio
22

3-
import verifiers as vf
4-
53
from prime_rl.orchestrator.utils import (
64
set_semaphore,
75
)
@@ -13,7 +11,7 @@
1311
setup_admin_clients,
1412
setup_clients,
1513
)
16-
from prime_rl.utils.logger import setup_logger
14+
from prime_rl.utils.logger import intercept_verifiers_logging, setup_logger
1715
from prime_rl.utils.pydantic_config import parse_argv
1816
from prime_rl.utils.utils import clean_exit, get_env_ids_to_install, install_env
1917

@@ -24,7 +22,7 @@ async def synthesize(config: SynthesizeConfig):
2422
logger = setup_logger(
2523
config.log.level, log_file=config.output_dir / "logs" / "synthesize.log" if config.log.file else None
2624
)
27-
vf.setup_logging(level=config.log.vf_level.upper())
25+
intercept_verifiers_logging(level=config.log.vf_level)
2826

2927
env_names = [env.name or env.id for env in config.env]
3028
logger.info(f"Starting synthetic data generation for {config.model.name} in environments {', '.join(env_names)}")

src/prime_rl/utils/config.py

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,13 @@ class LogConfig(BaseConfig):
7373
),
7474
] = True
7575

76+
env_worker_logs: Annotated[
77+
bool,
78+
Field(
79+
description="Whether env workers log to files. If True, workers write to logs/env_workers/{env_name}.log.",
80+
),
81+
] = True
82+
7683
log_data: Annotated[
7784
bool,
7885
Field(

0 commit comments

Comments
 (0)