From 351b22fe01b8e54209790be6383762582f9d6539 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 09:40:20 +0200 Subject: [PATCH 01/13] Optionally include infos in expert data --- scratch/etienne/intersimple/data/expert.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/scratch/etienne/intersimple/data/expert.py b/scratch/etienne/intersimple/data/expert.py index 56105af..e198502 100644 --- a/scratch/etienne/intersimple/data/expert.py +++ b/scratch/etienne/intersimple/data/expert.py @@ -1,4 +1,4 @@ -from intersim.envs.intersimple import Intersimple +from intersim.envs.intersimple import Intersimple, InfoFilter from stable_baselines3.common.policies import BasePolicy import gym import intersim.envs.intersimple @@ -117,6 +117,7 @@ def demonstrations(expert='NormalizedIntersimpleExpert', env='NRasterizedRandomA save_video(env, policy) path = path or (policy.__class__.__name__ + '_' + env.__class__.__name__ + '.pkl') + include_infos = isinstance(env, InfoFilter) rollout.rollout_and_save( path=path, @@ -125,7 +126,8 @@ def demonstrations(expert='NormalizedIntersimpleExpert', env='NRasterizedRandomA sample_until=rollout.make_sample_until( min_timesteps=min_timesteps, min_episodes=min_episodes, - ) + ), + exclude_infos=not include_infos, ) if __name__ == '__main__': From 71f69c43edc3804514389fd4bd37c2171bca4442 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 09:40:52 +0200 Subject: [PATCH 02/13] Fix typo in dataset folder name --- scratch/etienne/intersimple/gail_options_image_alltracks.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/scratch/etienne/intersimple/gail_options_image_alltracks.py b/scratch/etienne/intersimple/gail_options_image_alltracks.py index 71d4741..84b418b 100644 --- a/scratch/etienne/intersimple/gail_options_image_alltracks.py +++ b/scratch/etienne/intersimple/gail_options_image_alltracks.py @@ -331,7 +331,7 @@ if __name__ == '__main__': #env_class = NRasterized #env_settings = {'agent': 51, 'width': 36, 'height': 36, 'm_per_px': 2} - files = ['../../../expert_data/DR_USA_Roundabout_FT0/track%04i/expert.pkl'%(i) for i in range(5)] + files = ['../../../expert_data/DR_USA_Roundabout_FT/track%04i/expert.pkl'%(i) for i in range(5)] transitions=load_experts(files) generator = train( From 3857716cec078893d5a69ed5c59590b8faa03d38 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 09:42:27 +0200 Subject: [PATCH 03/13] Rename predict to forward This is done to be consistent with stable baselines interface. predict is then automatically defined. This is necessary to use stable baselines' evaluate_policy method --- src/gail/options.py | 2 +- src/policies/options.py | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/gail/options.py b/src/gail/options.py index 21dd7af..8eb5f16 100644 --- a/src/gail/options.py +++ b/src/gail/options.py @@ -36,7 +36,7 @@ class OptionsEnv(gym.Wrapper): self.episode_start = True self.m = available_actions(self.env, self.options) - self.ch, self.value, self.log_prob = generator.policy.predict({ + self.ch, self.value, self.log_prob = generator.policy.forward({ 'obs': torch.tensor(self.s).unsqueeze(0).to(generator.policy.device), 'mask': torch.tensor(self.m).unsqueeze(0).to(generator.policy.device), }) diff --git a/src/policies/options.py b/src/policies/options.py index 9bf71c9..851edd4 100644 --- a/src/policies/options.py +++ b/src/policies/options.py @@ -22,7 +22,7 @@ class OptionsCnnPolicy(stable_baselines3.common.policies.ActorCriticCnnPolicy): values = self.value_net(latent_vf) return values, distribution.distribution - def predict(self, obs, eps=1e-6): + def forward(self, obs, eps=1e-6): """ Will mask invalid states before making action selections Args: From d89c77409ceac1e6e12998bb4a8f09b9ebd168e2 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 09:43:55 +0200 Subject: [PATCH 04/13] Add some scratch --- scratch/johannes/evaluation.py | 90 ++++++++++++++++++++++++++++++ scratch/johannes/raytune_simple.py | 69 +++++++++++++++++++++++ 2 files changed, 159 insertions(+) create mode 100644 scratch/johannes/evaluation.py create mode 100644 scratch/johannes/raytune_simple.py diff --git a/scratch/johannes/evaluation.py b/scratch/johannes/evaluation.py new file mode 100644 index 0000000..67fec87 --- /dev/null +++ b/scratch/johannes/evaluation.py @@ -0,0 +1,90 @@ + +def evaluate_policy_simple( + model, + env: gym.Env, + n_eval_episodes: int = 10, + deterministic: bool = True, + render: bool = False, + callback = None, + reward_threshold = None, + return_episode_rewards: bool = False, + warn: bool = True, +): + """ + Runs policy for ``n_eval_episodes`` episodes and returns average reward. + If a vector env is passed in, this divides the episodes to evaluate onto the + different elements of the vector env. This static division of work is done to + remove bias. See https://github.com/DLR-RM/stable-baselines3/issues/402 for more + details and discussion. + + .. note:: + If environment has not been wrapped with ``Monitor`` wrapper, reward and + episode lengths are counted as it appears with ``env.step`` calls. If + the environment contains wrappers that modify rewards or episode lengths + (e.g. reward scaling, early episode reset), these will affect the evaluation + results as well. You can avoid this by wrapping environment with ``Monitor`` + wrapper before anything else. + + :param model: The RL agent you want to evaluate. + :param env: The gym environment or ``VecEnv`` environment. + :param n_eval_episodes: Number of episode to evaluate the agent + :param deterministic: Whether to use deterministic or stochastic actions + :param render: Whether to render the environment or not + :param callback: callback function to do additional checks, + called after each step. Gets locals() and globals() passed as parameters. + :param reward_threshold: Minimum expected reward per episode, + this will raise an error if the performance is not met + :param return_episode_rewards: If True, a list of rewards and episode lengths + per episode will be returned instead of the mean. + :param warn: If True (default), warns user about lack of a Monitor wrapper in the + evaluation environment. + :return: Mean reward per episode, std of reward per episode. + Returns ([float], [int]) when ``return_episode_rewards`` is True, first + list containing per-episode rewards and second containing per-episode lengths + (in number of steps). + """ + episode_rewards = [] + episode_lengths = [] + + episode_counts = 0 + + current_rewards = 0 + current_lengths = 0 + observations = env.reset() + states = None + while (episode_counts < n_eval_episodes): + actions, states = model.predict(observations, state=states, deterministic=deterministic) + observations, rewards, dones, infos = env.step(actions) + print(env._env.t) + current_rewards += rewards + current_lengths += 1 + + # unpack values so that the callback can access the local variables + reward = rewards + done = dones + info = infos + if info['collision']: + print("COLLISION") + + if callback is not None: + callback(locals(), globals()) + + if dones: + episode_rewards.append(current_rewards) + episode_lengths.append(current_lengths) + episode_counts += 1 + current_rewards = 0 + current_lengths = 0 + if states is not None: + states *= 0 + + if render: + env.render() + + mean_reward = np.mean(episode_rewards) + std_reward = np.std(episode_rewards) + if reward_threshold is not None: + assert mean_reward > reward_threshold, "Mean reward below threshold: " f"{mean_reward:.2f} < {reward_threshold:.2f}" + if return_episode_rewards: + return episode_rewards, episode_lengths + return mean_reward, std_reward diff --git a/scratch/johannes/raytune_simple.py b/scratch/johannes/raytune_simple.py new file mode 100644 index 0000000..5979c11 --- /dev/null +++ b/scratch/johannes/raytune_simple.py @@ -0,0 +1,69 @@ +"""This example demonstrates basic Ray Tune random search and grid search.""" +import time + +import ray +from ray import tune + + +def evaluation_fn(step, width, height): + time.sleep(0.1) + return (0.1 + width * step / 100)**(-1) + height * 0.1 + +def easy_objective(config): + # Hyperparameters + width, height = config["width"], config["height"] + + mydata = ray.get(ray_data) + print(mydata) + + for step in range(config["steps"]): + # Iterative training function - can be any arbitrary training procedure + intermediate_score = evaluation_fn(step, width, height) + # Feed the score back back to Tune. + tune.report(iterations=step, mean_loss=intermediate_score) + + +if __name__ == "__main__": + import argparse + + parser = argparse.ArgumentParser() + parser.add_argument( + "--smoke-test", action="store_true", help="Finish quickly for testing") + parser.add_argument( + "--server-address", + type=str, + default=None, + required=False, + help="The address of server to connect to if using " + "Ray Client.") + args, _ = parser.parse_known_args() + if args.server_address is not None: + ray.init(f"ray://{args.server_address}") + else: + ray.init(configure_logging=False) + + # This will do a grid search over the `activation` parameter. This means + # that each of the two values (`relu` and `tanh`) will be sampled once + # for each sample (`num_samples`). We end up with 2 * 50 = 100 samples. + # The `width` and `height` parameters are sampled randomly. + # `steps` is a constant parameter. + + import numpy as np + N = 3 + data = np.random.rand(N,N,N) + ray_data = ray.put(data) + + + analysis = tune.run( + easy_objective, + metric="mean_loss", + mode="min", + num_samples=5 if args.smoke_test else 50, + config={ + "steps": 5 if args.smoke_test else 100, + "width": tune.uniform(0, 20), + "height": tune.uniform(-100, 100), + "activation": tune.grid_search(["relu", "tanh"]) + }) + + print("Best hyperparameters found were: ", analysis.best_config) \ No newline at end of file From 09afee4e1d52d16e074799e8da3be74b4aad8253 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 09:55:20 +0200 Subject: [PATCH 05/13] Add first metrics --- scratch/johannes/gail_options_image.py | 190 +++++++++++++++++++++++++ 1 file changed, 190 insertions(+) create mode 100644 scratch/johannes/gail_options_image.py diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py new file mode 100644 index 0000000..f3b7f0c --- /dev/null +++ b/scratch/johannes/gail_options_image.py @@ -0,0 +1,190 @@ +# %% +# import sys +# sys.path.append('../../../') + +from src.discriminator import CnnDiscriminator, CnnDiscriminatorFlatAction +from imitation.algorithms import adversarial +import stable_baselines3 +from stable_baselines3.common.evaluation import evaluate_policy +import torch.utils.data +import numpy as np +from intersim.envs.intersimple import Intersimple, NRasterized, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward +import itertools +import functools +from torch.distributions import Categorical +import gym +import torch +import pickle +import imitation.data.rollout as rollout +import tempfile +import pathlib +from imitation.util import logger +from stable_baselines3.common.env_util import make_vec_env +from tqdm import tqdm +from src.policies.options import OptionsCnnPolicy +from src.gail.options import OptionsEnv, LLOptions, HLOptions, RenderOptions +from src.gail.train import train_discriminator, train_generator +from src.metrics import nanmean, divergence, visualize_distribution + +model_name = 'gail_options_image' +env_settings = {'width': 36, 'height': 36, 'm_per_px': 2} + +ALL_OPTIONS = [(v,t) for v in [0,2,4,6,8] for t in [5]] # option 0 is safe fallback + +def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, discount=0.99): + env = NRasterizedRandomAgent(**env_settings) + env.discount = discount + + tempdir = tempfile.TemporaryDirectory(prefix="quickstart") + tempdir_path = pathlib.Path(tempdir.name) + logger.configure(tempdir_path / "GAIL/") + print(f"All Tensorboards and logging are being written inside {tempdir_path}/.") + + venv = make_vec_env(NRasterized, n_envs=1, env_kwargs=env_settings) + discriminator = adversarial.GAIL( + expert_data=expert_data, + expert_batch_size=expert_batch_size, + discrim_kwargs={'discrim_net': CnnDiscriminatorFlatAction(venv)}, + #discrim_kwargs={'discrim_net': CnnDiscriminator(venv)}, + venv=venv, # unused + gen_algo=stable_baselines3.PPO("CnnPolicy", venv), # unused + ) + + generator = stable_baselines3.PPO( + OptionsCnnPolicy, + OptionsEnv(env, options=ALL_OPTIONS), + verbose=1, + n_steps=generator_steps, + ) + + # PPO.train requires logger as set up in + # PPO._setup_learn (called by PPO.learn) + generator._logger = stable_baselines3.common.utils.configure_logger( + generator.verbose, + generator.tensorboard_log, + ) + + for epoch in tqdm(range(epochs)): + train_discriminator(LLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=expert_batch_size) + train_generator(HLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=generator_steps) + + eval_env = NRasterizedRandomAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) + ev = Evaluation(eval_env, n_eval_episodes=10) + ev.evaluate(epoch, generator, discriminator, expert_data) + + return generator + +from stable_baselines3.common.vec_env import VecEnv +class Evaluation: + def __init__(self, eval_env, n_eval_episodes=10): + # if env is a VecEnv, the code needs to be adapted, since the callback will be called after each step, + # so transitions of different envs will be mixed and the total number of episodes could be larger than n_eval_episodes! + assert not isinstance(eval_env, VecEnv) + self.env = eval_env + self.n_eval_episodes = n_eval_episodes + self.reset() + + def reset(self): + self._n_collisions = 0 + self._trajectories = [] + self._episode_done = True + + def evaluate(self, epoch, generator, discriminator, expert_data): + self.reset() + info = {} + + episode_rewards, episode_lengths = evaluate_policy( + generator, + self.env, + n_eval_episodes=self.n_eval_episodes, + callback=self.evaluate_policy_callback, + return_episode_rewards=True + ) + + collision_rate = self._n_collisions / self.n_eval_episodes + info['collision_rate'] = collision_rate + + assert len(self._trajectories) >= self.n_eval_episodes + + # average velocity of each episode + # this first averages velocity over single trajectories and then averages over trajectories + # avg_velocities = [nanmean(torch.stack(t)[:,2]) for t in self._trajectories] + # avg_velocity = np.mean(avg_velocities) + + # velocities produced by generator + policy_velocities = torch.cat([torch.stack(t)[:,2] for t in self._trajectories]) + policy_velocities = policy_velocities[~torch.isnan(policy_velocities)] + + # expert velocities + extract_state = lambda info: info['projected_state'][info['agent']] + expert_velocities = torch.stack([extract_state(info) for info in expert_data.infos])[:,2] + expert_velocities = expert_velocities[~torch.isnan(expert_velocities)] + + info['avg_velocity_loss'] = expert_velocities.mean() - policy_velocities.mean() + + # divergence (true, policy) + + # same for acceleration + + print(info) + return info + + def evaluate_policy_callback(self, local_vars, global_vars): + info = local_vars['info'] + done = local_vars['done'] + env = local_vars['env'].envs[local_vars['i']] + assert isinstance(env, Intersimple) + + # Increase collision counter if episode terminated with a collision + if info['collision']: + assert done + self._n_collisions += 1 + + # if last episode is done, start new trajectory + if self._episode_done: + self._trajectories.append([]) + self._trajectories[-1].append(info['projected_state'][env._agent]) + self._episode_done = done + + + + +# class CAPolicy: +# def __init__(self, a): +# self.a = torch.tensor([a]) + +# def predict(self, obs, state=None, deterministic=False): +# return self.a, state + + +# %% +if __name__ == '__main__': + # %% + + # env = NRasterizedIncrementingAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) + # generator = CAPolicy(0.0) + + # ev = Evaluation(env, 10) + # ev.evaluate(1, generator, None, None) + + # exit() + + # %% + + with open("scratch/etienne/intersimple/data/NormalizedIntersimpleExpertMu.001N200_NRasterizedInfoAgent51w36h36mppx2.pkl", "rb") as f: + trajectories = pickle.load(f) + transitions = rollout.flatten_trajectories(trajectories) + generator = train(transitions) + + generator.save(model_name) + + # %% + model = stable_baselines3.PPO.load(model_name) + + env = RenderOptions(NRasterizedRandomAgent(**env_settings), options=ALL_OPTIONS) + + for s in env.sample_ll(model): + if s['dones']: + break + + env.close(filestr='render/'+model_name) From f217daf2514f483cf4cd06629036178fb5ad2842 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:10:44 +0200 Subject: [PATCH 06/13] update evaluation --- scratch/johannes/gail_options_image.py | 62 +++++++++++++++++--------- 1 file changed, 40 insertions(+), 22 deletions(-) diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py index f3b7f0c..f405db1 100644 --- a/scratch/johannes/gail_options_image.py +++ b/scratch/johannes/gail_options_image.py @@ -8,7 +8,7 @@ import stable_baselines3 from stable_baselines3.common.evaluation import evaluate_policy import torch.utils.data import numpy as np -from intersim.envs.intersimple import Intersimple, NRasterized, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward +from intersim.envs.intersimple import Intersimple, NormalizedActionSpace, NRasterized, NRasterizedInfo, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward import itertools import functools from torch.distributions import Categorical @@ -27,12 +27,13 @@ from src.gail.train import train_discriminator, train_generator from src.metrics import nanmean, divergence, visualize_distribution model_name = 'gail_options_image' +Env = NRasterized env_settings = {'width': 36, 'height': 36, 'm_per_px': 2} ALL_OPTIONS = [(v,t) for v in [0,2,4,6,8] for t in [5]] # option 0 is safe fallback def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, discount=0.99): - env = NRasterizedRandomAgent(**env_settings) + env = Env(**env_settings) env.discount = discount tempdir = tempfile.TemporaryDirectory(prefix="quickstart") @@ -68,7 +69,7 @@ def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, disc train_discriminator(LLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=expert_batch_size) train_generator(HLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=generator_steps) - eval_env = NRasterizedRandomAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) + eval_env = Env(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) ev = Evaluation(eval_env, n_eval_episodes=10) ev.evaluate(epoch, generator, discriminator, expert_data) @@ -88,10 +89,11 @@ class Evaluation: self._n_collisions = 0 self._trajectories = [] self._episode_done = True + self._accelerations = [] def evaluate(self, epoch, generator, discriminator, expert_data): self.reset() - info = {} + metrics = {} episode_rewards, episode_lengths = evaluate_policy( generator, @@ -102,7 +104,7 @@ class Evaluation: ) collision_rate = self._n_collisions / self.n_eval_episodes - info['collision_rate'] = collision_rate + metrics['collision_rate'] = collision_rate assert len(self._trajectories) >= self.n_eval_episodes @@ -113,6 +115,7 @@ class Evaluation: # velocities produced by generator policy_velocities = torch.cat([torch.stack(t)[:,2] for t in self._trajectories]) + # if episodes terminate without collisions, then the state is fully nan policy_velocities = policy_velocities[~torch.isnan(policy_velocities)] # expert velocities @@ -120,19 +123,28 @@ class Evaluation: expert_velocities = torch.stack([extract_state(info) for info in expert_data.infos])[:,2] expert_velocities = expert_velocities[~torch.isnan(expert_velocities)] - info['avg_velocity_loss'] = expert_velocities.mean() - policy_velocities.mean() + metrics['avg_velocity_loss'] = (expert_velocities.mean() - policy_velocities.mean()).item() + metrics['velocity_divergence'] = divergence(policy_velocities, expert_velocities, type='js') + - # divergence (true, policy) + # accelerations produced by generator + policy_accelerations = torch.tensor(self._accelerations) + # expert accelerations + extract_accel = lambda info: info['action_taken'][info['agent']] + expert_accelerations = torch.cat([extract_accel(info) for info in expert_data.infos]) - # same for acceleration + metrics['acceleration_divergence'] = divergence(policy_accelerations, expert_accelerations, type='js') + visualize_distribution(expert_accelerations, policy_accelerations, 'output/_action_viz{:02}'.format(epoch)) - print(info) - return info + print(metrics) + return metrics def evaluate_policy_callback(self, local_vars, global_vars): + venv_i = local_vars['i'] info = local_vars['info'] done = local_vars['done'] - env = local_vars['env'].envs[local_vars['i']] + _agent = info['agent'] + env = local_vars['env'].envs[venv_i] assert isinstance(env, Intersimple) # Increase collision counter if episode terminated with a collision @@ -143,7 +155,12 @@ class Evaluation: # if last episode is done, start new trajectory if self._episode_done: self._trajectories.append([]) - self._trajectories[-1].append(info['projected_state'][env._agent]) + agent_state = info['projected_state'][_agent] + self._trajectories[-1].append(agent_state) + # this does not work since agent_action are high level options + # agent_action = local_vars['actions'][venv_i] # normalized intersimple action + # acceleration = env._unnormalize(agent_action) if isinstance(env, NormalizedActionSpace) else agent_action + self._accelerations.append(info['action_taken'][_agent]) self._episode_done = done @@ -160,20 +177,21 @@ class Evaluation: # %% if __name__ == '__main__': # %% - - # env = NRasterizedIncrementingAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) - # generator = CAPolicy(0.0) - - # ev = Evaluation(env, 10) - # ev.evaluate(1, generator, None, None) - - # exit() - - # %% with open("scratch/etienne/intersimple/data/NormalizedIntersimpleExpertMu.001N200_NRasterizedInfoAgent51w36h36mppx2.pkl", "rb") as f: trajectories = pickle.load(f) transitions = rollout.flatten_trajectories(trajectories) + +### + # env = NRasterizedIncrementingAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) + # generator = CAPolicy(.5) + + # ev = Evaluation(env, 10) + # ev.evaluate(1, generator, None, transitions) + + # exit() +### + generator = train(transitions) generator.save(model_name) From 2218d144098dc265187c7a90b947532c5825b054 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:26:17 +0200 Subject: [PATCH 07/13] Move metrics --- src/__init__.py | 2 +- src/evaluation/__init__.py | 0 src/{ => evaluation}/metrics.py | 0 3 files changed, 1 insertion(+), 1 deletion(-) create mode 100644 src/evaluation/__init__.py rename src/{ => evaluation}/metrics.py (100%) diff --git a/src/__init__.py b/src/__init__.py index d44975d..aebb694 100644 --- a/src/__init__.py +++ b/src/__init__.py @@ -1,3 +1,3 @@ from src.data.expert_data import generate_expert_data, load_expert_data from src.data.data_utils import InteractionDatasetSingleAgent -from src.metrics import metrics \ No newline at end of file +from src.evaluation.metrics import metrics \ No newline at end of file diff --git a/src/evaluation/__init__.py b/src/evaluation/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/src/metrics.py b/src/evaluation/metrics.py similarity index 100% rename from src/metrics.py rename to src/evaluation/metrics.py From 7ffcc0b4b8bb8601502edab3617e79ef56bab259 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:29:00 +0200 Subject: [PATCH 08/13] Extract evaluation code to separate file --- scratch/johannes/gail_options_image.py | 95 +------------------------- src/evaluation/evaluation.py | 92 +++++++++++++++++++++++++ 2 files changed, 95 insertions(+), 92 deletions(-) create mode 100644 src/evaluation/evaluation.py diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py index f405db1..ac1676a 100644 --- a/scratch/johannes/gail_options_image.py +++ b/scratch/johannes/gail_options_image.py @@ -8,7 +8,7 @@ import stable_baselines3 from stable_baselines3.common.evaluation import evaluate_policy import torch.utils.data import numpy as np -from intersim.envs.intersimple import Intersimple, NormalizedActionSpace, NRasterized, NRasterizedInfo, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward +from intersim.envs.intersimple import Intersimple, NRasterized, NRasterizedInfo, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward import itertools import functools from torch.distributions import Categorical @@ -24,10 +24,9 @@ from tqdm import tqdm from src.policies.options import OptionsCnnPolicy from src.gail.options import OptionsEnv, LLOptions, HLOptions, RenderOptions from src.gail.train import train_discriminator, train_generator -from src.metrics import nanmean, divergence, visualize_distribution model_name = 'gail_options_image' -Env = NRasterized +Env = NRasterizedRandomAgent env_settings = {'width': 36, 'height': 36, 'm_per_px': 2} ALL_OPTIONS = [(v,t) for v in [0,2,4,6,8] for t in [5]] # option 0 is safe fallback @@ -41,7 +40,7 @@ def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, disc logger.configure(tempdir_path / "GAIL/") print(f"All Tensorboards and logging are being written inside {tempdir_path}/.") - venv = make_vec_env(NRasterized, n_envs=1, env_kwargs=env_settings) + venv = make_vec_env(Env, n_envs=1, env_kwargs=env_settings) discriminator = adversarial.GAIL( expert_data=expert_data, expert_batch_size=expert_batch_size, @@ -75,94 +74,6 @@ def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, disc return generator -from stable_baselines3.common.vec_env import VecEnv -class Evaluation: - def __init__(self, eval_env, n_eval_episodes=10): - # if env is a VecEnv, the code needs to be adapted, since the callback will be called after each step, - # so transitions of different envs will be mixed and the total number of episodes could be larger than n_eval_episodes! - assert not isinstance(eval_env, VecEnv) - self.env = eval_env - self.n_eval_episodes = n_eval_episodes - self.reset() - - def reset(self): - self._n_collisions = 0 - self._trajectories = [] - self._episode_done = True - self._accelerations = [] - - def evaluate(self, epoch, generator, discriminator, expert_data): - self.reset() - metrics = {} - - episode_rewards, episode_lengths = evaluate_policy( - generator, - self.env, - n_eval_episodes=self.n_eval_episodes, - callback=self.evaluate_policy_callback, - return_episode_rewards=True - ) - - collision_rate = self._n_collisions / self.n_eval_episodes - metrics['collision_rate'] = collision_rate - - assert len(self._trajectories) >= self.n_eval_episodes - - # average velocity of each episode - # this first averages velocity over single trajectories and then averages over trajectories - # avg_velocities = [nanmean(torch.stack(t)[:,2]) for t in self._trajectories] - # avg_velocity = np.mean(avg_velocities) - - # velocities produced by generator - policy_velocities = torch.cat([torch.stack(t)[:,2] for t in self._trajectories]) - # if episodes terminate without collisions, then the state is fully nan - policy_velocities = policy_velocities[~torch.isnan(policy_velocities)] - - # expert velocities - extract_state = lambda info: info['projected_state'][info['agent']] - expert_velocities = torch.stack([extract_state(info) for info in expert_data.infos])[:,2] - expert_velocities = expert_velocities[~torch.isnan(expert_velocities)] - - metrics['avg_velocity_loss'] = (expert_velocities.mean() - policy_velocities.mean()).item() - metrics['velocity_divergence'] = divergence(policy_velocities, expert_velocities, type='js') - - - # accelerations produced by generator - policy_accelerations = torch.tensor(self._accelerations) - # expert accelerations - extract_accel = lambda info: info['action_taken'][info['agent']] - expert_accelerations = torch.cat([extract_accel(info) for info in expert_data.infos]) - - metrics['acceleration_divergence'] = divergence(policy_accelerations, expert_accelerations, type='js') - visualize_distribution(expert_accelerations, policy_accelerations, 'output/_action_viz{:02}'.format(epoch)) - - print(metrics) - return metrics - - def evaluate_policy_callback(self, local_vars, global_vars): - venv_i = local_vars['i'] - info = local_vars['info'] - done = local_vars['done'] - _agent = info['agent'] - env = local_vars['env'].envs[venv_i] - assert isinstance(env, Intersimple) - - # Increase collision counter if episode terminated with a collision - if info['collision']: - assert done - self._n_collisions += 1 - - # if last episode is done, start new trajectory - if self._episode_done: - self._trajectories.append([]) - agent_state = info['projected_state'][_agent] - self._trajectories[-1].append(agent_state) - # this does not work since agent_action are high level options - # agent_action = local_vars['actions'][venv_i] # normalized intersimple action - # acceleration = env._unnormalize(agent_action) if isinstance(env, NormalizedActionSpace) else agent_action - self._accelerations.append(info['action_taken'][_agent]) - self._episode_done = done - diff --git a/src/evaluation/evaluation.py b/src/evaluation/evaluation.py new file mode 100644 index 0000000..882069a --- /dev/null +++ b/src/evaluation/evaluation.py @@ -0,0 +1,92 @@ +from stable_baselines3.common.vec_env import VecEnv +from stable_baselines3.common.evaluation import evaluate_policy + +from intersim.envs.intersimple import Intersimple +from src.evaluation.metrics import nanmean, divergence, visualize_distribution + +class Evaluation: + def __init__(self, eval_env, n_eval_episodes=10): + # if env is a VecEnv, the code needs to be adapted, since the callback will be called after each step, + # so transitions of different envs will be mixed and the total number of episodes could be larger than n_eval_episodes! + assert not isinstance(eval_env, VecEnv) + self.env = eval_env + self.n_eval_episodes = n_eval_episodes + self.reset() + + def reset(self): + self._n_collisions = 0 + self._trajectories = [] + self._episode_done = True + self._accelerations = [] + + def evaluate(self, epoch, generator, discriminator, expert_data): + self.reset() + metrics = {} + + episode_rewards, episode_lengths = evaluate_policy( + generator, + self.env, + n_eval_episodes=self.n_eval_episodes, + callback=self.evaluate_policy_callback, + return_episode_rewards=True + ) + + collision_rate = self._n_collisions / self.n_eval_episodes + metrics['collision_rate'] = collision_rate + + assert len(self._trajectories) >= self.n_eval_episodes + + # average velocity of each episode + # this first averages velocity over single trajectories and then averages over trajectories + # avg_velocities = [nanmean(torch.stack(t)[:,2]) for t in self._trajectories] + # avg_velocity = np.mean(avg_velocities) + + # velocities produced by generator + policy_velocities = torch.cat([torch.stack(t)[:,2] for t in self._trajectories]) + # if episodes terminate without collisions, then the state is fully nan + policy_velocities = policy_velocities[~torch.isnan(policy_velocities)] + + # expert velocities + extract_state = lambda info: info['projected_state'][info['agent']] + expert_velocities = torch.stack([extract_state(info) for info in expert_data.infos])[:,2] + expert_velocities = expert_velocities[~torch.isnan(expert_velocities)] + + metrics['avg_velocity_loss'] = (expert_velocities.mean() - policy_velocities.mean()).item() + metrics['velocity_divergence'] = divergence(policy_velocities, expert_velocities, type='js') + + + # accelerations produced by generator + policy_accelerations = torch.tensor(self._accelerations) + # expert accelerations + extract_accel = lambda info: info['action_taken'][info['agent']] + expert_accelerations = torch.cat([extract_accel(info) for info in expert_data.infos]) + + metrics['acceleration_divergence'] = divergence(policy_accelerations, expert_accelerations, type='js') + visualize_distribution(expert_accelerations, policy_accelerations, 'output/_action_viz{:02}'.format(epoch)) + + print(metrics) + return metrics + + def evaluate_policy_callback(self, local_vars, global_vars): + venv_i = local_vars['i'] + info = local_vars['info'] + done = local_vars['done'] + _agent = info['agent'] + env = local_vars['env'].envs[venv_i] + assert isinstance(env, Intersimple) + + # Increase collision counter if episode terminated with a collision + if info['collision']: + assert done + self._n_collisions += 1 + + # if last episode is done, start new trajectory + if self._episode_done: + self._trajectories.append([]) + agent_state = info['projected_state'][_agent] + self._trajectories[-1].append(agent_state) + # this does not work since agent_action are high level options + # agent_action = local_vars['actions'][venv_i] # normalized intersimple action + # acceleration = env._unnormalize(agent_action) if isinstance(env, NormalizedActionSpace) else agent_action + self._accelerations.append(info['action_taken'][_agent]) + self._episode_done = done From 1bab1aaab7dc6986797bd0a325477f44df271bef Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:31:08 +0200 Subject: [PATCH 09/13] Cleanup evaluation --- src/evaluation/evaluation.py | 15 +++++---------- 1 file changed, 5 insertions(+), 10 deletions(-) diff --git a/src/evaluation/evaluation.py b/src/evaluation/evaluation.py index 882069a..cbbd0da 100644 --- a/src/evaluation/evaluation.py +++ b/src/evaluation/evaluation.py @@ -1,3 +1,5 @@ +import torch +import numpy as np from stable_baselines3.common.vec_env import VecEnv from stable_baselines3.common.evaluation import evaluate_policy @@ -36,11 +38,6 @@ class Evaluation: assert len(self._trajectories) >= self.n_eval_episodes - # average velocity of each episode - # this first averages velocity over single trajectories and then averages over trajectories - # avg_velocities = [nanmean(torch.stack(t)[:,2]) for t in self._trajectories] - # avg_velocity = np.mean(avg_velocities) - # velocities produced by generator policy_velocities = torch.cat([torch.stack(t)[:,2] for t in self._trajectories]) # if episodes terminate without collisions, then the state is fully nan @@ -81,12 +78,10 @@ class Evaluation: self._n_collisions += 1 # if last episode is done, start new trajectory + # this is currently not necessary, only if velocity is to be averaged over individual trajectories first + # and then averaging over all trajectories if self._episode_done: self._trajectories.append([]) - agent_state = info['projected_state'][_agent] - self._trajectories[-1].append(agent_state) - # this does not work since agent_action are high level options - # agent_action = local_vars['actions'][venv_i] # normalized intersimple action - # acceleration = env._unnormalize(agent_action) if isinstance(env, NormalizedActionSpace) else agent_action + self._trajectories[-1].append(info['projected_state'][_agent]) self._accelerations.append(info['action_taken'][_agent]) self._episode_done = done From 8703b11deee847d423f2d93fb49810df8b14abe1 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:31:30 +0200 Subject: [PATCH 10/13] import evaluation --- scratch/johannes/gail_options_image.py | 1 + 1 file changed, 1 insertion(+) diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py index ac1676a..03849d5 100644 --- a/scratch/johannes/gail_options_image.py +++ b/scratch/johannes/gail_options_image.py @@ -24,6 +24,7 @@ from tqdm import tqdm from src.policies.options import OptionsCnnPolicy from src.gail.options import OptionsEnv, LLOptions, HLOptions, RenderOptions from src.gail.train import train_discriminator, train_generator +from src.evaluation.evaluation import Evaluation model_name = 'gail_options_image' Env = NRasterizedRandomAgent From c5b043c49ffa1b8b05c4db792f04b09d81fb2720 Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:47:54 +0200 Subject: [PATCH 11/13] Cleanup --- scratch/johannes/gail_options_image.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py index 03849d5..1d84807 100644 --- a/scratch/johannes/gail_options_image.py +++ b/scratch/johannes/gail_options_image.py @@ -5,10 +5,9 @@ from src.discriminator import CnnDiscriminator, CnnDiscriminatorFlatAction from imitation.algorithms import adversarial import stable_baselines3 -from stable_baselines3.common.evaluation import evaluate_policy import torch.utils.data import numpy as np -from intersim.envs.intersimple import Intersimple, NRasterized, NRasterizedInfo, NRasterizedIncrementingAgent, NRasterizedRandomAgent, speed_reward +from intersim.envs.intersimple import NRasterized, NRasterizedRandomAgent, speed_reward import itertools import functools from torch.distributions import Categorical From 070b8fc785f740ca501ca0564e434c974a1bfb6c Mon Sep 17 00:00:00 2001 From: Johannes Fischer Date: Thu, 28 Oct 2021 13:49:14 +0200 Subject: [PATCH 12/13] Cleanup --- scratch/johannes/gail_options_image.py | 47 +++++++------------------- 1 file changed, 12 insertions(+), 35 deletions(-) diff --git a/scratch/johannes/gail_options_image.py b/scratch/johannes/gail_options_image.py index 1d84807..1ea179d 100644 --- a/scratch/johannes/gail_options_image.py +++ b/scratch/johannes/gail_options_image.py @@ -1,13 +1,13 @@ # %% -# import sys -# sys.path.append('../../../') +import sys +sys.path.append('../../../') from src.discriminator import CnnDiscriminator, CnnDiscriminatorFlatAction from imitation.algorithms import adversarial import stable_baselines3 import torch.utils.data import numpy as np -from intersim.envs.intersimple import NRasterized, NRasterizedRandomAgent, speed_reward +from intersim.envs.intersimple import NRasterized, speed_reward import itertools import functools from torch.distributions import Categorical @@ -26,13 +26,12 @@ from src.gail.train import train_discriminator, train_generator from src.evaluation.evaluation import Evaluation model_name = 'gail_options_image' -Env = NRasterizedRandomAgent -env_settings = {'width': 36, 'height': 36, 'm_per_px': 2} +env_settings = {'agent': 51, 'width': 36, 'height': 36, 'm_per_px': 2} -ALL_OPTIONS = [(v,t) for v in [0,2,4,6,8] for t in [5]] # option 0 is safe fallback +ALL_OPTIONS = [(v,t) for v in [0,2,4,6,8] for t in [5, 10]] # option 0 is safe fallback -def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, discount=0.99): - env = Env(**env_settings) +def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=1024, discount=0.99): + env = NRasterized(**env_settings) env.discount = discount tempdir = tempfile.TemporaryDirectory(prefix="quickstart") @@ -40,7 +39,7 @@ def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, disc logger.configure(tempdir_path / "GAIL/") print(f"All Tensorboards and logging are being written inside {tempdir_path}/.") - venv = make_vec_env(Env, n_envs=1, env_kwargs=env_settings) + venv = make_vec_env(NRasterized, n_envs=1, env_kwargs=env_settings) discriminator = adversarial.GAIL( expert_data=expert_data, expert_batch_size=expert_batch_size, @@ -68,41 +67,19 @@ def train(expert_data, epochs=20, expert_batch_size=32, generator_steps=32, disc train_discriminator(LLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=expert_batch_size) train_generator(HLOptions(env, options=ALL_OPTIONS), generator, discriminator, num_samples=generator_steps) - eval_env = Env(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) - ev = Evaluation(eval_env, n_eval_episodes=10) + eval_env = env + ev = Evaluation(eval_env, n_eval_episodes=100) ev.evaluate(epoch, generator, discriminator, expert_data) return generator - - - -# class CAPolicy: -# def __init__(self, a): -# self.a = torch.tensor([a]) - -# def predict(self, obs, state=None, deterministic=False): -# return self.a, state - - # %% if __name__ == '__main__': # %% - with open("scratch/etienne/intersimple/data/NormalizedIntersimpleExpertMu.001N200_NRasterizedInfoAgent51w36h36mppx2.pkl", "rb") as f: + with open("data/NormalizedIntersimpleExpertMu.001_NRasterizedInfoAgent51w36h36mppx2.pkl", "rb") as f: trajectories = pickle.load(f) transitions = rollout.flatten_trajectories(trajectories) - -### - # env = NRasterizedIncrementingAgent(reward=functools.partial(speed_reward, collision_penalty=0.), **env_settings) - # generator = CAPolicy(.5) - - # ev = Evaluation(env, 10) - # ev.evaluate(1, generator, None, transitions) - - # exit() -### - generator = train(transitions) generator.save(model_name) @@ -110,7 +87,7 @@ if __name__ == '__main__': # %% model = stable_baselines3.PPO.load(model_name) - env = RenderOptions(NRasterizedRandomAgent(**env_settings), options=ALL_OPTIONS) + env = RenderOptions(NRasterized(**env_settings), options=ALL_OPTIONS) for s in env.sample_ll(model): if s['dones']: From b1740764e3d1f6c9a5ac9f7f6b47755a65b966c0 Mon Sep 17 00:00:00 2001 From: ebuehrle <43623224+ebuehrle@users.noreply.github.com> Date: Thu, 28 Oct 2021 17:58:34 +0200 Subject: [PATCH 13/13] Parameterize number of discriminator updates per epoch --- src/gail/train.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/gail/train.py b/src/gail/train.py index 31edd46..ee2715c 100644 --- a/src/gail/train.py +++ b/src/gail/train.py @@ -9,10 +9,11 @@ def flatten_transitions(transitions): 'dones': np.stack(list(t['dones'] for t in transitions), axis=0), } -def train_discriminator(env, generator, discriminator, num_samples): +def train_discriminator(env, generator, discriminator, num_samples, n_updates=1): transitions = list(itertools.islice(env.sample_ll(generator), num_samples)) generator_samples = flatten_transitions(transitions) - discriminator.train_disc(gen_samples=generator_samples) + for _ in range(n_updates): + discriminator.train_disc(gen_samples=generator_samples) def train_generator(env, generator, discriminator, num_samples): generator_samples = list(itertools.islice(env.sample_hl(generator, discriminator), num_samples+1))