diff --git a/cookbook/appworld/appworld_react_agent.py b/cookbook/appworld/appworld_react_agent.py index 635c9b71..b5ad4f61 100644 --- a/cookbook/appworld/appworld_react_agent.py +++ b/cookbook/appworld/appworld_react_agent.py @@ -35,6 +35,7 @@ class AppworldReactAgent: temperature: float = 0.9, max_interactions: int = 30, max_response_size: int = 2000, + num_runs: int = 1, use_experience: bool = False): self.index: int = index @@ -44,6 +45,7 @@ class AppworldReactAgent: self.temperature: float = temperature self.max_interactions: int = max_interactions self.max_response_size: int = max_response_size + self.num_runs: int = num_runs self.use_experience: bool = use_experience self.llm_client = OpenAI() @@ -104,39 +106,42 @@ class AppworldReactAgent: def execute(self): result = [] for task_index, task_id in enumerate(tqdm(self.task_ids, desc=f"ray_index={self.index}")): - with AppWorld(task_id=task_id, experiment_name=self.experiment_name) as world: - history = self.prompt_messages(world=world) - before_score = self.get_reward(world) - logger.info(f"ray_id={self.index} task_index={task_index} instruction={world.task.instruction} " - f"before_score={before_score:.4f}") + # Run each task num_runs times + for run_id in range(self.num_runs): + with AppWorld(task_id=task_id, experiment_name=f"{self.experiment_name}_run_{run_id}") as world: + history = self.prompt_messages(world=world) + before_score = self.get_reward(world) + logger.info(f"ray_id={self.index} task_index={task_index} run_id={run_id} " + f"instruction={world.task.instruction} before_score={before_score:.4f}") - for i in range(self.max_interactions): - code = self.call_llm(history) - history.append({"role": "assistant", "content": code}) + for i in range(self.max_interactions): + code = self.call_llm(history) + history.append({"role": "assistant", "content": code}) - output = world.execute(code) - if len(output) > self.max_response_size: - logger.warning(f"output exceed max size={len(output)}") - output = output[:self.max_response_size] - history.append({"role": "user", "content": output}) + output = world.execute(code) + if len(output) > self.max_response_size: + logger.warning(f"output exceed max size={len(output)}") + output = output[:self.max_response_size] + history.append({"role": "user", "content": output}) - logger.info(f"ray_id={self.index} task_index={task_index} step={i} complete~") + logger.info(f"ray_id={self.index} task_index={task_index} run_id={run_id} step={i} complete~") - if world.task_completed(): - break + if world.task_completed(): + break - after_score = self.get_reward(world) - uplift_score = after_score - before_score - t_result = { - "task_id": world.task_id, - "experiment_name": self.experiment_name, - "task_completed": world.task_completed(), - "before_score": before_score, - "after_score": after_score, - "uplift_score": uplift_score, - "task_history": history, - } - result.append(t_result) + after_score = self.get_reward(world) + uplift_score = after_score - before_score + t_result = { + "task_id": world.task_id, + "run_id": run_id, # Add run_id field + "experiment_name": self.experiment_name, + "task_completed": world.task_completed(), + "before_score": before_score, + "after_score": after_score, + "uplift_score": uplift_score, + "task_history": history, + } + result.append(t_result) return result @@ -164,10 +169,10 @@ class AppworldReactAgent: def main(): dataset_name = "train" task_ids = load_task_ids(dataset_name) - agent = AppworldReactAgent(index=0, task_ids=task_ids[0:1], experiment_name=f"jinli_{dataset_name}") + agent = AppworldReactAgent(index=0, task_ids=task_ids[0:1], experiment_name=dataset_name, num_runs=4) result = agent.execute() logger.info(f"result={json.dumps(result)}") if __name__ == "__main__": - main() + main() \ No newline at end of file diff --git a/cookbook/appworld/run_appworld.py b/cookbook/appworld/run_appworld.py index 6b07579d..afd4021f 100644 --- a/cookbook/appworld/run_appworld.py +++ b/cookbook/appworld/run_appworld.py @@ -17,7 +17,7 @@ from appworld import load_task_ids from appworld_react_agent import AppworldReactAgent -def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_experience: bool=False): +def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, num_runs: int = 1, use_experience: bool = False): experiment_name = dataset_name + "_" + experiment_suffix path: Path = Path(f"./exp_result") path.mkdir(parents=True, exist_ok=True) @@ -33,9 +33,12 @@ def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_e if max_workers > 1: future_list: list = [] for i in range(max_workers): + # Assign tasks to each worker, ensuring each task runs num_runs times + worker_task_ids = task_ids[i::max_workers] actor = AppworldReactAgent.remote(index=i, - task_ids=task_ids[i::max_workers], + task_ids=worker_task_ids, experiment_name=experiment_name, + num_runs=num_runs, use_experience=use_experience) future = actor.execute.remote() future_list.append(future) @@ -50,23 +53,32 @@ def run_agent(dataset_name: str, experiment_suffix: str, max_workers: int, use_e else: result.append(t_result) - logger.info(f"{i + 1}/{len(task_ids)} complete") + logger.info(f"worker {i + 1}/{max_workers} complete") dump_file() else: for index, task_id in enumerate(task_ids): - agent = AppworldReactAgent(index=index, task_ids=[task_id], experiment_name=experiment_name, use_experience=use_experience) - result.append(agent.execute()) + agent = AppworldReactAgent(index=index, + task_ids=[task_id], + experiment_name=experiment_name, + num_runs=num_runs, + use_experience=use_experience) + task_results = agent.execute() + if isinstance(task_results, list): + result.extend(task_results) + else: + result.append(task_results) dump_file() def main(): - max_workers = 2 + max_workers = 4 + num_runs = 4 # Run each task 4 times if max_workers > 1: - ray.init(num_cpus=max_workers) - # run_agent(dataset_name="train", experiment_suffix="v2", max_workers=max_workers) - run_agent(dataset_name="dev", experiment_suffix="add-exp", max_workers=max_workers, use_experience=True) + ray.init(num_cpus=4) + # run_agent(dataset_name="train", experiment_suffix="v2", max_workers=max_workers, num_runs=num_runs) + run_agent(dataset_name="dev", experiment_suffix="add-exp", max_workers=max_workers, num_runs=num_runs, use_experience=True) if __name__ == "__main__": - main() + main() \ No newline at end of file diff --git a/cookbook/appworld/run_exp_statistic.py b/cookbook/appworld/run_exp_statistic.py index 1cafe731..e00be7b9 100644 --- a/cookbook/appworld/run_exp_statistic.py +++ b/cookbook/appworld/run_exp_statistic.py @@ -1,45 +1,138 @@ import json from pathlib import Path +from collections import defaultdict +import pandas as pd from loguru import logger +def calculate_best_at_k(scores: list, k: int) -> float: + """ + Calculate best@k + Divide scores into groups of size k, take the maximum value in each group, + then average these maximum values + + Args: + scores: List of after_score values for all runs of a task + k: Group size + + Returns: + best@k value + """ + if len(scores) % k != 0: + raise ValueError(f"Length of scores ({len(scores)}) must be divisible by k ({k})") + + group_maxs = [] + for i in range(0, len(scores), k): + group = scores[i:i + k] + group_maxs.append(max(group)) + + return sum(group_maxs) / len(group_maxs) + + +def get_possible_k_values(total_runs: int) -> list: + """ + Get all possible k values (factors of total_runs) + + Args: + total_runs: Total number of runs + + Returns: + List of k values in descending order + """ + k_values = [] + for k in range(1, total_runs + 1): + if total_runs % k == 0: + k_values.append(k) + return sorted(k_values, reverse=True) # Sort from large to small + + def run_exp_statistic(): path: Path = Path(f"./exp_result") + # Store results for all experiments + all_results = {} + for file in path.glob("*.jsonl"): + # Group results by task_id + task_results = defaultdict(list) + with open(file, "r") as f: - task_completed_list = [] - before_score_list = [] - after_score_list = [] - task_success_list = [] for line in f: if not line.strip(): continue data = json.loads(line) + if isinstance(data, list): for part_data in data: - task_completed_list.append(1 if part_data["task_completed"] is True else 0) - before_score_list.append(part_data["before_score"]) - after_score_list.append(part_data["after_score"]) - task_success_list.append(part_data["after_score"] > 0.9) - + task_id = part_data["task_id"] + after_score = part_data["after_score"] + task_results[task_id].append(after_score) else: - task_completed_list.append(1 if data["task_completed"] is True else 0) - before_score_list.append(data["before_score"]) - after_score_list.append(data["after_score"]) - task_success_list.append(data["after_score"] > 0.9) + task_id = data["task_id"] + after_score = data["after_score"] + task_results[task_id].append(after_score) - task_completed_ratio = sum(task_completed_list) / len(task_completed_list) - before_score_ratio = sum(before_score_list) / len(before_score_list) - after_score_ratio = sum(after_score_list) / len(after_score_list) - task_success_ratio = sum(task_success_list) / len(task_success_list) - logger.info(f"file={file} " - f"task_completed_ratio={task_completed_ratio:.2f} " - f"before_score_ratio={before_score_ratio:.2f} " - f"after_score_ratio={after_score_ratio:.2f} " - f"task_success_ratio={task_success_ratio:.2f}") + if not task_results: + logger.warning(f"No valid data found in file {file}") + continue + + # Check if each task has consistent number of runs + run_counts = [len(scores) for scores in task_results.values()] + if len(set(run_counts)) > 1: + logger.warning(f"Inconsistent number of runs for different tasks in file {file}: {set(run_counts)}") + continue + + num_runs = run_counts[0] + logger.info(f"File {file}: {len(task_results)} tasks, {num_runs} runs per task") + + # Get all possible k values + k_values = get_possible_k_values(num_runs) + logger.info(f"Calculable best@k values: {k_values}") + + # Calculate various best@k values + file_results = {"file": file.name} + + for k in k_values: + best_at_k_scores = [] + for task_id, scores in task_results.items(): + try: + best_k_score = calculate_best_at_k(scores, k) + best_at_k_scores.append(best_k_score) + except ValueError as e: + logger.error(f"Error calculating best@{k} for task {task_id}: {e}") + continue + + if best_at_k_scores: + avg_best_at_k = sum(best_at_k_scores) / len(best_at_k_scores) + file_results[f"best@{k}"] = avg_best_at_k + logger.info(f"file={file.name} best@{k}={avg_best_at_k:.4f}") + + all_results[file.name] = file_results + + # Create and display table + if all_results: + df = pd.DataFrame(list(all_results.values())) + df = df.set_index('file') + + # Sort columns by the number in column name (best@8, best@4, best@2, best@1) + best_columns = [col for col in df.columns if col.startswith('best@')] + best_columns.sort(key=lambda x: int(x.split('@')[1]), reverse=True) + df = df[best_columns] + + print("\n" + "=" * 80) + print("Experiment Results Summary Table") + print("=" * 80) + print(df.round(4)) + print("=" * 80) + + # Save table to CSV + output_path = path / "experiment_summary.csv" + df.to_csv(output_path) + logger.info(f"Results table saved to: {output_path}") + else: + logger.warning("No valid experiment results found") if __name__ == "__main__": - run_exp_statistic() + run_exp_statistic() \ No newline at end of file