From a535bb2d8dc035ddc085252c30e802d4213aad4d Mon Sep 17 00:00:00 2001 From: "jinli.yl" Date: Fri, 25 Jul 2025 15:19:33 +0800 Subject: [PATCH] change dump/load from vectornode to experience --- doc/future_roadmap.md | 7 ++-- .../op/vector_store/vector_store_action_op.py | 15 ++++++- experiencemaker/schema/experience.py | 39 +++++++++++++++---- .../vector_store/base_vector_store.py | 30 +++++++++----- 4 files changed, 69 insertions(+), 22 deletions(-) diff --git a/doc/future_roadmap.md b/doc/future_roadmap.md index 22b06121..472337b5 100644 --- a/doc/future_roadmap.md +++ b/doc/future_roadmap.md @@ -64,11 +64,12 @@ Democratize AI experience sharing by making curated experience libraries publicl - [ ] cook_book-appworld code & readme @jiaji - [ ] cook_book-bfcl-v3 op @zouyin delay 0730 - [x] fix multi-process bug @jinli -- [ ] logo optimize @jiaji +- [x] logo optimize @jiaji - [x] Ready-made Experience Store @jinli, add appworld/bfcl-v3 default experience store @jiaji -- [ ] op config make up @jiaji +- [ ] op config make up @jinli - [x] config make up, easy to understand @jinli - [x] refine readme @jinli +- [ ] bugfix dump/load experience @jinli dump@jiaji - [x] integrate into beyond-agent @jinli -- [ ] rm workspace_id in code @jinli +- [x] rm workspace_id in code @jinli diff --git a/experiencemaker/op/vector_store/vector_store_action_op.py b/experiencemaker/op/vector_store/vector_store_action_op.py index c12f442a..331b3c51 100644 --- a/experiencemaker/op/vector_store/vector_store_action_op.py +++ b/experiencemaker/op/vector_store/vector_store_action_op.py @@ -1,7 +1,9 @@ from experiencemaker.op import OP_REGISTRY from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import vector_node_to_experience, dict_to_experience, BaseExperience from experiencemaker.schema.request import VectorStoreRequest from experiencemaker.schema.response import VectorStoreResponse +from experiencemaker.schema.vector_node import VectorNode @OP_REGISTRY.register() @@ -19,12 +21,21 @@ class VectorStoreActionOp(BaseOp): result = self.vector_store.delete_workspace(workspace_id=request.workspace_id) elif request.action == "dump": + def node_to_experience(node: VectorNode) -> dict: + return vector_node_to_experience(node).model_dump() + result = self.vector_store.dump_workspace(workspace_id=request.workspace_id, - path=request.path) + path=request.path, + callback_fn=node_to_experience) elif request.action == "load": + def experience_dict_to_node(experience_dict: dict) -> VectorNode: + experience: BaseExperience = dict_to_experience(experience_dict=experience_dict) + return experience.to_vector_node() + result = self.vector_store.load_workspace(workspace_id=request.workspace_id, - path=request.path) + path=request.path, + callback_fn=experience_dict_to_node) else: raise ValueError(f"invalid action={request.action}") diff --git a/experiencemaker/schema/experience.py b/experiencemaker/schema/experience.py index f8573555..7d93643a 100644 --- a/experiencemaker/schema/experience.py +++ b/experiencemaker/schema/experience.py @@ -30,6 +30,17 @@ class BaseExperience(BaseModel, ABC): score: float | None = Field(default=None) metadata: ExperienceMeta = Field(default_factory=ExperienceMeta) + def to_vector_node(self) -> VectorNode: + raise NotImplementedError + + @classmethod + def from_vector_node(cls, node: VectorNode): + raise NotImplementedError + + +class TextExperience(BaseExperience): + experience_type: str = Field(default="text") + def to_vector_node(self) -> VectorNode: return VectorNode(unique_id=self.experience_id, workspace_id=self.workspace_id, @@ -41,14 +52,6 @@ class BaseExperience(BaseModel, ABC): "metadata": self.metadata.model_dump(), }) - @classmethod - def from_vector_node(cls, node: VectorNode): - raise NotImplementedError - - -class TextExperience(BaseExperience): - experience_type: str = Field(default="text") - @classmethod def from_vector_node(cls, node: VectorNode): return cls(workspace_id=node.workspace_id, @@ -106,6 +109,26 @@ def vector_node_to_experience(node: VectorNode) -> BaseExperience: logger.warning(f"experience type {experience_type} not supported") return TextExperience.from_vector_node(node) + +def dict_to_experience(experience_dict: dict) -> BaseExperience: + experience_type = experience_dict.get("experience_type", "text") + if experience_type == "text": + return TextExperience(**experience_dict) + + elif experience_type == "function": + return FuncExperience(**experience_dict) + + elif experience_type == "personal": + return PersonalExperience(**experience_dict) + + elif experience_type == "knowledge": + return KnowledgeExperience(**experience_dict) + + else: + logger.warning(f"experience type {experience_type} not supported") + return TextExperience(**experience_dict) + + if __name__ == "__main__": e1 = TextExperience( workspace_id="w_1024", diff --git a/experiencemaker/vector_store/base_vector_store.py b/experiencemaker/vector_store/base_vector_store.py index 22f5e70f..f16f2707 100644 --- a/experiencemaker/vector_store/base_vector_store.py +++ b/experiencemaker/vector_store/base_vector_store.py @@ -17,7 +17,7 @@ class BaseVectorStore(BaseModel, ABC): batch_size: int = Field(default=1024) @staticmethod - def _load_from_path(workspace_id: str, path: str | Path, **kwargs) -> Iterable[VectorNode]: + def _load_from_path(workspace_id: str, path: str | Path, callback_fn=None, **kwargs) -> Iterable[VectorNode]: workspace_path = Path(path) / f"{workspace_id}.jsonl" if not workspace_path.exists(): logger.warning(f"workspace_path={workspace_path} is not exists!") @@ -28,7 +28,11 @@ class BaseVectorStore(BaseModel, ABC): try: for line in tqdm(f, desc="load from path"): if line.strip(): - node = VectorNode(**json.loads(line.strip(), **kwargs)) + node_dict = json.loads(line.strip()) + if callback_fn: + node = callback_fn(node_dict) + else: + node = VectorNode(**node_dict, **kwargs) node.workspace_id = workspace_id yield node @@ -36,7 +40,7 @@ class BaseVectorStore(BaseModel, ABC): fcntl.flock(f, fcntl.LOCK_UN) @staticmethod - def _dump_to_path(nodes: Iterable[VectorNode], workspace_id: str, path: str | Path = "", + def _dump_to_path(nodes: Iterable[VectorNode], workspace_id: str, path: str | Path = "", callback_fn=None, ensure_ascii: bool = False, **kwargs): dump_path: Path = Path(path) dump_path.mkdir(parents=True, exist_ok=True) @@ -48,7 +52,12 @@ class BaseVectorStore(BaseModel, ABC): try: for node in tqdm(nodes, desc="dump to path"): node.workspace_id = workspace_id - f.write(json.dumps(node.model_dump(), ensure_ascii=ensure_ascii, **kwargs)) + if callback_fn: + node_dict = callback_fn(node) + else: + node_dict = node.model_dump() + assert isinstance(node_dict, dict) + f.write(json.dumps(node_dict, ensure_ascii=ensure_ascii, **kwargs)) f.write("\n") count += 1 @@ -68,16 +77,19 @@ class BaseVectorStore(BaseModel, ABC): def _iter_workspace_nodes(self, workspace_id: str, **kwargs) -> Iterable[VectorNode]: raise NotImplementedError - def dump_workspace(self, workspace_id: str, path: str | Path = "", **kwargs): + def dump_workspace(self, workspace_id: str, path: str | Path = "", callback_fn=None, **kwargs): if not self.exist_workspace(workspace_id=workspace_id, **kwargs): logger.warning(f"workspace_id={workspace_id} is not exist!") return {} return self._dump_to_path(nodes=self._iter_workspace_nodes(workspace_id=workspace_id, **kwargs), - workspace_id=workspace_id, - path=path, **kwargs) + workspace_id=workspace_id, + path=path, + callback_fn=callback_fn, + **kwargs) - def load_workspace(self, workspace_id: str, path: str | Path = "", nodes: List[VectorNode] = None, **kwargs): + def load_workspace(self, workspace_id: str, path: str | Path = "", nodes: List[VectorNode] = None, callback_fn=None, + **kwargs): if self.exist_workspace(workspace_id, **kwargs): self.delete_workspace(workspace_id=workspace_id, **kwargs) logger.info(f"delete workspace_id={workspace_id}") @@ -87,7 +99,7 @@ class BaseVectorStore(BaseModel, ABC): all_nodes: List[VectorNode] = [] if nodes: all_nodes.extend(nodes) - for node in self._load_from_path(path=path, workspace_id=workspace_id, **kwargs): + for node in self._load_from_path(path=path, workspace_id=workspace_id, callback_fn=callback_fn, **kwargs): all_nodes.append(node) self.insert(nodes=all_nodes, workspace_id=workspace_id, **kwargs) return {"size": len(all_nodes)}