From ffc85936f1d4026cdb5ac23742e27af1e324793f Mon Sep 17 00:00:00 2001 From: you-n-g Date: Sun, 16 Feb 2025 01:40:44 +0800 Subject: [PATCH] refactor: refactor core framework to better propogate feedbacks (#599) * refactor: Update type annotations and remove unused class in evolving modules * refactor: Simplify evolving agent and feedback handling in CoSTEER module * lint & CI * mypy * ruff for core * mypy * refactor: remove unnecessary comments and update feedback handling logic * refactor: Add prev_task_feedback parameter to evolving strategies * feat: Clear folder before extracting zip file in DockerEnv * fix: Correct retrieval of last experiment from history --- rdagent/components/coder/CoSTEER/__init__.py | 54 ++++++++++++---- .../components/coder/CoSTEER/evaluators.py | 41 +++++++++--- .../coder/CoSTEER/evolving_agent.py | 30 --------- .../coder/CoSTEER/evolving_strategy.py | 42 ++++++++++--- .../coder/data_science/ensemble/__init__.py | 8 ++- .../coder/data_science/feature/__init__.py | 8 ++- .../coder/data_science/model/__init__.py | 8 ++- .../data_science/raw_data_loader/__init__.py | 8 ++- .../coder/data_science/workflow/__init__.py | 8 ++- .../components/coder/factor_coder/__init__.py | 9 +++ .../coder/factor_coder/evolving_strategy.py | 2 + .../coder/model_coder/evolving_strategy.py | 2 + rdagent/core/developer.py | 10 ++- rdagent/core/evaluation.py | 29 +++++---- rdagent/core/evolving_agent.py | 63 +++++++++---------- rdagent/core/evolving_framework.py | 8 +-- rdagent/core/experiment.py | 13 +++- rdagent/core/proposal.py | 6 +- rdagent/oai/backend/deprec.py | 4 +- .../data_science/dev/runner/__init__.py | 11 +++- .../data_science/experiment/experiment.py | 5 ++ .../data_science/proposal/exp_gen.py | 3 +- .../scenarios/qlib/developer/factor_runner.py | 31 +++++---- rdagent/utils/env.py | 6 +- 24 files changed, 261 insertions(+), 148 deletions(-) delete mode 100644 rdagent/components/coder/CoSTEER/evolving_agent.py diff --git a/rdagent/components/coder/CoSTEER/__init__.py b/rdagent/components/coder/CoSTEER/__init__.py index a6c53a65..9c66ba96 100644 --- a/rdagent/components/coder/CoSTEER/__init__.py +++ b/rdagent/components/coder/CoSTEER/__init__.py @@ -2,8 +2,8 @@ import pickle from pathlib import Path from rdagent.components.coder.CoSTEER.config import CoSTEERSettings +from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiFeedback from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem -from rdagent.components.coder.CoSTEER.evolving_agent import FilterFailedRAGEvoAgent from rdagent.components.coder.CoSTEER.knowledge_management import ( CoSTEERKnowledgeBaseV1, CoSTEERKnowledgeBaseV2, @@ -11,8 +11,9 @@ from rdagent.components.coder.CoSTEER.knowledge_management import ( CoSTEERRAGStrategyV2, ) from rdagent.core.developer import Developer -from rdagent.core.evaluation import Evaluator -from rdagent.core.evolving_agent import EvolvingStrategy +from rdagent.core.evaluation import Evaluator, Feedback +from rdagent.core.evolving_agent import EvolvingStrategy, RAGEvoAgent +from rdagent.core.exception import CoderError from rdagent.core.experiment import Experiment from rdagent.log import rdagent_logger as logger @@ -83,9 +84,9 @@ class CoSTEER(Developer[Experiment]): def develop(self, exp: Experiment) -> Experiment: # init intermediate items - experiment = EvolvingItem.from_experiment(exp) + evo_exp = EvolvingItem.from_experiment(exp) - self.evolve_agent = FilterFailedRAGEvoAgent( + self.evolve_agent = RAGEvoAgent( max_loop=self.max_loop, evolving_strategy=self.evolving_strategy, rag=self.rag, @@ -94,16 +95,43 @@ class CoSTEER(Developer[Experiment]): knowledge_self_gen=self.knowledge_self_gen, ) - experiment = self.evolve_agent.multistep_evolve( - experiment, - self.evaluator, - filter_final_evo=self.filter_final_evo, - ) + for evo_exp in self.evolve_agent.multistep_evolve(evo_exp, self.evaluator): + assert isinstance(evo_exp, Experiment) # multiple inheritance + logger.log_object(evo_exp.sub_workspace_list, tag="evolving code") + for sw in evo_exp.sub_workspace_list: + logger.info(f"evolving code workspace: {sw}") + + if self.with_feedback and self.filter_final_evo: + evo_exp = self._exp_postprocess_by_feedback(evo_exp, self.evolve_agent.evolving_trace[-1].feedback) # save new knowledge base if self.new_knowledge_base_path is not None: - pickle.dump(self.knowledge_base, open(self.new_knowledge_base_path, "wb")) + with self.new_knowledge_base_path.open("wb") as f: + pickle.dump(self.knowledge_base, f) logger.info(f"New knowledge base saved to {self.new_knowledge_base_path}") - exp.sub_workspace_list = experiment.sub_workspace_list - exp.experiment_workspace = experiment.experiment_workspace + exp.sub_workspace_list = evo_exp.sub_workspace_list + exp.experiment_workspace = evo_exp.experiment_workspace return exp + + def _exp_postprocess_by_feedback(self, evo: Experiment, feedback: CoSTEERMultiFeedback) -> Experiment: + """ + Responsibility: + - Raise Error if it failed to handle the develop task + - + """ + assert isinstance(evo, Experiment) + assert isinstance(feedback, CoSTEERMultiFeedback) + assert len(evo.sub_workspace_list) == len(feedback) + + # FIXME: when whould the feedback be None? + failed_feedbacks = [ + f"- feedback{index + 1:02d}:\n - execution: {f.execution}\n - return_checking: {f.return_checking}\n - code: {f.code}" + for index, f in enumerate(feedback) + if f is not None and not f.final_decision + ] + + if len(failed_feedbacks) == len(feedback): + feedback_summary = "\n".join(failed_feedbacks) + raise CoderError(f"All tasks are failed:\n{feedback_summary}") + + return evo diff --git a/rdagent/components/coder/CoSTEER/evaluators.py b/rdagent/components/coder/CoSTEER/evaluators.py index 91bbb7a7..612b6895 100644 --- a/rdagent/components/coder/CoSTEER/evaluators.py +++ b/rdagent/components/coder/CoSTEER/evaluators.py @@ -1,6 +1,6 @@ from abc import abstractmethod from dataclasses import dataclass -from typing import List +from typing import TYPE_CHECKING, List from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem from rdagent.core.conf import RD_AGENT_SETTINGS @@ -10,6 +10,9 @@ from rdagent.core.experiment import Task, Workspace from rdagent.core.utils import multiprocessing_wrapper from rdagent.log import rdagent_logger as logger +if TYPE_CHECKING: + from rdagent.core.scenario import Scenario + # TODO: # 1. It seems logically sound, but we currently lack a scenario to apply it. # 2. If it proves to be useful, relocate it to a more general location. @@ -113,14 +116,35 @@ This implementation is {'SUCCESS' if self.final_decision else 'FAIL'}. """ -class CoSTEERMultiFeedback( - Feedback, - List[CoSTEERSingleFeedback], -): +class CoSTEERMultiFeedback(Feedback): """Feedback contains a list, each element is the corresponding feedback for each factor implementation.""" + def __init__(self, feedback_list: List[CoSTEERSingleFeedback]) -> None: + self.feedback_list = feedback_list + + def __getitem__(self, index: int) -> CoSTEERSingleFeedback: + return self.feedback_list[index] + + def __len__(self) -> int: + return len(self.feedback_list) + + def append(self, feedback: CoSTEERSingleFeedback) -> None: + self.feedback_list.append(feedback) + + def __iter__(self): + return iter(self.feedback_list) + + def __bool__(self): + return all(feedback.final_decision for feedback in self.feedback_list) + class CoSTEEREvaluator(Evaluator): + def __init__( + self, + scen: "Scenario", + ) -> None: + self.scen = scen + # TODO: # I think we should have unified interface for all evaluates, for examples. # So we should adjust the interface of other factors @@ -135,7 +159,7 @@ class CoSTEEREvaluator(Evaluator): raise NotImplementedError("Please implement the `evaluator` method") -class CoSTEERMultiEvaluator(Evaluator): +class CoSTEERMultiEvaluator(CoSTEEREvaluator): """This is for evaluation of experiment. Due to we have multiple tasks, so we will return a list of evaluation feebacks""" def __init__(self, single_evaluator: CoSTEEREvaluator, *args, **kwargs) -> None: @@ -164,9 +188,6 @@ class CoSTEERMultiEvaluator(Evaluator): n=RD_AGENT_SETTINGS.multi_proc_n, ) - for index in range(len(evo.sub_tasks)): - evo.sub_workspace_list[index].feedback = multi_implementation_feedback[index] - final_decision = [ None if single_feedback is None else single_feedback.final_decision for single_feedback in multi_implementation_feedback @@ -177,4 +198,4 @@ class CoSTEERMultiEvaluator(Evaluator): if final_decision[index]: evo.sub_tasks[index].factor_implementation = True - return multi_implementation_feedback + return CoSTEERMultiFeedback(multi_implementation_feedback) diff --git a/rdagent/components/coder/CoSTEER/evolving_agent.py b/rdagent/components/coder/CoSTEER/evolving_agent.py deleted file mode 100644 index d4d9e55c..00000000 --- a/rdagent/components/coder/CoSTEER/evolving_agent.py +++ /dev/null @@ -1,30 +0,0 @@ -from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem -from rdagent.core.evolving_agent import RAGEvoAgent -from rdagent.core.evolving_framework import EvolvableSubjects -from rdagent.core.exception import CoderError - - -class FilterFailedRAGEvoAgent(RAGEvoAgent): - - def filter_evolvable_subjects_by_feedback(self, evo: EvolvableSubjects, feedback: list) -> EvolvableSubjects: - assert isinstance(evo, EvolvingItem) - # FIXME: the list does not align with the annotation; It should be MultipleFeedback instead of a list of feedbacks - assert isinstance(feedback, list) - assert len(evo.sub_workspace_list) == len(feedback) - - for index in range(len(evo.sub_workspace_list)): - evo.sub_workspace_list[index].feedback = None - if evo.sub_workspace_list[index] is not None and feedback[index] is not None and not feedback[index]: - evo.sub_workspace_list[index].clear() - - failed_feedbacks = [ - f"- feedback{index + 1:02d}:\n - execution: {f.execution}\n - return_checking: {f.return_checking}\n - code: {f.code}" - for index, f in enumerate(feedback) - if f is not None and not f.final_decision - ] - - if len(failed_feedbacks) == len(feedback): - feedback_summary = "\n".join(failed_feedbacks) - raise CoderError(f"All tasks are failed:\n{feedback_summary}") - - return evo diff --git a/rdagent/components/coder/CoSTEER/evolving_strategy.py b/rdagent/components/coder/CoSTEER/evolving_strategy.py index 131becb6..a7bbdfcf 100644 --- a/rdagent/components/coder/CoSTEER/evolving_strategy.py +++ b/rdagent/components/coder/CoSTEER/evolving_strategy.py @@ -4,13 +4,17 @@ from abc import abstractmethod from pathlib import Path from rdagent.components.coder.CoSTEER.config import CoSTEERSettings +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiFeedback, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem from rdagent.components.coder.CoSTEER.knowledge_management import ( CoSTEERQueriedKnowledge, ) from rdagent.components.coder.CoSTEER.scheduler import random_select from rdagent.core.conf import RD_AGENT_SETTINGS -from rdagent.core.evolving_framework import EvolvingStrategy, QueriedKnowledge +from rdagent.core.evolving_framework import EvolvingStrategy, EvoStep, QueriedKnowledge from rdagent.core.experiment import FBWorkspace, Task from rdagent.core.prompts import Prompts from rdagent.core.scenario import Scenario @@ -28,14 +32,27 @@ class MultiProcessEvolvingStrategy(EvolvingStrategy): def implement_one_task( self, target_task: Task, - queried_knowledge: QueriedKnowledge = None, + queried_knowledge: QueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: # FIXME: fix interface of previous implement """ This method will input the task & current workspace, and output the modification to applied to the workspace. (i.e. replace the content with ) + Parameters + ---------- + target_task : Task + + queried_knowledge : QueriedKnowledge | None + + workspace : FBWorkspace | None + + prev_task_feedback : CoSTEERSingleFeedback | None + task feedback for previous evolving step + None indicate it is the first loop. + Return ------ The new files {: } to update the workspace. @@ -54,10 +71,13 @@ class MultiProcessEvolvingStrategy(EvolvingStrategy): return random_select(to_be_finished_task_index, evo, selected_num, queried_knowledge, scen) @abstractmethod - def assign_code_list_to_evo(self, code_list: list, evo: EvolvingItem) -> None: + def assign_code_list_to_evo(self, code_list: list[dict], evo: EvolvingItem) -> None: """ Assign the code list to the evolving item. + Due to the implement_one_task take `workspace` as input and output the `modification`. + We should apply implmentation to evo + The code list is aligned with the evolving item's sub-tasks. If a task is not implemented, put a None in the list. """ @@ -68,6 +88,7 @@ class MultiProcessEvolvingStrategy(EvolvingStrategy): *, evo: EvolvingItem, queried_knowledge: CoSTEERQueriedKnowledge | None = None, + evolving_trace: list[EvoStep] = [], **kwargs, ) -> EvolvingItem: # 1.找出需要evolve的task @@ -93,11 +114,20 @@ class MultiProcessEvolvingStrategy(EvolvingStrategy): to_be_finished_task_index, evo, self.settings.select_threshold, queried_knowledge, self.scen ) + last_feedback = None + if len(evolving_trace) > 0: + last_feedback = evolving_trace[-1].feedback + assert isinstance(last_feedback, CoSTEERMultiFeedback) result = multiprocessing_wrapper( [ ( self.implement_one_task, - (evo.sub_tasks[target_index], queried_knowledge, evo.experiment_workspace), + ( + evo.sub_tasks[target_index], + queried_knowledge, + evo.experiment_workspace, + None if last_feedback is None else last_feedback[target_index], + ), ) for target_index in to_be_finished_task_index ], @@ -110,8 +140,4 @@ class MultiProcessEvolvingStrategy(EvolvingStrategy): evo = self.assign_code_list_to_evo(code_list, evo) evo.corresponding_selection = to_be_finished_task_index - # After implementation, the feedback should be reset - for workspace in evo.sub_workspace_list: - workspace.feedback = None - return evo diff --git a/rdagent/components/coder/data_science/ensemble/__init__.py b/rdagent/components/coder/data_science/ensemble/__init__.py index 8080e3cf..4f180716 100644 --- a/rdagent/components/coder/data_science/ensemble/__init__.py +++ b/rdagent/components/coder/data_science/ensemble/__init__.py @@ -15,7 +15,10 @@ import json from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -37,6 +40,7 @@ class EnsembleMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: EnsembleTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: # Get task information for knowledge querying ensemble_information_str = target_task.get_task_information() @@ -74,7 +78,7 @@ class EnsembleMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): user_prompt = T(".prompts:ensemble_coder.user").r( ensemble_spec=workspace.file_dict["spec/ensemble.md"], latest_code=workspace.file_dict.get("ensemble.py"), - latest_code_feedback=workspace.feedback, + latest_code_feedback=prev_task_feedback, ) for _ in range(5): diff --git a/rdagent/components/coder/data_science/feature/__init__.py b/rdagent/components/coder/data_science/feature/__init__.py index 98b13056..b3547c5e 100644 --- a/rdagent/components/coder/data_science/feature/__init__.py +++ b/rdagent/components/coder/data_science/feature/__init__.py @@ -2,7 +2,10 @@ import json from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -24,6 +27,7 @@ class FeatureMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: FeatureTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: # return a workspace with "load_data.py", "spec/load_data.md" inside # assign the implemented code to the new workspace. @@ -59,7 +63,7 @@ class FeatureMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): user_prompt = T(".prompts:feature_coder.user").r( feature_spec=workspace.file_dict["spec/feature.md"], latest_code=workspace.file_dict.get("feature.py"), - latest_code_feedback=workspace.feedback, + latest_code_feedback=prev_task_feedback, ) for _ in range(5): diff --git a/rdagent/components/coder/data_science/model/__init__.py b/rdagent/components/coder/data_science/model/__init__.py index feea5d0e..3042d2df 100644 --- a/rdagent/components/coder/data_science/model/__init__.py +++ b/rdagent/components/coder/data_science/model/__init__.py @@ -5,7 +5,10 @@ from jinja2 import Environment, StrictUndefined from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -30,6 +33,7 @@ class ModelMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: ModelTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: model_information_str = target_task.get_task_information() @@ -74,7 +78,7 @@ class ModelMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): latest_model_code=workspace.get_codes( r"^model_(?!test)\w+\.py$" ), # TODO: If we have high failure rate here, we should clean this step with less information. - latest_code_feedback=workspace.feedback, + latest_code_feedback=prev_task_feedback, ) for _ in range(5): diff --git a/rdagent/components/coder/data_science/raw_data_loader/__init__.py b/rdagent/components/coder/data_science/raw_data_loader/__init__.py index 376fba74..c54132a6 100644 --- a/rdagent/components/coder/data_science/raw_data_loader/__init__.py +++ b/rdagent/components/coder/data_science/raw_data_loader/__init__.py @@ -26,7 +26,10 @@ import json from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -51,6 +54,7 @@ class DataLoaderMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: DataLoaderTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: # return a workspace with "load_data.py", "spec/load_data.md" inside # assign the implemented code to the new workspace. @@ -134,7 +138,7 @@ class DataLoaderMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): data_loader_spec=data_loader_spec, folder_spec=data_folder_info, latest_code=workspace.file_dict.get("load_data.py"), - latest_code_feedback=workspace.feedback, + latest_code_feedback=prev_task_feedback, ) for _ in range(5): diff --git a/rdagent/components/coder/data_science/workflow/__init__.py b/rdagent/components/coder/data_science/workflow/__init__.py index f41655b9..53d6ad01 100644 --- a/rdagent/components/coder/data_science/workflow/__init__.py +++ b/rdagent/components/coder/data_science/workflow/__init__.py @@ -2,7 +2,10 @@ import json from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -26,6 +29,7 @@ class WorkflowMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: WorkflowTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: # competition_info = self.scen.competition_descriptions workflow_information_str = target_task.get_task_information() @@ -64,7 +68,7 @@ class WorkflowMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): ensemble_code=workspace.file_dict["ensemble.py"], latest_code=workspace.file_dict.get("main.py"), workflow_spec=workspace.file_dict["spec/workflow.md"], - latest_code_feedback=workspace.feedback, + latest_code_feedback=prev_task_feedback, ) for _ in range(5): diff --git a/rdagent/components/coder/factor_coder/__init__.py b/rdagent/components/coder/factor_coder/__init__.py index be80b121..e90ca970 100644 --- a/rdagent/components/coder/factor_coder/__init__.py +++ b/rdagent/components/coder/factor_coder/__init__.py @@ -5,6 +5,7 @@ from rdagent.components.coder.factor_coder.evaluators import FactorEvaluatorForC from rdagent.components.coder.factor_coder.evolving_strategy import ( FactorMultiProcessEvolvingStrategy, ) +from rdagent.core.experiment import Experiment from rdagent.core.scenario import Scenario @@ -20,3 +21,11 @@ class FactorCoSTEER(CoSTEER): es = FactorMultiProcessEvolvingStrategy(scen=scen, settings=FACTOR_COSTEER_SETTINGS) super().__init__(*args, settings=setting, eva=eva, es=es, evolving_version=2, scen=scen, **kwargs) + + def develop(self, exp: Experiment) -> Experiment: + try: + exp = super().develop(exp) + finally: + es = self.evolve_agent.evolving_trace[-1] + exp.prop_dev_feedback = es.feedback + return exp diff --git a/rdagent/components/coder/factor_coder/evolving_strategy.py b/rdagent/components/coder/factor_coder/evolving_strategy.py index 5ddc938d..37b132e7 100644 --- a/rdagent/components/coder/factor_coder/evolving_strategy.py +++ b/rdagent/components/coder/factor_coder/evolving_strategy.py @@ -5,6 +5,7 @@ from pathlib import Path from jinja2 import Environment, StrictUndefined +from rdagent.components.coder.CoSTEER.evaluators import CoSTEERSingleFeedback from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -74,6 +75,7 @@ class FactorMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: FactorTask, queried_knowledge: CoSTEERQueriedKnowledge, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> str: target_factor_task_information = target_task.get_task_information() diff --git a/rdagent/components/coder/model_coder/evolving_strategy.py b/rdagent/components/coder/model_coder/evolving_strategy.py index 83b7afa3..aa22ac7e 100644 --- a/rdagent/components/coder/model_coder/evolving_strategy.py +++ b/rdagent/components/coder/model_coder/evolving_strategy.py @@ -4,6 +4,7 @@ from pathlib import Path from jinja2 import Environment, StrictUndefined from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS +from rdagent.components.coder.CoSTEER.evaluators import CoSTEERSingleFeedback from rdagent.components.coder.CoSTEER.evolving_strategy import ( MultiProcessEvolvingStrategy, ) @@ -30,6 +31,7 @@ class ModelMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: ModelTask, queried_knowledge: CoSTEERQueriedKnowledge = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> str: model_information_str = target_task.get_task_information() diff --git a/rdagent/core/developer.py b/rdagent/core/developer.py index 2ec4bc6b..f87eb7a7 100644 --- a/rdagent/core/developer.py +++ b/rdagent/core/developer.py @@ -14,13 +14,21 @@ class Developer(ABC, Generic[ASpecificExp]): self.scen: Scenario = scen @abstractmethod - def develop(self, exp: ASpecificExp) -> ASpecificExp: + def develop(self, exp: ASpecificExp) -> ASpecificExp: # TODO: remove return value """ Task Generator should take in an experiment. Because the schedule of different tasks is crucial for the final performance due to it affects the learning process. + Current constraints: + - The developer should **inplace** edit the exp instead of returning value; + - because we have a lot of use cases to raise errors, but we need the intermediate results in exp. + - So we should remove the return value in the future. + + Responsibilities: + - Generate a new experiment after developing on it. + - If it tries to deliver message for future development, it should set a ExperimentFeedback """ error_message = "generate method is not implemented." raise NotImplementedError(error_message) diff --git a/rdagent/core/evaluation.py b/rdagent/core/evaluation.py index 4536e9cc..e49ceca7 100644 --- a/rdagent/core/evaluation.py +++ b/rdagent/core/evaluation.py @@ -1,9 +1,8 @@ -import typing -from abc import ABC, abstractmethod +""" +It is expected to be shared among different frameworks. +""" -if typing.TYPE_CHECKING: - from rdagent.core.experiment import Task, Workspace - from rdagent.core.scenario import Scenario +from abc import ABC, abstractmethod class Feedback: @@ -17,6 +16,15 @@ class Feedback: return True +class EvaluableObj: + """ + A set of information that is evaluable. Following things can be included. + - Task + - Solution + - Ground Truth + """ + + class Evaluator(ABC): """ Design Principle: @@ -27,18 +35,9 @@ class Evaluator(ABC): 2. advanced/summarized feedback information. (evaluate will handle this) """ - def __init__( - self, - scen: "Scenario", - ) -> None: - self.scen = scen - @abstractmethod def evaluate( self, - target_task: "Task", - implementation: "Workspace", - gt_implementation: "Workspace", - **kwargs: object, + eo: EvaluableObj, ) -> Feedback: raise NotImplementedError diff --git a/rdagent/core/evolving_agent.py b/rdagent/core/evolving_agent.py index 4d7e8e5c..3012d563 100644 --- a/rdagent/core/evolving_agent.py +++ b/rdagent/core/evolving_agent.py @@ -1,20 +1,23 @@ from __future__ import annotations from abc import ABC, abstractmethod -from typing import TYPE_CHECKING, Any +from collections.abc import Generator +from typing import TYPE_CHECKING, Any, Generic, TypeVar from tqdm import tqdm if TYPE_CHECKING: - from rdagent.core.evaluation import Evaluator from rdagent.core.evolving_framework import EvolvableSubjects -from rdagent.core.evaluation import Feedback +from rdagent.core.evaluation import EvaluableObj, Evaluator, Feedback from rdagent.core.evolving_framework import EvolvingStrategy, EvoStep from rdagent.log import rdagent_logger as logger +ASpecificEvaluator = TypeVar("ASpecificEvaluator", bound=Evaluator) + + +class EvoAgent(ABC, Generic[ASpecificEvaluator]): -class EvoAgent(ABC): def __init__(self, max_loop: int, evolving_strategy: EvolvingStrategy) -> None: self.max_loop = max_loop self.evolving_strategy = evolving_strategy @@ -23,24 +26,32 @@ class EvoAgent(ABC): def multistep_evolve( self, evo: EvolvableSubjects, - eva: Evaluator | Feedback, - filter_final_evo: bool = False, - ) -> EvolvableSubjects: ... + eva: ASpecificEvaluator | Feedback, + ) -> Generator[EvolvableSubjects, None, None]: + """ + yield EvolvableSubjects for caller for easier process control and logging. + """ + + +class RAGEvaluator(Evaluator): @abstractmethod - def filter_evolvable_subjects_by_feedback( + def evaluate( self, - evo: EvolvableSubjects, - feedback: Feedback | list[Feedback] | None, - ) -> EvolvableSubjects: ... + eo: EvaluableObj, + queried_knowledge: object = None, + ) -> Feedback: + raise NotImplementedError -class RAGEvoAgent(EvoAgent): +class RAGEvoAgent(EvoAgent[RAGEvaluator]): + def __init__( self, max_loop: int, evolving_strategy: EvolvingStrategy, rag: Any, + *, with_knowledge: bool = False, with_feedback: bool = True, knowledge_self_gen: bool = False, @@ -55,9 +66,8 @@ class RAGEvoAgent(EvoAgent): def multistep_evolve( self, evo: EvolvableSubjects, - eva: Evaluator | Feedback, - filter_final_evo: bool = False, - ) -> EvolvableSubjects: + eva: RAGEvaluator | Feedback, + ) -> Generator[EvolvableSubjects, None, None]: for evo_loop_id in tqdm(range(self.max_loop), "Implementing"): with logger.tag(f"evo_loop_{evo_loop_id}"): # 1. knowledge self-evolving @@ -75,10 +85,7 @@ class RAGEvoAgent(EvoAgent): evolving_trace=self.evolving_trace, queried_knowledge=queried_knowledge, ) - # TODO: Due to design issues, we have chosen to ignore this mypy error. - logger.log_object(evo.sub_workspace_list, tag="evolving code") # type: ignore[attr-defined] - for sw in evo.sub_workspace_list: # type: ignore[attr-defined] - logger.info(f"evolving code workspace: {sw}") + yield evo # yield the control to caller for process control and logging. # 4. Pack evolve results es = EvoStep(evo, queried_knowledge) @@ -86,11 +93,7 @@ class RAGEvoAgent(EvoAgent): # 5. Evaluation if self.with_feedback: es.feedback = ( - # TODO: Due to the irregular design of rdagent.core.evaluation.Evaluator, - # it fails mypy's test here, so we'll ignore this error for now. - eva - if isinstance(eva, Feedback) - else eva.evaluate(evo, queried_knowledge=queried_knowledge) # type: ignore[arg-type, call-arg] + eva if isinstance(eva, Feedback) else eva.evaluate(evo, queried_knowledge=queried_knowledge) ) logger.log_object(es.feedback, tag="evolving feedback") @@ -98,12 +101,6 @@ class RAGEvoAgent(EvoAgent): self.evolving_trace.append(es) # 7. check if all tasks are completed - if self.with_feedback: - all_completed = all(es.feedback) if isinstance(es.feedback, list) else es.feedback - if all_completed: - logger.info("All tasks in evolving subject have been completed.") - break - - if self.with_feedback and filter_final_evo: - evo = self.filter_evolvable_subjects_by_feedback(evo, self.evolving_trace[-1].feedback) - return evo + if self.with_feedback and es.feedback: + logger.info("All tasks in evolving subject have been completed.") + break diff --git a/rdagent/core/evolving_framework.py b/rdagent/core/evolving_framework.py index 9ff87448..e8917bd2 100644 --- a/rdagent/core/evolving_framework.py +++ b/rdagent/core/evolving_framework.py @@ -5,6 +5,7 @@ from abc import ABC, abstractmethod from dataclasses import dataclass from typing import TYPE_CHECKING, Any +from rdagent.core.evaluation import EvaluableObj from rdagent.core.knowledge_base import KnowledgeBase if TYPE_CHECKING: @@ -28,16 +29,13 @@ class EvolvingKnowledgeBase(KnowledgeBase): raise NotImplementedError -class EvolvableSubjects: +class EvolvableSubjects(EvaluableObj): """The target object to be evolved""" def clone(self) -> EvolvableSubjects: return copy.deepcopy(self) -class QlibEvolvableSubjects(EvolvableSubjects): ... - - @dataclass class EvoStep: """At a specific step, @@ -52,7 +50,7 @@ class EvoStep: evolvable_subjects: EvolvableSubjects queried_knowledge: QueriedKnowledge | None = None - feedback: Feedback | list[Feedback] | None = None + feedback: Feedback | None = None class EvolvingStrategy(ABC): diff --git a/rdagent/core/experiment.py b/rdagent/core/experiment.py index 680e712e..96f08352 100644 --- a/rdagent/core/experiment.py +++ b/rdagent/core/experiment.py @@ -18,7 +18,7 @@ from rdagent.utils import filter_progress_bar from rdagent.utils.fmt import shrink_text if typing.TYPE_CHECKING: - from rdagent.core.proposal import Hypothesis + from rdagent.core.proposal import ExperimentFeedback, Hypothesis from rdagent.utils.env import Env """ @@ -280,6 +280,16 @@ class Experiment( # If we implement the whole workflow, we don't have to use it, then we remove it. self.based_experiments: Sequence[ASpecificWSForExperiment] = based_experiments + self.experiment_workspace: ASpecificWSForExperiment | None = None + + # The experiment may be developed by different developers. + # Last feedback is used to propagate info to the next developer. + # Life cycle: + # - Developer assigns feedback for next component; + # - Workflow control clears feedback. + self.prop_dev_feedback: Feedback | None = None + + # TODO: (xiao) I think this is too concrete; we should move it into # NOTE: Assumption # - only runner will assign this variable # - We will always create a new Experiment without copying previous results when we goto the next new loop. @@ -287,7 +297,6 @@ class Experiment( self.sub_results: dict[str, float] = ( {} ) # TODO: in Kaggle, now sub results are all saved in self.result, remove this in the future. - self.experiment_workspace: ASpecificWSForExperiment | None = None ASpecificExp = TypeVar("ASpecificExp", bound=Experiment) diff --git a/rdagent/core/proposal.py b/rdagent/core/proposal.py index 2c341640..5aba6e66 100644 --- a/rdagent/core/proposal.py +++ b/rdagent/core/proposal.py @@ -58,8 +58,9 @@ class Hypothesis: class ExperimentFeedback(Feedback): def __init__( self, - decision: bool, reason: str, + *, + decision: bool, exception: Exception | None = None, ) -> None: self.decision = decision @@ -91,9 +92,10 @@ class HypothesisFeedback(ExperimentFeedback): hypothesis_evaluation: str, new_hypothesis: str, reason: str, + *, decision: bool, ) -> None: - super().__init__(decision, reason) + super().__init__(reason, decision=decision) self.observations = observations self.hypothesis_evaluation = hypothesis_evaluation self.new_hypothesis = new_hypothesis diff --git a/rdagent/oai/backend/deprec.py b/rdagent/oai/backend/deprec.py index 501c7add..e5601d49 100644 --- a/rdagent/oai/backend/deprec.py +++ b/rdagent/oai/backend/deprec.py @@ -12,7 +12,7 @@ import urllib.request import uuid from copy import deepcopy from pathlib import Path -from typing import Any, Optional +from typing import Any, Optional, cast import numpy as np import openai @@ -161,7 +161,7 @@ class SQliteLazyCache(SingletonBaseClass): def message_get(self, conversation_id: str) -> list[dict[str, Any]]: self.c.execute("SELECT message FROM message_cache WHERE conversation_id=?", (conversation_id,)) result = self.c.fetchone() - return [] if result is None else json.loads(result[0]) + return [] if result is None else cast(list[dict[str, Any]], json.loads(result[0])) def message_set(self, conversation_id: str, message_value: list[dict[str, Any]]) -> None: self.c.execute( diff --git a/rdagent/scenarios/data_science/dev/runner/__init__.py b/rdagent/scenarios/data_science/dev/runner/__init__.py index e83a1f4a..0feef719 100644 --- a/rdagent/scenarios/data_science/dev/runner/__init__.py +++ b/rdagent/scenarios/data_science/dev/runner/__init__.py @@ -6,7 +6,10 @@ from rdagent.app.data_science.conf import DS_RD_SETTING from rdagent.components.coder import CoSTEER from rdagent.components.coder.CoSTEER import CoSTEER from rdagent.components.coder.CoSTEER.config import CoSTEER_SETTINGS -from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiEvaluator +from rdagent.components.coder.CoSTEER.evaluators import ( + CoSTEERMultiEvaluator, + CoSTEERSingleFeedback, +) from rdagent.components.coder.CoSTEER.evolvable_subjects import FBWorkspace from rdagent.components.coder.CoSTEER.evolving_strategy import ( CoSTEERQueriedKnowledge, @@ -29,8 +32,10 @@ class DSRunnerMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): target_task: CoSTEERTask, queried_knowledge: CoSTEERQueriedKnowledge | None = None, workspace: FBWorkspace | None = None, + prev_task_feedback: CoSTEERSingleFeedback | None = None, ) -> dict[str, str]: - if workspace.feedback is None: + if prev_task_feedback is None: + # if no prev_tak_feedback, it is the first loop; we do not make any changes and goto evaluators directly. return {} task_information_str = target_task.get_task_information() @@ -41,7 +46,7 @@ class DSRunnerMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy): ) user_prompt = T(".prompts:DSCoSTEER_debugger.user").r( code=workspace.all_codes, - feedback=workspace.feedback, + feedback=prev_task_feedback, ) batch_edit = BatchEditOut.extract_output( diff --git a/rdagent/scenarios/data_science/experiment/experiment.py b/rdagent/scenarios/data_science/experiment/experiment.py index 4f19f390..68fe18e6 100644 --- a/rdagent/scenarios/data_science/experiment/experiment.py +++ b/rdagent/scenarios/data_science/experiment/experiment.py @@ -11,6 +11,11 @@ COMPONENT = Literal["DataLoadSpec", "FeatureEng", "Model", "Ensemble", "Workflow class DSExperiment(Experiment[Task, FBWorkspace, FBWorkspace]): def __init__(self, pending_tasks_list: list, *args, **kwargs) -> None: super().__init__(sub_tasks=[], *args, **kwargs) + # Status + # - Initial: blank; + # - Injecting from SOTA code; + # - New version no matter successful or not + # the initial workspace or the successful new version after coding self.experiment_workspace = FBWorkspace() self.pending_tasks_list = pending_tasks_list self.format_check_result = None diff --git a/rdagent/scenarios/data_science/proposal/exp_gen.py b/rdagent/scenarios/data_science/proposal/exp_gen.py index c6f23b03..317b15e9 100644 --- a/rdagent/scenarios/data_science/proposal/exp_gen.py +++ b/rdagent/scenarios/data_science/proposal/exp_gen.py @@ -274,8 +274,7 @@ class DSExpGen(ExpGen): # - Extra RAG sota_exp = trace.sota_experiment() assert sota_exp is not None, "SOTA experiment is not provided." - exp_and_feedback = trace.last_runnable_exp_fb() - assert exp_and_feedback is not None, "Last runnable experiment is not provided." + exp_and_feedback = trace.hist[-1] last_exp = exp_and_feedback[0] # Step 1: Generate component diff --git a/rdagent/scenarios/qlib/developer/factor_runner.py b/rdagent/scenarios/qlib/developer/factor_runner.py index d307a0ee..ade70b1c 100644 --- a/rdagent/scenarios/qlib/developer/factor_runner.py +++ b/rdagent/scenarios/qlib/developer/factor_runner.py @@ -5,6 +5,7 @@ from typing import List import pandas as pd from pandarallel import pandarallel +from rdagent.components.coder.CoSTEER.evaluators import CoSTEERMultiFeedback from rdagent.core.conf import RD_AGENT_SETTINGS from rdagent.core.utils import cache_with_pickle, multiprocessing_wrapper @@ -133,17 +134,25 @@ class QlibFactorRunner(CachedRunner[QlibFactorExperiment]): # Collect all exp's dataframes for exp in exp_or_list: - # Iterate over sub-implementations and execute them to get each factor data - message_and_df_list = multiprocessing_wrapper( - [(implementation.execute, ("All",)) for implementation in exp.sub_workspace_list if implementation], - n=RD_AGENT_SETTINGS.multi_proc_n, - ) - for message, df in message_and_df_list: - # Check if factor generation was successful - if df is not None and "datetime" in df.index.names: - time_diff = df.index.get_level_values("datetime").to_series().diff().dropna().unique() - if pd.Timedelta(minutes=1) not in time_diff: - factor_dfs.append(df) + if len(exp.sub_tasks) > 0: + # if it has no sub_tasks, the experiment is results from template project. + # otherwise, it is developed with designed task. So it should have feedback. + assert isinstance(exp.prop_dev_feedback, CoSTEERMultiFeedback) + # Iterate over sub-implementations and execute them to get each factor data + message_and_df_list = multiprocessing_wrapper( + [ + (implementation.execute, ("All",)) + for implementation, fb in zip(exp.sub_workspace_list, exp.prop_dev_feedback) + if implementation and fb + ], # only execute successfully feedback + n=RD_AGENT_SETTINGS.multi_proc_n, + ) + for message, df in message_and_df_list: + # Check if factor generation was successful + if df is not None and "datetime" in df.index.names: + time_diff = df.index.get_level_values("datetime").to_series().diff().dropna().unique() + if pd.Timedelta(minutes=1) not in time_diff: + factor_dfs.append(df) # Combine all successful factor data if factor_dfs: diff --git a/rdagent/utils/env.py b/rdagent/utils/env.py index b885096a..7f3741cb 100644 --- a/rdagent/utils/env.py +++ b/rdagent/utils/env.py @@ -417,7 +417,11 @@ class DockerEnv(Env[DockerConf]): """ Unzip a file into a folder, use zipfile instead of subprocess """ - shutil.rmtree(folder_path, ignore_errors=True) + # Clear folder_path before extracting + if os.path.exists(folder_path): + shutil.rmtree(folder_path) + os.makedirs(folder_path) + with zipfile.ZipFile(zip_file_path, "r") as z: z.extractall(folder_path)