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 = {} 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