diff --git a/pufferlib/config/puffer_drive.yaml b/pufferlib/config/puffer_drive.yaml index e2dc186f23..e78430f310 100644 --- a/pufferlib/config/puffer_drive.yaml +++ b/pufferlib/config/puffer_drive.yaml @@ -128,6 +128,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 78a1dee3a4..26504c2caf 100644 --- a/pufferlib/config_schema.py +++ b/pufferlib/config_schema.py @@ -136,6 +136,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 c8fe42d850..a6d41411c2 100644 --- a/pufferlib/ocean/drive/binding.c +++ b/pufferlib/ocean/drive/binding.c @@ -26,14 +26,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; +} + // Seeds are 63-bit non-negative so they survive int64 round-trips (numpy, pandas, CSV). static int unpack_seed(PyObject *kwargs, uint64_t *seed_out) { PyObject *seed_obj = PyDict_GetItemString(kwargs, "seed"); @@ -1925,6 +1940,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"); @@ -2007,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); @@ -2025,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 12170cc006..62ffd1952a 100644 --- a/pufferlib/ocean/drive/drive.h +++ b/pufferlib/ocean/drive/drive.h @@ -99,6 +99,9 @@ struct Log { float multi_lane_score; float total_distance_travelled; float total_infractions; + // 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) float no_at_fault; float no_offroad; @@ -313,6 +316,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 { @@ -1937,9 +1944,39 @@ 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 +// 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) { + // 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)); + 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; + env->log.total_distance_travelled = interval_distance_travelled; + env->log.total_infractions = interval_infraction_count; +} + 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; @@ -1947,51 +1984,55 @@ 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; + 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; } + 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; - 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( @@ -1999,41 +2040,48 @@ 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]; + } + } } + + prepare_log(env); env->log_episode_seed = env->episode_seed; } @@ -2868,6 +2916,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); } @@ -2929,6 +2979,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 83cf98201c..9cb76f9071 100644 --- a/pufferlib/ocean/drive/drive.py +++ b/pufferlib/ocean/drive/drive.py @@ -99,6 +99,7 @@ def __init__( max_scenarios_per_batch=None, 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", @@ -232,6 +233,9 @@ 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 + 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 self.terminate_on_goal = terminate_on_goal @@ -558,6 +562,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, @@ -631,6 +636,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 67c2357fb9..cdbba8cff6 100644 --- a/pufferlib/ocean/env_binding.h +++ b/pufferlib/ocean/env_binding.h @@ -723,8 +723,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}; @@ -737,8 +735,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/smoke_tests/data/drive_rollout_golden.json b/tests/smoke_tests/data/drive_rollout_golden.json index 8c55869b2c..4e6c779cca 100644 --- a/tests/smoke_tests/data/drive_rollout_golden.json +++ b/tests/smoke_tests/data/drive_rollout_golden.json @@ -1,33 +1,36 @@ { "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, + "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, + "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..dc04f5366d 100644 --- a/tests/smoke_tests/data/drive_smoke_golden.json +++ b/tests/smoke_tests/data/drive_smoke_golden.json @@ -1,36 +1,39 @@ { "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.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, 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..33fb09ef87 --- /dev/null +++ b/tests/unit_tests/test_drive_per_agent_logging.py @@ -0,0 +1,215 @@ +"""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 (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 (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 T4 both fail against the completion-sum +implementation. +""" + +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 _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 + 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_emissions_identical_when_no_new_completions(): + """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() + 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() + + +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(): + """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 + # 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))}; " + "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()