add feedback to workspace and ds runner base on costeer

This commit is contained in:
Xu Yang
2025-02-11 20:50:19 +08:00
committed by GitHub
parent f96c710c23
commit 203967ae9a
25 changed files with 383 additions and 231 deletions
-35
View File
@@ -10,41 +10,6 @@ class DataScienceBasePropSetting(KaggleBasePropSetting):
scen: str = "rdagent.scenarios.data_science.scen.KaggleScen"
"""Scenario class for data mining model"""
## proposal
exp_gen: str = "rdagent.scenarios.data_science.proposal.exp_gen.DSExpGen"
# exp_gen_init_kwargs: dict = {"max_trace_hist": 3} # TODO: to be configurable
# the two below should be used in ExpGen
# hypothesis_gen: str = "rdagent.scenarios.kaggle.proposal.proposal.KGHypothesisGen"
# """Hypothesis generation class"""
#
# hypothesis2experiment: str = "rdagent.scenarios.kaggle.proposal.proposal.KGHypothesis2Experiment"
# """Hypothesis to experiment class"""
## dev/coder
data_loader_coder: str = "rdagent.components.coder.data_science.raw_data_loader.DataLoaderCoSTEER"
"""Data Loader CoSTEER"""
# feature_coder: str = "rdagent.scenarios.kaggle.developer.coder.KGFactorCoSTEER"
# """Feature Coder class"""
# model_feature_selection_coder: str = "rdagent.scenarios.kaggle.developer.coder.KGModelFeatureSelectionCoder"
# """Model Feature Selection Coder class"""
# model_coder: str = "rdagent.scenarios.kaggle.developer.coder.KGModelCoSTEER"
# """Model Coder class"""
## dev/runner
feature_runner: str = "rdagent.scenarios.kaggle.developer.runner.KGFactorRunner"
"""Feature Runner class"""
model_runner: str = "rdagent.scenarios.kaggle.developer.runner.KGModelRunner"
"""Model Runner class"""
## feedback
summarizer: str = "rdagent.scenarios.kaggle.developer.feedback.KGExperiment2Feedback"
"""Summarizer class"""
## Workflow Related
consecutive_errors: int = 5
+3 -3
View File
@@ -12,12 +12,12 @@ from rdagent.components.coder.data_science.workflow import WorkflowCoSTEER
from rdagent.components.workflow.conf import BasePropSetting
from rdagent.components.workflow.rd_loop import RDLoop
from rdagent.core.exception import CoderError, RunnerError
from rdagent.core.proposal import ExperimentFeedback, HypothesisFeedback
from rdagent.core.proposal import ExperimentFeedback
from rdagent.core.scenario import Scenario
from rdagent.core.utils import import_class
from rdagent.log import rdagent_logger as logger
from rdagent.scenarios.data_science.dev.feedback import DSExperiment2Feedback
from rdagent.scenarios.data_science.dev.runner import DSRunner
from rdagent.scenarios.data_science.dev.runner import DSCoSTEERRunner
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
from rdagent.scenarios.data_science.proposal.exp_gen import DSExpGen, DSTrace
from rdagent.scenarios.kaggle.kaggle_crawler import download_data
@@ -49,7 +49,7 @@ class DataScienceRDLoop(RDLoop):
self.ensemble_coder = EnsembleCoSTEER(scen)
self.workflow_coder = WorkflowCoSTEER(scen)
self.runner = DSRunner(scen)
self.runner = DSCoSTEERRunner(scen)
# self.summarizer: Experiment2Feedback = import_class(PROP_SETTING.summarizer)(scen)
# logger.log_object(self.summarizer, tag="summarizer")
@@ -6,8 +6,7 @@ from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem
from rdagent.core.conf import RD_AGENT_SETTINGS
from rdagent.core.evaluation import Evaluator, Feedback
from rdagent.core.evolving_framework import QueriedKnowledge
from rdagent.core.experiment import Workspace
from rdagent.core.scenario import Task
from rdagent.core.experiment import Task, Workspace
from rdagent.core.utils import multiprocessing_wrapper
from rdagent.log import rdagent_logger as logger
@@ -165,6 +164,9 @@ 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
@@ -1,6 +1,5 @@
from rdagent.core.evolving_framework import EvolvableSubjects
from rdagent.core.experiment import Experiment, FBWorkspace
from rdagent.core.scenario import Task
from rdagent.core.experiment import Experiment, FBWorkspace, Task
from rdagent.log import rdagent_logger as logger
@@ -1,4 +1,3 @@
from rdagent.components.coder.CoSTEER.evaluators import CoSTEERSingleFeedbackDeprecated
from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem
from rdagent.core.evolving_agent import RAGEvoAgent
from rdagent.core.evolving_framework import EvolvableSubjects
@@ -7,15 +6,14 @@ from rdagent.core.exception import CoderError
class FilterFailedRAGEvoAgent(RAGEvoAgent):
def filter_evolvable_subjects_by_feedback(
self, evo: EvolvableSubjects, feedback: CoSTEERSingleFeedbackDeprecated
) -> EvolvableSubjects:
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()
@@ -10,11 +10,10 @@ from rdagent.components.coder.CoSTEER.knowledge_management import (
)
from rdagent.components.coder.CoSTEER.scheduler import random_select
from rdagent.core.conf import RD_AGENT_SETTINGS
from rdagent.core.evaluation import Scenario
from rdagent.core.evolving_framework import EvolvingStrategy, QueriedKnowledge
from rdagent.core.experiment import FBWorkspace
from rdagent.core.experiment import FBWorkspace, Task
from rdagent.core.prompts import Prompts
from rdagent.core.scenario import Task
from rdagent.core.scenario import Scenario
from rdagent.core.utils import multiprocessing_wrapper
implement_prompts = Prompts(file_path=Path(__file__).parent / "prompts.yaml")
@@ -111,4 +110,8 @@ 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
@@ -25,9 +25,8 @@ from rdagent.core.evolving_framework import (
QueriedKnowledge,
RAGStrategy,
)
from rdagent.core.experiment import FBWorkspace
from rdagent.core.experiment import FBWorkspace, Task
from rdagent.core.prompts import Prompts
from rdagent.core.scenario import Task
from rdagent.log import rdagent_logger as logger
from rdagent.oai.llm_utils import (
APIBackend,
@@ -4,7 +4,7 @@ from rdagent.components.coder.CoSTEER.evolvable_subjects import EvolvingItem
from rdagent.components.coder.CoSTEER.knowledge_management import (
CoSTEERQueriedKnowledge,
)
from rdagent.core.evaluation import Scenario
from rdagent.core.scenario import Scenario
from rdagent.log import rdagent_logger as logger
@@ -52,21 +52,14 @@ class EnsembleMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
if queried_knowledge is not None
else []
)
latest_code_feedback = [
knowledge.feedback
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("ensemble.py") is not None
and knowledge.implementation.file_dict.get("ensemble.py") == workspace.file_dict.get("ensemble.py")
]
if len(latest_code_feedback) > 0:
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("ensemble.py") != workspace.file_dict.get("ensemble.py")
],
queried_former_failed_knowledge[1],
)
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("ensemble.py") != workspace.file_dict.get("ensemble.py")
],
queried_former_failed_knowledge[1],
)
# Generate code with knowledge integration
competition_info = self.scen.get_scenario_all_desc()
@@ -81,7 +74,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=latest_code_feedback[0] if len(latest_code_feedback) > 0 else None,
latest_code_feedback=workspace.feedback,
)
for _ in range(5):
@@ -40,21 +40,14 @@ class FeatureMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
if queried_knowledge is not None
else []
)
latest_code_feedback = [
knowledge.feedback
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("feature.py") is not None
and knowledge.implementation.file_dict.get("feature.py") == workspace.file_dict.get("feature.py")
]
if len(latest_code_feedback) > 0:
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("feature.py") != workspace.file_dict.get("feature.py")
],
queried_former_failed_knowledge[1],
)
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("feature.py") != workspace.file_dict.get("feature.py")
],
queried_former_failed_knowledge[1],
)
# 2. code
system_prompt = T(".prompts:feature_coder.system").r(
@@ -66,7 +59,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=latest_code_feedback[0] if len(latest_code_feedback) > 0 else None,
latest_code_feedback=workspace.feedback,
)
for _ in range(5):
@@ -44,23 +44,15 @@ class ModelMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
if queried_knowledge is not None
else []
)
latest_code_feedback = [
knowledge.feedback
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get(f"{target_task.name}.py") is not None
and knowledge.implementation.file_dict.get(f"{target_task.name}.py")
== workspace.file_dict.get(f"{target_task.name}.py")
]
if len(latest_code_feedback) > 0:
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get(f"{target_task.name}.py")
!= workspace.file_dict.get(f"{target_task.name}.py")
],
queried_former_failed_knowledge[1],
)
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get(f"{target_task.name}.py")
!= workspace.file_dict.get(f"{target_task.name}.py")
],
queried_former_failed_knowledge[1],
)
# 2. code
system_prompt = T(".prompts:model_coder.system").r(
@@ -82,7 +74,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=latest_code_feedback[0] if len(latest_code_feedback) > 0 else None,
latest_code_feedback=workspace.feedback,
)
for _ in range(5):
@@ -69,21 +69,14 @@ class DataLoaderMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
if queried_knowledge is not None
else []
)
latest_code_feedback = [
knowledge.feedback
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("load_data.py") is not None
and knowledge.implementation.file_dict.get("load_data.py") == workspace.file_dict.get("load_data.py")
]
if len(latest_code_feedback) > 0:
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("load_data.py") != workspace.file_dict.get("load_data.py")
],
queried_former_failed_knowledge[1],
)
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("load_data.py") != workspace.file_dict.get("load_data.py")
],
queried_former_failed_knowledge[1],
)
# 1. specifications
# TODO: We may move spec into a separated COSTEER task
@@ -141,7 +134,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=latest_code_feedback[0] if len(latest_code_feedback) > 0 else None,
latest_code_feedback=workspace.feedback,
)
for _ in range(5):
@@ -41,21 +41,14 @@ class WorkflowMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
if queried_knowledge is not None
else []
)
latest_code_feedback = [
knowledge.feedback
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("main.py") is not None
and knowledge.implementation.file_dict.get("main.py") == workspace.file_dict.get("main.py")
]
if len(latest_code_feedback) > 0:
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("main.py") != workspace.file_dict.get("main.py")
],
queried_former_failed_knowledge[1],
)
queried_former_failed_knowledge = (
[
knowledge
for knowledge in queried_former_failed_knowledge[0]
if knowledge.implementation.file_dict.get("main.py") != workspace.file_dict.get("main.py")
],
queried_former_failed_knowledge[1],
)
# 2. code
system_prompt = T(".prompts:workflow_coder.system").r(
@@ -71,7 +64,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=latest_code_feedback[0] if len(latest_code_feedback) > 0 else None,
latest_code_feedback=workspace.feedback,
)
for _ in range(5):
@@ -97,21 +97,21 @@ class WorkflowGeneralCaseSpecEvaluator(CoSTEEREvaluator):
if not submission_fp.exists():
stdout += "\nSubmission file (submission.csv) is not generated."
else:
base_check_code = (DIRNAME / "eval_tests" / "submission_check.txt").read_text()
implementation.inject_files(**{"submission_check.py": base_check_code})
base_check_code = (DIRNAME / "eval_tests" / "submission_format_test.txt").read_text()
implementation.inject_files(**{"submission_format_test.py": base_check_code})
# stdout += "----Submission Check 1-----\n"
stdout += implementation.execute(env=de, entry="python submission_check.py")
stdout += implementation.execute(env=de, entry="python submission_format_test.py")
# MLEBench Check
# !!! Since we are running on a sampled dataset, mlebench check is not required.
# mle_check_code = (
# (DIRNAME / "eval_tests" / "mle_submission_check.txt")
# (DIRNAME / "eval_tests" / "mle_submission_format_test.txt")
# .read_text()
# .replace("<competition_id>", self.scen.competition)
# )
# implementation.inject_files(**{"mle_submission_check.py": mle_check_code})
# implementation.inject_files(**{"mle_submission_format_test.py": mle_check_code})
# stdout += "----Submission Check 2-----\n"
# stdout += implementation.execute(env=mde, entry=f"python mle_submission_check.py")
# stdout += implementation.execute(env=mde, entry=f"python mle_submission_format_test.py")
system_prompt = T(".prompts:workflow_eval.system").r(
scenario=self.scen.get_scenario_all_desc(),
-5
View File
@@ -1,8 +1,3 @@
import pickle
import shutil
from pathlib import Path
from typing import Any, Tuple
from rdagent.core.developer import Developer
from rdagent.core.experiment import ASpecificExp, Experiment
from rdagent.oai.llm_utils import md5_hash
+5 -6
View File
@@ -1,10 +1,9 @@
import typing
from abc import ABC, abstractmethod
from rdagent.core.scenario import Scenario
if typing.TYPE_CHECKING:
from rdagent.core.experiment import Task, Workspace
from rdagent.core.scenario import Scenario
class Feedback:
@@ -23,14 +22,14 @@ class Evaluator(ABC):
Design Principle:
It should cover the building process of feedback from raw information.
Typically the buiilding of feedback will be two phases.
1. raw information including stdout & workspace (feeedback itself will handle this)
2. advanced/summaried feedback information. (evaluate will handle this)
Typically the building of feedback will be two phases.
1. raw information including stdout & workspace (feedback itself will handle this)
2. advanced/summarized feedback information. (evaluate will handle this)
"""
def __init__(
self,
scen: Scenario,
scen: "Scenario",
) -> None:
self.scen = scen
+4 -1
View File
@@ -13,6 +13,7 @@ from pathlib import Path
from typing import Any, Generic, TypeVar
from rdagent.core.conf import RD_AGENT_SETTINGS
from rdagent.core.evaluation import Feedback
from rdagent.utils import filter_progress_bar
from rdagent.utils.fmt import shrink_text
@@ -55,9 +56,10 @@ class Task(AbsTask):
ASpecificTask = TypeVar("ASpecificTask", bound=Task)
ASpecificFeedback = TypeVar("ASpecificFeedback", bound=Feedback)
class Workspace(ABC, Generic[ASpecificTask]):
class Workspace(ABC, Generic[ASpecificTask, ASpecificFeedback]):
"""
A workspace is a place to store the task implementation. It evolves as the developer implements the task.
To get a snapshot of the workspace, make sure call `copy` to get a copy of the workspace.
@@ -65,6 +67,7 @@ class Workspace(ABC, Generic[ASpecificTask]):
def __init__(self, target_task: ASpecificTask | None = None) -> None:
self.target_task: ASpecificTask | None = target_task
self.feedback: ASpecificFeedback | None = None
@abstractmethod
def execute(self, *args: Any, **kwargs: Any) -> object | None:
@@ -1,76 +0,0 @@
import json
import os
from pathlib import Path
import pandas as pd
from rdagent.app.data_science.conf import DS_RD_SETTING
from rdagent.core.developer import Developer
from rdagent.core.exception import RunnerError
from rdagent.log import rdagent_logger as logger
from rdagent.scenarios.data_science.experiment.experiment import DSExperiment
from rdagent.utils.env import DockerEnv, DSDockerConf, MLEBDockerConf
class DSRunner(Developer[DSExperiment]):
def develop(self, exp: DSExperiment) -> DSExperiment:
ds_docker_conf = DSDockerConf()
ds_docker_conf.extra_volumes = {f"{DS_RD_SETTING.local_data_path}/{self.scen.competition}": "/kaggle/input"}
ds_docker_conf.running_timeout_period = DS_RD_SETTING.full_timeout
de = DockerEnv(conf=ds_docker_conf)
stdout = exp.experiment_workspace.execute(
env=de, entry=f"rm submission.csv scores.csv"
) # Remove previous submission and scores files generated by worklfow.
# execute workflow
stdout = exp.experiment_workspace.execute(env=de, entry="coverage run main.py")
score_fp = exp.experiment_workspace.workspace_path / "scores.csv"
if not score_fp.exists():
logger.error("Metrics file (scores.csv) is not generated.")
raise RunnerError(f"Metrics file (scores.csv) is not generated, log is:\n{stdout}")
submission_fp = exp.experiment_workspace.workspace_path / "submission.csv"
if not submission_fp.exists():
logger.error("Submission file (submission.csv) is not generated.")
raise RunnerError(f"Submission file (submission.csv) is not generated, log is:\n{stdout}")
else:
# DockerEnv for MLEBench submission validation
mle_de_conf = MLEBDockerConf()
mle_de_conf.extra_volumes = {
f"{DS_RD_SETTING.local_data_path}/zip_files": "/mle/data",
}
mde = DockerEnv(conf=mle_de_conf)
mde.prepare()
# MLEBench Check
mle_check_code = (
(Path(__file__).absolute().resolve().parent / "eval_tests" / "mle_submission_check.txt")
.read_text()
.replace("<competition_id>", self.scen.competition)
)
exp.experiment_workspace.inject_files(**{"mle_submission_check.py": mle_check_code})
exp.format_check_result = exp.experiment_workspace.execute(env=mde, entry=f"python mle_submission_check.py")
exp.result = pd.read_csv(score_fp, index_col=0)
# remove unused files
stdout = exp.experiment_workspace.execute(env=de, entry="coverage json -o coverage.json")
if Path(exp.experiment_workspace.workspace_path / "coverage.json").exists():
with open(exp.experiment_workspace.workspace_path / "coverage.json") as f:
used_files = set(json.load(f)["files"].keys()) | {"submission_check.py", "mle_submission_check.py"}
logger.info("All used scripts: {}".format(used_files))
all_python_files = set(Path(exp.experiment_workspace.workspace_path).rglob("*.py"))
unused_files = [
py_file
for py_file in all_python_files
if not (py_file.name in used_files or py_file.name.endswith("test.py"))
]
if unused_files:
logger.warning(f"Unused scripts: {unused_files}")
exp.experiment_workspace.inject_files(
**{file_path.name: exp.experiment_workspace.DEL_KEY for file_path in unused_files}
)
os.remove(exp.experiment_workspace.workspace_path / "coverage.json")
return exp
@@ -0,0 +1,125 @@
from pathlib import Path
import pandas as pd
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.evolvable_subjects import FBWorkspace
from rdagent.components.coder.CoSTEER.evolving_strategy import (
CoSTEERQueriedKnowledge,
MultiProcessEvolvingStrategy,
)
from rdagent.components.coder.CoSTEER.task import CoSTEERTask
from rdagent.core.exception import RunnerError
from rdagent.core.scenario import Scenario
from rdagent.log import rdagent_logger as logger
from rdagent.scenarios.data_science.dev.runner.eval import DSCoSTEERCoSTEEREvaluator
from rdagent.utils import APIBackend
from rdagent.utils.agent.ret import BatchEditOut
from rdagent.utils.agent.tpl import T
from rdagent.utils.env import DockerEnv, MLEBDockerConf
class DSRunnerMultiProcessEvolvingStrategy(MultiProcessEvolvingStrategy):
def implement_one_task(
self,
target_task: CoSTEERTask,
queried_knowledge: CoSTEERQueriedKnowledge | None = None,
workspace: FBWorkspace | None = None,
) -> dict[str, str]:
if workspace.feedback is None:
return {}
task_information_str = target_task.get_task_information()
# 1. code
system_prompt = T(".prompts:DSCoSTEER_debugger.system").r(
task_desc=task_information_str,
out_spec=BatchEditOut.get_spec(with_del=False),
)
user_prompt = T(".prompts:DSCoSTEER_debugger.user").r(
code=workspace.all_codes,
feedback=workspace.feedback,
)
batch_edit = BatchEditOut.extract_output(
APIBackend().build_messages_and_create_chat_completion(
user_prompt=user_prompt,
system_prompt=system_prompt,
json_mode=BatchEditOut.json_mode,
)
)
batch_edit = {k: v for k, v in batch_edit.items() if k in workspace.file_dict.keys()}
return batch_edit
def assign_code_list_to_evo(self, code_list: list[dict[str, str]], evo):
"""
Assign the code list to the evolving item.
The code list is aligned with the evolving item's sub-tasks.
If a task is not implemented, put a None in the list.
"""
for index in range(len(evo.sub_tasks)):
if code_list[index] is None:
continue
if evo.sub_workspace_list[index] is None:
# evo.sub_workspace_list[index] = FBWorkspace(target_task=evo.sub_tasks[index])
evo.sub_workspace_list[index] = evo.experiment_workspace
evo.sub_workspace_list[index].inject_files(**code_list[index])
return evo
class DSCoSTEERRunner(CoSTEER):
def __init__(
self,
scen: Scenario,
*args,
**kwargs,
) -> None:
eva = CoSTEERMultiEvaluator(
DSCoSTEERCoSTEEREvaluator(scen=scen), scen=scen
) # Please specify whether you agree running your eva in parallel or not
es = DSRunnerMultiProcessEvolvingStrategy(scen=scen, settings=CoSTEER_SETTINGS)
super().__init__(*args, settings=CoSTEER_SETTINGS, eva=eva, es=es, evolving_version=2, scen=scen, **kwargs)
def develop(self, exp):
bak_sub_tasks = exp.sub_tasks
exp.sub_tasks = [
CoSTEERTask(
name="Debug running solution",
description="The whole workflow of the solution has finished with some execution error, please check the error message and debug the whole code repo.",
)
]
exp = super().develop(exp)
exp.sub_tasks = bak_sub_tasks
score_fp = exp.experiment_workspace.workspace_path / "scores.csv"
if not score_fp.exists():
logger.error("Metrics file (scores.csv) is not generated.")
raise RunnerError(f"Metrics file (scores.csv) is not generated")
exp.result = pd.read_csv(score_fp, index_col=0)
# DockerEnv for MLEBench submission validation
mle_de_conf = MLEBDockerConf()
mle_de_conf.extra_volumes = {
f"{DS_RD_SETTING.local_data_path}/zip_files": "/mle/data",
}
mde = DockerEnv(conf=mle_de_conf)
mde.prepare()
# MLEBench Check
mle_check_code = (
(Path(__file__).absolute().resolve().parent / "eval_tests" / "mle_submission_format_test.txt")
.read_text()
.replace("<competition_id>", self.scen.competition)
)
exp.experiment_workspace.inject_files(**{"mle_submission_format_test.py": mle_check_code})
exp.format_check_result = exp.experiment_workspace.execute(
env=mde, entry=f"python mle_submission_format_test.py"
)
return exp
@@ -0,0 +1,107 @@
import json
import os
from pathlib import Path
from rdagent.app.data_science.conf import DS_RD_SETTING
from rdagent.components.coder.CoSTEER.evaluators import (
CoSTEEREvaluator,
CoSTEERSingleFeedback,
)
from rdagent.core.evolving_framework import QueriedKnowledge
from rdagent.core.experiment import FBWorkspace, Task
from rdagent.log import rdagent_logger as logger
from rdagent.oai.llm_utils import APIBackend
from rdagent.utils.agent.tpl import T
from rdagent.utils.agent.workflow import build_cls_from_json_with_retry
from rdagent.utils.env import DockerEnv, DSDockerConf, MLEBDockerConf
from rdagent.utils.fmt import shrink_text
DIRNAME = Path(__file__).absolute().resolve().parent
DSCoSTEEREvalFeedback = CoSTEERSingleFeedback
class DSCoSTEERCoSTEEREvaluator(CoSTEEREvaluator):
def evaluate(
self,
target_task: Task,
implementation: FBWorkspace,
gt_implementation: FBWorkspace,
queried_knowledge: QueriedKnowledge = None,
**kwargs,
) -> DSCoSTEEREvalFeedback:
ds_docker_conf = DSDockerConf()
ds_docker_conf.extra_volumes = {f"{DS_RD_SETTING.local_data_path}/{self.scen.competition}": "/kaggle/input"}
ds_docker_conf.running_timeout_period = DS_RD_SETTING.full_timeout
de = DockerEnv(conf=ds_docker_conf)
stdout = implementation.execute(
env=de, entry=f"rm submission.csv scores.csv"
) # Remove previous submission and scores files generated by worklfow.
# execute workflow
stdout = implementation.execute(env=de, entry="coverage run main.py")
score_fp = implementation.workspace_path / "scores.csv"
if not score_fp.exists():
stdout += "\n Metrics file (scores.csv) is not generated!"
else:
stdout += "\n Metrics file (scores.csv) is generated."
submission_fp = implementation.workspace_path / "submission.csv"
if not submission_fp.exists():
stdout += "\n Submission file (submission.csv) is not generated!"
else:
# DockerEnv for MLEBench submission validation
mle_de_conf = MLEBDockerConf()
mle_de_conf.extra_volumes = {
f"{DS_RD_SETTING.local_data_path}/zip_files": "/mle/data",
}
mde = DockerEnv(conf=mle_de_conf)
mde.prepare()
# MLEBench Check
mle_check_code = (
(Path(__file__).absolute().resolve().parent / "eval_tests" / "mle_submission_format_test.txt")
.read_text()
.replace("<competition_id>", self.scen.competition)
)
implementation.inject_files(**{"mle_submission_format_test.py": mle_check_code})
stdout += f"\n MLEBench submission check:"
stdout += implementation.execute(env=mde, entry="python mle_submission_format_test.py")
# remove unused files
implementation.execute(env=de, entry="coverage json -o coverage.json")
if Path(implementation.workspace_path / "coverage.json").exists():
with open(implementation.workspace_path / "coverage.json") as f:
used_files = set(json.load(f)["files"].keys()) | {
"submission_format_test.py",
"mle_submission_format_test.py",
}
logger.info("All used scripts: {}".format(used_files))
all_python_files = set(Path(implementation.workspace_path).rglob("*.py"))
unused_files = [
py_file
for py_file in all_python_files
if not (py_file.name in used_files or py_file.name.endswith("test.py"))
]
if unused_files:
logger.warning(f"Unused scripts: {unused_files}")
implementation.inject_files(
**{file_path.name: implementation.DEL_KEY for file_path in unused_files}
)
os.remove(implementation.workspace_path / "coverage.json")
system_prompt = T(".prompts:DSCoSTEER_eval.system").r(
scenario=self.scen.get_scenario_all_desc(),
task_desc=target_task.get_task_information(),
)
user_prompt = T(".prompts:DSCoSTEER_eval.user").r(
code=implementation.all_codes,
stdout=shrink_text(stdout),
)
return build_cls_from_json_with_retry(
DSCoSTEEREvalFeedback, system_prompt=system_prompt, user_prompt=user_prompt
)
@@ -0,0 +1,65 @@
DSCoSTEER_eval:
system: |-
You are a data scientist responsible for evaluating all the code.
## Task Description
The user is trying to build a data science solution in the following scenario:
{{ scenario }}
The task is as follows:
{{ task_desc }}
The whole workflow includes multiple stages, such as:
- Data loading
- Feature engineering
- Model training
- Ensembling
The user will provide you the whole code base, some logs generated during the execution of the whole workflow. Your evaluation scope includes whether the workflow code:
1. Executes successfully, correctly organizing components and generating a final submission.
2. Generates predictions in the correct format, ensuring they align with the **sample submission** structure!
Please respond with your feedback in the following JSON format and order
```json
{
"execution": "Describe whether the whole code base executed successfully and generating the final submission. Include any errors or issues encountered, and retain all error messages and traceback details.",
"return_checking": "Verify the generated files, particularly the submission file. Ensure that its format matches the sample submission",
"code": "Provide feedback on code quality, readability, and adherence to the given specifications.",
"final_decision": <true/false>
}
```
user: |-
--------- code base ---------
{{ code }}
--------- test stdout ---------
{{ stdout }}
DSCoSTEER_debugger:
system: |-
You are a world-class data scientist and machine learning engineer with deep expertise in statistics, mathematics, and computer science.
You have finished the implementation of the whole workflow which has executed well on a sampled dataset. However, the user has reported that the workflow failed to execute on the full dataset.
Your current job is to debug the whole code base, try to correct the errors, and ensure that the workflow can execute successfully on the full dataset.
The user will provide your the whole code base and some feedback generated during the execution of the whole workflow. Please identify the issues and provide the corrected code.
Task description:
{{ task_desc }}
Your modified code should follow the minimal changes principle. You should only modify the code that is necessary to fix the issues but not affect any other parts of the code. Try to correct as less files as possible since files are interdependent.
## Output Format
{% if out_spec %}
{{ out_spec }}
{% else %}
Please response the code in the following json format. Here is an example structure for the JSON output:
{
"code": "The Python code as a string."
}
{% endif %}
user: |-
--------- code base ---------
{{ code }}
--------- feedback ---------
{{ feedback }}
+2 -2
View File
@@ -41,8 +41,8 @@ class BatchEditOut(AgentOut):
json_mode: bool = True
@classmethod
def get_spec(cls):
return T(".tpl:BatchEditOut").r()
def get_spec(cls, with_del=True):
return T(".tpl:BatchEditOut").r(with_del=with_del)
@classmethod
def extract_output(cls, resp: str):
+5 -1
View File
@@ -12,6 +12,10 @@ BatchEditOut: |-
For example:
Inject the code into the folder. Your file name should always contain the suffix. Your file name keys should be unique to avoid delete or replace conflicts.
{
<file name1>: "<code>", // indicate writing <code> into <file name> (create new file or replace existing file)
<file name1>: "<code>", // indicate writing <code> into <file name1> (create new file or replace existing file)
{% if with_del %}
<file name2>: "__DEL__" // indicate removing file name2. When we want to replace a file to a new one, we usually use this
{% else %}
<file name2>(optional): "<code>" // indicate writing <code> into <file name2> (create new file or replace existing file)
{% endif %}
}