From e908109bb5bb652f6ef8097ef121ec807eac6fd9 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Sun, 26 Jul 2026 18:07:10 -0400 Subject: [PATCH 01/12] vec_log per-agent population mean: EMA slots instead of completion sums MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Rebased onto current 3.0 (the original branch predated the goal-system rewrite and metrics refactor). add_log no longer sums completed episodes into env->log — that weighted agents by completion frequency, so short-episode agents (early crashes) dominated the reported means and curves only looked clean after a resample-forced synchronized reset. Each agent slot now keeps an EMA of its completed-episode Log (seeded on first completion, then slot = alpha*slot + (1-alpha)*episode); prepare_log sums the seeded slots into env->log on demand without resetting them, and vec_log's divide-by-n yields a population mean with one weight per agent. binding.vec_prepare_log rebuilds the slots right before each vec_log call, inside drive.py's existing not-eval_mode guard (prepare_log overwrites env->log, which eval reads per episode). The vec_log gate becomes n >= 1 (the old n >= num_agents throttle made no sense for persistent per-agent state). log_ema_alpha is exposed end-to-end (drive.h, binding, drive.py, config_schema, puffer_drive.yaml; default 0.707 ~ 2-episode half-life). expert_static_car_count and static_car_count move into the per-agent episode log, so the divide-by-n recovers the true per-env count rather than count/num_agents. Their logged scale changes accordingly. Behavioral tests cover the emission gate, steady-state n == num_agents, and EMA persistence across emissions without new completions. Smoke goldens are not included; they are only bit-reproducible in the pinned CI image and are regenerated there. Co-Authored-By: Claude Fable 5 --- pufferlib/config/puffer_drive.yaml | 5 + pufferlib/config_schema.py | 1 + pufferlib/ocean/drive/binding.c | 18 ++- pufferlib/ocean/drive/drive.h | 143 +++++++++++------- pufferlib/ocean/drive/drive.py | 5 + pufferlib/ocean/env_binding.h | 6 +- .../test_drive_per_agent_logging.py | 121 +++++++++++++++ 7 files changed, 242 insertions(+), 57 deletions(-) create mode 100644 tests/unit_tests/test_drive_per_agent_logging.py diff --git a/pufferlib/config/puffer_drive.yaml b/pufferlib/config/puffer_drive.yaml index 5cfd5bdd8a..4949001a44 100644 --- a/pufferlib/config/puffer_drive.yaml +++ b/pufferlib/config/puffer_drive.yaml @@ -124,6 +124,11 @@ env: reward_timestep: 2.5e-05 reward_overspeed: 0.05 reward_ade: 0.0 + # --- Logging --- + # EMA coefficient for per-agent log aggregation: on each completion, + # slot = alpha*slot + (1-alpha)*episode. 0 keeps only the newest episode; + # 0.707 has a half-life of ~2 episodes per agent. + log_ema_alpha: 0.707 # --- Map --- # Path to map used for training map_dir: pufferlib/resources/drive/binaries/carla diff --git a/pufferlib/config_schema.py b/pufferlib/config_schema.py index 1d2d174997..99c14d19da 100644 --- a/pufferlib/config_schema.py +++ b/pufferlib/config_schema.py @@ -135,6 +135,7 @@ class DriveEnvConfig: reward_timestep: float = MISSING reward_overspeed: float = MISSING reward_ade: float = MISSING + log_ema_alpha: float = MISSING map_dir: str = MISSING num_maps: int = MISSING obs_slots_lane_n: int = MISSING diff --git a/pufferlib/ocean/drive/binding.c b/pufferlib/ocean/drive/binding.c index 170f43a8ff..f4eb4f6af9 100644 --- a/pufferlib/ocean/drive/binding.c +++ b/pufferlib/ocean/drive/binding.c @@ -25,14 +25,29 @@ static PyObject *map_cache_live_count_py( return PyLong_FromLong(live); } +static void prepare_log(Drive *env); +static PyObject *vec_prepare_log_py(PyObject *self __attribute__((unused)), PyObject *args); + // clang-format off #define MY_METHODS \ {"map_cache_size", map_cache_size_py, METH_NOARGS, "Map cache slot count."}, \ - {"map_cache_live_count", map_cache_live_count_py, METH_NOARGS, "Map cache live count."} + {"map_cache_live_count", map_cache_live_count_py, METH_NOARGS, "Map cache live count."}, \ + {"vec_prepare_log", vec_prepare_log_py, METH_VARARGS, "Aggregate per-agent log EMAs into each env->log."} // clang-format on #include "../env_binding.h" +static PyObject *vec_prepare_log_py(PyObject *self __attribute__((unused)), PyObject *args) { + VecEnv *vec = unpack_vecenv(args); + if (!vec) { + return NULL; + } + for (int i = 0; i < vec->num_envs; i++) { + prepare_log((Drive *) vec->envs[i]); + } + Py_RETURN_NONE; +} + static int my_put(Env *env, PyObject *args, PyObject *kwargs) { PyObject *obs = PyDict_GetItemString(kwargs, "observations"); if (!PyObject_TypeCheck(obs, &PyArray_Type)) { @@ -1893,6 +1908,7 @@ static int my_init(Env *env, PyObject *args, PyObject *kwargs) { env->reward_randomization = (bool) unpack(kwargs, "reward_randomization"); env->compute_eval_metrics = (bool) unpack(kwargs, "compute_eval_metrics"); env->eval_mode = (int) unpack(kwargs, "eval_mode"); + env->log_ema_alpha = (float) unpack(kwargs, "log_ema_alpha"); env->obs_norm_goal_offset_m = (float) unpack(kwargs, "obs_norm_goal_offset_m"); env->obs_norm_xy_offset_m = (float) unpack(kwargs, "obs_norm_xy_offset_m"); env->obs_norm_veh_length_m = (float) unpack(kwargs, "obs_norm_veh_length_m"); diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index 84fcc1f610..525a8530c9 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -308,6 +308,10 @@ struct Drive { // Rendering char video_suffix[64]; char resource_root[512]; + // One EMA slot per active controllable agent; sized in init() from active_agent_count. + Log *per_agent_log_ema; + int *per_agent_has_data; + float log_ema_alpha; }; typedef struct { @@ -1926,9 +1930,33 @@ static float calculate_puffer_score(Log *agent_log, float duration_steps, float return agent_log->puffer_score; } +// Sum each agent's current EMA slot into env->log and set env->log.n to the +// number of slots that have ever been seeded. vec_log divides by aggregate.n +// downstream, producing a cross-agent population mean with one weight per +// agent regardless of that agent's completion frequency. Does not reset the +// slots: emissions between completions repeat the previous values. +static void prepare_log(Drive *env) { + memset(&env->log, 0, sizeof(Log)); + int num_keys = sizeof(Log) / sizeof(float); + int num_with_data = 0; + for (int i = 0; i < env->active_agent_count; i++) { + if (!env->per_agent_has_data[i]) { + continue; + } + float *slot = (float *) &env->per_agent_log_ema[i]; + float *dst = (float *) &env->log; + for (int j = 0; j < num_keys; j++) { + dst[j] += slot[j]; + } + num_with_data++; + } + env->log.n = (float) num_with_data; +} + static void add_log(Drive *env) { int safe_timestep = (env->timestep > 0) ? env->timestep : 1; - Log episode_log = {0}; + float alpha = env->log_ema_alpha; + int num_log_keys = sizeof(Log) / sizeof(float); for (int i = 0; i < env->active_agent_count; i++) { Agent *agent = &env->agents[env->active_agent_indices[i]]; float episode_duration_s = env->logs[i].episode_length * env->dt; @@ -1936,51 +1964,53 @@ static void add_log(Drive *env) { reference_progress_distance = fmaxf(reference_progress_distance, 1.0f); env->logs[i].progress_ratio = agent->distance_since_spawn / reference_progress_distance; + Log episode_log = {0}; + int offroad = env->logs[i].offroad_rate; - episode_log.offroad_rate += offroad; + episode_log.offroad_rate = offroad; int collided = env->logs[i].collision_rate; - episode_log.collision_rate += collided; + episode_log.collision_rate = collided; int red_light_violations = env->logs[i].red_light_violation_rate; - episode_log.red_light_violation_rate += red_light_violations; + episode_log.red_light_violation_rate = red_light_violations; int total_infractions = (offroad || collided || red_light_violations) ? 1 : 0; float avg_speed_per_agent = env->logs[i].avg_speed_per_agent; - episode_log.avg_speed_per_agent += avg_speed_per_agent / safe_timestep; + episode_log.avg_speed_per_agent = avg_speed_per_agent / safe_timestep; int num_goals_reached = env->logs[i].num_goals_reached; - episode_log.num_goals_reached += num_goals_reached; + episode_log.num_goals_reached = num_goals_reached; // Score: 1 per agent that reached its full goal set without being removed/stopped. if (num_goals_reached >= env->num_goals && !agent->removed && !agent->stopped) { - episode_log.score += 1.0f; + episode_log.score = 1.0f; } if (!offroad && !collided && !red_light_violations && num_goals_reached < 1) { - episode_log.dnf_rate += 1.0f; + episode_log.dnf_rate = 1.0f; } - episode_log.total_distance_travelled += agent->distance_since_spawn; + episode_log.total_distance_travelled = agent->distance_since_spawn; if (total_infractions > 0) { - episode_log.total_infractions += 1.0f; + episode_log.total_infractions = 1.0f; } float displacement_error = env->logs[i].avg_displacement_error; - episode_log.avg_displacement_error += displacement_error; - episode_log.episode_length += env->logs[i].episode_length; - episode_log.episode_return += env->logs[i].episode_return; + episode_log.avg_displacement_error = displacement_error; + episode_log.episode_length = env->logs[i].episode_length; + episode_log.episode_return = env->logs[i].episode_return; // Per-component reward sums (mirrors compute_rewards' env->rewards[i]+= sites). - episode_log.reward_collision += env->logs[i].reward_collision; - episode_log.reward_offroad += env->logs[i].reward_offroad; - episode_log.reward_red_light += env->logs[i].reward_red_light; - episode_log.reward_goal += env->logs[i].reward_goal; - episode_log.reward_lane_align += env->logs[i].reward_lane_align; - episode_log.reward_lane_center += env->logs[i].reward_lane_center; - episode_log.reward_comfort += env->logs[i].reward_comfort; - episode_log.reward_velocity += env->logs[i].reward_velocity; - episode_log.reward_timestep += env->logs[i].reward_timestep; - episode_log.reward_reverse += env->logs[i].reward_reverse; - episode_log.reward_overspeed += env->logs[i].reward_overspeed; - episode_log.reward_ade += env->logs[i].reward_ade; + episode_log.reward_collision = env->logs[i].reward_collision; + episode_log.reward_offroad = env->logs[i].reward_offroad; + episode_log.reward_red_light = env->logs[i].reward_red_light; + episode_log.reward_goal = env->logs[i].reward_goal; + episode_log.reward_lane_align = env->logs[i].reward_lane_align; + episode_log.reward_lane_center = env->logs[i].reward_lane_center; + episode_log.reward_comfort = env->logs[i].reward_comfort; + episode_log.reward_velocity = env->logs[i].reward_velocity; + episode_log.reward_timestep = env->logs[i].reward_timestep; + episode_log.reward_reverse = env->logs[i].reward_reverse; + episode_log.reward_overspeed = env->logs[i].reward_overspeed; + episode_log.reward_ade = env->logs[i].reward_ade; // Comfort and velocity metrics (normalized per timestep) - episode_log.comfort_violation_count += env->logs[i].comfort_violation_count / safe_timestep; - episode_log.velocity_progress_sum += env->logs[i].velocity_progress_sum / safe_timestep; + episode_log.comfort_violation_count = env->logs[i].comfort_violation_count / safe_timestep; + episode_log.velocity_progress_sum = env->logs[i].velocity_progress_sum / safe_timestep; // Lane metrics (normalized per timestep for average per episode) - episode_log.lane_center_rate += env->logs[i].lane_center_rate / safe_timestep; - episode_log.lane_heading_aligned_rate += env->logs[i].lane_heading_aligned_rate / safe_timestep; + episode_log.lane_center_rate = env->logs[i].lane_center_rate / safe_timestep; + episode_log.lane_heading_aligned_rate = env->logs[i].lane_heading_aligned_rate / safe_timestep; if (env->compute_eval_metrics) { env->logs[i].progress_ratio = agent->distance_since_spawn / reference_progress_distance; env->logs[i].comfort_score = calculate_duration_scaled_violation_score( @@ -1988,40 +2018,45 @@ static void add_log(Drive *env) { env->logs[i].episode_length, env->dt); calculate_puffer_score(&env->logs[i], env->logs[i].episode_length, env->dt); - episode_log.at_fault_collision_rate += env->logs[i].at_fault_collision_rate; - episode_log.ttc_within_bound_rate += env->logs[i].ttc_within_bound_rate; - episode_log.wrong_way_distance += env->logs[i].wrong_way_distance; - episode_log.speed_violation_sum += env->logs[i].speed_violation_sum; - episode_log.progress_ratio += env->logs[i].progress_ratio; - episode_log.comfort_score += env->logs[i].comfort_score; - episode_log.ttc_violations += env->logs[i].ttc_violations; - episode_log.ttc_samples += env->logs[i].ttc_samples; - episode_log.multi_lane_time += env->logs[i].multi_lane_time; - episode_log.multi_lane_score += env->logs[i].multi_lane_score; + episode_log.at_fault_collision_rate = env->logs[i].at_fault_collision_rate; + episode_log.ttc_within_bound_rate = env->logs[i].ttc_within_bound_rate; + episode_log.wrong_way_distance = env->logs[i].wrong_way_distance; + episode_log.speed_violation_sum = env->logs[i].speed_violation_sum; + episode_log.progress_ratio = env->logs[i].progress_ratio; + episode_log.comfort_score = env->logs[i].comfort_score; + episode_log.ttc_violations = env->logs[i].ttc_violations; + episode_log.ttc_samples = env->logs[i].ttc_samples; + episode_log.multi_lane_time = env->logs[i].multi_lane_time; + episode_log.multi_lane_score = env->logs[i].multi_lane_score; float wrong_dist = env->logs[i].wrong_way_distance; float direction_score = (wrong_dist <= 2.0f) ? 1.0f : (wrong_dist <= 6.0f) ? 0.5f : 0.0f; - episode_log.driving_direction_score += direction_score; + episode_log.driving_direction_score = direction_score; float safe_duration_s = safe_timestep * env->dt; float speed_compliance = fmaxf(0.0f, 1.0f - env->logs[i].speed_violation_sum / fmaxf(safe_duration_s, 1e-3f)); - episode_log.speed_limit_compliance += speed_compliance; + episode_log.speed_limit_compliance = speed_compliance; float making_progress = (env->logs[i].progress_ratio > 0.2f) ? 1.0f : 0.0f; - episode_log.making_progress_rate += making_progress; - episode_log.puffer_score += env->logs[i].puffer_score; + episode_log.making_progress_rate = making_progress; + episode_log.puffer_score = env->logs[i].puffer_score; } - episode_log.n += 1; - } - // Log composition counts per agent so vec_log averaging recovers the per-env value - episode_log.expert_static_car_count += env->expert_static_agent_count; - episode_log.static_car_count += env->static_agent_count; + episode_log.expert_static_car_count = env->expert_static_agent_count; + episode_log.static_car_count = env->static_agent_count; - // Fold this episode into the cumulative log that vec_log averages and resets. - int num_log_fields = sizeof(Log) / sizeof(float); - for (int field_idx = 0; field_idx < num_log_fields; field_idx++) { - ((float *) &env->log)[field_idx] += ((float *) &episode_log)[field_idx]; + // First completion seeds the slot; later ones decay toward the newest episode by (1 - alpha). + Log *slot = &env->per_agent_log_ema[i]; + if (!env->per_agent_has_data[i]) { + *slot = episode_log; + env->per_agent_has_data[i] = 1; + } else { + float *slot_f = (float *) slot; + float *episode_f = (float *) &episode_log; + for (int j = 0; j < num_log_keys; j++) { + slot_f[j] = alpha * slot_f[j] + (1.0f - alpha) * episode_f[j]; + } + } } env->log_episode_seed = env->episode_seed; @@ -2853,6 +2888,8 @@ void init(Drive *env) { } set_active_agents(env); env->logs_capacity = env->active_agent_count; + env->per_agent_log_ema = (Log *) calloc(env->active_agent_count, sizeof(Log)); + env->per_agent_has_data = (int *) calloc(env->active_agent_count, sizeof(int)); if (env->simulation_mode == SIMULATION_MODE_REPLAY) { remove_bad_trajectories(env); } @@ -2914,6 +2951,8 @@ void c_close(Drive *env) { free(env->traffic_elements); free(env->active_agent_indices); free(env->logs); + free(env->per_agent_log_ema); + free(env->per_agent_has_data); if (env->shared_map != NULL) { // Geometry is borrowed from the cache. Release our reference; free the // entry only on the last reference, and only in the process that built it. diff --git a/pufferlib/ocean/drive/drive.py b/pufferlib/ocean/drive/drive.py index 1046edca66..673edd152a 100644 --- a/pufferlib/ocean/drive/drive.py +++ b/pufferlib/ocean/drive/drive.py @@ -97,6 +97,7 @@ def __init__( num_eval_scenarios=16, eval_map_indices=None, eval_scenario_seeds=None, + log_ema_alpha=0.707, init_mode="create_all_valid", control_mode="control_vehicles", sdc_controller="policy", @@ -222,6 +223,7 @@ def __init__( if self.eval_scenario_seeds is None or len(self.eval_scenario_seeds) != len(self.eval_map_indices): raise ValueError("eval_scenario_seeds must have one seed per eval_map_indices entry") self.use_exact_episode_seed = bool(eval_mode) and self.eval_scenario_seeds is not None + self.log_ema_alpha = log_ema_alpha self.termination_mode = termination_mode self.inactive_agent_threshold = inactive_agent_threshold self.terminate_on_goal = terminate_on_goal @@ -553,6 +555,7 @@ def _env_init_kwargs(self, map_file, max_agents): "compute_eval_metrics": self.compute_eval_metrics, "eval_mode": self.eval_mode, "use_exact_episode_seed": int(self.use_exact_episode_seed), + "log_ema_alpha": self.log_ema_alpha, "obs_norm_goal_offset_m": self.obs_norm_goal_offset_m, "obs_norm_xy_offset_m": self.obs_norm_xy_offset_m, "obs_norm_veh_length_m": self.obs_norm_veh_length_m, @@ -608,6 +611,8 @@ def step(self, actions): # vec_log is the training aggregate; it resets env->log, which eval reads # per episode, so it must not run in eval mode. if not self.eval_mode and self.tick % self.report_interval == 0: + # Rebuilds env->log from the per-agent EMA slots so vec_log means over agents, not completions. + binding.vec_prepare_log(self.c_envs) log = binding.vec_log(self.c_envs, self.num_agents) if log: info.append(log) diff --git a/pufferlib/ocean/env_binding.h b/pufferlib/ocean/env_binding.h index 907a46754c..9c626f0d6d 100644 --- a/pufferlib/ocean/env_binding.h +++ b/pufferlib/ocean/env_binding.h @@ -697,8 +697,6 @@ static PyObject *vec_log(PyObject *self, PyObject *args) { // Iterates over logs one float at a time. Will break // horribly if Log has non-float data. - PyObject *num_agents_arg = PyTuple_GetItem(args, 1); - float num_agents = (float) PyLong_AsLong(num_agents_arg); int num_keys = sizeof(Log) / sizeof(float); Log aggregate = {0}; @@ -711,8 +709,8 @@ static PyObject *vec_log(PyObject *self, PyObject *args) { PyObject *dict = PyDict_New(); - // Only log if we have at least num_agents worth of data - if (aggregate.n < num_agents) { + // aggregate.n counts agent slots with data, not completions; n=0 would divide by zero below. + if (aggregate.n < 1) { return dict; } diff --git a/tests/unit_tests/test_drive_per_agent_logging.py b/tests/unit_tests/test_drive_per_agent_logging.py new file mode 100644 index 0000000000..078cf6efd2 --- /dev/null +++ b/tests/unit_tests/test_drive_per_agent_logging.py @@ -0,0 +1,121 @@ +"""Behavioral tests for the per-agent EMA logging path in Drive's vec_log. + +Covered: +- T1 (gate): vec_log emits nothing while no agent has completed an episode; + the first emission appears exactly when the scenario truncates. +- T2 (steady-state n): after enough steps for every agent to complete at + least one episode, dict["n"] equals num_agents. +- T3 (state persistence): two consecutive emissions with no new completions + between them produce identical metric values, proving prepare_log no + longer resets the per-agent EMAs. +""" + +from pathlib import Path + +import numpy as np +import pytest + +from pufferlib.ocean.drive.drive import Drive + +MAP_DIR = Path(__file__).resolve().parents[2] / "pufferlib/resources/drive/binaries/carla" + +NUM_AGENTS = 32 +SCENARIO_LENGTH = 5 + + +def _log_dicts(info): + """Filter info for the vec_log aggregate dict (carries both n and episode_return).""" + return [d for d in info if isinstance(d, dict) and "n" in d and "episode_return" in d] + + +def _make_env(): + if not MAP_DIR.is_dir() or not any(MAP_DIR.glob("*.bin")): + pytest.skip(f"Drive map binaries not available at {MAP_DIR}") + return Drive( + num_agents=NUM_AGENTS, + num_maps=1, + min_agents_per_env=1, + max_agents_per_env=8, + scenario_length=SCENARIO_LENGTH, + report_interval=1, + map_dir=str(MAP_DIR), + log_ema_alpha=0.95, + ) + + +def test_gate_skips_emissions_until_first_completion(): + """T1: vec_log returns no log dict until at least one agent has completed an + episode. With synced agents under null actions, completions first happen + when timestep reaches scenario_length.""" + env = _make_env() + env.reset(seed=0) + + # Steps 1..(L-1): no agent has truncated yet, gate sees aggregate.n == 0. + for step_idx in range(SCENARIO_LENGTH - 1): + _, _, _, _, info = env.step(np.zeros_like(env.actions)) + assert not _log_dicts(info), ( + f"Unexpected emission at step {step_idx + 1} (timestep={step_idx + 1}); " + "gate should skip before any agent has completed." + ) + + # Step L: all agents truncate at timestep == scenario_length. + _, _, _, _, info = env.step(np.zeros_like(env.actions)) + logs = _log_dicts(info) + assert len(logs) == 1, f"Expected exactly one log dict at scenario end, got {len(logs)}" + assert logs[0]["n"] >= 1, f"Expected n>=1 at first emission, got n={logs[0]['n']}" + + env.close() + + +def test_steady_state_n_equals_num_agents(): + """T2: once every agent has completed at least one episode, the emitted + n is the full population. With synced agents this is true from the first + emission onward.""" + env = _make_env() + env.reset(seed=0) + + # Run multiple scenarios so every agent has contributed. + last_log = None + for _ in range(4 * SCENARIO_LENGTH): + _, _, _, _, info = env.step(np.zeros_like(env.actions)) + logs = _log_dicts(info) + if logs: + last_log = logs[-1] + + assert last_log is not None, "Expected at least one emission across 4 scenarios" + assert last_log["n"] == NUM_AGENTS, ( + f"Expected steady-state n={NUM_AGENTS}, got n={last_log['n']} (some agents missing from the population mean)" + ) + + env.close() + + +def test_emissions_identical_when_no_new_completions(): + """T3: prepare_log preserves per-agent EMA state across emissions. Two + consecutive emissions with no intervening completions must produce + bit-for-bit identical metric values.""" + env = _make_env() + env.reset(seed=0) + + # Reach the first completion at timestep == scenario_length. + for _ in range(SCENARIO_LENGTH): + env.step(np.zeros_like(env.actions)) + + # Next two emissions sit between completions (no agent truncates between + # timestep=1 and timestep=scenario_length-1 of the next scenario). + _, _, _, _, info_a = env.step(np.zeros_like(env.actions)) + _, _, _, _, info_b = env.step(np.zeros_like(env.actions)) + + log_a = _log_dicts(info_a) + log_b = _log_dicts(info_b) + assert log_a and log_b, "Both consecutive steps should emit" + log_a, log_b = log_a[-1], log_b[-1] + + assert log_a["n"] == log_b["n"], f"n changed without new completions: {log_a['n']} vs {log_b['n']}" + for key in ("episode_return", "episode_length", "collision_rate", "offroad_rate"): + assert log_a[key] == pytest.approx(log_b[key], rel=0, abs=1e-6), ( + f"{key} changed without new completions: {log_a[key]} vs {log_b[key]} " + "(prepare_log may be resetting per-agent state)" + ) + + env.close() From 4fa16c39b6db2a6f0e2ee563871c06e10a588a25 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Mon, 10 Aug 2026 02:36:21 +0000 Subject: [PATCH 02/12] Regenerate smoke goldens --- .../data/drive_rollout_golden.json | 46 ++++++++++--------- .../smoke_tests/data/drive_smoke_golden.json | 44 +++++++++--------- 2 files changed, 47 insertions(+), 43 deletions(-) diff --git a/tests/smoke_tests/data/drive_rollout_golden.json b/tests/smoke_tests/data/drive_rollout_golden.json index 8e96952f46..8f26c25eaa 100644 --- a/tests/smoke_tests/data/drive_rollout_golden.json +++ b/tests/smoke_tests/data/drive_rollout_golden.json @@ -1,31 +1,33 @@ { "env": { - "avg_distance_per_infraction": 12.06708054339632, - "avg_speed_per_agent": 1.360816847770772, - "collision_rate": 0.023936170212765957, - "comfort_violation_count": 0.6436155717423622, - "dnf_rate": 0.535904255319149, - "episode_length": 13.48936170212766, - "episode_return": -0.7790729422518547, - "lane_center_rate": 0.6959362321711601, + "avg_distance_per_infraction": 11.880148284621052, + "avg_speed_per_agent": 1.356827261676262, + "collision_rate": 0.02294035734988253, + "comfort_violation_count": 0.6380039383838703, + "dnf_rate": 0.5361869363718993, + "episode_length": 13.349792389126566, + "episode_return": -0.7652764916903787, + "lane_center_rate": 0.7025819145046271, "n": 16.0, - "num_goals_reached": 0.00997340425531915, - "offroad_rate": 0.4115691489361702, - "red_light_violation_rate": 0.025930851063829786, + "num_goals_reached": 0.009789912615118525, + "offroad_rate": 0.4116112170620011, + "red_light_violation_rate": 0.02665171274707831, "reward_components/ade": 0.0, - "reward_components/collision": -0.042797685620632575, - "reward_components/comfort": -0.09395423511716913, - "reward_components/goal": 0.003324468085106383, - "reward_components/lane_align": -0.01126101781227725, - "reward_components/lane_center": -0.0026710387282172575, - "reward_components/offroad": -0.6115261237037942, + "reward_components/collision": -0.04310367152574109, + "reward_components/comfort": -0.09134016255243355, + "reward_components/goal": 0.003110341709638211, + "reward_components/lane_align": -0.01109200311815608, + "reward_components/lane_center": -0.0027703031133663805, + "reward_components/offroad": -0.6000605714688827, "reward_components/overspeed": 0.0, - "reward_components/red_light": -0.014523565179688181, - "reward_components/reverse": -0.005763038997835619, - "reward_components/timestep": -8.58909574007873e-05, - "reward_components/velocity": 0.00018516867541930356, + "reward_components/red_light": -0.014580663242999104, + "reward_components/reverse": -0.005549692918089908, + "reward_components/timestep": -8.510924470857543e-05, + "reward_components/velocity": 0.00019534487165850356, "score": 0.0, - "velocity_progress_sum": 0.018232046272308428 + "total_distance_travelled_sum": 87.34569526647593, + "total_infraction_count": 7.3792525982701935, + "velocity_progress_sum": 0.019855013023079193 }, "meta": { "bptt_horizon": 64, diff --git a/tests/smoke_tests/data/drive_smoke_golden.json b/tests/smoke_tests/data/drive_smoke_golden.json index 9077394e95..1df258df1a 100644 --- a/tests/smoke_tests/data/drive_smoke_golden.json +++ b/tests/smoke_tests/data/drive_smoke_golden.json @@ -1,34 +1,36 @@ { "env": { - "avg_distance_per_infraction": 12.355366627375284, - "avg_speed_per_agent": 1.3996734884050157, - "collision_rate": 0.013888888888888888, - "comfort_violation_count": 0.6618798706266615, - "dnf_rate": 0.5416666666666666, - "episode_length": 13.055555555555555, - "episode_return": -0.7164439492755466, - "lane_center_rate": 0.6980418860912323, + "avg_distance_per_infraction": 11.652713274145599, + "avg_speed_per_agent": 1.3747441177738124, + "collision_rate": 0.011199567458945615, + "comfort_violation_count": 0.6526833839457611, + "dnf_rate": 0.5479880072947206, + "episode_length": 12.521783195692917, + "episode_return": -0.6689560803873785, + "lane_center_rate": 0.725223951298615, "n": 16.0, - "num_goals_reached": 0.020833333333333332, + "num_goals_reached": 0.010593278654188656, "obs/max": 49.0, "obs/mean": 0.24742739112116396, "obs/min": -1.720494793727994, - "offroad_rate": 0.40625, - "red_light_violation_rate": 0.03125, + "offroad_rate": 0.40545966244977094, + "red_light_violation_rate": 0.02961260791675284, "reward_components/ade": 0.0, - "reward_components/collision": -0.02127074119117525, - "reward_components/comfort": -0.09980892741845714, + "reward_components/collision": -0.01718186237849295, + "reward_components/comfort": -0.08829760264027221, "reward_components/goal": 0.0, - "reward_components/lane_align": -0.011237355065532029, - "reward_components/lane_center": -0.002885955537850451, - "reward_components/offroad": -0.5649742384751638, + "reward_components/lane_align": -0.011195697115156156, + "reward_components/lane_center": -0.002516856037204732, + "reward_components/offroad": -0.5360425002872944, "reward_components/overspeed": 0.0, - "reward_components/red_light": -0.011040630271761782, - "reward_components/reverse": -0.00537483084998611, - "reward_components/timestep": -8.41927073148933e-05, - "reward_components/velocity": 0.0002329377817255186, + "reward_components/red_light": -0.009031652451838078, + "reward_components/reverse": -0.0048044504957064854, + "reward_components/timestep": -8.012855822267613e-05, + "reward_components/velocity": 0.00019470069981190167, "score": 0.0, - "velocity_progress_sum": 0.022617154185556702 + "total_distance_travelled": 19303.431358337402, + "total_infractions": 1656.561086177826, + "velocity_progress_sum": 0.019867269547078115 }, "losses": { "approx_kl": 0.001297435179973642, From 71aa850d59210b1e62fa27980aaaeba09d3f6824 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Sun, 9 Aug 2026 22:42:45 -0400 Subject: [PATCH 03/12] Refresh env->log from the EMA slots at episode end add_log wrote only the per-agent EMA slots, leaving env->log populated solely by prepare_log via the Python vec_prepare_log binding. Consumers that drive the env directly in C never reach that path, so tests/drive/test_drive_env_smoke.c saw env.log.n == 0 and failed both its saw_log and truncation assertions. Calling prepare_log at the end of add_log keeps env->log valid for any C-side reader. drive.py's per-report vec_prepare_log call stays: vec_log zeroes env->log after each emission, so envs with no completion since the last emission must still be rebuilt from their slots to contribute to the cross-agent mean. --- pufferlib/ocean/drive/drive.h | 2 ++ 1 file changed, 2 insertions(+) diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index 525a8530c9..477b6f27f9 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -2059,6 +2059,8 @@ static void add_log(Drive *env) { } } + // Keeps env->log valid for C-side readers, which never reach the Python vec_prepare_log. + prepare_log(env); env->log_episode_seed = env->episode_seed; } From b107f4aa6e8708fd1523a8316c908f8222c3f319 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Sun, 9 Aug 2026 23:26:09 -0400 Subject: [PATCH 04/12] Keep distance and infraction totals as interval accumulators [update-golden] The per-agent EMA turned every Log field into a gauge, but total_distance_travelled and total_infractions feed total_distance_travelled_sum / total_infraction_count, which reduce_environment_metrics and the eval report sum across emissions to build avg_distance_per_infraction from global totals. Summing a gauge double-counted them once per emission rather than once per completed episode. Both fields now accumulate raw per-episode values directly into env->log and are excluded from the EMA slots; prepare_log carries them across its memset so only vec_log's post-emission zeroing clears them. That restores reset-on-consumption, so emissions with no completed episode contribute zero instead of re-reporting the running total. Measured over 60 steps at scenario_length=10: 51 emissions, of which 6 carry data, matching the 6 episode ends. The summed distance drops from 3590.4 to 422.4. --- pufferlib/ocean/drive/drive.h | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index 477b6f27f9..5bbebb712c 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -1936,6 +1936,10 @@ static float calculate_puffer_score(Log *agent_log, float duration_steps, float // agent regardless of that agent's completion frequency. Does not reset the // slots: emissions between completions repeat the previous values. static void prepare_log(Drive *env) { + // Raw interval totals, not per-agent gauges: they accumulate until vec_log emits and zeroes them, + // so summing them across emissions yields the global totals avg_distance_per_infraction needs. + float interval_distance_travelled = env->log.total_distance_travelled; + float interval_infraction_count = env->log.total_infractions; memset(&env->log, 0, sizeof(Log)); int num_keys = sizeof(Log) / sizeof(float); int num_with_data = 0; @@ -1951,6 +1955,8 @@ static void prepare_log(Drive *env) { num_with_data++; } env->log.n = (float) num_with_data; + env->log.total_distance_travelled = interval_distance_travelled; + env->log.total_infractions = interval_infraction_count; } static void add_log(Drive *env) { @@ -1984,9 +1990,9 @@ static void add_log(Drive *env) { if (!offroad && !collided && !red_light_violations && num_goals_reached < 1) { episode_log.dnf_rate = 1.0f; } - episode_log.total_distance_travelled = agent->distance_since_spawn; + env->log.total_distance_travelled += agent->distance_since_spawn; if (total_infractions > 0) { - episode_log.total_infractions = 1.0f; + env->log.total_infractions += 1.0f; } float displacement_error = env->logs[i].avg_displacement_error; episode_log.avg_displacement_error = displacement_error; From 29be94427167e33bbac9b57475b410c6f2ee1bea Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Mon, 10 Aug 2026 03:39:55 +0000 Subject: [PATCH 05/12] Regenerate smoke goldens --- tests/smoke_tests/data/drive_rollout_golden.json | 6 +++--- tests/smoke_tests/data/drive_smoke_golden.json | 6 +++--- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/tests/smoke_tests/data/drive_rollout_golden.json b/tests/smoke_tests/data/drive_rollout_golden.json index 8f26c25eaa..9d24b747d8 100644 --- a/tests/smoke_tests/data/drive_rollout_golden.json +++ b/tests/smoke_tests/data/drive_rollout_golden.json @@ -1,6 +1,6 @@ { "env": { - "avg_distance_per_infraction": 11.880148284621052, + "avg_distance_per_infraction": 0.9207025739279661, "avg_speed_per_agent": 1.356827261676262, "collision_rate": 0.02294035734988253, "comfort_violation_count": 0.6380039383838703, @@ -25,8 +25,8 @@ "reward_components/timestep": -8.510924470857543e-05, "reward_components/velocity": 0.00019534487165850356, "score": 0.0, - "total_distance_travelled_sum": 87.34569526647593, - "total_infraction_count": 7.3792525982701935, + "total_distance_travelled_sum": 6.756095629233818, + "total_infraction_count": 0.5633116883116883, "velocity_progress_sum": 0.019855013023079193 }, "meta": { diff --git a/tests/smoke_tests/data/drive_smoke_golden.json b/tests/smoke_tests/data/drive_smoke_golden.json index 1df258df1a..f0e36ccbd6 100644 --- a/tests/smoke_tests/data/drive_smoke_golden.json +++ b/tests/smoke_tests/data/drive_smoke_golden.json @@ -1,6 +1,6 @@ { "env": { - "avg_distance_per_infraction": 11.652713274145599, + "avg_distance_per_infraction": 12.26463147676908, "avg_speed_per_agent": 1.3747441177738124, "collision_rate": 0.011199567458945615, "comfort_violation_count": 0.6526833839457611, @@ -28,8 +28,8 @@ "reward_components/timestep": -8.012855822267613e-05, "reward_components/velocity": 0.00019470069981190167, "score": 0.0, - "total_distance_travelled": 19303.431358337402, - "total_infractions": 1656.561086177826, + "total_distance_travelled": 1594.4020919799805, + "total_infractions": 130.0, "velocity_progress_sum": 0.019867269547078115 }, "losses": { From 3fb448b6f02c0332bb6a8bf6dcf8b528b2fca0c1 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Fri, 14 Aug 2026 07:33:32 -0400 Subject: [PATCH 06/12] Add a discriminating logging test and validate log_ema_alpha [update-golden] T1-T3 run every agent in lockstep, which is the one regime where completion-weighted and agent-weighted aggregation agree: with the base implementation spliced back in, two of the three still passed. T4 gives sub-envs independent truncation (termination_mode=1, infractions remove) so completions arrive staggered, then asserts the emitted population climbs to num_agents and stays there. It fails on 3.0 and passes here. log_ema_alpha reaches the EMA update straight from config with no range check. At 1.0 the (1 - alpha) term vanishes and every metric stays pinned to its first completed episode for the whole run; above 1.0 the slots diverge. Both are silent, so reject them at construction. Co-Authored-By: Claude Opus 5 (1M context) --- pufferlib/ocean/drive/drive.py | 4 ++ .../test_drive_per_agent_logging.py | 68 +++++++++++++++++++ 2 files changed, 72 insertions(+) diff --git a/pufferlib/ocean/drive/drive.py b/pufferlib/ocean/drive/drive.py index 23979bc5c4..529f38bbe7 100644 --- a/pufferlib/ocean/drive/drive.py +++ b/pufferlib/ocean/drive/drive.py @@ -233,6 +233,10 @@ def __init__( if self.eval_scenario_seeds is None or len(self.eval_scenario_seeds) != len(self.eval_map_indices): raise ValueError("eval_scenario_seeds must have one seed per eval_map_indices entry") self.use_exact_episode_seed = bool(eval_mode) and self.eval_scenario_seeds is not None + # Outside [0, 1) the slot update stops being a convex combination: 1.0 pins every + # metric to its first episode, above 1.0 the slots diverge. + if not 0.0 <= log_ema_alpha < 1.0: + raise ValueError(f"log_ema_alpha must be in [0.0, 1.0). Got: {log_ema_alpha}") self.log_ema_alpha = log_ema_alpha self.termination_mode = termination_mode self.inactive_agent_threshold = inactive_agent_threshold diff --git a/tests/unit_tests/test_drive_per_agent_logging.py b/tests/unit_tests/test_drive_per_agent_logging.py index 078cf6efd2..74e8a8e3db 100644 --- a/tests/unit_tests/test_drive_per_agent_logging.py +++ b/tests/unit_tests/test_drive_per_agent_logging.py @@ -8,6 +8,9 @@ - T3 (state persistence): two consecutive emissions with no new completions between them produce identical metric values, proving prepare_log no longer resets the per-agent EMAs. +- T4 (unequal completion rates): with sub-envs truncating at different + timesteps, every emission still weights one term per agent. This is the + case T1-T3 cannot distinguish, since they run all agents in lockstep. """ from pathlib import Path @@ -43,6 +46,29 @@ def _make_env(): ) +def _make_staggered_env(): + """Sub-envs truncate on their own inactive-agent ratio, so they complete at + different timesteps instead of in lockstep.""" + if not MAP_DIR.is_dir() or not any(MAP_DIR.glob("*.bin")): + pytest.skip(f"Drive map binaries not available at {MAP_DIR}") + return Drive( + num_agents=NUM_AGENTS, + num_maps=4, + min_agents_per_env=1, + max_agents_per_env=4, + scenario_length=40, + report_interval=1, + map_dir=str(MAP_DIR), + log_ema_alpha=0.95, + termination_mode=1, + offroad_behavior="remove", + collision_behavior="remove", + # Resampling rebuilds the envs and so discards the slots; keep it out of + # the measured window so only completion weighting is under test. + resample_frequency=100_000, + ) + + def test_gate_skips_emissions_until_first_completion(): """T1: vec_log returns no log dict until at least one agent has completed an episode. With synced agents under null actions, completions first happen @@ -119,3 +145,45 @@ def test_emissions_identical_when_no_new_completions(): ) env.close() + + +def test_population_size_is_agent_count_under_unequal_completion_rates(): + """T4: the discriminating case. Sub-envs truncate independently here, so + completions arrive staggered rather than all at once. Completion-weighted + aggregation would make the emitted n track how many agents completed in + that interval; agent-weighted aggregation pins it to the number of agents + that have ever completed, which saturates at num_agents and stays there.""" + env = _make_staggered_env() + env.reset(seed=0) + + population_sizes = [] + for _ in range(240): + _, _, _, _, info = env.step(np.zeros_like(env.actions)) + logs = _log_dicts(info) + if logs: + population_sizes.append(logs[-1]["n"]) + + assert population_sizes, "Expected emissions once sub-envs began completing" + + # Guard against the test silently degenerating into the lockstep regime T1-T3 + # already cover: staggered completions make n climb through intermediate values. + warmup = [size for size in population_sizes if size < NUM_AGENTS] + assert len(set(warmup)) > 1, ( + f"Expected staggered completions to grow the population gradually, saw {sorted(set(warmup))}; " + "sub-envs are completing in lockstep so this test is not exercising unequal completion rates" + ) + + assert max(population_sizes) == NUM_AGENTS, ( + f"Population never reached the full agent count: max n={max(population_sizes)} of {NUM_AGENTS}" + ) + + # Monotonic and saturating: once an agent has completed it contributes to + # every later emission, whether or not it completed again. + saturation_idx = population_sizes.index(NUM_AGENTS) + after_saturation = population_sizes[saturation_idx:] + assert set(after_saturation) == {NUM_AGENTS}, ( + f"n fell back below {NUM_AGENTS} after saturating: {sorted(set(after_saturation))}; " + "the emitted mean is weighted by completions in the interval, not by agent" + ) + + env.close() From cc0b7ac7337d32de794c6afe06738f2c7ad57c9e Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 14 Aug 2026 11:47:59 +0000 Subject: [PATCH 07/12] Regenerate smoke goldens --- .../data/drive_rollout_golden.json | 48 +++++++++---------- .../smoke_tests/data/drive_smoke_golden.json | 42 ++++++++-------- 2 files changed, 45 insertions(+), 45 deletions(-) diff --git a/tests/smoke_tests/data/drive_rollout_golden.json b/tests/smoke_tests/data/drive_rollout_golden.json index 8c55869b2c..9ba4f12e3b 100644 --- a/tests/smoke_tests/data/drive_rollout_golden.json +++ b/tests/smoke_tests/data/drive_rollout_golden.json @@ -1,33 +1,33 @@ { "env": { - "avg_distance_per_infraction": 12.31972300351321, - "avg_speed_per_agent": 1.3578476041227907, - "collision_rate": 0.025412087912087912, - "comfort_violation_count": 0.6445396886422083, - "dnf_rate": 0.5322802197802198, - "episode_length": 13.747252747252746, - "episode_return": -1.1555632654127184, - "lane_center_rate": 0.6866458556154272, + "avg_distance_per_infraction": 0.912943642768487, + "avg_speed_per_agent": 1.359724880915123, + "collision_rate": 0.026141917114206725, + "comfort_violation_count": 0.6377636792217093, + "dnf_rate": 0.5323061052867566, + "episode_length": 13.7684256323774, + "episode_return": -1.171682134853124, + "lane_center_rate": 0.6847131364896942, "n": 16.0, - "num_goals_reached": 0.015796703296703296, - "offroad_rate": 0.4024725274725275, - "red_light_violation_rate": 0.029532967032967032, + "num_goals_reached": 0.015526922819062806, + "offroad_rate": 0.40039702070458316, + "red_light_violation_rate": 0.03061665838026117, "reward_components/ade": 0.0, - "reward_components/collision": -0.03617091595635309, - "reward_components/comfort": -0.4257013105101638, - "reward_components/goal": 0.0027472527472527475, - "reward_components/lane_align": -0.026205073141462202, - "reward_components/lane_center": -0.004745660013273604, - "reward_components/offroad": -0.642226483140673, + "reward_components/collision": -0.03978979088827143, + "reward_components/comfort": -0.4286592380911209, + "reward_components/goal": 0.0022560067138109804, + "reward_components/lane_align": -0.026883024588858266, + "reward_components/lane_center": -0.004865653432923876, + "reward_components/offroad": -0.6491618211281027, "reward_components/overspeed": 0.0, - "reward_components/red_light": -0.013046603813603685, - "reward_components/reverse": -0.010339878637671144, - "reward_components/timestep": -8.684237648164956e-05, - "reward_components/velocity": 0.00021225640647481759, + "reward_components/red_light": -0.01431413074702795, + "reward_components/reverse": -0.010398259812265222, + "reward_components/timestep": -8.704732119221224e-05, + "reward_components/velocity": 0.00022083869197386843, "score": 0.0, - "total_distance_travelled_sum": 89.77611935793699, - "total_infraction_count": 7.318681318681318, - "velocity_progress_sum": 0.02037534874782048 + "total_distance_travelled_sum": 6.65279060388621, + "total_infraction_count": 0.5423452768729642, + "velocity_progress_sum": 0.021379507935186167 }, "meta": { "bptt_horizon": 64, diff --git a/tests/smoke_tests/data/drive_smoke_golden.json b/tests/smoke_tests/data/drive_smoke_golden.json index 5f6bd7b663..f5b85520b5 100644 --- a/tests/smoke_tests/data/drive_smoke_golden.json +++ b/tests/smoke_tests/data/drive_smoke_golden.json @@ -1,36 +1,36 @@ { "env": { "avg_distance_per_infraction": 12.032092877288363, - "avg_speed_per_agent": 1.3549120161268446, - "collision_rate": 0.006944444444444444, - "comfort_violation_count": 0.6202899664640427, - "dnf_rate": 0.5243055555555556, - "episode_length": 13.722222222222221, - "episode_return": -1.181912683778339, - "lane_center_rate": 0.716822452015347, + "avg_speed_per_agent": 1.3546471775605762, + "collision_rate": 0.0009471982203680893, + "comfort_violation_count": 0.6261929119455403, + "dnf_rate": 0.5358763861245123, + "episode_length": 13.164048149667938, + "episode_return": -1.215755915847318, + "lane_center_rate": 0.7255455193848446, "n": 16.0, - "num_goals_reached": 0.017361111111111112, + "num_goals_reached": 0.01153640180087552, "obs/max": 49.0, "obs/mean": 0.24155639042146504, "obs/min": -1.0741224959492683, - "offroad_rate": 0.4236111111111111, - "red_light_violation_rate": 0.034722222222222224, + "offroad_rate": 0.4215648413218301, + "red_light_violation_rate": 0.035249389097865284, "reward_components/ade": 0.0, - "reward_components/collision": -0.009900637798839144, - "reward_components/comfort": -0.4281758947504891, - "reward_components/goal": 0.003472222222222222, - "reward_components/lane_align": -0.027358024432841275, - "reward_components/lane_center": -0.004938648845483031, - "reward_components/offroad": -0.6899763031138314, + "reward_components/collision": -0.00135041285177757, + "reward_components/comfort": -0.42865003529807616, + "reward_components/goal": 0.002803422402237253, + "reward_components/lane_align": -0.026865696004623997, + "reward_components/lane_center": -0.004688873391293375, + "reward_components/offroad": -0.7286143891256431, "reward_components/overspeed": 0.0, - "reward_components/red_light": -0.014334955770108435, - "reward_components/reverse": -0.010830783155850239, - "reward_components/timestep": -8.736979240590396e-05, - "reward_components/velocity": 0.00021771855729942521, + "reward_components/red_light": -0.017992268673722344, + "reward_components/reverse": -0.010515955479510513, + "reward_components/timestep": -8.355726609367812e-05, + "reward_components/velocity": 0.00020186775057889713, "score": 0.0, "total_distance_travelled": 1612.3004455566406, "total_infractions": 134.0, - "velocity_progress_sum": 0.020827788549164932 + "velocity_progress_sum": 0.02026809492662292 }, "losses": { "approx_kl": 0.004087306525824326, From b7116c46ac6497e32636233603c051bc1f015fef Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Fri, 14 Aug 2026 07:42:39 -0400 Subject: [PATCH 08/12] Drop the redundant steady-state population test MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Its assertion — n == num_agents at the last emission with agents in lockstep — is a strict subset of what the staggered test already asserts, and it passes against the completion-sum implementation, so it guards nothing the other two don't. The gate test stays: it is the only one that pins the no-emission-before-any-completion behavior. Co-Authored-By: Claude Opus 5 (1M context) --- .../test_drive_per_agent_logging.py | 45 +++++-------------- 1 file changed, 12 insertions(+), 33 deletions(-) diff --git a/tests/unit_tests/test_drive_per_agent_logging.py b/tests/unit_tests/test_drive_per_agent_logging.py index 74e8a8e3db..bcf7d2a905 100644 --- a/tests/unit_tests/test_drive_per_agent_logging.py +++ b/tests/unit_tests/test_drive_per_agent_logging.py @@ -3,14 +3,16 @@ Covered: - T1 (gate): vec_log emits nothing while no agent has completed an episode; the first emission appears exactly when the scenario truncates. -- T2 (steady-state n): after enough steps for every agent to complete at - least one episode, dict["n"] equals num_agents. -- T3 (state persistence): two consecutive emissions with no new completions +- T2 (state persistence): two consecutive emissions with no new completions between them produce identical metric values, proving prepare_log no longer resets the per-agent EMAs. -- T4 (unequal completion rates): with sub-envs truncating at different - timesteps, every emission still weights one term per agent. This is the - case T1-T3 cannot distinguish, since they run all agents in lockstep. +- T3 (unequal completion rates): with sub-envs truncating at different + timesteps, every emission still weights one term per agent. + +T1 runs all agents in lockstep, the regime where completion-weighted and +agent-weighted aggregation agree, so it holds either way; it is here to pin +the emission gate. T2 and T3 both fail against the completion-sum +implementation. """ from pathlib import Path @@ -93,31 +95,8 @@ def test_gate_skips_emissions_until_first_completion(): env.close() -def test_steady_state_n_equals_num_agents(): - """T2: once every agent has completed at least one episode, the emitted - n is the full population. With synced agents this is true from the first - emission onward.""" - env = _make_env() - env.reset(seed=0) - - # Run multiple scenarios so every agent has contributed. - last_log = None - for _ in range(4 * SCENARIO_LENGTH): - _, _, _, _, info = env.step(np.zeros_like(env.actions)) - logs = _log_dicts(info) - if logs: - last_log = logs[-1] - - assert last_log is not None, "Expected at least one emission across 4 scenarios" - assert last_log["n"] == NUM_AGENTS, ( - f"Expected steady-state n={NUM_AGENTS}, got n={last_log['n']} (some agents missing from the population mean)" - ) - - env.close() - - def test_emissions_identical_when_no_new_completions(): - """T3: prepare_log preserves per-agent EMA state across emissions. Two + """T2: prepare_log preserves per-agent EMA state across emissions. Two consecutive emissions with no intervening completions must produce bit-for-bit identical metric values.""" env = _make_env() @@ -148,7 +127,7 @@ def test_emissions_identical_when_no_new_completions(): def test_population_size_is_agent_count_under_unequal_completion_rates(): - """T4: the discriminating case. Sub-envs truncate independently here, so + """T3: the discriminating case. Sub-envs truncate independently here, so completions arrive staggered rather than all at once. Completion-weighted aggregation would make the emitted n track how many agents completed in that interval; agent-weighted aggregation pins it to the number of agents @@ -165,8 +144,8 @@ def test_population_size_is_agent_count_under_unequal_completion_rates(): assert population_sizes, "Expected emissions once sub-envs began completing" - # Guard against the test silently degenerating into the lockstep regime T1-T3 - # already cover: staggered completions make n climb through intermediate values. + # Guard against the test silently degenerating into the lockstep regime T1 + # already covers: staggered completions make n climb through intermediate values. warmup = [size for size in population_sizes if size < NUM_AGENTS] assert len(set(warmup)) > 1, ( f"Expected staggered completions to grow the population gradually, saw {sorted(set(warmup))}; " From 5503f2bb8681c20948cd5aa445707f886c3c25da Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Fri, 14 Aug 2026 16:08:16 -0400 Subject: [PATCH 09/12] Report distance-per-infraction under both weightings [update-golden] MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit avg_distance_per_infraction pools every episode completed in the interval, so it answers "how far does the fleet travel per infraction" — dominated by whichever agents drive the most. That is the right question for fleet safety, but it is the completion-weighted scheme the rest of this PR moves away from, so on its own it no longer matches the other curves. Add the same ratio under per-agent weighting. agent_weighted_distance_ travelled and agent_weighted_infractions are ordinary EMA slots, so their ratio counts one episode per agent and answers "how far does a typical agent go between infractions". Both are now reported. The two are the same quantity when agents complete equally often and diverge only when completion rates differ: measured over 400 steps at alpha=0, lockstep gives 30.601 vs 30.621, while sub-envs truncating independently give 1339.776 vs 76.729. T3 pins the lockstep identity. The pooled pair resets on consumption and so must be summed across emissions; the agent-weighted pair persists and must be averaged. The emitted names avoid _sum/_count for the latter to keep that straight. Co-Authored-By: Claude Opus 5 (1M context) --- pufferlib/ocean/drive/binding.c | 7 +++ pufferlib/ocean/drive/drive.h | 8 +++ .../test_drive_per_agent_logging.py | 53 +++++++++++++++++-- 3 files changed, 65 insertions(+), 3 deletions(-) diff --git a/pufferlib/ocean/drive/binding.c b/pufferlib/ocean/drive/binding.c index ec942521ed..a6d41411c2 100644 --- a/pufferlib/ocean/drive/binding.c +++ b/pufferlib/ocean/drive/binding.c @@ -2023,6 +2023,10 @@ static int my_log(PyObject *dict, Env *env, Log *log, float n) { float total_distance_travelled = log->total_distance_travelled * n; float total_infractions = log->total_infractions * n; float avg_distance_per_infraction = total_distance_travelled / fmaxf(1.0f, total_infractions); + float agent_weighted_distance_travelled = log->agent_weighted_distance_travelled * n; + float agent_weighted_infractions = log->agent_weighted_infractions * n; + float agent_weighted_distance_per_infraction + = agent_weighted_distance_travelled / fmaxf(1.0f, agent_weighted_infractions); assign_to_dict(dict, "n", log->n); assign_to_dict(dict, "offroad_rate", log->offroad_rate); @@ -2041,6 +2045,9 @@ static int my_log(PyObject *dict, Env *env, Log *log, float n) { assign_to_dict(dict, "avg_distance_per_infraction", avg_distance_per_infraction); assign_to_dict(dict, "total_distance_travelled_sum", total_distance_travelled); assign_to_dict(dict, "total_infraction_count", total_infractions); + assign_to_dict(dict, "agent_weighted_distance_per_infraction", agent_weighted_distance_per_infraction); + assign_to_dict(dict, "agent_weighted_distance_travelled", agent_weighted_distance_travelled); + assign_to_dict(dict, "agent_weighted_infractions", agent_weighted_infractions); assign_to_dict(dict, "reward_components/collision", log->reward_collision); assign_to_dict(dict, "reward_components/offroad", log->reward_offroad); assign_to_dict(dict, "reward_components/red_light", log->reward_red_light); diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index 64ff0d6371..154b89f457 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -97,8 +97,14 @@ struct Log { float ttc_samples; float multi_lane_time; float multi_lane_score; + // Pooled over every episode completed in the interval; prepare_log keeps them + // out of the EMA so vec_log's zeroing makes them reset-on-consumption. float total_distance_travelled; float total_infractions; + // Same quantities as ordinary EMA slots, so their ratio weights one episode + // per agent instead of one per completion. + float agent_weighted_distance_travelled; + float agent_weighted_infractions; // Agent-only puffer display fields (for serialization) float no_at_fault; float no_offroad; @@ -2005,6 +2011,8 @@ static void add_log(Drive *env) { if (total_infractions > 0) { env->log.total_infractions += 1.0f; } + episode_log.agent_weighted_distance_travelled = agent->distance_since_spawn; + episode_log.agent_weighted_infractions = (total_infractions > 0) ? 1.0f : 0.0f; float displacement_error = env->logs[i].avg_displacement_error; episode_log.avg_displacement_error = displacement_error; episode_log.episode_length = env->logs[i].episode_length; diff --git a/tests/unit_tests/test_drive_per_agent_logging.py b/tests/unit_tests/test_drive_per_agent_logging.py index bcf7d2a905..33fb09ef87 100644 --- a/tests/unit_tests/test_drive_per_agent_logging.py +++ b/tests/unit_tests/test_drive_per_agent_logging.py @@ -6,12 +6,14 @@ - T2 (state persistence): two consecutive emissions with no new completions between them produce identical metric values, proving prepare_log no longer resets the per-agent EMAs. -- T3 (unequal completion rates): with sub-envs truncating at different +- T3 (two weightings agree in lockstep): the pooled fleet ratio and the + agent-weighted ratio coincide when agents complete equally often. +- T4 (unequal completion rates): with sub-envs truncating at different timesteps, every emission still weights one term per agent. T1 runs all agents in lockstep, the regime where completion-weighted and agent-weighted aggregation agree, so it holds either way; it is here to pin -the emission gate. T2 and T3 both fail against the completion-sum +the emission gate. T2 and T4 both fail against the completion-sum implementation. """ @@ -126,8 +128,53 @@ def test_emissions_identical_when_no_new_completions(): env.close() +def test_pooled_and_agent_weighted_ratios_agree_in_lockstep(): + """T3: avg_distance_per_infraction pools every completed episode; + agent_weighted_distance_per_infraction weights one episode per agent. Those + are the same quantity exactly when agents complete equally often, so in + lockstep they must agree. alpha=0 removes the EMA's time smoothing, leaving + the weighting as the only difference.""" + if not MAP_DIR.is_dir() or not any(MAP_DIR.glob("*.bin")): + pytest.skip(f"Drive map binaries not available at {MAP_DIR}") + env = Drive( + num_agents=NUM_AGENTS, + num_maps=1, + min_agents_per_env=1, + max_agents_per_env=8, + scenario_length=40, + report_interval=1, + map_dir=str(MAP_DIR), + log_ema_alpha=0.0, + resample_frequency=100_000, + ) + env.reset(seed=0) + + pooled_distance = 0.0 + pooled_infractions = 0.0 + agent_weighted_ratios = [] + for _ in range(400): + _, _, _, _, info = env.step(np.zeros_like(env.actions)) + for log in _log_dicts(info): + # The pooled pair resets on every emission, so summing recovers the + # interval totals; the agent-weighted ratio is a gauge, so it averages. + pooled_distance += log["total_distance_travelled_sum"] + pooled_infractions += log["total_infraction_count"] + agent_weighted_ratios.append(log["agent_weighted_distance_per_infraction"]) + + assert pooled_infractions > 0, "Expected at least one infraction to exercise the ratio" + pooled_ratio = pooled_distance / pooled_infractions + agent_weighted_ratio = float(np.mean(agent_weighted_ratios)) + + assert agent_weighted_ratio == pytest.approx(pooled_ratio, rel=0.02), ( + f"lockstep ratios diverged: pooled={pooled_ratio} agent_weighted={agent_weighted_ratio}; " + "with equal completion counts the two weightings must coincide" + ) + + env.close() + + def test_population_size_is_agent_count_under_unequal_completion_rates(): - """T3: the discriminating case. Sub-envs truncate independently here, so + """T4: the discriminating case. Sub-envs truncate independently here, so completions arrive staggered rather than all at once. Completion-weighted aggregation would make the emitted n track how many agents completed in that interval; agent-weighted aggregation pins it to the number of agents From 277870e0da7fda8252c4f96d41f439fe544b3dd7 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Fri, 14 Aug 2026 20:16:18 +0000 Subject: [PATCH 10/12] Regenerate smoke goldens --- tests/smoke_tests/data/drive_rollout_golden.json | 3 +++ tests/smoke_tests/data/drive_smoke_golden.json | 3 +++ 2 files changed, 6 insertions(+) diff --git a/tests/smoke_tests/data/drive_rollout_golden.json b/tests/smoke_tests/data/drive_rollout_golden.json index 9ba4f12e3b..4e6c779cca 100644 --- a/tests/smoke_tests/data/drive_rollout_golden.json +++ b/tests/smoke_tests/data/drive_rollout_golden.json @@ -1,5 +1,8 @@ { "env": { + "agent_weighted_distance_per_infraction": 12.329273682852133, + "agent_weighted_distance_travelled": 90.10730401700793, + "agent_weighted_infractions": 7.314489596441437, "avg_distance_per_infraction": 0.912943642768487, "avg_speed_per_agent": 1.359724880915123, "collision_rate": 0.026141917114206725, diff --git a/tests/smoke_tests/data/drive_smoke_golden.json b/tests/smoke_tests/data/drive_smoke_golden.json index f5b85520b5..dc04f5366d 100644 --- a/tests/smoke_tests/data/drive_smoke_golden.json +++ b/tests/smoke_tests/data/drive_smoke_golden.json @@ -1,5 +1,8 @@ { "env": { + "agent_weighted_distance_per_infraction": 11.734331180309427, + "agent_weighted_distance_travelled": 85.89256253735772, + "agent_weighted_infractions": 7.324183032430452, "avg_distance_per_infraction": 12.032092877288363, "avg_speed_per_agent": 1.3546471775605762, "collision_rate": 0.0009471982203680893, From cf862d9fc63e46ea614992bb5684336a88b70cb5 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Fri, 14 Aug 2026 22:07:21 -0400 Subject: [PATCH 11/12] Trim the Log comments on the distance/infraction fields Co-Authored-By: Claude Opus 5 (1M context) --- pufferlib/ocean/drive/drive.h | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index 154b89f457..fc23646c00 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -97,12 +97,9 @@ struct Log { float ttc_samples; float multi_lane_time; float multi_lane_score; - // Pooled over every episode completed in the interval; prepare_log keeps them - // out of the EMA so vec_log's zeroing makes them reset-on-consumption. float total_distance_travelled; float total_infractions; - // Same quantities as ordinary EMA slots, so their ratio weights one episode - // per agent instead of one per completion. + // Used to compute average distance per infraction weighted per agent slot rather than per episode float agent_weighted_distance_travelled; float agent_weighted_infractions; // Agent-only puffer display fields (for serialization) From c1c9b6dae6621a62715405a46afdc4b0dd6d8641 Mon Sep 17 00:00:00 2001 From: Eugene Vinitsky Date: Fri, 14 Aug 2026 22:15:30 -0400 Subject: [PATCH 12/12] Reword the prepare_log comments and drop the log_ema_alpha one Co-Authored-By: Claude Opus 5 (1M context) --- pufferlib/ocean/drive/drive.h | 9 ++++----- pufferlib/ocean/drive/drive.py | 2 -- 2 files changed, 4 insertions(+), 7 deletions(-) diff --git a/pufferlib/ocean/drive/drive.h b/pufferlib/ocean/drive/drive.h index fc23646c00..62ffd1952a 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -1946,12 +1946,12 @@ static float calculate_puffer_score(Log *agent_log, float duration_steps, float // Sum each agent's current EMA slot into env->log and set env->log.n to the // number of slots that have ever been seeded. vec_log divides by aggregate.n -// downstream, producing a cross-agent population mean with one weight per -// agent regardless of that agent's completion frequency. Does not reset the +// to compute the per-agent logged values. Does not reset the // slots: emissions between completions repeat the previous values. static void prepare_log(Drive *env) { - // Raw interval totals, not per-agent gauges: they accumulate until vec_log emits and zeroes them, - // so summing them across emissions yields the global totals avg_distance_per_infraction needs. + // this specific metric is also computed in an episodic way + // in addition to a per-agent way so this one metric needs + // to be treated separately. float interval_distance_travelled = env->log.total_distance_travelled; float interval_infraction_count = env->log.total_infractions; memset(&env->log, 0, sizeof(Log)); @@ -2081,7 +2081,6 @@ static void add_log(Drive *env) { } } - // Keeps env->log valid for C-side readers, which never reach the Python vec_prepare_log. prepare_log(env); env->log_episode_seed = env->episode_seed; } diff --git a/pufferlib/ocean/drive/drive.py b/pufferlib/ocean/drive/drive.py index 529f38bbe7..9cb76f9071 100644 --- a/pufferlib/ocean/drive/drive.py +++ b/pufferlib/ocean/drive/drive.py @@ -233,8 +233,6 @@ def __init__( if self.eval_scenario_seeds is None or len(self.eval_scenario_seeds) != len(self.eval_map_indices): raise ValueError("eval_scenario_seeds must have one seed per eval_map_indices entry") self.use_exact_episode_seed = bool(eval_mode) and self.eval_scenario_seeds is not None - # Outside [0, 1) the slot update stops being a convex combination: 1.0 pins every - # metric to its first episode, above 1.0 the slots diverge. if not 0.0 <= log_ema_alpha < 1.0: raise ValueError(f"log_ema_alpha must be in [0.0, 1.0). Got: {log_ema_alpha}") self.log_ema_alpha = log_ema_alpha