mirror of
https://github.com/NicolasBohn/NexQuant.git
synced 2026-08-01 17:37:43 +00:00
Implement model (and some factor) coder with evolving (#52)
* store code into FBImplementation * fix path related bugs * fix a bug * fix factor related small bugs * re-submit all model related code * new code to model coder * finish the model evolving code --------- Co-authored-by: xuyang1 <xuyang1@microsoft.com>
This commit is contained in:
@@ -0,0 +1,86 @@
|
||||
import pickle
|
||||
from pathlib import Path
|
||||
|
||||
from rdagent.components.coder.model_coder.conf import MODEL_IMPL_SETTINGS
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evaluators import (
|
||||
ModelCoderMultiEvaluator,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evolvable_subjects import (
|
||||
ModelEvolvingItem,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evolving_strategy import (
|
||||
ModelCoderEvolvingStrategy,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.CoSTEER.knowledge_management import (
|
||||
ModelKnowledgeBase,
|
||||
ModelRAGStrategy,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.model import ModelExperiment
|
||||
from rdagent.core.evolving_agent import RAGEvoAgent
|
||||
from rdagent.core.task_generator import TaskGenerator
|
||||
|
||||
|
||||
class ModelCoSTEER(TaskGenerator[ModelExperiment]):
|
||||
def __init__(
|
||||
self,
|
||||
*args,
|
||||
with_knowledge: bool = True,
|
||||
with_feedback: bool = True,
|
||||
knowledge_self_gen: bool = True,
|
||||
**kwargs,
|
||||
) -> None:
|
||||
super().__init__(*args, **kwargs)
|
||||
self.max_loop = MODEL_IMPL_SETTINGS.max_loop
|
||||
self.knowledge_base_path = (
|
||||
Path(MODEL_IMPL_SETTINGS.knowledge_base_path)
|
||||
if MODEL_IMPL_SETTINGS.knowledge_base_path is not None
|
||||
else None
|
||||
)
|
||||
self.new_knowledge_base_path = (
|
||||
Path(MODEL_IMPL_SETTINGS.new_knowledge_base_path)
|
||||
if MODEL_IMPL_SETTINGS.new_knowledge_base_path is not None
|
||||
else None
|
||||
)
|
||||
self.with_knowledge = with_knowledge
|
||||
self.with_feedback = with_feedback
|
||||
self.knowledge_self_gen = knowledge_self_gen
|
||||
self.evolving_strategy = ModelCoderEvolvingStrategy(scen=self.scen)
|
||||
self.model_evaluator = ModelCoderMultiEvaluator(scen=self.scen)
|
||||
|
||||
def load_or_init_knowledge_base(self, former_knowledge_base_path: Path = None, component_init_list: list = []):
|
||||
if former_knowledge_base_path is not None and former_knowledge_base_path.exists():
|
||||
model_knowledge_base = pickle.load(open(former_knowledge_base_path, "rb"))
|
||||
if not isinstance(model_knowledge_base, ModelKnowledgeBase):
|
||||
raise ValueError("The former knowledge base is not compatible with the current version")
|
||||
else:
|
||||
model_knowledge_base = ModelKnowledgeBase()
|
||||
|
||||
return model_knowledge_base
|
||||
|
||||
def generate(self, exp: ModelExperiment) -> ModelExperiment:
|
||||
# init knowledge base
|
||||
model_knowledge_base = self.load_or_init_knowledge_base(
|
||||
former_knowledge_base_path=self.knowledge_base_path,
|
||||
component_init_list=[],
|
||||
)
|
||||
# init rag method
|
||||
self.rag = ModelRAGStrategy(model_knowledge_base)
|
||||
|
||||
# init intermediate items
|
||||
model_experiment = ModelEvolvingItem(sub_tasks=exp.sub_tasks)
|
||||
|
||||
self.evolve_agent = RAGEvoAgent(max_loop=self.max_loop, evolving_strategy=self.evolving_strategy, rag=self.rag)
|
||||
|
||||
model_experiment = self.evolve_agent.multistep_evolve(
|
||||
model_experiment,
|
||||
self.model_evaluator,
|
||||
with_knowledge=self.with_knowledge,
|
||||
with_feedback=self.with_feedback,
|
||||
knowledge_self_gen=self.knowledge_self_gen,
|
||||
)
|
||||
|
||||
# save new knowledge base
|
||||
if self.new_knowledge_base_path is not None:
|
||||
pickle.dump(model_knowledge_base, open(self.new_knowledge_base_path, "wb"))
|
||||
self.knowledge_base = model_knowledge_base
|
||||
return model_experiment
|
||||
@@ -0,0 +1,333 @@
|
||||
import json
|
||||
import random
|
||||
from pathlib import Path
|
||||
from typing import List, Tuple
|
||||
|
||||
import numpy as np
|
||||
import torch
|
||||
from jinja2 import Environment, StrictUndefined
|
||||
|
||||
from rdagent.components.coder.model_coder.conf import MODEL_IMPL_SETTINGS
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evolvable_subjects import (
|
||||
ModelEvolvingItem,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.model import ModelImplementation, ModelTask
|
||||
from rdagent.core.conf import RD_AGENT_SETTINGS
|
||||
from rdagent.core.evaluation import Evaluator
|
||||
from rdagent.core.evolving_framework import QueriedKnowledge
|
||||
from rdagent.core.experiment import Implementation, Task
|
||||
from rdagent.core.log import RDAgentLog
|
||||
from rdagent.core.prompts import Prompts
|
||||
from rdagent.core.utils import multiprocessing_wrapper
|
||||
from rdagent.oai.llm_utils import APIBackend
|
||||
|
||||
evaluate_prompts = Prompts(file_path=Path(__file__).parent.parent / "prompts.yaml")
|
||||
|
||||
|
||||
def shape_evaluator(prediction: torch.Tensor, target_shape: Tuple = None) -> Tuple[str, bool]:
|
||||
if target_shape is None or prediction is None:
|
||||
return "No output generated from the model. No shape evaluation conducted.", False
|
||||
pre_shape = prediction.shape
|
||||
|
||||
if pre_shape == target_shape:
|
||||
return "The shape of the output is correct.", True
|
||||
else:
|
||||
return f"The shape of the output is incorrect. Expected {target_shape}, but got {pre_shape}.", False
|
||||
|
||||
|
||||
def reshape_tensor(original_tensor, target_shape):
|
||||
new_tensor = torch.zeros(target_shape)
|
||||
for i, dim in enumerate(original_tensor.shape):
|
||||
new_tensor = new_tensor.narrow(i, 0, dim).copy_(original_tensor)
|
||||
|
||||
return new_tensor
|
||||
|
||||
|
||||
def value_evaluator(
|
||||
prediction: torch.Tensor,
|
||||
target: torch.Tensor,
|
||||
) -> Tuple[torch.Tensor, bool]:
|
||||
if target is None or prediction is None:
|
||||
return "No output generated from the model. No value evaluation conducted.", False
|
||||
else:
|
||||
# Calculate the mean absolute difference
|
||||
diff = torch.mean(torch.abs(target - prediction)).item()
|
||||
return (
|
||||
f"The value of the output is correct. The mean absolute difference is {diff}.",
|
||||
diff < 0.1,
|
||||
)
|
||||
|
||||
|
||||
class ModelCodeEvaluator(Evaluator):
|
||||
def evaluate(
|
||||
self,
|
||||
target_task: Task,
|
||||
implementation: Implementation,
|
||||
gt_implementation: Implementation,
|
||||
model_execution_feedback: str = "",
|
||||
model_value_feedback: str = "",
|
||||
):
|
||||
assert isinstance(target_task, ModelTask)
|
||||
assert isinstance(implementation, ModelImplementation)
|
||||
if gt_implementation is not None:
|
||||
assert isinstance(gt_implementation, ModelImplementation)
|
||||
|
||||
model_task_information = target_task.get_information()
|
||||
code = implementation.code
|
||||
|
||||
system_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(evaluate_prompts["evaluator_code_feedback"]["system"])
|
||||
.render(scenario=self.scen.get_scenario_all_desc() if self.scen is not None else "No scenario description.")
|
||||
)
|
||||
|
||||
execution_feedback_to_render = model_execution_feedback
|
||||
for _ in range(10): # 10 times to split the content is enough
|
||||
user_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(
|
||||
evaluate_prompts["evaluator_code_feedback"]["user"],
|
||||
)
|
||||
.render(
|
||||
model_information=model_task_information,
|
||||
code=code,
|
||||
model_execution_feedback=execution_feedback_to_render,
|
||||
model_value_feedback=model_value_feedback,
|
||||
gt_code=gt_implementation.code if gt_implementation else None,
|
||||
)
|
||||
)
|
||||
if (
|
||||
APIBackend().build_messages_and_calculate_token(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
)
|
||||
> RD_AGENT_SETTINGS.chat_token_limit
|
||||
):
|
||||
execution_feedback_to_render = execution_feedback_to_render[len(execution_feedback_to_render) // 2 :]
|
||||
else:
|
||||
break
|
||||
|
||||
critic_response = APIBackend().build_messages_and_create_chat_completion(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
json_mode=False,
|
||||
)
|
||||
|
||||
return critic_response, None
|
||||
|
||||
|
||||
class ModelFinalEvaluator(Evaluator):
|
||||
def evaluate(
|
||||
self,
|
||||
target_task: Task,
|
||||
implementation: Implementation,
|
||||
gt_implementation: Implementation,
|
||||
model_execution_feedback: str,
|
||||
model_value_feedback: str,
|
||||
model_code_feedback: str,
|
||||
):
|
||||
assert isinstance(target_task, ModelTask)
|
||||
assert isinstance(implementation, ModelImplementation)
|
||||
if gt_implementation is not None:
|
||||
assert isinstance(gt_implementation, ModelImplementation)
|
||||
|
||||
system_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(evaluate_prompts["evaluator_final_feedback"]["system"])
|
||||
.render(scenario=self.scen.get_scenario_all_desc() if self.scen is not None else "No scenario description.")
|
||||
)
|
||||
|
||||
execution_feedback_to_render = model_execution_feedback
|
||||
|
||||
for _ in range(10): # 10 times to split the content is enough
|
||||
user_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(
|
||||
evaluate_prompts["evaluator_final_feedback"]["user"],
|
||||
)
|
||||
.render(
|
||||
model_information=target_task.get_information(),
|
||||
model_execution_feedback=execution_feedback_to_render,
|
||||
model_code_feedback=model_code_feedback,
|
||||
model_value_feedback=model_value_feedback,
|
||||
)
|
||||
)
|
||||
if (
|
||||
APIBackend().build_messages_and_calculate_token(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
)
|
||||
> RD_AGENT_SETTINGS.chat_token_limit
|
||||
):
|
||||
execution_feedback_to_render = execution_feedback_to_render[len(execution_feedback_to_render) // 2 :]
|
||||
else:
|
||||
break
|
||||
|
||||
final_evaluation_dict = json.loads(
|
||||
APIBackend().build_messages_and_create_chat_completion(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
json_mode=True,
|
||||
),
|
||||
)
|
||||
if isinstance(final_evaluation_dict["final_decision"], str) and final_evaluation_dict[
|
||||
"final_decision"
|
||||
].lower() in ("true", "false"):
|
||||
final_evaluation_dict["final_decision"] = bool(final_evaluation_dict["final_decision"])
|
||||
return (
|
||||
final_evaluation_dict["final_feedback"],
|
||||
final_evaluation_dict["final_decision"],
|
||||
)
|
||||
|
||||
|
||||
class ModelCoderFeedback:
|
||||
"""This feedback includes all the content to the model coder"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
execution_feedback: str,
|
||||
shape_feedback: str,
|
||||
value_feedback: str,
|
||||
code_feedback: str,
|
||||
final_feedback: str,
|
||||
final_decision: bool,
|
||||
):
|
||||
self.execution_feedback: str = execution_feedback
|
||||
self.shape_feedback: str = shape_feedback
|
||||
self.value_feedback: str = value_feedback
|
||||
self.code_feedback: str = code_feedback
|
||||
self.final_feedback: str = final_feedback
|
||||
self.final_decision: str = final_decision
|
||||
|
||||
def __str__(self) -> str:
|
||||
return f"""------------------Model Execution Feedback------------------
|
||||
{self.execution_feedback}
|
||||
------------------Model Shape Feedback------------------
|
||||
{self.shape_feedback}
|
||||
------------------Model Value Feedback------------------
|
||||
{self.value_feedback}
|
||||
------------------Model Code Feedback------------------
|
||||
{self.code_feedback}
|
||||
------------------Model Final Feedback------------------
|
||||
{self.final_feedback}
|
||||
------------------Model Final Decision------------------
|
||||
This implementation is {'SUCCESS' if self.final_decision else 'FAIL'}.
|
||||
"""
|
||||
|
||||
|
||||
class ModelCoderEvaluator(Evaluator):
|
||||
def evaluate(
|
||||
self,
|
||||
target_task: Task,
|
||||
implementation: Implementation,
|
||||
gt_implementation: Implementation,
|
||||
queried_knowledge: QueriedKnowledge = None,
|
||||
**kwargs,
|
||||
) -> ModelCoderFeedback:
|
||||
target_task_information = target_task.get_information()
|
||||
if (
|
||||
queried_knowledge is not None
|
||||
and target_task_information in queried_knowledge.success_task_to_knowledge_dict
|
||||
):
|
||||
return queried_knowledge.success_task_to_knowledge_dict[target_task_information].feedback
|
||||
elif queried_knowledge is not None and target_task_information in queried_knowledge.failed_task_info_set:
|
||||
return ModelCoderFeedback(
|
||||
execution_feedback="This task has failed too many times, skip implementation.",
|
||||
shape_feedback="This task has failed too many times, skip implementation.",
|
||||
value_feedback="This task has failed too many times, skip implementation.",
|
||||
code_feedback="This task has failed too many times, skip implementation.",
|
||||
final_feedback="This task has failed too many times, skip implementation.",
|
||||
final_decision=False,
|
||||
)
|
||||
assert isinstance(target_task, ModelTask)
|
||||
|
||||
batch_size, num_features, num_timesteps = (
|
||||
random.randint(6, 10),
|
||||
random.randint(6, 10),
|
||||
random.randint(6, 10),
|
||||
)
|
||||
input_value, param_init_value = random.random(), random.random()
|
||||
|
||||
assert isinstance(implementation, ModelImplementation)
|
||||
model_execution_feedback, gen_tensor = implementation.execute(
|
||||
batch_size=batch_size,
|
||||
num_features=num_features,
|
||||
num_timesteps=num_timesteps,
|
||||
input_value=input_value,
|
||||
param_init_value=param_init_value,
|
||||
)
|
||||
if gt_implementation is not None:
|
||||
assert isinstance(gt_implementation, ModelImplementation)
|
||||
_, gt_tensor = gt_implementation.execute(
|
||||
batch_size=batch_size,
|
||||
num_features=num_features,
|
||||
num_timesteps=num_timesteps,
|
||||
input_value=input_value,
|
||||
param_init_value=param_init_value,
|
||||
)
|
||||
else:
|
||||
gt_tensor = None
|
||||
|
||||
shape_feedback, shape_decision = shape_evaluator(gen_tensor, (batch_size, 1))
|
||||
value_feedback, value_decision = value_evaluator(gt_tensor, gen_tensor)
|
||||
code_feedback, _ = ModelCodeEvaluator(scen=self.scen).evaluate(
|
||||
target_task=target_task,
|
||||
implementation=implementation,
|
||||
gt_implementation=gt_implementation,
|
||||
model_execution_feedback=model_execution_feedback,
|
||||
model_value_feedback="\n".join([shape_feedback, value_feedback]),
|
||||
)
|
||||
final_feedback, final_decision = ModelFinalEvaluator(scen=self.scen).evaluate(
|
||||
target_task=target_task,
|
||||
implementation=implementation,
|
||||
gt_implementation=gt_implementation,
|
||||
model_execution_feedback=model_execution_feedback,
|
||||
model_value_feedback=value_feedback,
|
||||
model_code_feedback=code_feedback,
|
||||
)
|
||||
|
||||
return ModelCoderFeedback(
|
||||
execution_feedback=model_execution_feedback,
|
||||
shape_feedback=shape_feedback,
|
||||
value_feedback=value_feedback,
|
||||
code_feedback=code_feedback,
|
||||
final_feedback=final_feedback,
|
||||
final_decision=final_decision,
|
||||
)
|
||||
|
||||
|
||||
class ModelCoderMultiEvaluator(Evaluator):
|
||||
def evaluate(
|
||||
self,
|
||||
evo: ModelEvolvingItem,
|
||||
queried_knowledge: QueriedKnowledge = None,
|
||||
**kwargs,
|
||||
) -> List[ModelCoderFeedback]:
|
||||
multi_implementation_feedback = []
|
||||
|
||||
calls = []
|
||||
for index in range(len(evo.sub_tasks)):
|
||||
corresponding_implementation = evo.sub_implementations[index]
|
||||
corresponding_gt_implementation = (
|
||||
evo.sub_gt_implementations[index] if evo.sub_gt_implementations is not None else None
|
||||
)
|
||||
calls.append(
|
||||
(
|
||||
ModelCoderEvaluator(scen=self.scen).evaluate,
|
||||
(
|
||||
evo.sub_tasks[index],
|
||||
corresponding_implementation,
|
||||
corresponding_gt_implementation,
|
||||
queried_knowledge,
|
||||
),
|
||||
),
|
||||
)
|
||||
multi_implementation_feedback = multiprocessing_wrapper(calls, n=MODEL_IMPL_SETTINGS.evo_multi_proc_n)
|
||||
|
||||
final_decision = [
|
||||
None if single_feedback is None else single_feedback.final_decision
|
||||
for single_feedback in multi_implementation_feedback
|
||||
]
|
||||
RDAgentLog().info(f"Final decisions: {final_decision} True count: {final_decision.count(True)}")
|
||||
|
||||
return multi_implementation_feedback
|
||||
@@ -0,0 +1,29 @@
|
||||
from rdagent.components.coder.model_coder.model import (
|
||||
ModelExperiment,
|
||||
ModelImplementation,
|
||||
ModelTask,
|
||||
)
|
||||
from rdagent.core.evolving_framework import EvolvableSubjects
|
||||
from rdagent.core.log import RDAgentLog
|
||||
|
||||
|
||||
class ModelEvolvingItem(ModelExperiment, EvolvableSubjects):
|
||||
"""
|
||||
Intermediate item of model implementation.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
sub_tasks: list[ModelTask],
|
||||
sub_gt_implementations: list[ModelImplementation] = None,
|
||||
):
|
||||
ModelExperiment.__init__(self, sub_tasks=sub_tasks)
|
||||
if sub_gt_implementations is not None and len(
|
||||
sub_gt_implementations,
|
||||
) != len(self.sub_tasks):
|
||||
self.sub_gt_implementations = None
|
||||
RDAgentLog().warning(
|
||||
"The length of sub_gt_implementations is not equal to the length of sub_tasks, set sub_gt_implementations to None",
|
||||
)
|
||||
else:
|
||||
self.sub_gt_implementations = sub_gt_implementations
|
||||
@@ -0,0 +1,145 @@
|
||||
import json
|
||||
from copy import deepcopy
|
||||
from pathlib import Path
|
||||
|
||||
from jinja2 import Environment, StrictUndefined
|
||||
|
||||
from rdagent.components.coder.model_coder.conf import MODEL_IMPL_SETTINGS
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evolvable_subjects import (
|
||||
ModelEvolvingItem,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.CoSTEER.knowledge_management import (
|
||||
ModelQueriedKnowledge,
|
||||
)
|
||||
from rdagent.components.coder.model_coder.model import ModelImplementation, ModelTask
|
||||
from rdagent.core.conf import RD_AGENT_SETTINGS
|
||||
from rdagent.core.evolving_framework import EvolvingStrategy
|
||||
from rdagent.core.prompts import Prompts
|
||||
from rdagent.core.utils import multiprocessing_wrapper
|
||||
from rdagent.oai.llm_utils import APIBackend
|
||||
|
||||
coder_prompts = Prompts(file_path=Path(__file__).parent.parent / "prompts.yaml")
|
||||
|
||||
|
||||
class ModelCoderEvolvingStrategy(EvolvingStrategy):
|
||||
def implement_one_model(
|
||||
self,
|
||||
target_task: ModelTask,
|
||||
queried_knowledge: ModelQueriedKnowledge = None,
|
||||
) -> ModelImplementation:
|
||||
model_information_str = target_task.get_information()
|
||||
|
||||
if queried_knowledge is not None and model_information_str in queried_knowledge.success_task_to_knowledge_dict:
|
||||
return queried_knowledge.success_task_to_knowledge_dict[model_information_str].implementation
|
||||
elif queried_knowledge is not None and model_information_str in queried_knowledge.failed_task_info_set:
|
||||
return None
|
||||
else:
|
||||
queried_similar_successful_knowledge = (
|
||||
queried_knowledge.working_task_to_similar_successful_knowledge_dict[model_information_str]
|
||||
if queried_knowledge is not None
|
||||
else []
|
||||
)
|
||||
queried_former_failed_knowledge = (
|
||||
queried_knowledge.working_task_to_former_failed_knowledge_dict[model_information_str]
|
||||
if queried_knowledge is not None
|
||||
else []
|
||||
)
|
||||
|
||||
queried_former_failed_knowledge_to_render = queried_former_failed_knowledge
|
||||
|
||||
system_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(
|
||||
coder_prompts["evolving_strategy_model_coder"]["system"],
|
||||
)
|
||||
.render(
|
||||
scenario=self.scen.get_scenario_all_desc(),
|
||||
queried_former_failed_knowledge=queried_former_failed_knowledge_to_render,
|
||||
)
|
||||
)
|
||||
|
||||
queried_similar_successful_knowledge_to_render = queried_similar_successful_knowledge
|
||||
for _ in range(10): # max attempt to reduce the length of user_prompt
|
||||
user_prompt = (
|
||||
Environment(undefined=StrictUndefined)
|
||||
.from_string(
|
||||
coder_prompts["evolving_strategy_model_coder"]["user"],
|
||||
)
|
||||
.render(
|
||||
model_information_str=model_information_str,
|
||||
queried_similar_successful_knowledge=queried_similar_successful_knowledge_to_render,
|
||||
queried_former_failed_knowledge=queried_former_failed_knowledge_to_render,
|
||||
)
|
||||
.strip("\n")
|
||||
)
|
||||
if (
|
||||
APIBackend().build_messages_and_calculate_token(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
)
|
||||
< RD_AGENT_SETTINGS.chat_token_limit
|
||||
):
|
||||
break
|
||||
elif len(queried_former_failed_knowledge_to_render) > 1:
|
||||
queried_former_failed_knowledge_to_render = queried_former_failed_knowledge_to_render[1:]
|
||||
elif len(queried_similar_successful_knowledge_to_render) > 1:
|
||||
queried_similar_successful_knowledge_to_render = queried_similar_successful_knowledge_to_render[1:]
|
||||
|
||||
code = json.loads(
|
||||
APIBackend(use_chat_cache=True).build_messages_and_create_chat_completion(
|
||||
user_prompt=user_prompt,
|
||||
system_prompt=system_prompt,
|
||||
json_mode=True,
|
||||
),
|
||||
)["code"]
|
||||
# ast.parse(code)
|
||||
model_implementation = ModelImplementation(
|
||||
target_task,
|
||||
)
|
||||
model_implementation.prepare()
|
||||
model_implementation.inject_code(**{"model.py": code})
|
||||
|
||||
return model_implementation
|
||||
|
||||
def evolve(
|
||||
self,
|
||||
*,
|
||||
evo: ModelEvolvingItem,
|
||||
queried_knowledge: ModelQueriedKnowledge | None = None,
|
||||
**kwargs,
|
||||
) -> ModelEvolvingItem:
|
||||
new_evo = deepcopy(evo)
|
||||
|
||||
# 1.找出需要evolve的model
|
||||
to_be_finished_task_index = []
|
||||
for index, target_model_task in enumerate(new_evo.sub_tasks):
|
||||
target_model_task_desc = target_model_task.get_information()
|
||||
if target_model_task_desc in queried_knowledge.success_task_to_knowledge_dict:
|
||||
new_evo.sub_implementations[index] = queried_knowledge.success_task_to_knowledge_dict[
|
||||
target_model_task_desc
|
||||
].implementation
|
||||
elif (
|
||||
target_model_task_desc not in queried_knowledge.success_task_to_knowledge_dict
|
||||
and target_model_task_desc not in queried_knowledge.failed_task_info_set
|
||||
):
|
||||
to_be_finished_task_index.append(index)
|
||||
|
||||
result = multiprocessing_wrapper(
|
||||
[
|
||||
(self.implement_one_model, (new_evo.sub_tasks[target_index], queried_knowledge))
|
||||
for target_index in to_be_finished_task_index
|
||||
],
|
||||
n=MODEL_IMPL_SETTINGS.evo_multi_proc_n,
|
||||
)
|
||||
|
||||
for index, target_index in enumerate(to_be_finished_task_index):
|
||||
new_evo.sub_implementations[target_index] = result[index]
|
||||
|
||||
# for target_index in to_be_finished_task_index:
|
||||
# new_evo.sub_implementations[target_index] = self.implement_one_model(
|
||||
# new_evo.sub_tasks[target_index], queried_knowledge
|
||||
# )
|
||||
|
||||
new_evo.corresponding_selection = to_be_finished_task_index
|
||||
|
||||
return new_evo
|
||||
@@ -0,0 +1,167 @@
|
||||
from rdagent.components.coder.model_coder.conf import MODEL_IMPL_SETTINGS
|
||||
from rdagent.components.coder.model_coder.CoSTEER.evaluators import ModelCoderFeedback
|
||||
from rdagent.components.coder.model_coder.model import ModelTask
|
||||
from rdagent.core.evolving_framework import (
|
||||
EvolvableSubjects,
|
||||
EvoStep,
|
||||
Knowledge,
|
||||
KnowledgeBase,
|
||||
QueriedKnowledge,
|
||||
RAGStrategy,
|
||||
)
|
||||
from rdagent.core.experiment import Implementation
|
||||
from rdagent.oai.llm_utils import calculate_embedding_distance_between_str_list
|
||||
|
||||
|
||||
class ModelKnowledge(Knowledge):
|
||||
def __init__(
|
||||
self,
|
||||
target_task: ModelTask,
|
||||
implementation: Implementation,
|
||||
feedback: ModelCoderFeedback,
|
||||
) -> None:
|
||||
"""
|
||||
Initialize a ModelKnowledge object. The ModelKnowledge object is used to store a model implementation without the ground truth code and value.
|
||||
|
||||
Args:
|
||||
model (Model): The model object associated with the KnowledgeManagement.
|
||||
|
||||
Returns:
|
||||
None
|
||||
"""
|
||||
self.target_task = target_task
|
||||
self.implementation = implementation
|
||||
self.feedback = feedback
|
||||
|
||||
def get_implementation_and_feedback_str(self) -> str:
|
||||
return f"""------------------Model implementation code:------------------
|
||||
{self.implementation.code}
|
||||
------------------Model implementation feedback:------------------
|
||||
{self.feedback!s}
|
||||
"""
|
||||
|
||||
|
||||
class ModelQueriedKnowledge(QueriedKnowledge):
|
||||
def __init__(self, success_task_to_knowledge_dict: dict = {}, failed_task_info_set: set = set()) -> None:
|
||||
self.success_task_to_knowledge_dict = success_task_to_knowledge_dict
|
||||
self.failed_task_info_set = failed_task_info_set
|
||||
self.working_task_to_former_failed_knowledge_dict = dict()
|
||||
self.working_task_to_similar_successful_knowledge_dict = dict()
|
||||
|
||||
|
||||
class ModelKnowledgeBase(KnowledgeBase):
|
||||
def __init__(self) -> None:
|
||||
self.implementation_trace: dict[str, ModelKnowledge] = dict()
|
||||
self.success_task_info_set: set[str] = set()
|
||||
|
||||
self.task_to_embedding = dict()
|
||||
|
||||
def query(self) -> QueriedKnowledge | None:
|
||||
"""
|
||||
Query the knowledge base to get the queried knowledge. So far is handled in RAG strategy.
|
||||
"""
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
class ModelRAGStrategy(RAGStrategy):
|
||||
def __init__(self, knowledgebase: ModelKnowledgeBase) -> None:
|
||||
super().__init__(knowledgebase)
|
||||
self.current_generated_trace_count = 0
|
||||
|
||||
def generate_knowledge(
|
||||
self,
|
||||
evolving_trace: list[EvoStep],
|
||||
*,
|
||||
return_knowledge: bool = False,
|
||||
) -> Knowledge | None:
|
||||
if len(evolving_trace) == self.current_generated_trace_count:
|
||||
return
|
||||
else:
|
||||
for trace_index in range(
|
||||
self.current_generated_trace_count,
|
||||
len(evolving_trace),
|
||||
):
|
||||
evo_step = evolving_trace[trace_index]
|
||||
implementations = evo_step.evolvable_subjects
|
||||
feedback = evo_step.feedback
|
||||
for task_index in range(len(implementations.sub_tasks)):
|
||||
target_task = implementations.sub_tasks[task_index]
|
||||
target_task_information = target_task.get_information()
|
||||
implementation = implementations.sub_implementations[task_index]
|
||||
single_feedback = feedback[task_index]
|
||||
if single_feedback is None:
|
||||
continue
|
||||
single_knowledge = ModelKnowledge(
|
||||
target_task=target_task,
|
||||
implementation=implementation,
|
||||
feedback=single_feedback,
|
||||
)
|
||||
if target_task_information not in self.knowledgebase.success_task_info_set:
|
||||
self.knowledgebase.implementation_trace.setdefault(
|
||||
target_task_information,
|
||||
[],
|
||||
).append(single_knowledge)
|
||||
|
||||
if single_feedback.final_decision == True:
|
||||
self.knowledgebase.success_task_info_set.add(
|
||||
target_task_information,
|
||||
)
|
||||
self.current_generated_trace_count = len(evolving_trace)
|
||||
|
||||
def query(
|
||||
self,
|
||||
evo: EvolvableSubjects,
|
||||
evolving_trace: list[EvoStep],
|
||||
) -> QueriedKnowledge | None:
|
||||
query_former_trace_limit = MODEL_IMPL_SETTINGS.query_former_trace_limit
|
||||
query_similar_success_limit = MODEL_IMPL_SETTINGS.query_similar_success_limit
|
||||
fail_task_trial_limit = MODEL_IMPL_SETTINGS.fail_task_trial_limit
|
||||
|
||||
queried_knowledge = ModelQueriedKnowledge()
|
||||
for target_model_task in evo.sub_tasks:
|
||||
target_model_task_information = target_model_task.get_information()
|
||||
if target_model_task_information in self.knowledgebase.success_task_info_set:
|
||||
queried_knowledge.success_task_to_knowledge_dict[target_model_task_information] = (
|
||||
self.knowledgebase.implementation_trace[target_model_task_information][-1]
|
||||
)
|
||||
elif (
|
||||
len(
|
||||
self.knowledgebase.implementation_trace.setdefault(
|
||||
target_model_task_information,
|
||||
[],
|
||||
),
|
||||
)
|
||||
>= fail_task_trial_limit
|
||||
):
|
||||
queried_knowledge.failed_task_info_set.add(target_model_task_information)
|
||||
else:
|
||||
queried_knowledge.working_task_to_former_failed_knowledge_dict[target_model_task_information] = (
|
||||
self.knowledgebase.implementation_trace.setdefault(
|
||||
target_model_task_information,
|
||||
[],
|
||||
)[-query_former_trace_limit:]
|
||||
)
|
||||
|
||||
knowledge_base_success_task_list = list(
|
||||
self.knowledgebase.success_task_info_set,
|
||||
)
|
||||
similarity = calculate_embedding_distance_between_str_list(
|
||||
[target_model_task_information],
|
||||
knowledge_base_success_task_list,
|
||||
)[0]
|
||||
similar_indexes = sorted(
|
||||
range(len(similarity)),
|
||||
key=lambda i: similarity[i],
|
||||
reverse=True,
|
||||
)[:query_similar_success_limit]
|
||||
similar_successful_knowledge = [
|
||||
self.knowledgebase.implementation_trace.setdefault(
|
||||
knowledge_base_success_task_list[index],
|
||||
[],
|
||||
)[-1]
|
||||
for index in similar_indexes
|
||||
]
|
||||
queried_knowledge.working_task_to_similar_successful_knowledge_dict[target_model_task_information] = (
|
||||
similar_successful_knowledge
|
||||
)
|
||||
return queried_knowledge
|
||||
Reference in New Issue
Block a user