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
This commit is contained in:
you-n-g
2025-02-16 01:40:44 +08:00
committed by GitHub
parent 241c05ccac
commit ffc85936f1
24 changed files with 261 additions and 148 deletions
+41 -13
View File
@@ -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
+31 -10
View File
@@ -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)
@@ -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
@@ -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 <filename> with <content>)
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 {<filename>: <content>} 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
@@ -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):
@@ -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):
@@ -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):
@@ -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):
@@ -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):
@@ -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
@@ -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()
@@ -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()
+9 -1
View File
@@ -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)
+14 -15
View File
@@ -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
+30 -33
View File
@@ -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
+3 -5
View File
@@ -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):
+11 -2
View File
@@ -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)
+4 -2
View File
@@ -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
+2 -2
View File
@@ -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(
@@ -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(
@@ -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
@@ -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
@@ -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:
+5 -1
View File
@@ -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)