diff --git a/docs/sphinx/source/en/2-user_guide/1-training/3-logging.md b/docs/sphinx/source/en/2-user_guide/1-training/3-logging.md index 64be3b465..6078aa3a6 100644 --- a/docs/sphinx/source/en/2-user_guide/1-training/3-logging.md +++ b/docs/sphinx/source/en/2-user_guide/1-training/3-logging.md @@ -8,8 +8,11 @@ persistent backend, which keeps each unsmoothed iteration value. This page first covers the log directory shared by all algorithms and the Manager-Based reward metric contract, then documents the off-policy terminal used by -SAC / FlashSAC / WarpSAC and APPO. Every terminal field in the tables maps directly to one -backend key; an `_ms` suffix always means milliseconds. +SAC / FlashSAC / WarpSAC and APPO. Persisted backend tags follow the canonical schema +owned by `uni_rl.logging.metric_schema`: millisecond timing fields end in `_ms`, +while `Perf/learning_time`, `Perf/collection_time`, and `Perf/iteration_time` are +seconds. A few terminal fields (`Other`, phase percentages, collector cycle total, +`Rows/s`) are derived in memory for the display and are intentionally not persisted. ## Log Directory and Backend @@ -80,8 +83,9 @@ The bottom of the terminal has three columns: The panel-border title reports `GPUs N`. In multi-GPU training, only rank 0 owns the terminal and persistent logger. Learner metrics and timings are averaged across ranks first and then across the terminal's two-second window. `Steps/s` and -`Samples/s` are exceptions to the rank average: per-rank collector step rates and -learner sample rates are summed, so both fields are total job throughput. The +`Rows/s` are exceptions to the rank average: per-rank collector step rates and +learner replay-row rates are summed, so both fields are total job throughput. Only +`Steps/s` is persisted (as `Perf/total_fps`); `Rows/s` is terminal-only. The `Avg 2s (n=...)` header field gives the number of already rank-reduced learner samples in the current time window. @@ -99,22 +103,22 @@ rows nor get persisted for other algorithms; for example, `Replay Stage` and | Terminal field | TensorBoard / W&B key | Path | Meaning | | --- | --- | --- | --- | -| Collector Wait | `timing/learner_collector_wait_ms` | All | Learner-main-thread wait to reach this iteration's update boundary; SAC-like paths serve inference requests and continue until replay is ready and the `env_steps_per_sync` tick count is met, while APPO waits for a rollout in the ring | -| Inference | `timing/learner_inference_ms` | Learner-owned inference | Total wall time for observation H2D, actor forward, and action D2H; nested details are listed below | -| Collector Release | `timing/learner_collector_release_ms` | Learner-owned inference | Publishing the action response token; normally short, while blocking means the response queue has not drained | -| Replay Batch Wait | `timing/learner_replay_batch_wait_ms` | Device replay | Waiting for the prefetched device batch to finish ingress commit and gather; near zero on a prefetch hit | -| Replay Stage | `timing/learner_replay_stage_ms` | APPO | Sequentially materializing newly arrived NumPy rollouts from the ring into the learner staging pool; this is an exclusive learner-main-thread phase | -| Replay Sample | `timing/learner_replay_sample_ms` | All | Acquiring the ready batch; usually hot/cold swap plus views on CUDA, with possible slot-event waiting on MPS | -| Train | `timing/learner_train_ms` | All | Learner update-phase wall time | -| Weight Publish | `timing/learner_weight_publish_ms` | APPO | Writing fresh actor / critic weights to shared memory | -| Other | `timing/learner_other_ms` | All | Residual after subtracting the phases above from `Iter Wall`, including metrics drain, reward stats, and loop bookkeeping | -| Iter Wall | `perf/iter_ms` | All | Wall time from learner-loop iteration start through update completion; always shown as 100% | - -Backends also record `perf/learner_train_pct`, `perf/learner_accounted_pct`, and -`perf/learner_other_pct`. Accounted time contains only the mutually exclusive +| Collector Wait | `Perf/learner_collector_wait_ms` | All | Learner-main-thread wait to reach this iteration's update boundary; SAC-like paths serve inference requests and continue until replay is ready and the `env_steps_per_sync` tick count is met, while APPO waits for a rollout in the ring | +| Inference | `Perf/learner_inference_ms` | Learner-owned inference | Total wall time for observation H2D, actor forward, and action D2H; nested details are listed below | +| Collector Release | `Perf/learner_collector_release_ms` | Learner-owned inference | Publishing the action response token; normally short, while blocking means the response queue has not drained | +| Replay Batch Wait | `Perf/learner_replay_batch_wait_ms` | Device replay | Waiting for the prefetched device batch to finish ingress commit and gather; near zero on a prefetch hit | +| Replay Stage | `Perf/learner_replay_stage_ms` | APPO | Sequentially materializing newly arrived NumPy rollouts from the ring into the learner staging pool; this is an exclusive learner-main-thread phase | +| Replay Sample | `Perf/learner_replay_sample_ms` | All | Acquiring the ready batch; usually hot/cold swap plus views on CUDA, with possible slot-event waiting on MPS | +| Learning | `Perf/learning_time` (seconds) | All | Learner update-phase wall time | +| Weight Publish | `Perf/learner_weight_publish_ms` | APPO | Writing fresh actor / critic weights to shared memory | +| Other | terminal-only | All | Residual after subtracting the phases above from `Iter Wall`, including metrics drain, reward stats, and loop bookkeeping | +| Iter Wall | `Perf/iteration_time` (seconds) | All | Wall time from learner-loop iteration start through update completion; always shown as 100%; emitted only when the runner measured the complete iteration wall time | + +`Other` and the per-row percentages are derived in memory for the terminal and are +not persisted as backend charts. Accounted time contains only the mutually exclusive main-thread phases above, never nested or background work. Normally `accounted + other = 100%`. If a clock anomaly or a future overlapping timer makes -`accounted > Iter Wall`, `Other` is clamped to zero and `accounted_pct > 100%` is a +`accounted > Iter Wall`, `Other` is clamped to zero and accounted exceeding 100% is a contract violation to fix, not an interpretable parallel-work percentage. ### Learner Nested and Background Diagnostics @@ -124,10 +128,10 @@ slices: | TensorBoard / W&B key | Parent or execution thread | Meaning | | --- | --- | --- | -| `timing/learner_inference_h2d_ms` | Child of `Inference` | Observation copy from the shared CPU slot to the learner device | -| `timing/learner_inference_forward_ms` | Child of `Inference` | `learner.actor` inference, including the current device synchronization | -| `timing/learner_inference_d2h_ms` | Child of `Inference` | Action copy into the shared CPU slot | -| `timing/replay_ingress_h2d_submit_ms` | Replay ingress; CUDA daemon or MPS learner thread | CPU-side duration of the latest transition-span submission into the authoritative device ring; it can occur inside any main phase and must not be added to learner percentages | +| `Perf/learner_inference_h2d_ms` | Child of `Inference` | Observation copy from the shared CPU slot to the learner device | +| `Perf/learner_inference_forward_ms` | Child of `Inference` | `learner.actor` inference, including the current device synchronization | +| `Perf/learner_inference_d2h_ms` | Child of `Inference` | Action copy into the shared CPU slot | +| `Perf/replay_ingress_h2d_submit_ms` | Replay ingress; CUDA daemon or MPS learner thread | CPU-side duration of the latest transition-span submission into the authoritative device ring; it can occur inside any main phase and must not be added to learner percentages | The three inference details should approximately compose `Inference`; small gaps come from Python work between timers. `Replay H2D Submit` overlaps the learner timeline: @@ -137,33 +141,49 @@ copy / gather interval instead of treating submit wall time as extra iteration s ### Tags in Existing Runs -Historical TensorBoard event files are not rewritten. New runs use the canonical -tags below, while an existing run continues to show its old names: +Historical TensorBoard event files are not rewritten. Runs made with unilab-rl 1.4.1 +or later use the canonical tags below (source of truth: +`uni_rl.logging.metric_schema`, with the full migration table in the unilab_rl +repo's `docs/metrics.md`), while an older run continues to show its retired names: -| Old tag | New tag | Reason | +| Old tag | New tag | Note | | --- | --- | --- | -| `timing/inference_total_ms` | `timing/learner_inference_ms` | Match terminal `Inference` and identify the owner | -| `timing/inference_{h2d,forward,d2h}_ms` | `timing/learner_inference_{h2d,forward,d2h}_ms` | Put all three nested items in the learner namespace | -| `timing/learner_incremental_h2d_ms` (SAC-like) | `timing/replay_ingress_h2d_submit_ms` | Identify a potentially parallel submit diagnostic rather than a learner main phase | -| `timing/learner_incremental_h2d_ms` (APPO) | `timing/learner_replay_stage_ms` | Identify synchronous staging that is part of `Iter Wall` | -| `timing/collector_inference_wait_ms` | `timing/collector_learner_action_wait_ms` | The wait can include the remaining learner update, not only inference latency | - -`perf/learner_pipeline_ms` was removed because it mixed exclusive main phases with a -background H2D submission. Use `perf/iter_ms` for the main timeline and reconcile it -with `perf/learner_accounted_pct` plus `perf/learner_other_pct`. +| `perf/steps_per_sec` | `Perf/total_fps` | Aggregate collector env-step throughput (cross-rank sum in DP) | +| `reward/mean`, `reward/mean_ep100` | `Train/mean_reward` | Collector's latest 100-episode mean return; the runner's 10-report smoothing is checkpoint state only | +| `episode/timeout_rate` | `Episode/timeout_rate` | Omitted until the first completed episode | +| `perf/iter_ms` | `Perf/iteration_time` | Milliseconds → seconds; emitted only for a measured complete iteration wall time | +| `timing/learner_train_ms` | `Perf/learning_time` | Milliseconds → seconds | +| `timing/learner__ms` | `Perf/learner__ms` | Same millisecond phase timings, canonical namespace | +| `timing/inference_total_ms`, `timing/inference_{h2d,forward,d2h}_ms` | `Perf/learner_inference_ms`, `Perf/learner_inference_{h2d,forward,d2h}_ms` | Pre-canonical names of the same fields | +| `timing/learner_incremental_h2d_ms` (SAC-like), `timing/replay_ingress_h2d_submit_ms` | `Perf/replay_ingress_h2d_submit_ms` | Potentially parallel submit diagnostic, not a learner main phase | +| `timing/learner_incremental_h2d_ms` (APPO) | `Perf/learner_replay_stage_ms` | Synchronous staging that is part of `Iter Wall` | +| `timing/collector__ms` | `Perf/collector__ms` | Same millisecond collector phases, canonical namespace | +| `timing/collector_inference_wait_ms` | `Perf/collector_learner_action_wait_ms` | The wait can include the remaining learner update, not only inference latency | +| `timing/collector_rollout_ms` | `Perf/collection_time` | Milliseconds → seconds; APPO's whole-rollout wall time | +| `train/dp_sync_time` | `Perf/dp_gradient_sync_ms_per_rank` | Seconds → milliseconds; explicitly per rank | +| `train/dp_gradient_sync_calls` | `Perf/dp_gradient_sync_calls_per_rank` | Explicitly per rank | + +Retired without a replacement tag: `perf/effective_samples_per_sec` / +`perf/learner_samples_per_sec` (derive replay rows from the run configuration and +`Perf/iteration_time`), `perf/collector_active_steps_per_sec` (derive from collector +timing and the run configuration), `perf/collector_cycle_ms`, `perf/learner_*_pct`, +`timing/learner_other_ms`, and `perf/learner_pipeline_ms` — all of these are +derivable from the canonical fields above or are terminal-only derived values. The +extraction helper `scripts/benchmark/rl/extract_offpolicy_metrics.py` accepts both +schemas, preferring the canonical tag and scaling units to match. ### Collector Timeline SAC / FlashSAC record four mutually exclusive hot-path phases per vectorized -env tick. Terminal percentages use `perf/collector_cycle_ms`, the sum of these four -phases: +env tick. Terminal percentages use the sum of these four phases (the collector +cycle), which is derived in memory and not persisted as a backend chart: | Terminal field | TensorBoard / W&B key | Meaning | | --- | --- | --- | -| Inference Request | `timing/collector_inference_request_ms` | Publish observations / dones to the shared slot and notify the learner | -| Learner Action Wait | `timing/collector_learner_action_wait_ms` | Barrier wall time from request publication until the learner publishes this tick's action | -| Env Step | `timing/collector_env_step_ms` | `env.step()` wall time | -| Replay Write | `timing/collector_replay_write_ms` | Transition post-processing, packing, and bounded-ingress write | +| Inference Request | `Perf/collector_inference_request_ms` | Publish observations / dones to the shared slot and notify the learner | +| Learner Action Wait | `Perf/collector_learner_action_wait_ms` | Barrier wall time from request publication until the learner publishes this tick's action | +| Env Step | `Perf/collector_env_step_ms` | `env.step()` wall time | +| Replay Write | `Perf/collector_replay_write_ms` | Transition post-processing, packing, and bounded-ingress write | `Learner Action Wait` is deliberately not named “Inference Wait”: it is not pure inference latency. If the collector finishes `Env Step + Replay Write` and submits its @@ -172,11 +192,12 @@ update, a small scheduling part of the next learner `Collector Wait`, and the ne `Inference + Collector Release`. A long value therefore agrees with parallel execution: it means the collector reached the next barrier before the learner. -The persistent `perf/collector_active_steps_per_sec` diagnostic is calculated as -`num_envs / (Inference Request + Env Step + Replay Write)`. It intentionally excludes -`Learner Action Wait`; low-value sub-millisecond episode / metrics bookkeeping is no -longer timed separately. The terminal `Steps/s` field instead reports total -synchronized collector throughput. The indented Backend Step / Update State / Reset +The collector active-throughput diagnostic (`num_envs / (Inference Request + +Env Step + Replay Write)`, intentionally excluding `Learner Action Wait`) is no +longer persisted as a backend chart; it can be derived from the canonical collector +timing fields and the run configuration when needed. The terminal `Steps/s` field +instead reports total synchronized collector throughput, persisted as +`Perf/total_fps`. The indented Backend Step / Update State / Reset Done rows are nested `Env Step` details. They do not enter the cycle sum, though their displayed percentages use the same collector-cycle denominator. @@ -184,13 +205,14 @@ APPO has a different collection contract and reports: | Terminal field | TensorBoard / W&B key | Basis | | --- | --- | --- | -| MLP Infer | `timing/collector_mlp_infer_ms` | Per-step policy-inference EMA | -| Env Step | `timing/collector_env_step_ms` | Single-`env.step()` EMA | -| Rollout Wall | `timing/collector_rollout_ms` | Whole `steps_per_env` rollout wall-time EMA | +| MLP Infer | `Perf/collector_mlp_infer_ms` | Per-step policy-inference EMA | +| Env Step | `Perf/collector_env_step_ms` | Single-`env.step()` EMA | +| Rollout Wall | `Perf/collection_time` (seconds) | Whole `steps_per_env` rollout wall time | These are not one percentage breakdown: the first two are per-step EMAs and `Rollout Wall` is a whole-rollout total, so the terminal shows milliseconds only. The backend -active-throughput diagnostic is `(num_envs * steps_per_env) / Rollout Wall`. +active-throughput diagnostic `(num_envs * steps_per_env) / Rollout Wall` is likewise +no longer persisted and can be derived from the fields above. ## FastSAC Dual Timeline @@ -286,7 +308,7 @@ adjacent, additive learner phases, not background H2D diagnostics. | Observation | Direct meaning | Check first | | --- | --- | --- | | Learner `Collector Wait` is high | The collector's data / request is not ready when the learner reaches iteration start | Env step, transition post-processing, collector liveness, and IPC | -| Collector `Learner Action Wait` is high while learner `Collector Wait` is low | Collector reaches the next barrier first and waits for update completion plus inference | `Train`, `updates_per_step`, and batch size; then inference details | +| Collector `Learner Action Wait` is high while learner `Collector Wait` is low | Collector reaches the next barrier first and waits for update completion plus inference | `Learning`, `updates_per_step`, and batch size; then inference details | | Learner `Inference` is high | The learner-owned action path itself is slow | H2D / forward / D2H children | | Learner `Replay Batch Wait` is high | Device-replay prefetch misses the consumption point | Ingress commit, side-stream gather, and GPU contention | | Collector `Replay Write` is high | Bounded-ingress write or transition post-processing slows down | Exhausted ingress slots and delayed device commit | diff --git a/docs/sphinx/source/en/2-user_guide/2-algorithms/2-appo.md b/docs/sphinx/source/en/2-user_guide/2-algorithms/2-appo.md index bde402557..45da2b5b6 100644 --- a/docs/sphinx/source/en/2-user_guide/2-algorithms/2-appo.md +++ b/docs/sphinx/source/en/2-user_guide/2-algorithms/2-appo.md @@ -60,10 +60,10 @@ gantt Weight Publish → collector :crit, l4, 28000, 30000 section Iter Wall - perf/iter_ms (learner loop only) : l5, 12000, 30000 + Perf/iteration_time (learner loop only) : l5, 12000, 30000 ``` -> The axis is schematic (relative, not real-ms). The collector subprocess produces rollouts through the 4-slot ring buffer in parallel with the learner, so **Collector Wait ≈ 0** in steady state. `perf/iter_ms` counts only this learner loop (it includes Collector Wait but not the collector's parallel rollout compute); the red Weight Publish marks the end of the iteration when fresh weights are published to the collector. Field meanings are on the [logging page](../1-training/3-logging.md). +> The axis is schematic (relative, not real-ms). The collector subprocess produces rollouts through the 4-slot ring buffer in parallel with the learner, so **Collector Wait ≈ 0** in steady state. `Perf/iteration_time` counts only this learner loop (it includes Collector Wait but not the collector's parallel rollout compute); the red Weight Publish marks the end of the iteration when fresh weights are published to the collector. Field meanings are on the [logging page](../1-training/3-logging.md). ## Key Fields diff --git a/docs/sphinx/source/zh_CN/2-user_guide/1-training/3-logging.md b/docs/sphinx/source/zh_CN/2-user_guide/1-training/3-logging.md index 4ab64c279..7b390cc36 100644 --- a/docs/sphinx/source/zh_CN/2-user_guide/1-training/3-logging.md +++ b/docs/sphinx/source/zh_CN/2-user_guide/1-training/3-logging.md @@ -6,8 +6,11 @@ training logger 提交一次指标。终端面板按固定 2 Hz 时钟刷新, backend 保留每个 iteration 未经时间平滑的值。 本文先说明所有算法共用的日志目录和 Manager-Based reward 指标契约,再详细说明 -SAC / FlashSAC / WarpSAC 与 APPO 共用的 off-policy 终端视图。表中的“终端字段”与 -backend key 一一对应;后缀 `_ms` 均为毫秒。 +SAC / FlashSAC / WarpSAC 与 APPO 共用的 off-policy 终端视图。持久化 backend tag +遵循 `uni_rl.logging.metric_schema` 持有的 canonical schema:毫秒计时字段以 +`_ms` 结尾,而 `Perf/learning_time`、`Perf/collection_time` 与 +`Perf/iteration_time` 单位为秒。少数终端字段(`Other`、各阶段百分比、collector +cycle 合计、`Rows/s`)是在内存中为显示推导的值,有意不持久化。 ## 日志目录与 backend @@ -72,9 +75,10 @@ key 的契约如下: - `System`:buffer 大小、timeout rate、env 数和每 rank batch 大小。 面板边框标题显示 `GPUs N`。多卡训练中只有 rank 0 持有终端与持久化 logger;learner -指标和计时先在 rank 间取平均,再进入终端的两秒时间窗口。`Steps/s` 与 `Samples/s` -不取 rank 平均:前者把各 rank collector step rate 求和,后者把各 rank learner sample -rate 求和,因此两者都是整个训练任务的总吞吐。标题中的 `Avg 2s (n=...)` 表示当前窗口 +指标和计时先在 rank 间取平均,再进入终端的两秒时间窗口。`Steps/s` 与 `Rows/s` +不取 rank 平均:前者把各 rank collector step rate 求和,后者把各 rank learner +replay 行速率求和,因此两者都是整个训练任务的总吞吐。其中只有 `Steps/s` 会持久化 +(即 `Perf/total_fps`),`Rows/s` 仅终端显示。标题中的 `Avg 2s (n=...)` 表示当前窗口 包含多少个已经完成 rank 聚合的 learner 样本。 两列时间不能横向相加。只有 learner 列从 `Collector Wait` 到 `Other` 的行是互斥主线程 @@ -88,21 +92,21 @@ APPO 出现。 | 终端字段 | TensorBoard / W&B key | 适用路径 | 含义 | | --- | --- | --- | --- | -| Collector Wait | `timing/learner_collector_wait_ms` | 全部 | learner 主线程等待达到本轮 update 边界;SAC 类路径会服务 inference request,并继续等 replay ready 与 `env_steps_per_sync` tick 数满足,APPO 等 ring 中 rollout | -| Inference | `timing/learner_inference_ms` | learner-owned inference | observation H2D、actor forward、action D2H 的总墙钟;三项嵌套明细见下表 | -| Collector Release | `timing/learner_collector_release_ms` | learner-owned inference | 将 action response token 发给 collector;正常应很短,阻塞表示 response queue 尚未腾空 | -| Replay Batch Wait | `timing/learner_replay_batch_wait_ms` | device replay | 等已预取的 device batch 完成 ingress commit 与 gather;预取命中时接近 0 | -| Replay Stage | `timing/learner_replay_stage_ms` | APPO | 将 ring buffer 中本轮新到的 NumPy rollout 顺序 materialize 到 learner staging pool;这是 learner 主线程的独占阶段 | -| Replay Sample | `timing/learner_replay_sample_ms` | 全部 | 取得 ready batch;CUDA device replay 通常只是 hot/cold swap 与 view,MPS 还可能等待 slot event | -| Train | `timing/learner_train_ms` | 全部 | learner update 阶段的墙钟 | -| Weight Publish | `timing/learner_weight_publish_ms` | APPO | 将新 actor / critic 权重写入共享内存 | -| Other | `timing/learner_other_ms` | 全部 | `Iter Wall` 减去上述互斥阶段的 residual,例如 metrics drain、reward stats 与 loop bookkeeping | -| Iter Wall | `perf/iter_ms` | 全部 | 从本轮 learner loop 开始到 update 完成的墙钟,固定显示 100% | - -backend 还记录 `perf/learner_train_pct`、`perf/learner_accounted_pct` 与 -`perf/learner_other_pct`。`accounted` 只包含上表的互斥主线程阶段,不包含任何嵌套或后台 +| Collector Wait | `Perf/learner_collector_wait_ms` | 全部 | learner 主线程等待达到本轮 update 边界;SAC 类路径会服务 inference request,并继续等 replay ready 与 `env_steps_per_sync` tick 数满足,APPO 等 ring 中 rollout | +| Inference | `Perf/learner_inference_ms` | learner-owned inference | observation H2D、actor forward、action D2H 的总墙钟;三项嵌套明细见下表 | +| Collector Release | `Perf/learner_collector_release_ms` | learner-owned inference | 将 action response token 发给 collector;正常应很短,阻塞表示 response queue 尚未腾空 | +| Replay Batch Wait | `Perf/learner_replay_batch_wait_ms` | device replay | 等已预取的 device batch 完成 ingress commit 与 gather;预取命中时接近 0 | +| Replay Stage | `Perf/learner_replay_stage_ms` | APPO | 将 ring buffer 中本轮新到的 NumPy rollout 顺序 materialize 到 learner staging pool;这是 learner 主线程的独占阶段 | +| Replay Sample | `Perf/learner_replay_sample_ms` | 全部 | 取得 ready batch;CUDA device replay 通常只是 hot/cold swap 与 view,MPS 还可能等待 slot event | +| Learning | `Perf/learning_time`(秒) | 全部 | learner update 阶段的墙钟 | +| Weight Publish | `Perf/learner_weight_publish_ms` | APPO | 将新 actor / critic 权重写入共享内存 | +| Other | 仅终端 | 全部 | `Iter Wall` 减去上述互斥阶段的 residual,例如 metrics drain、reward stats 与 loop bookkeeping | +| Iter Wall | `Perf/iteration_time`(秒) | 全部 | 从本轮 learner loop 开始到 update 完成的墙钟,固定显示 100%;仅在 runner 测得完整迭代墙钟时写出 | + +`Other` 与各行百分比是在内存中为终端推导的值,不作为 backend chart 持久化。 +`accounted` 只包含上表的互斥主线程阶段,不包含任何嵌套或后台 诊断。正常采样下 `accounted + other = 100%`;如果系统时钟异常或未来埋点意外重叠使 -`accounted > Iter Wall`,`Other` 会钳制为 0,而 `accounted_pct > 100%` 是应修复的 contract +`accounted > Iter Wall`,`Other` 会钳制为 0,而 accounted 超过 100% 是应修复的 contract 告警,不是可解释的并行占比。 ### Learner 嵌套与后台诊断 @@ -111,10 +115,10 @@ backend 还记录 `perf/learner_train_pct`、`perf/learner_accounted_pct` 与 | TensorBoard / W&B key | 父级或执行线程 | 含义 | | --- | --- | --- | -| `timing/learner_inference_h2d_ms` | `Inference` 子项 | observation 从共享 CPU slot 拷到 learner device | -| `timing/learner_inference_forward_ms` | `Inference` 子项 | `learner.actor` 推理;包含当前实现中的 device synchronize | -| `timing/learner_inference_d2h_ms` | `Inference` 子项 | action 写回共享 CPU slot | -| `timing/replay_ingress_h2d_submit_ms` | replay ingress;CUDA daemon 或 MPS learner 线程 | 最近一次 transition span 提交到 authoritative device ring 的 CPU 侧耗时;它可能落在任意主阶段内,因此绝不能再加到 learner 百分比中 | +| `Perf/learner_inference_h2d_ms` | `Inference` 子项 | observation 从共享 CPU slot 拷到 learner device | +| `Perf/learner_inference_forward_ms` | `Inference` 子项 | `learner.actor` 推理;包含当前实现中的 device synchronize | +| `Perf/learner_inference_d2h_ms` | `Inference` 子项 | action 写回共享 CPU slot | +| `Perf/replay_ingress_h2d_submit_ms` | replay ingress;CUDA daemon 或 MPS learner 线程 | 最近一次 transition span 提交到 authoritative device ring 的 CPU 侧耗时;它可能落在任意主阶段内,因此绝不能再加到 learner 百分比中 | 三项 inference 明细应近似组成 `Inference`;计时调用之间的 Python 开销可能造成小差值。 `Replay H2D Submit` 则与 learner 主时间线重叠:CUDA 由 @@ -124,32 +128,46 @@ backend 还记录 `perf/learner_train_pct`、`perf/learner_accounted_pct` 与 ### 旧 run 的 tag 对照 -历史 TensorBoard event 不会被重写。新 run 使用下列 canonical tag;打开旧 run 时仍会看到 -旧名: +历史 TensorBoard event 不会被重写。unilab-rl 1.4.1 及以后的 run 使用下列 +canonical tag(权威定义:`uni_rl.logging.metric_schema`,完整迁移表见 unilab_rl +仓库的 `docs/metrics.md`);打开旧 run 时仍会看到已退役的旧名: -| 旧 tag | 新 tag | 变化原因 | +| 旧 tag | 新 tag | 说明 | | --- | --- | --- | -| `timing/inference_total_ms` | `timing/learner_inference_ms` | 与终端 `Inference` 统一,并明确 owner | -| `timing/inference_{h2d,forward,d2h}_ms` | `timing/learner_inference_{h2d,forward,d2h}_ms` | 三个嵌套项统一放入 learner namespace | -| `timing/learner_incremental_h2d_ms`(SAC 类) | `timing/replay_ingress_h2d_submit_ms` | 明确它是可能并行的 submit 诊断,不是 learner 主阶段 | -| `timing/learner_incremental_h2d_ms`(APPO) | `timing/learner_replay_stage_ms` | 明确它是可计入 `Iter Wall` 的同步 staging 阶段 | -| `timing/collector_inference_wait_ms` | `timing/collector_learner_action_wait_ms` | 等待范围还包含剩余 learner update,不等于 inference latency | - -`perf/learner_pipeline_ms` 已移除:它曾把互斥主阶段与后台 H2D submit 混加。主时间线请使用 -`perf/iter_ms`,完整性请对照 `perf/learner_accounted_pct` 与 -`perf/learner_other_pct`。 +| `perf/steps_per_sec` | `Perf/total_fps` | 聚合 collector env-step 吞吐(DP 下跨 rank 求和) | +| `reward/mean`、`reward/mean_ep100` | `Train/mean_reward` | collector 最近 100 个 episode 的 return 均值;runner 的 10-report 平滑仅是 checkpoint 状态 | +| `episode/timeout_rate` | `Episode/timeout_rate` | 首个完成 episode 之前不写出 | +| `perf/iter_ms` | `Perf/iteration_time` | 毫秒 → 秒;仅在测得完整迭代墙钟时写出 | +| `timing/learner_train_ms` | `Perf/learning_time` | 毫秒 → 秒 | +| `timing/learner__ms` | `Perf/learner__ms` | 同样的毫秒阶段计时,canonical 命名空间 | +| `timing/inference_total_ms`、`timing/inference_{h2d,forward,d2h}_ms` | `Perf/learner_inference_ms`、`Perf/learner_inference_{h2d,forward,d2h}_ms` | 同一批字段的更早期名称 | +| `timing/learner_incremental_h2d_ms`(SAC 类)、`timing/replay_ingress_h2d_submit_ms` | `Perf/replay_ingress_h2d_submit_ms` | 可能并行的 submit 诊断,不是 learner 主阶段 | +| `timing/learner_incremental_h2d_ms`(APPO) | `Perf/learner_replay_stage_ms` | 可计入 `Iter Wall` 的同步 staging 阶段 | +| `timing/collector__ms` | `Perf/collector__ms` | 同样的毫秒 collector 阶段,canonical 命名空间 | +| `timing/collector_inference_wait_ms` | `Perf/collector_learner_action_wait_ms` | 等待范围还包含剩余 learner update,不等于 inference latency | +| `timing/collector_rollout_ms` | `Perf/collection_time` | 毫秒 → 秒;APPO 的整条 rollout 墙钟 | +| `train/dp_sync_time` | `Perf/dp_gradient_sync_ms_per_rank` | 秒 → 毫秒;明确为每 rank 值 | +| `train/dp_gradient_sync_calls` | `Perf/dp_gradient_sync_calls_per_rank` | 明确为每 rank 值 | + +已退役且无替代 tag:`perf/effective_samples_per_sec` / +`perf/learner_samples_per_sec`(可由 run 配置与 `Perf/iteration_time` 推导 replay +行数)、`perf/collector_active_steps_per_sec`(可由 collector 计时与 run 配置推导)、 +`perf/collector_cycle_ms`、`perf/learner_*_pct`、`timing/learner_other_ms` 与 +`perf/learner_pipeline_ms`——它们都可以从上面的 canonical 字段推导,或本来就是仅终端 +显示的推导值。提取工具 `scripts/benchmark/rl/extract_offpolicy_metrics.py` 两种 schema +都认,优先使用 canonical tag 并换算单位。 ### Collector 自有时间线 SAC / FlashSAC 每个 vectorized env tick 记录四个热路径互斥阶段,终端百分比使用 -`perf/collector_cycle_ms`(这四项之和)作分母: +这四项之和(collector cycle)作分母;该合计在内存中推导,不作为 backend chart 持久化: | 终端字段 | TensorBoard / W&B key | 含义 | | --- | --- | --- | -| Inference Request | `timing/collector_inference_request_ms` | 发布 observation / dones 到共享 slot,并通知 learner | -| Learner Action Wait | `timing/collector_learner_action_wait_ms` | request 发出后,等 learner 发布当前 tick action 的屏障墙钟 | -| Env Step | `timing/collector_env_step_ms` | `env.step()` 墙钟 | -| Replay Write | `timing/collector_replay_write_ms` | transition 后处理、打包并写入 bounded ingress | +| Inference Request | `Perf/collector_inference_request_ms` | 发布 observation / dones 到共享 slot,并通知 learner | +| Learner Action Wait | `Perf/collector_learner_action_wait_ms` | request 发出后,等 learner 发布当前 tick action 的屏障墙钟 | +| Env Step | `Perf/collector_env_step_ms` | `env.step()` 墙钟 | +| Replay Write | `Perf/collector_replay_write_ms` | transition 后处理、打包并写入 bounded ingress | `Learner Action Wait` 特意不叫 “Inference Wait”:它不是纯 inference latency。如果 collector 在 learner update 期间先完成 `Env Step + Replay Write` 并提交下一 request,这一项会包含 @@ -157,10 +175,10 @@ SAC / FlashSAC 每个 vectorized env tick 记录四个热路径互斥阶段, Collector Release`。所以它很长并不与“两侧并行”矛盾,反而说明 collector 比 learner update 更早到达下一屏障。 -持久化指标 `perf/collector_active_steps_per_sec` 按 collector 活跃路径计算: -`num_envs / (Inference Request + Env Step + Replay Write)` 计算;它有意排除 -`Learner Action Wait`;诊断价值较低的亚毫秒 episode / metrics bookkeeping 不再单独计时。 -终端 `Steps/s` 则报告同步 collector 的总吞吐。`Env Step` 下缩进的 Backend Step / +collector 活跃吞吐诊断(`num_envs / (Inference Request + Env Step + Replay Write)`, +有意排除 `Learner Action Wait`)不再作为 backend chart 持久化;需要时可由 canonical +collector 计时字段与 run 配置推导。终端 `Steps/s` 则报告同步 collector 的总吞吐, +并持久化为 `Perf/total_fps`。`Env Step` 下缩进的 Backend Step / Update State / Reset Done 是父项的嵌套明细,不参与 cycle 求和,但百分比仍使用同一 collector cycle 分母。 @@ -168,13 +186,13 @@ APPO 使用不同的采集 contract,因此 collector 只上报: | 终端字段 | TensorBoard / W&B key | 口径 | | --- | --- | --- | -| MLP Infer | `timing/collector_mlp_infer_ms` | 单步策略推理 EMA | -| Env Step | `timing/collector_env_step_ms` | 单次 `env.step()` EMA | -| Rollout Wall | `timing/collector_rollout_ms` | 完整 `steps_per_env` 步 rollout 的墙钟 EMA | +| MLP Infer | `Perf/collector_mlp_infer_ms` | 单步策略推理 EMA | +| Env Step | `Perf/collector_env_step_ms` | 单次 `env.step()` EMA | +| Rollout Wall | `Perf/collection_time`(秒) | 完整 `steps_per_env` 步 rollout 的墙钟 | APPO 的三个值不是一组百分比分解:前两个是单步 EMA,`Rollout Wall` 是整条 rollout 总量, -所以终端只显示 ms。backend 的活跃吞吐诊断按 -`(num_envs * steps_per_env) / Rollout Wall` 计算。 +所以终端只显示 ms。backend 的活跃吞吐诊断 `(num_envs * steps_per_env) / Rollout Wall` +同样不再持久化,可由上述字段推导。 ## FastSAC 双时间线 @@ -264,7 +282,7 @@ learner `Iter Wall`,取决于 ring backlog 和两侧吞吐。 | 现象 | 直接含义 | 优先检查 | | --- | --- | --- | | learner `Collector Wait` 高 | learner 到达迭代开头后,collector 数据/request 尚未就绪 | env step、transition 后处理、collector 存活与 IPC | -| collector `Learner Action Wait` 高,同时 learner `Collector Wait` 低 | collector 先到下一屏障,在等 learner 完成 update 并服务 inference | `Train`、`updates_per_step`、batch size;再看 inference 三项明细 | +| collector `Learner Action Wait` 高,同时 learner `Collector Wait` 低 | collector 先到下一屏障,在等 learner 完成 update 并服务 inference | `Learning`、`updates_per_step`、batch size;再看 inference 三项明细 | | learner `Inference` 高 | learner-owned action 路径本身慢 | H2D / forward / D2H 三项子指标 | | learner `Replay Batch Wait` 高 | device replay 预取未赶上消费 | ingress commit、side-stream gather 与 GPU 竞争 | | collector `Replay Write` 高 | bounded ingress 写入或 transition 后处理变慢 | ingress 槽是否耗尽、device commit 是否落后 | diff --git a/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/2-appo.md b/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/2-appo.md index ceb2e7da3..61bdf7d67 100644 --- a/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/2-appo.md +++ b/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/2-appo.md @@ -58,10 +58,10 @@ gantt Weight Publish 写回 collector :crit, l4, 28000, 30000 section Iter Wall - perf/iter_ms(仅 learner 这圈) : l5, 12000, 30000 + Perf/iteration_time(仅 learner 这圈) : l5, 12000, 30000 ``` -> 横轴为示意相对时长(非真实 ms 比例)。collector 子进程经 4 槽 ring buffer 与 learner 并行产出 rollout,稳态下 **Collector Wait ≈ 0**。`perf/iter_ms` 仅计 learner 这一圈(含 Collector Wait,但不含 collector 的并行采集计算);红色 Weight Publish 标志该轮迭代结束、向 collector 发布新权重。各字段含义见[日志页](../1-training/3-logging.md)。 +> 横轴为示意相对时长(非真实 ms 比例)。collector 子进程经 4 槽 ring buffer 与 learner 并行产出 rollout,稳态下 **Collector Wait ≈ 0**。`Perf/iteration_time` 仅计 learner 这一圈(含 Collector Wait,但不含 collector 的并行采集计算);红色 Weight Publish 标志该轮迭代结束、向 collector 发布新权重。各字段含义见[日志页](../1-training/3-logging.md)。 ## 关键字段 diff --git a/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/3-sac.md b/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/3-sac.md index e0094420d..1c33d9ac5 100644 --- a/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/3-sac.md +++ b/docs/sphinx/source/zh_CN/2-user_guide/2-algorithms/3-sac.md @@ -66,9 +66,9 @@ uv run train --algo sac --task g1_walk_flat --sim mujoco \ rank 0 独占终端与 TensorBoard/W&B logger;其他 learner 不刷新终端,也不创建独立 tfevents。rank 0 在每个 learner iteration 汇总跨 rank 标量:loss、reward 和 timing -取均值,计数与并行吞吐求和。`perf/steps_per_sec` / 终端 `Steps/s` 表示所有 -collector 的聚合 env-step 吞吐,`perf/effective_samples_per_sec` / 终端 `Samples/s` -表示所有 learner 的聚合有效样本吞吐。checkpoint 同样由 rank 0 独占:每个保存间隔和 +取均值,计数与并行吞吐求和。`Perf/total_fps` / 终端 `Steps/s` 表示所有 +collector 的聚合 env-step 吞吐;终端 `Rows/s` 表示所有 learner 的聚合 replay 行 +吞吐(仅终端显示,不持久化;可由 run 配置与 `Perf/iteration_time` 推导)。checkpoint 同样由 rank 0 独占:每个保存间隔和 训练结束只在 canonical run 目录写一份模型;其他 rank 复用该路径完成进程协调,不创建 rank 子目录或任何日志文件。 自动生成的多卡 run 目录以 `_gpuxN` 结尾(例如 `_gpux2`);单卡目录和显式 diff --git a/pyproject.rocm.toml b/pyproject.rocm.toml index a491a1064..2b249cc9b 100644 --- a/pyproject.rocm.toml +++ b/pyproject.rocm.toml @@ -72,12 +72,12 @@ mujoco = [ "mjbatch-uni~=0.2.1", ] motrix = ["unisim-core[motrix]>=1.7.8"] -uni_rl = ["unilab-rl==1.4.0"] +uni_rl = ["unilab-rl==1.4.1"] [dependency-groups] # Dev/CI exercises the APPO/off-policy/multi-GPU paths, so the optional # unilab-rl runtime stays installed in the dev environment. -dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.0"] +dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.1"] [[tool.uv.index]] name = "pytorch-rocm72" diff --git a/pyproject.toml b/pyproject.toml index 1f51fb0f1..396336b05 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -109,7 +109,7 @@ newton = [ ] motrix = ["unisim-core[motrix]>=1.7.8"] genesis = ["genesis-world==1.3.3"] -uni_rl = ["unilab-rl==1.4.0"] +uni_rl = ["unilab-rl==1.4.1"] # Wheels are CPython 3.12/3.13 Linux x86_64 only; the marker keeps uv lock # resolvable for the other required-environments. superdex = [ @@ -119,7 +119,7 @@ superdex = [ [dependency-groups] # Dev/CI exercises the APPO/off-policy/multi-GPU paths, so the optional # unilab-rl runtime stays installed in the dev environment. -dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.0"] +dev = ["pytest", "pytest-cov", "ruff", "mypy", "pyright>=1.1.408", "unilab-rl==1.4.1"] [[tool.uv.index]] name = "r2-cu130" diff --git a/scripts/benchmark/rl/benchmark_offpolicy_dp_scaling.py b/scripts/benchmark/rl/benchmark_offpolicy_dp_scaling.py index daa6fc447..ddcb8f454 100644 --- a/scripts/benchmark/rl/benchmark_offpolicy_dp_scaling.py +++ b/scripts/benchmark/rl/benchmark_offpolicy_dp_scaling.py @@ -12,17 +12,22 @@ Measurement conventions: - Collector throughput (Steps/s): mean of the last 50% of canonical rank-0 - ``perf/steps_per_sec`` samples. In DP runs the runner already sums the - per-rank collector rates before logging. -- Learner throughput (Samples/s): the same tail mean over - ``perf/effective_samples_per_sec``. This counts effective learner samples - (including configured sample multipliers), summed across ranks. + ``Perf/total_fps`` samples (legacy ``perf/steps_per_sec`` for pre-1.4.1 + event files). In DP runs the runner already sums the per-rank collector + rates before logging. +- Learner throughput (Samples/s): the same tail mean over the legacy + ``perf/effective_samples_per_sec`` tag. unilab-rl 1.4.1 no longer persists + learner replay throughput, so runs without that tag derive it per iteration + as ``world_size * algo.batch_size * algo.updates_per_step / + Perf/iteration_time`` (run configuration from ``run_config.json``). This + counts effective learner replay rows (including configured sample + multipliers), summed across ranks. - Both fields report their own N-way / N=1 scaling ratio. The roadmap verdict remains attached to collector Steps/s: ``pass`` at >= 1.7 (``SCALING_PASS_THRESHOLD``), otherwise ``below threshold``. The verdict is data only and does not affect the exit code. - The tail uses ceil(n/2) points (with n samples, ``n - n//2``), skipping - collector warm-up and replay prefill. Missing either throughput tag is a + collector warm-up and replay prefill. Missing either throughput series is a hard error, never silently skipped. - Exit code is non-zero only when a run itself failed (subprocess error, non-completed summary, or missing artifacts). @@ -65,10 +70,23 @@ # --sim mujoco` (see src/unilab/cli.py build_route for off-policy algos). ROUTE_OVERRIDES = ("task=g1_walk_flat/mujoco",) -STEPS_PER_SEC_TAG = "perf/steps_per_sec" -SAMPLES_PER_SEC_TAG = "perf/effective_samples_per_sec" -REWARD_TAG = "reward/mean" -DP_SYNC_TIME_TAG = "train/dp_sync_time" +# Scalar tag candidates are ``(tag, scale)`` pairs ordered newest schema +# first. unilab-rl 1.4.1 introduced the canonical metric schema +# (``uni_rl.logging.metric_schema``; migration table in unilab_rl +# ``docs/metrics.md``) without rewriting historical event files, so the first +# tag present in a run wins and its values are scaled into the units the +# report columns already use (per-second rates, seconds for durations). +STEPS_PER_SEC_TAGS = (("Perf/total_fps", 1.0), ("perf/steps_per_sec", 1.0)) +REWARD_TAGS = (("Train/mean_reward", 1.0), ("reward/mean", 1.0)) +DP_SYNC_TIME_TAGS = (("Perf/dp_gradient_sync_ms_per_rank", 0.001), ("train/dp_sync_time", 1.0)) +ITERATION_TIME_TAGS = (("Perf/iteration_time", 1.0), ("perf/iter_ms", 0.001)) +# unilab-rl 1.4.1 removed the persisted learner replay throughput chart with +# no replacement tag; post-1.4.1 runs derive it from the run configuration +# and ``Perf/iteration_time`` (see derive_learner_samples_per_sec). +SAMPLES_PER_SEC_TAGS = ( + ("perf/effective_samples_per_sec", 1.0), + ("perf/learner_samples_per_sec", 1.0), +) DEFAULT_ITERATIONS = 300 DEFAULT_DEVICES = "0,1" @@ -117,8 +135,8 @@ def find_event_files(rank_dir: Path) -> list[Path]: return sorted(rank_dir.glob("events.out.tfevents.*")) -def read_scalar_series(rank_dir: Path, tag: str) -> list[float]: - """Scalar values of ``tag`` from a rank's tfevents (in event order). +def read_scalar_series(rank_dir: Path, candidates: Sequence[tuple[str, float]]) -> list[float]: + """Scaled values of the first present tag among ``candidates``. An absent tag yields ``[]``; a rank directory without any tfevents file raises ``RunParseError`` so missing rank data is never silent. @@ -134,9 +152,37 @@ def read_scalar_series(rank_dir: Path, tag: str) -> list[float]: accumulator = event_accumulator.EventAccumulator(str(event_files[0])) accumulator.Reload() - if tag not in accumulator.Tags()["scalars"]: + tags = accumulator.Tags()["scalars"] + for tag, scale in candidates: + if tag in tags: + return [float(event.value) * scale for event in accumulator.Scalars(tag)] + return [] + + +def derive_learner_samples_per_sec(run_dir: Path, world_size: int) -> list[float]: + """Reconstruct aggregate learner replay rows/s for post-1.4.1 runs. + + unilab-rl 1.4.1 no longer persists learner replay throughput; per + unilab_rl ``docs/metrics.md`` it derives from the run configuration — + ``algo.batch_size * algo.updates_per_step`` replay rows per rank per + iteration, summed across ranks — divided by the measured + ``Perf/iteration_time``. + """ + iteration_times = read_scalar_series(run_dir, ITERATION_TIME_TAGS) + config_path = Path(run_dir) / "run_config.json" + if not iteration_times or not config_path.is_file(): + return [] + config = json.loads(config_path.read_text(encoding="utf-8")).get("config") or {} + algo = config.get("algo") or {} + try: + rows_per_rank_per_iter = int(algo["batch_size"]) * int(algo["updates_per_step"]) + except (KeyError, TypeError, ValueError): return [] - return [float(event.value) for event in accumulator.Scalars(tag)] + return [ + world_size * rows_per_rank_per_iter / iteration_time + for iteration_time in iteration_times + if iteration_time > 0 + ] def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]: @@ -160,14 +206,24 @@ def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]: f"run {run_dir} did not complete: status={status!r} error={summary.get('error')!r}" ) - collector_series = read_scalar_series(run_dir, STEPS_PER_SEC_TAG) + collector_series = read_scalar_series(run_dir, STEPS_PER_SEC_TAGS) if not collector_series: - raise RunParseError(f"run has no {STEPS_PER_SEC_TAG!r} samples under {run_dir}") - learner_series = read_scalar_series(run_dir, SAMPLES_PER_SEC_TAG) + raise RunParseError( + f"run has no collector throughput samples " + f"({' or '.join(tag for tag, _ in STEPS_PER_SEC_TAGS)}) under {run_dir}" + ) + learner_series = read_scalar_series(run_dir, SAMPLES_PER_SEC_TAGS) + learner_throughput_source = "tfevents" if not learner_series: - raise RunParseError(f"run has no {SAMPLES_PER_SEC_TAG!r} samples under {run_dir}") - dp_sync_samples = read_scalar_series(run_dir, DP_SYNC_TIME_TAG) - reward_series = read_scalar_series(run_dir, REWARD_TAG) + learner_series = derive_learner_samples_per_sec(run_dir, world_size) + learner_throughput_source = "derived" + if not learner_series: + raise RunParseError( + f"run has no 'perf/effective_samples_per_sec' samples and learner " + f"throughput cannot be derived under {run_dir}" + ) + dp_sync_samples = read_scalar_series(run_dir, DP_SYNC_TIME_TAGS) + reward_series = read_scalar_series(run_dir, REWARD_TAGS) return { "run_dir": str(run_dir), "world_size": world_size, @@ -176,6 +232,7 @@ def parse_run(run_dir: Path, world_size: int) -> dict[str, Any]: "training_wall_time_sec": summary.get("training_wall_time_sec"), "num_collector_throughput_samples": len(collector_series), "num_learner_throughput_samples": len(learner_series), + "learner_throughput_source": learner_throughput_source, "steady_state_collector_steps_per_s": steady_state_mean(collector_series), "steady_state_learner_samples_per_s": steady_state_mean(learner_series), "final_mean_reward": reward_series[-1] if reward_series else None, @@ -468,8 +525,11 @@ def main(argv: list[str] | None = None) -> int: "scaling_pass_threshold": SCALING_PASS_THRESHOLD, "steady_state_tail_fraction": STEADY_STATE_TAIL_FRACTION, "throughput_tags": { - "collector_steps_per_sec": STEPS_PER_SEC_TAG, - "learner_samples_per_sec": SAMPLES_PER_SEC_TAG, + "collector_steps_per_sec": [tag for tag, _ in STEPS_PER_SEC_TAGS], + "learner_samples_per_sec": [tag for tag, _ in SAMPLES_PER_SEC_TAGS], + "learner_samples_per_sec_derived": ( + "world_size * algo.batch_size * algo.updates_per_step / Perf/iteration_time" + ), }, }, }, diff --git a/scripts/benchmark/rl/extract_offpolicy_metrics.py b/scripts/benchmark/rl/extract_offpolicy_metrics.py index 0a1677d2a..f8071eb2e 100644 --- a/scripts/benchmark/rl/extract_offpolicy_metrics.py +++ b/scripts/benchmark/rl/extract_offpolicy_metrics.py @@ -1,5 +1,18 @@ #!/usr/bin/env python3 -"""Extract end-to-end off-policy benchmark metrics from TensorBoard event files.""" +"""Extract end-to-end off-policy benchmark metrics from TensorBoard event files. + +unilab-rl 1.4.1 replaced the ``timing/*`` and lowercase ``perf/*`` tags with the +canonical schema in ``uni_rl.logging.metric_schema`` (see ``docs/metrics.md`` in +the unilab_rl repo). Historical event files are immutable, so every row below +lists ``(tag, scale)`` candidates newest-first: the first tag present in the +file wins and its values are multiplied by ``scale`` to keep each row's +long-standing unit (milliseconds / per-second rates). Two canonical fields are +seconds where the retired tags were milliseconds: ``Perf/iteration_time`` +(was ``perf/iter_ms``) and ``Perf/learning_time`` (was +``timing/learner_train_ms``). Rows whose source chart was retired without a +replacement (collector active throughput, collector cycle total) only carry +their legacy tag and report NaN on post-1.4.1 runs. +""" from __future__ import annotations @@ -9,20 +22,44 @@ from tensorboard.backend.event_processing import event_accumulator -TAGS = { - "iter_ms": "perf/iter_ms", - "steps_per_sec": "perf/steps_per_sec", - "collector_active_steps_per_sec": "perf/collector_active_steps_per_sec", - "collector_cycle_ms": "perf/collector_cycle_ms", - "collector_learner_action_wait_ms": "timing/collector_learner_action_wait_ms", - "collector_replay_write_ms": "timing/collector_replay_write_ms", - "learner_collector_wait_ms": "timing/learner_collector_wait_ms", - "learner_inference_ms": "timing/learner_inference_ms", - "learner_replay_batch_wait_ms": "timing/learner_replay_batch_wait_ms", - "learner_replay_sample_ms": "timing/learner_replay_sample_ms", - "replay_ingress_h2d_submit_ms": "timing/replay_ingress_h2d_submit_ms", - "learner_collector_release_ms": "timing/learner_collector_release_ms", - "learner_train_ms": "timing/learner_train_ms", +TAGS: dict[str, tuple[tuple[str, float], ...]] = { + "iter_ms": (("Perf/iteration_time", 1000.0), ("perf/iter_ms", 1.0)), + "steps_per_sec": (("Perf/total_fps", 1.0), ("perf/steps_per_sec", 1.0)), + "collector_active_steps_per_sec": (("perf/collector_active_steps_per_sec", 1.0),), + "collector_cycle_ms": (("perf/collector_cycle_ms", 1.0),), + "collector_learner_action_wait_ms": ( + ("Perf/collector_learner_action_wait_ms", 1.0), + ("timing/collector_learner_action_wait_ms", 1.0), + ), + "collector_replay_write_ms": ( + ("Perf/collector_replay_write_ms", 1.0), + ("timing/collector_replay_write_ms", 1.0), + ), + "learner_collector_wait_ms": ( + ("Perf/learner_collector_wait_ms", 1.0), + ("timing/learner_collector_wait_ms", 1.0), + ), + "learner_inference_ms": ( + ("Perf/learner_inference_ms", 1.0), + ("timing/learner_inference_ms", 1.0), + ), + "learner_replay_batch_wait_ms": ( + ("Perf/learner_replay_batch_wait_ms", 1.0), + ("timing/learner_replay_batch_wait_ms", 1.0), + ), + "learner_replay_sample_ms": ( + ("Perf/learner_replay_sample_ms", 1.0), + ("timing/learner_replay_sample_ms", 1.0), + ), + "replay_ingress_h2d_submit_ms": ( + ("Perf/replay_ingress_h2d_submit_ms", 1.0), + ("timing/replay_ingress_h2d_submit_ms", 1.0), + ), + "learner_collector_release_ms": ( + ("Perf/learner_collector_release_ms", 1.0), + ("timing/learner_collector_release_ms", 1.0), + ), + "learner_train_ms": (("Perf/learning_time", 1000.0), ("timing/learner_train_ms", 1.0)), } @@ -41,6 +78,18 @@ def _average_last(scalars: list[event_accumulator.ScalarEvent], n: int) -> float return sum(values[-n:]) / len(values[-n:]) +def _extract_row( + ea: event_accumulator.EventAccumulator, + candidates: tuple[tuple[str, float], ...], + last: int, +) -> float: + available = set(ea.Tags()["scalars"]) + for tag, scale in candidates: + if tag in available: + return _average_last(ea.Scalars(tag), last) * scale + return float("nan") + + def main(argv: list[str]) -> int: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("log_dir", type=Path) @@ -56,17 +105,7 @@ def main(argv: list[str]) -> int: ea = event_accumulator.EventAccumulator(str(event_file)) ea.Reload() - available_tags = {tag: False for tag in TAGS.values()} - for tag in ea.Tags()["scalars"]: - available_tags[tag] = True - - rows = [] - for label, tag in TAGS.items(): - if not available_tags.get(tag, False): - rows.append((label, float("nan"))) - continue - scalars = ea.Scalars(tag) - rows.append((label, _average_last(scalars, args.last))) + rows = [(label, _extract_row(ea, candidates, args.last)) for label, candidates in TAGS.items()] if args.json: import json diff --git a/tests/algos/test_offpolicy_dp_sync.py b/tests/algos/test_offpolicy_dp_sync.py index 25f7179a5..c3f624252 100644 --- a/tests/algos/test_offpolicy_dp_sync.py +++ b/tests/algos/test_offpolicy_dp_sync.py @@ -102,7 +102,10 @@ def test_runner_collects_per_iteration_gradient_sync_metrics(): runner = _runner_with(_SyncLearner(), _FakeDpSync()) metrics = defaultdict(list) runner._collect_dp_sync_metrics(metrics) - assert metrics == {"dp_sync_time": [0.25], "dp_gradient_sync_calls": [3.0]} + assert metrics == { + "Perf/dp_gradient_sync_ms_per_rank": [250.0], + "Perf/dp_gradient_sync_calls_per_rank": [3.0], + } def test_runner_requires_gradient_sync_contract(): @@ -175,16 +178,15 @@ def test_log_statistics_mean_scalars_and_sum_concurrent_throughput(): logger._total_steps = 100 logger._buffer_size = 50 logger._buffer_target = 200 - logger._collector_active_steps_per_sec = 1_000.0 logger._mean_ep_length = 25.0 logger._collector_timing = {"env_step_ms": 2.0} payload = runner._aggregate_log_statistics( logger, - metrics={"critic_loss": 2.0}, - reward=3.0, - reward_metrics={"mean_ep100": 4.0}, - reward_components={"reward/tracking": 5.0}, + metrics={"Loss/critic": 2.0}, + checkpoint_return_mean_reports10=None, + return_mean_ep100=3.0, + reward_components={"tracking": 5.0}, train_time=0.4, collector_wait_time=0.1, replay_batch_wait_time=0.01, @@ -198,39 +200,37 @@ def test_log_statistics_mean_scalars_and_sum_concurrent_throughput(): iteration_time=0.5, extra_info={ "throughput_steps": 100, - "collector_active_steps_per_sec": 1_000.0, "batch_size_per_rank": 64, "effective_batch_size": 64, - "replay_samples_per_iter": 128, - "learner_samples_per_iter": 256, + "learner_replay_rows_per_iter": 256, }, ) - assert payload["metrics"] == {"critic_loss": pytest.approx(2.0)} - assert payload["reward"] == pytest.approx(3.0) + assert payload["metrics"] == {"Loss/critic": pytest.approx(2.0)} + assert payload["return_mean_ep100"] == pytest.approx(3.0) extra_info = payload["extra_info"] assert isinstance(extra_info, dict) - assert extra_info["steps_per_sec"] == pytest.approx(400.0) - assert extra_info["learner_samples_per_sec"] == pytest.approx(1_024.0) - assert extra_info["collector_active_steps_per_sec"] == pytest.approx(2_000.0) + assert extra_info["env_steps_per_sec"] == pytest.approx(400.0) + assert extra_info["learner_replay_rows_per_sec"] == pytest.approx(1_024.0) assert extra_info["effective_batch_size"] == 128 - assert extra_info["learner_samples_per_iter"] == 512 + assert extra_info["learner_replay_rows_per_iter"] == 512 assert logger._total_steps == 200 assert logger._buffer_size == 100 assert logger._collector_timing == {"env_step_ms": pytest.approx(2.0)} - logger.log_step(iteration=1, **payload) - assert logger._get_iter_steps_per_sec() == pytest.approx(400.0) - assert logger._get_effective_samples_per_sec() == pytest.approx(1_024.0) + log_step_payload = dict(payload) + del log_step_payload["checkpoint_return_mean_reports10"] + logger.log_step(iteration=1, **log_step_payload) + assert logger._get_iter_env_steps_per_sec() == pytest.approx(400.0) + assert logger._get_learner_replay_rows_per_sec() == pytest.approx(1_024.0) header = logger._build_compact_header(include_status=False).plain assert "Steps/s 400" in header - assert "Samples/s 1,024" in header + assert "Rows/s 1,024" in header assert "Collector/s" not in header assert "GPUs 2" in logger._build_display().title.plain runner._restore_local_logger_statistics(logger) assert logger._total_steps == 100 assert logger._buffer_size == 50 - assert logger._collector_active_steps_per_sec == pytest.approx(1_000.0) def test_only_rank_zero_owns_terminal_and_tensorboard_backend(): diff --git a/tests/benchmark/test_extract_offpolicy_metrics.py b/tests/benchmark/test_extract_offpolicy_metrics.py new file mode 100644 index 000000000..befb6a7f0 --- /dev/null +++ b/tests/benchmark/test_extract_offpolicy_metrics.py @@ -0,0 +1,81 @@ +from __future__ import annotations + +import json +import math +from pathlib import Path + +import pytest +from scripts.benchmark.rl import extract_offpolicy_metrics as extract + +torch = pytest.importorskip("torch") +from torch.utils.tensorboard import SummaryWriter # noqa: E402 + + +def _write_tfevents(log_dir: Path, series: dict[str, list[float]]) -> None: + log_dir.mkdir(parents=True, exist_ok=True) + writer = SummaryWriter(log_dir=str(log_dir)) + try: + for tag, values in series.items(): + for step, value in enumerate(values): + writer.add_scalar(tag, value, step) + finally: + writer.close() + + +def _run_json(log_dir: Path, capsys: pytest.CaptureFixture) -> dict[str, float]: + assert extract.main([str(log_dir), "--json", "--last", "20"]) == 0 + return json.loads(capsys.readouterr().out) + + +def test_extract_new_schema_tags_with_unit_conversion(tmp_path: Path, capsys) -> None: + _write_tfevents( + tmp_path, + { + "Perf/iteration_time": [1.0, 1.5], # seconds + "Perf/learning_time": [0.4, 0.6], # seconds + "Perf/total_fps": [100.0, 300.0], + "Perf/learner_collector_wait_ms": [5.0, 7.0], + }, + ) + rows = _run_json(tmp_path, capsys) + assert rows["iter_ms"] == pytest.approx(1_250.0) + assert rows["learner_train_ms"] == pytest.approx(500.0) + assert rows["steps_per_sec"] == pytest.approx(200.0) + assert rows["learner_collector_wait_ms"] == pytest.approx(6.0) + # Retired charts without a canonical replacement stay NaN on new runs. + assert math.isnan(rows["collector_active_steps_per_sec"]) + assert math.isnan(rows["collector_cycle_ms"]) + + +def test_extract_legacy_tags_fall_back_unscaled(tmp_path: Path, capsys) -> None: + _write_tfevents( + tmp_path, + { + "perf/iter_ms": [1_000.0, 1_500.0], + "timing/learner_train_ms": [400.0, 600.0], + "perf/steps_per_sec": [100.0, 300.0], + "perf/collector_active_steps_per_sec": [50.0, 70.0], + }, + ) + rows = _run_json(tmp_path, capsys) + assert rows["iter_ms"] == pytest.approx(1_250.0) + assert rows["learner_train_ms"] == pytest.approx(500.0) + assert rows["steps_per_sec"] == pytest.approx(200.0) + assert rows["collector_active_steps_per_sec"] == pytest.approx(60.0) + + +def test_extract_prefers_new_tag_when_both_exist(tmp_path: Path, capsys) -> None: + _write_tfevents( + tmp_path, + { + "Perf/total_fps": [999.0], + "perf/steps_per_sec": [111.0], + }, + ) + rows = _run_json(tmp_path, capsys) + assert rows["steps_per_sec"] == pytest.approx(999.0) + + +def test_extract_missing_event_file_is_an_error(tmp_path: Path, capsys) -> None: + assert extract.main([str(tmp_path)]) == 1 + assert "No event file found" in capsys.readouterr().err diff --git a/tests/benchmark/test_offpolicy_dp_scaling_benchmark.py b/tests/benchmark/test_offpolicy_dp_scaling_benchmark.py index 4414f3a9e..358ff6afc 100644 --- a/tests/benchmark/test_offpolicy_dp_scaling_benchmark.py +++ b/tests/benchmark/test_offpolicy_dp_scaling_benchmark.py @@ -34,6 +34,15 @@ def _write_run_summary(run_dir: Path, **overrides: object) -> None: (run_dir / "run_summary.json").write_text(json.dumps(summary), encoding="utf-8") +def _write_run_config(run_dir: Path, *, batch_size: int, updates_per_step: int) -> None: + run_dir.mkdir(parents=True, exist_ok=True) + payload = { + "run": {}, + "config": {"algo": {"batch_size": batch_size, "updates_per_step": updates_per_step}}, + } + (run_dir / "run_config.json").write_text(json.dumps(payload), encoding="utf-8") + + def _make_run_dir( run_dir: Path, *, @@ -42,15 +51,40 @@ def _make_run_dir( reward: list[float] | None = None, dp_sync_time: list[float] | None = None, ) -> None: + """A pre-1.4.1 run directory carrying only the legacy scalar tags.""" _write_run_summary(run_dir) rank0_series = { - bench.STEPS_PER_SEC_TAG: steps_per_sec or [], - bench.SAMPLES_PER_SEC_TAG: samples_per_sec or [], + "perf/steps_per_sec": steps_per_sec or [], + "perf/effective_samples_per_sec": samples_per_sec or [], } if reward is not None: - rank0_series[bench.REWARD_TAG] = reward + rank0_series["reward/mean"] = reward if dp_sync_time is not None: - rank0_series[bench.DP_SYNC_TIME_TAG] = dp_sync_time + rank0_series["train/dp_sync_time"] = dp_sync_time + _write_tfevents(run_dir, rank0_series) + + +def _make_canonical_run_dir( + run_dir: Path, + *, + total_fps: list[float], + iteration_time: list[float], + batch_size: int = 8_192, + updates_per_step: int = 4, + reward: list[float] | None = None, + dp_sync_ms_per_rank: list[float] | None = None, +) -> None: + """A post-1.4.1 run directory: canonical tags only, no legacy tags.""" + _write_run_summary(run_dir) + _write_run_config(run_dir, batch_size=batch_size, updates_per_step=updates_per_step) + rank0_series = { + "Perf/total_fps": total_fps, + "Perf/iteration_time": iteration_time, + } + if reward is not None: + rank0_series["Train/mean_reward"] = reward + if dp_sync_ms_per_rank is not None: + rank0_series["Perf/dp_gradient_sync_ms_per_rank"] = dp_sync_ms_per_rank _write_tfevents(run_dir, rank0_series) @@ -136,6 +170,54 @@ def test_parse_run_missing_learner_throughput_is_a_hard_error(tmp_path: Path) -> bench.parse_run(run_dir, world_size=2) +def test_parse_run_canonical_tags_derive_learner_throughput(tmp_path: Path) -> None: + run_dir = tmp_path / "n2" + _make_canonical_run_dir( + run_dir, + total_fps=[100.0, 200.0, 300.0, 400.0], + iteration_time=[0.5, 0.5, 1.0, 1.0], + batch_size=1_000, + updates_per_step=2, + reward=[0.5, 1.5], + ) + metrics = bench.parse_run(run_dir, world_size=2) + assert metrics["steady_state_collector_steps_per_s"] == pytest.approx(350.0) + # Derived per iteration: world_size * batch_size * updates_per_step / + # iteration_time = 2 * 1000 * 2 / 0.5 = 8000, 8000, 4000, 4000. + assert metrics["steady_state_learner_samples_per_s"] == pytest.approx(4_000.0) + assert metrics["learner_throughput_source"] == "derived" + assert metrics["final_mean_reward"] == pytest.approx(1.5) + + +def test_parse_run_canonical_dp_sync_ms_is_converted_to_seconds(tmp_path: Path) -> None: + run_dir = tmp_path / "n2" + _make_canonical_run_dir( + run_dir, + total_fps=[100.0], + iteration_time=[0.5], + dp_sync_ms_per_rank=[20.0, 40.0], + ) + metrics = bench.parse_run(run_dir, world_size=2) + assert metrics["mean_dp_sync_time_sec"] == pytest.approx(0.03) + + +def test_parse_run_prefers_canonical_tags_over_legacy(tmp_path: Path) -> None: + run_dir = tmp_path / "n1" + _write_run_summary(run_dir) + _write_tfevents( + run_dir, + { + "Perf/total_fps": [999.0], + "perf/steps_per_sec": [111.0], + "perf/effective_samples_per_sec": [500.0], + }, + ) + metrics = bench.parse_run(run_dir, world_size=1) + assert metrics["steady_state_collector_steps_per_s"] == pytest.approx(999.0) + assert metrics["steady_state_learner_samples_per_s"] == pytest.approx(500.0) + assert metrics["learner_throughput_source"] == "tfevents" + + def test_parse_run_rejects_non_completed_status(tmp_path: Path) -> None: run_dir = tmp_path / "n1" _make_run_dir(run_dir, steps_per_sec=[100.0], samples_per_sec=[200.0]) diff --git a/tests/utils/test_experiment_tracking.py b/tests/utils/test_experiment_tracking.py index e0d5ec0bf..24c616d99 100644 --- a/tests/utils/test_experiment_tracking.py +++ b/tests/utils/test_experiment_tracking.py @@ -127,7 +127,7 @@ def _fake_time() -> float: logger.log_step(iteration=1, train_time=0.01, collector_wait_time=0.0) assert update_refresh_values == [] - logger.log_collector(total_steps=128, buffer_size=128, mean_reward=2.0) + logger.log_collector(total_steps=128, buffer_size=128) logger.log_status("Collector metrics updated") logger.log_save("/tmp/model_2.pt") assert update_refresh_values == [] @@ -154,12 +154,12 @@ def test_onpolicy_logger_uses_offpolicy_terminal_layout(): logger.start() logger.log_step( iteration=1, - metrics={"surrogate": 0.1, "value_loss": 0.2}, - reward=3.0, + metrics={"Loss/surrogate": 0.1, "Loss/value": 0.2}, + return_mean_ep100=3.0, collect_time=0.02, train_time=0.03, ) - logger.update_ep_length(12.0) + logger.update_mean_episode_length(12.0) console = Console(record=True, width=120) console.print(logger._build_display()) @@ -188,8 +188,7 @@ def test_offpolicy_logger_terminal_keeps_core_bottleneck_timing_rows(): ) logger.log_step( iteration=1, - metrics={"qf_loss": 0.1}, - reward=1.0, + metrics={"Loss/critic": 0.1}, train_time=0.5, collector_wait_time=1.0, replay_batch_wait_time=0.005, @@ -206,8 +205,7 @@ def test_offpolicy_logger_terminal_keeps_core_bottleneck_timing_rows(): "throughput_steps": 8, "batch_size_per_rank": 8, "effective_batch_size": 8, - "replay_samples_per_iter": 64, - "learner_samples_per_iter": 64, + "learner_replay_rows_per_iter": 64, }, ) @@ -218,7 +216,7 @@ def test_offpolicy_logger_terminal_keeps_core_bottleneck_timing_rows(): assert "Replay Wait" not in output assert "Collector Wait" in output assert "Inference" in output - assert "Train" in output + assert "Learning" in output assert "Iter Wall" in output assert "Replay Batch Wait" in output assert "Replay Sample" in output @@ -241,7 +239,7 @@ def test_offpolicy_logger_terminal_keeps_core_bottleneck_timing_rows(): assert "Batch/Rank" in output assert "Batch/Update" not in output assert "Samples/Iter" not in output - assert "Samples/s" in output + assert "Rows/s" in output logger.close() @@ -286,7 +284,7 @@ def test_offpolicy_logger_terminal_shows_replay_rows_and_effective_batch(): ) logger.log_step( iteration=1, - metrics={"qf_loss": 0.1}, + metrics={"Loss/critic": 0.1}, train_time=0.5, collector_wait_time=0.1, iteration_time=1.0, @@ -294,8 +292,7 @@ def test_offpolicy_logger_terminal_shows_replay_rows_and_effective_batch(): "throughput_steps": 4, "batch_size_per_rank": 16, "effective_batch_size": 16, - "replay_samples_per_iter": 8, - "learner_samples_per_iter": 16, + "learner_replay_rows_per_iter": 16, }, ) @@ -306,7 +303,7 @@ def test_offpolicy_logger_terminal_shows_replay_rows_and_effective_batch(): assert "Batch/Rank" in output assert "Replay/Iter" not in output assert "Samples/Iter" not in output - assert "Samples/s" in output + assert "Rows/s" in output assert "Rank Barrier" not in output assert "Param Sync" not in output assert "Other Loop" not in output @@ -530,57 +527,46 @@ def test_offpolicy_logger_logs_wait_and_iter_throughput(monkeypatch): metrics={}, train_time=0.75, collector_wait_time=10.0, - inference_time=0.03, - sync_coordination_time=0.04, - replay_batch_wait_time=0.06, learner_replay_stage_time=0.02, - replay_ingress_h2d_submit_time=0.07, weight_sync_time=0.05, iteration_time=10.9, extra_info={ "throughput_steps": 8, - "collector_active_steps_per_sec": 1234.5, "batch_size_per_rank": 8, "effective_batch_size": 8, - "replay_samples_per_iter": 32, - "learner_samples_per_iter": 32, + "learner_replay_rows_per_iter": 32, }, ) payload, step = fake_wandb.log_calls[-1] assert step == 1 - assert payload["timing/learner_collector_wait_ms"] == 10_000.0 + assert payload["Perf/learner_collector_wait_ms"] == 10_000.0 assert "timing/learner_collect_ms" not in payload assert "timing/learner_replay_wait_ms" not in payload - assert payload["timing/learner_replay_stage_ms"] == 20.0 - assert payload["timing/learner_replay_sample_ms"] == 0.0 - assert payload["timing/learner_train_ms"] == 750.0 + assert payload["Perf/learner_replay_stage_ms"] == 20.0 + assert payload["Perf/learner_replay_sample_ms"] == 0.0 + assert payload["Perf/learning_time"] == pytest.approx(0.75) assert "timing/learner_param_sync_ms" not in payload - assert payload["timing/learner_weight_publish_ms"] == 50.0 + assert payload["Perf/learner_weight_publish_ms"] == 50.0 for key in ( - "timing/learner_inference_ms", - "timing/learner_collector_release_ms", - "timing/learner_replay_batch_wait_ms", - "timing/learner_inference_h2d_ms", - "timing/learner_inference_forward_ms", - "timing/learner_inference_d2h_ms", - "timing/replay_ingress_h2d_submit_ms", + "Perf/learner_inference_ms", + "Perf/learner_collector_release_ms", + "Perf/learner_replay_batch_wait_ms", + "Perf/learner_inference_h2d_ms", + "Perf/learner_inference_forward_ms", + "Perf/learner_inference_d2h_ms", + "Perf/replay_ingress_h2d_submit_ms", ): assert key not in payload - assert payload["timing/learner_other_ms"] == pytest.approx(80.0) + assert "timing/learner_other_ms" not in payload assert "perf/learner_pipeline_ms" not in payload - assert payload["perf/iter_ms"] == pytest.approx(10_900.0) + assert payload["Perf/iteration_time"] == pytest.approx(10.9) + assert "perf/iter_ms" not in payload assert "perf/iter_unaccounted_ms" not in payload - assert payload["perf/learner_train_pct"] == pytest.approx(0.75 / 10.9 * 100) - assert payload["perf/learner_other_pct"] == pytest.approx(0.08 / 10.9 * 100) - assert payload["perf/learner_accounted_pct"] == pytest.approx(10.82 / 10.9 * 100) - assert payload["perf/steps_per_sec"] == pytest.approx(8.0 / 10.9) - assert payload["perf/collector_active_steps_per_sec"] == pytest.approx(1234.5) - assert payload["perf/effective_samples_per_sec"] == pytest.approx(32.0 / 10.9) - for key in ("axis/iteration", "axis/env_steps_total"): + assert payload["Perf/total_fps"] == pytest.approx(8.0 / 10.9) + for key in ("axis/iteration", "axis/env_steps_total", "Train/iteration"): assert key not in payload assert not any(key.startswith("distributed/") for key in payload) - assert "perf/effective_samples_per_sec_smoothed" not in payload assert "perf/collect_train_ratio" not in payload logger.finish() @@ -599,7 +585,7 @@ def test_offpolicy_logger_logs_collector_phase_timing_to_backends(monkeypatch): wandb_logger.log_step(iteration=3, metrics={}, train_time=0.1) payload, _ = fake_wandb.log_calls[-1] - assert payload["timing/collector_replay_write_ms"] == pytest.approx(1.25) + assert payload["Perf/collector_replay_write_ms"] == pytest.approx(1.25) wandb_logger.finish() tb_writer = _FakeTensorBoardWriter() @@ -612,7 +598,7 @@ def test_offpolicy_logger_logs_collector_phase_timing_to_backends(monkeypatch): tb_logger.update_collector_timing({"replay_write_ms": 2.5}) tb_logger.log_step(iteration=4, metrics={}, train_time=0.1) - assert ("timing/collector_replay_write_ms", 2.5, 4) in tb_writer.scalars + assert ("Perf/collector_replay_write_ms", 2.5, 4) in tb_writer.scalars tb_logger.finish() @@ -657,23 +643,22 @@ def test_offpolicy_logger_uses_same_canonical_timing_names_in_terminal_and_backe "Collector Release", "Replay Batch Wait", "Replay Sample", - "Train", + "Learning", "Other", "Iter Wall", ] for key in ( - "timing/learner_collector_wait_ms", - "timing/learner_inference_ms", - "timing/learner_collector_release_ms", - "timing/learner_replay_batch_wait_ms", - "timing/learner_replay_sample_ms", - "timing/learner_train_ms", - "timing/learner_other_ms", - "perf/iter_ms", + "Perf/learner_collector_wait_ms", + "Perf/learner_inference_ms", + "Perf/learner_collector_release_ms", + "Perf/learner_replay_batch_wait_ms", + "Perf/learner_replay_sample_ms", + "Perf/learning_time", + "Perf/iteration_time", ): assert key in payload - assert "timing/learner_replay_stage_ms" not in payload - assert "timing/learner_weight_publish_ms" not in payload + assert "Perf/learner_replay_stage_ms" not in payload + assert "Perf/learner_weight_publish_ms" not in payload assert collector_labels[:4] == [ "Inference Request", "Learner Action Wait", @@ -681,15 +666,14 @@ def test_offpolicy_logger_uses_same_canonical_timing_names_in_terminal_and_backe "Replay Write", ] for key in ( - "timing/collector_inference_request_ms", - "timing/collector_learner_action_wait_ms", - "timing/collector_env_step_ms", - "timing/collector_replay_write_ms", - "perf/collector_cycle_ms", + "Perf/collector_inference_request_ms", + "Perf/collector_learner_action_wait_ms", + "Perf/collector_env_step_ms", + "Perf/collector_replay_write_ms", ): assert key in payload + assert "perf/collector_cycle_ms" not in payload assert "timing/collector_bookkeeping_ms" not in payload - assert payload["perf/collector_cycle_ms"] == pytest.approx(10.0) logger.finish() @@ -718,45 +702,41 @@ def test_offpolicy_logger_tensorboard_logs_wall_clock_without_axis_scalars(): iteration_time=2.15, extra_info={ "throughput_steps": 16, - "collector_active_steps_per_sec": 4321.0, "batch_size_per_rank": 16, "effective_batch_size": 16, - "replay_samples_per_iter": 64, - "learner_samples_per_iter": 64, + "learner_replay_rows_per_iter": 64, }, ) scalars = {tag: value for tag, value, _ in tb_writer.scalars} - assert scalars["timing/learner_collector_wait_ms"] == pytest.approx(1_000.0) + assert scalars["Perf/learner_collector_wait_ms"] == pytest.approx(1_000.0) assert "timing/learner_replay_wait_ms" not in scalars - assert scalars["timing/learner_replay_batch_wait_ms"] == pytest.approx(20.0) - assert scalars["timing/replay_ingress_h2d_submit_ms"] == pytest.approx(50.0) - assert scalars["timing/learner_replay_sample_ms"] == pytest.approx(30.0) - assert scalars["timing/learner_collector_release_ms"] == pytest.approx(40.0) - assert scalars["timing/learner_inference_h2d_ms"] == pytest.approx(10.0) - assert scalars["timing/learner_inference_forward_ms"] == pytest.approx(180.0) - assert scalars["timing/learner_inference_d2h_ms"] == pytest.approx(10.0) - assert scalars["timing/learner_inference_ms"] == pytest.approx(200.0) + assert scalars["Perf/learner_replay_batch_wait_ms"] == pytest.approx(20.0) + assert scalars["Perf/replay_ingress_h2d_submit_ms"] == pytest.approx(50.0) + assert scalars["Perf/learner_replay_sample_ms"] == pytest.approx(30.0) + assert scalars["Perf/learner_collector_release_ms"] == pytest.approx(40.0) + assert scalars["Perf/learner_inference_h2d_ms"] == pytest.approx(10.0) + assert scalars["Perf/learner_inference_forward_ms"] == pytest.approx(180.0) + assert scalars["Perf/learner_inference_d2h_ms"] == pytest.approx(10.0) + assert scalars["Perf/learner_inference_ms"] == pytest.approx(200.0) + assert scalars["Perf/learning_time"] == pytest.approx(0.7) assert "timing/learner_param_sync_ms" not in scalars - assert "timing/learner_replay_stage_ms" not in scalars - assert "timing/learner_weight_publish_ms" not in scalars - assert scalars["timing/learner_other_ms"] == pytest.approx(160.0) + assert "Perf/learner_replay_stage_ms" not in scalars + assert "Perf/learner_weight_publish_ms" not in scalars + assert "timing/learner_other_ms" not in scalars assert "perf/learner_pipeline_ms" not in scalars - assert scalars["perf/iter_ms"] == pytest.approx(2_150.0) - assert scalars["perf/learner_train_pct"] == pytest.approx(0.7 / 2.15 * 100) - assert scalars["perf/learner_other_pct"] == pytest.approx(0.16 / 2.15 * 100) - assert scalars["perf/collector_active_steps_per_sec"] == pytest.approx(4321.0) - assert scalars["perf/effective_samples_per_sec"] == pytest.approx(64.0 / 2.15) - assert scalars["episode/timeout_rate"] == pytest.approx(0.0) + assert scalars["Perf/iteration_time"] == pytest.approx(2.15) + assert scalars["Perf/total_fps"] == pytest.approx(16.0 / 2.15) + assert "Episode/timeout_rate" not in scalars + assert "episode/timeout_rate" not in scalars assert "episode/terminated_rate" not in scalars for key in ("axis/iteration", "axis/env_steps_total"): assert key not in scalars assert not any(key.startswith("distributed/") for key in scalars) - assert "perf/effective_samples_per_sec_smoothed" not in scalars logger.finish() -def test_offpolicy_logger_logs_reward_comparison_metrics(monkeypatch): +def test_offpolicy_logger_logs_episode_return_and_reward_terms(monkeypatch): fake_wandb = _FakeWandb() monkeypatch.setitem(sys.modules, "wandb", fake_wandb) @@ -769,14 +749,16 @@ def test_offpolicy_logger_logs_reward_comparison_metrics(monkeypatch): logger.log_step( iteration=2, metrics={}, - reward=3.0, - reward_metrics={"mean_ep100": 2.0}, + return_mean_ep100=2.0, + reward_components={"tracking": 3.0}, ) payload, step = fake_wandb.log_calls[-1] assert step == 128 - assert payload["reward/mean"] == 3.0 - assert payload["reward/mean_ep100"] == 2.0 + assert payload["Train/mean_reward"] == 2.0 + assert payload["reward/tracking"] == 3.0 + assert "reward/mean" not in payload + assert "reward/mean_ep100" not in payload assert "reward/mean_unilab_100x100" not in payload logger.finish() @@ -791,8 +773,6 @@ def test_offpolicy_logger_omits_iteration_extra_fields_when_not_supplied(monkeyp env_name="Go2JoystickFlat", log_backend="wandb", ) - logger._start_time = 1.0 - monkeypatch.setattr(common_module.time, "time", lambda: 2.0) logger.log_collector(total_steps=8, buffer_size=8) logger.log_step( iteration=1, @@ -805,9 +785,12 @@ def test_offpolicy_logger_omits_iteration_extra_fields_when_not_supplied(monkeyp payload, _ = fake_wandb.log_calls[-1] assert "timing/learner_collect_ms" not in payload - assert "perf/steps_per_sec" not in payload - assert "timing/learner_weight_publish_ms" not in payload - assert payload["perf/iter_ms"] == pytest.approx(1_750.0) + assert "Perf/total_fps" not in payload + assert "Perf/iteration_time" not in payload + assert "Perf/learner_weight_publish_ms" not in payload + assert payload["Perf/learning_time"] == pytest.approx(0.75) + assert payload["Perf/learner_collector_wait_ms"] == pytest.approx(1_000.0) + assert payload["Perf/replay_ingress_h2d_submit_ms"] == pytest.approx(20.0) logger.finish() diff --git a/uv.lock b/uv.lock index 2da768893..c0718b513 100644 --- a/uv.lock +++ b/uv.lock @@ -5268,7 +5268,7 @@ requires-dist = [ { name = "tqdm" }, { name = "trimesh", specifier = ">=3.21.7" }, { name = "typing-extensions" }, - { name = "unilab-rl", marker = "extra == 'uni-rl'", specifier = "==1.4.0" }, + { name = "unilab-rl", marker = "extra == 'uni-rl'", specifier = "==1.4.1" }, { name = "unisim-core", specifier = ">=1.7.8" }, { name = "unisim-core", extras = ["motrix"], marker = "extra == 'motrix'", specifier = ">=1.7.8" }, { name = "unisim-core", extras = ["superdex"], marker = "python_full_version >= '3.12' and platform_machine == 'x86_64' and sys_platform == 'linux' and extra == 'superdex'", specifier = ">=1.7.8" }, @@ -5286,12 +5286,12 @@ dev = [ { name = "pytest" }, { name = "pytest-cov" }, { name = "ruff" }, - { name = "unilab-rl", specifier = "==1.4.0" }, + { name = "unilab-rl", specifier = "==1.4.1" }, ] [[package]] name = "unilab-rl" -version = "1.4.0" +version = "1.4.1" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "hydra-core" }, @@ -5307,9 +5307,9 @@ dependencies = [ { name = "torch", version = "2.14.0+cu130", source = { registry = "https://download-r2.pytorch.org/whl/cu130" }, marker = "sys_platform == 'linux' or sys_platform == 'win32'" }, { name = "wandb" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/58/64/859b40ae7ae5202b392573c37cf74e81e3f6c5d2cef05ed4b5be8a9e2c23/unilab_rl-1.4.0.tar.gz", hash = "sha256:52b638ed40ad58291a0c28709d32e84b272c8b370d9a6ce9dbf7e3fd0a9c4508", size = 131925, upload-time = "2026-09-25T15:06:01.671Z" } +sdist = { url = "https://files.pythonhosted.org/packages/af/90/2b8905acce6a5fdca7bab1ac5bcfccf7017fa3b286dadc307f110f8249e1/unilab_rl-1.4.1.tar.gz", hash = "sha256:440566af3422e1915587b62f7a47cc1c0e40615dff7fd94b1ff040f9725d5d3a", size = 142073, upload-time = "2026-09-26T18:24:13.396Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/73/b6/265b3dc612659c08afd31f3583ffb98e76eec1845d0bc823ab87edd08e96/unilab_rl-1.4.0-py3-none-any.whl", hash = "sha256:68ced237a3b42e4efc43602cc598d6cc319e2bb0d7863a8c84bfdd3bfab7658d", size = 163240, upload-time = "2026-09-25T15:06:00.401Z" }, + { url = "https://files.pythonhosted.org/packages/77/29/b33c7e3ec8b9f0e1efeda53c755f8ffe973d8a0b8c8c13beb8d122e5a89f/unilab_rl-1.4.1-py3-none-any.whl", hash = "sha256:5efa10f5b6a2705509ef51f3c44e4705d97c835144779d9851b93999575a15ac", size = 174061, upload-time = "2026-09-26T18:24:12.088Z" }, ] [[package]] diff --git a/uv.rocm.lock b/uv.rocm.lock index 1ab535eb4..900e0dad7 100644 --- a/uv.rocm.lock +++ b/uv.rocm.lock @@ -3511,7 +3511,7 @@ requires-dist = [ { name = "trimesh", specifier = ">=3.21.7" }, { name = "triton-rocm", marker = "platform_machine == 'x86_64' and sys_platform == 'linux'", specifier = "==3.8.0", index = "https://download.pytorch.org/whl/rocm7.2" }, { name = "typing-extensions" }, - { name = "unilab-rl", marker = "extra == 'uni-rl'", specifier = "==1.4.0" }, + { name = "unilab-rl", marker = "extra == 'uni-rl'", specifier = "==1.4.1" }, { name = "unisim-core", specifier = ">=1.7.8" }, { name = "unisim-core", extras = ["motrix"], marker = "extra == 'motrix'", specifier = ">=1.7.8" }, { name = "viser", specifier = ">=1.0.26" }, @@ -3526,12 +3526,12 @@ dev = [ { name = "pytest" }, { name = "pytest-cov" }, { name = "ruff" }, - { name = "unilab-rl", specifier = "==1.4.0" }, + { name = "unilab-rl", specifier = "==1.4.1" }, ] [[package]] name = "unilab-rl" -version = "1.4.0" +version = "1.4.1" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "hydra-core" }, @@ -3548,9 +3548,9 @@ dependencies = [ { name = "torch", version = "2.14.0+rocm7.2", source = { registry = "https://download.pytorch.org/whl/rocm7.2" }, marker = "platform_machine == 'x86_64' and sys_platform == 'linux'" }, { name = "wandb" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/58/64/859b40ae7ae5202b392573c37cf74e81e3f6c5d2cef05ed4b5be8a9e2c23/unilab_rl-1.4.0.tar.gz", hash = "sha256:52b638ed40ad58291a0c28709d32e84b272c8b370d9a6ce9dbf7e3fd0a9c4508", size = 131925, upload-time = "2026-09-25T15:06:01.671Z" } +sdist = { url = "https://files.pythonhosted.org/packages/af/90/2b8905acce6a5fdca7bab1ac5bcfccf7017fa3b286dadc307f110f8249e1/unilab_rl-1.4.1.tar.gz", hash = "sha256:440566af3422e1915587b62f7a47cc1c0e40615dff7fd94b1ff040f9725d5d3a", size = 142073, upload-time = "2026-09-26T18:24:13.396Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/73/b6/265b3dc612659c08afd31f3583ffb98e76eec1845d0bc823ab87edd08e96/unilab_rl-1.4.0-py3-none-any.whl", hash = "sha256:68ced237a3b42e4efc43602cc598d6cc319e2bb0d7863a8c84bfdd3bfab7658d", size = 163240, upload-time = "2026-09-25T15:06:00.401Z" }, + { url = "https://files.pythonhosted.org/packages/77/29/b33c7e3ec8b9f0e1efeda53c755f8ffe973d8a0b8c8c13beb8d122e5a89f/unilab_rl-1.4.1-py3-none-any.whl", hash = "sha256:5efa10f5b6a2705509ef51f3c44e4705d97c835144779d9851b93999575a15ac", size = 174061, upload-time = "2026-09-26T18:24:12.088Z" }, ] [[package]]