change dump/load from vectornode to experience

This commit is contained in:
jinli.yl 2025-07-25 15:19:33 +08:00
parent 906d3d680a
commit a535bb2d8d
4 changed files with 69 additions and 22 deletions

View file

@ -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

View file

@ -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}")

View file

@ -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",

View file

@ -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)}