From 411a6ee3e0e0fe5a852c93f40774a3b4af4b4b0e Mon Sep 17 00:00:00 2001 From: TATP-233 Date: Thu, 27 Aug 2026 14:32:13 +0800 Subject: [PATCH 1/2] fix(logging): make collector reward reporting timely Reward displays (tensorboard reward/mean and the terminal logger) lagged badly on off-policy and APPO runs: - collectors sent metrics only every num_envs * 10 env steps, so the reported reward changed just once per ~10 learner iterations; - runners then averaged the last 100 (off-policy) or 50 (APPO) reports, each already a rolling 100-episode mean, delaying the visible curve by ~1000 iterations. Report metrics every collector cycle, keep the runner-side window at the last 10 reports, and bound the per-worker episode reward/length buffers with deque(maxlen=100) instead of lists that grew for the whole run. Co-Authored-By: Claude Fable 5 --- src/unilab/algos/appo/runner.py | 9 +++++---- src/unilab/algos/appo/worker.py | 16 +++++++++------ src/unilab/algos/hora/appo_runner.py | 9 +++++---- src/unilab/algos/hora/appo_worker.py | 16 +++++++++------ .../algos/offpolicy/double_buffer_runner.py | 5 ++++- src/unilab/algos/offpolicy/worker.py | 20 ++++++++++--------- 6 files changed, 45 insertions(+), 30 deletions(-) diff --git a/src/unilab/algos/appo/runner.py b/src/unilab/algos/appo/runner.py index 4473b46ac..370ad4b2b 100644 --- a/src/unilab/algos/appo/runner.py +++ b/src/unilab/algos/appo/runner.py @@ -312,7 +312,10 @@ def learn( ) logger_started = False - reward_history: deque = deque(maxlen=200) + # Recent collector reports; each entry is already the collector's + # rolling 100-episode mean, so a short window keeps the logged + # reward timely without losing smoothing. + reward_history: deque = deque(maxlen=10) latest_reward_components: dict = {} staging_pool = RolloutStagingPool( @@ -396,9 +399,7 @@ def learn( logger.update_staging_pool(staging_pool.active_count, staging_pool.capacity) mean_reward = ( - sum(list(reward_history)[-50:]) / max(len(list(reward_history)[-50:]), 1) - if reward_history - else 0.0 + sum(reward_history) / max(len(reward_history), 1) if reward_history else 0.0 ) last_mean_reward = float(mean_reward) best_mean_reward = max(best_mean_reward, last_mean_reward) diff --git a/src/unilab/algos/appo/worker.py b/src/unilab/algos/appo/worker.py index 6375aed69..35eed27d3 100644 --- a/src/unilab/algos/appo/worker.py +++ b/src/unilab/algos/appo/worker.py @@ -8,7 +8,7 @@ import statistics import sys import time -from collections import defaultdict +from collections import defaultdict, deque from queue import Empty, Full from typing import Any, Dict @@ -235,8 +235,10 @@ def to_float32_np(x): obs_td = TensorDict({"policy": obs_torch}, batch_size=num_envs, device=collector_device) total_steps = 0 - ep_rewards = [] - ep_lengths = [] + # Bounded rolling window of the most recent completed episodes; an + # unbounded list here grows for the entire run. + ep_rewards: deque[float] = deque(maxlen=100) + ep_lengths: deque[int] = deque(maxlen=100) current_ep_rewards = np.zeros(num_envs, dtype=np.float32) current_ep_lengths = np.zeros(num_envs, dtype=np.int32) ep_reward_components = defaultdict(list) @@ -353,15 +355,17 @@ def to_float32_np(x): if k.startswith("reward/"): ep_reward_components[k].append(v) - if metrics_queue is not None and total_steps % (num_envs * 10) == 0: + # Report every env step so learner-side reward and throughput + # displays track the current policy without extra lag. + if metrics_queue is not None: try: msg: dict[str, Any] = { "total_steps": total_steps, } if ep_rewards: - msg["mean_ep_reward"] = statistics.mean(ep_rewards[-100:]) + msg["mean_ep_reward"] = statistics.mean(ep_rewards) msg["mean_ep_length"] = ( - statistics.mean(ep_lengths[-100:]) if ep_lengths else 0.0 + statistics.mean(ep_lengths) if ep_lengths else 0.0 ) if ep_completions > 0: msg["timeout_rate"] = ep_timeouts / ep_completions diff --git a/src/unilab/algos/hora/appo_runner.py b/src/unilab/algos/hora/appo_runner.py index 252591023..3bfb1b965 100644 --- a/src/unilab/algos/hora/appo_runner.py +++ b/src/unilab/algos/hora/appo_runner.py @@ -303,7 +303,10 @@ def learn( f"epochs={learner.num_learning_epochs})" ) - reward_history: deque = deque(maxlen=200) + # Recent collector reports; each entry is already the collector's + # rolling 100-episode mean, so a short window keeps the logged + # reward timely without losing smoothing. + reward_history: deque = deque(maxlen=10) latest_reward_components: dict = {} staging_pool = RolloutStagingPool( capacity=self.staging_pool_size, @@ -379,9 +382,7 @@ def learn( logger.update_staging_pool(staging_pool.active_count, staging_pool.capacity) mean_reward = ( - sum(list(reward_history)[-50:]) / max(len(list(reward_history)[-50:]), 1) - if reward_history - else 0.0 + sum(reward_history) / max(len(reward_history), 1) if reward_history else 0.0 ) last_mean_reward = float(mean_reward) best_mean_reward = max(best_mean_reward, last_mean_reward) diff --git a/src/unilab/algos/hora/appo_worker.py b/src/unilab/algos/hora/appo_worker.py index a5cccb432..0d5ec8002 100644 --- a/src/unilab/algos/hora/appo_worker.py +++ b/src/unilab/algos/hora/appo_worker.py @@ -5,7 +5,7 @@ import statistics import sys import time -from collections import defaultdict +from collections import defaultdict, deque from typing import Any, Dict import numpy as np @@ -227,8 +227,10 @@ def to_float32_np(x): ) total_steps = 0 - ep_rewards = [] - ep_lengths = [] + # Bounded rolling window of the most recent completed episodes; an + # unbounded list here grows for the entire run. + ep_rewards: deque[float] = deque(maxlen=100) + ep_lengths: deque[int] = deque(maxlen=100) current_ep_rewards = np.zeros(num_envs, dtype=np.float32) current_ep_lengths = np.zeros(num_envs, dtype=np.int32) ep_reward_components = defaultdict(list) @@ -365,15 +367,17 @@ def to_float32_np(x): if k.startswith("reward/"): ep_reward_components[k].append(v) - if metrics_queue is not None and total_steps % (num_envs * 10) == 0: + # Report every env step so learner-side reward and throughput + # displays track the current policy without extra lag. + if metrics_queue is not None: try: msg: dict[str, Any] = { "total_steps": total_steps, } if ep_rewards: - msg["mean_ep_reward"] = statistics.mean(ep_rewards[-100:]) + msg["mean_ep_reward"] = statistics.mean(ep_rewards) msg["mean_ep_length"] = ( - statistics.mean(ep_lengths[-100:]) if ep_lengths else 0.0 + statistics.mean(ep_lengths) if ep_lengths else 0.0 ) if ep_completions > 0: msg["timeout_rate"] = ep_timeouts / ep_completions diff --git a/src/unilab/algos/offpolicy/double_buffer_runner.py b/src/unilab/algos/offpolicy/double_buffer_runner.py index f689811ba..ce506b23b 100644 --- a/src/unilab/algos/offpolicy/double_buffer_runner.py +++ b/src/unilab/algos/offpolicy/double_buffer_runner.py @@ -969,7 +969,10 @@ def learn( time.sleep(0.5) - reward_history: deque = deque(maxlen=100) + # Recent collector reports; each entry is already the collector's + # rolling 100-episode mean, so a short window keeps the logged + # reward timely without losing smoothing. + reward_history: deque = deque(maxlen=10) latest_reward_components: dict[str, float] = {} has_logged_reward = False last_buf_log = 0 diff --git a/src/unilab/algos/offpolicy/worker.py b/src/unilab/algos/offpolicy/worker.py index 7bb8e748b..07fa38fca 100644 --- a/src/unilab/algos/offpolicy/worker.py +++ b/src/unilab/algos/offpolicy/worker.py @@ -243,12 +243,15 @@ def _run_collector( replay_buffer.trace_recorder = trace_recorder replay_buffer.trace_thread_time = trace_thread_time replay_buffer.attach_stop_event(stop_event) + from collections import defaultdict, deque + total_steps = 0 - ep_rewards = [] - ep_lengths = [] + # Bounded rolling window of the most recent completed episodes; an + # unbounded list here grows for the entire run. + ep_rewards: deque[float] = deque(maxlen=100) + ep_lengths: deque[int] = deque(maxlen=100) current_ep_rewards = np.zeros(num_envs, dtype=np.float32) current_ep_lengths = np.zeros(num_envs, dtype=np.int32) - from collections import defaultdict ep_reward_components = defaultdict(list) timing_accum_ms: defaultdict[str, float] = defaultdict(float) @@ -454,8 +457,9 @@ def _run_collector( if k.startswith("reward/"): ep_reward_components[k].append(v) - # Send metrics periodically - if metrics_queue is not None and total_steps % (num_envs * 10) == 0: + # Send metrics every collector cycle so learner-side reward and + # throughput displays track the current policy without extra lag. + if metrics_queue is not None: import statistics try: @@ -464,10 +468,8 @@ def _run_collector( "buffer_size": int(replay_buffer.size[0]), } if ep_rewards: - msg["mean_ep_reward"] = statistics.mean(ep_rewards[-100:]) - msg["mean_ep_length"] = ( - statistics.mean(ep_lengths[-100:]) if ep_lengths else 0.0 - ) + msg["mean_ep_reward"] = statistics.mean(ep_rewards) + msg["mean_ep_length"] = statistics.mean(ep_lengths) if ep_lengths else 0.0 # Add mean reward components if ep_reward_components: components_mean = {} From 22fc439ea546d4b38b3add57b50770291aa73811 Mon Sep 17 00:00:00 2001 From: TATP-233 Date: Thu, 27 Aug 2026 15:36:13 +0800 Subject: [PATCH 2/2] fix(env): keep per-step reward log entries through autoreset MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ManagerBasedRlEnv.reset() replaced state.info["log"] with the reset-only extras (Episode_Reward/*), wiping the fresh per-step reward/* entries that _update_state_in_read_phase() had just computed for the current transition. On any step where at least one env resets — with thousands of envs, nearly every step — collectors therefore saw no reward/* keys at all, so the per-term reward components in tensorboard and the terminal logger stayed frozen at one stale value for thousands of iterations (observed as long flat staircases on reward/motion_* etc.). Merge instead of replace on the autoreset path: the pre-reset per-step entries stay, reset extras layer on top. Standalone (non-autoreset) resets are unchanged. Co-Authored-By: Claude Fable 5 --- src/unilab/envs/manager_based_rl_env.py | 8 ++++++++ 1 file changed, 8 insertions(+) diff --git a/src/unilab/envs/manager_based_rl_env.py b/src/unilab/envs/manager_based_rl_env.py index 18369bcfa..07257827f 100644 --- a/src/unilab/envs/manager_based_rl_env.py +++ b/src/unilab/envs/manager_based_rl_env.py @@ -579,6 +579,14 @@ def reset( if self._state is not None: for name, values in reset_obs.items(): self._state.obs[name][ids] = values + if self._autoreset_reset_active: + # Autoreset runs at the tail of step(): keep this step's + # per-step log entries (reward/* etc., computed pre-reset) and + # layer the reset extras (Episode_Reward/* etc.) on top, so + # consumers still see the transition's reward breakdown. + step_log = self._state.info.get("log") + if step_log: + log = {**step_log, **log} self._state.info["log"] = log if not self._autoreset_reset_active: self._state.terminated[ids] = False