refactor(reme_ai): consolidate personal memory and update related ops

- Rename and restructure personal memory consolidation flow
- Update memory handling in context and operations
- Refactor memory schema to use time_created and time_modified fields
- Improve error handling and logging in memory operations
- Update test cases for new memory consolidation flow
This commit is contained in:
jinli.yl 2025-08-28 14:31:03 +08:00
parent 527e04140c
commit 074634d884
10 changed files with 175 additions and 104 deletions

View file

@ -1,10 +1,13 @@
# 代码框架
1. flowllm: 通过op pipeline的配置实现mcp接口的生成。
2. 重写Remy readme.
3. 迁移memoryscope/experiencemaker到flowllm的框架下
4. 迁移到新的op框架下
5. 重写memoryscope的cli-chat前端
6. 迁移memoryscope的文档
7. 迁移experiencemaker的文档
#
1. library 转化 @zouyin
2. index.html @jinli
3. reme_ai两个personal的调通
4. doc
1. readme @jiaji
2. experience maker @jiaji
3. personal @jinli
5. 新增op @zouyin
6. cookbook
1. appworld @jiaji P2
2. bfcl @zouyin P1
3. frozenlake @jiaji
4. simple_demo @jinli

View file

@ -18,7 +18,7 @@ http:
flow:
retrieve_task_memory:
flow_content: build_query_op >> recall_vector_store_op >> rerank_memory_op >> rewrite_memory_op
description: "Retrieve the most relevant top_k memory experience from historical memory based on the query to help solve tasks better now"
description: "Retrieves the most relevant top-k memory experiences from historical data based on the current query to enhance task-solving capabilities"
input_schema:
query:
type: "str"
@ -27,7 +27,7 @@ flow:
summary_task_memory:
flow_content: trajectory_preprocess_op >> (success_extraction_op|failure_extraction_op|comparative_extraction_op) >> memory_validation_op >> update_vector_store_op
description: "Summarize trajectories or messages into memories"
description: "Summarizes conversation trajectories or messages into structured memory representations for long-term storage"
input_schema:
trajectories:
type: "list"
@ -36,7 +36,7 @@ flow:
retrieve_task_memory_simple:
flow_content: build_query_op >> recall_vector_store_op >> merge_memory_op
description: "Retrieve the most relevant top_k memory experience from historical memory based on the query to help solve tasks better now"
description: "Retrieves the most relevant top-k memory experiences from historical data based on the current query with simplified processing"
input_schema:
query:
type: "str"
@ -45,7 +45,7 @@ flow:
summary_task_memory_simple:
flow_content: simple_summary_op >> update_vector_store_op
description: "Summarize trajectories or messages into memories"
description: "Summarizes conversation trajectories or messages into memories using a simplified approach"
input_schema:
trajectories:
type: "list"
@ -54,7 +54,7 @@ flow:
vector_store:
flow_content: vector_store_action_op
description: "directly operate the vector store."
description: "Directly operates on the vector store with various management actions"
input_schema:
action:
type: "str"
@ -64,16 +64,16 @@ flow:
retrieve_personal_memory:
flow_content: set_query_op >> (extract_time_op | (retrieve_memory_op >> semantic_rank_op)) >> fuse_rerank_op
description: "Retrieve the most relevant memories from historical memory based on the query to help answer better now."
description: "Retrieves the most relevant personal memories from historical data based on the query to enhance response quality"
input_schema:
query:
type: "str"
description: "user query"
required: true
consolidate_personal_memory:
summary_personal_memory:
flow_content: info_filter_op >> (get_observation_op | get_observation_with_time_op | load_today_memory_op) >> contra_repeat_op >> update_vector_store_op
description: "summary user's observation memory"
description: "Consolidates user observations and memories by filtering information and removing redundancies for efficient storage"
input_schema:
messages:
type: "list"
@ -88,7 +88,8 @@ flow:
llm:
default:
backend: openai_compatible
model_name: qwen3-30b-a3b-thinking-2507
# model_name: qwen3-30b-a3b-thinking-2507
model_name: qwen3-30b-a3b-instruct-2507
params:
temperature: 0.6

View file

@ -61,13 +61,12 @@ class SemanticRankOp(BaseLLMOp):
memory_list = ranked_memories
# Sort by score
memory_list = sorted(memory_list, key=lambda m: getattr(m, 'score', 0.0), reverse=True)
memory_list = sorted(memory_list, key=lambda m: m.score, reverse=True)
# Log top ranked memories
logger.info(f"Semantic ranking completed for query: {query[:50]}...")
for i, memory in enumerate(memory_list[:5]): # Log top 5
score = getattr(memory, 'score', 0.0)
logger.info(f"Top {i + 1}: Score={score:.3f}, Content={memory.content[:80]}...")
logger.info(f"Top {i + 1}: Score={memory.score:.3f}, Content={memory.content[:80]}...")
# Save ranked memories back to context
self.context.response.metadata["memory_list"] = memory_list

View file

@ -13,16 +13,16 @@ class BaseMemory(BaseModel, ABC):
when_to_use: str = Field(default="")
content: str | bytes = Field(default="")
score: float | None = Field(default=None)
score: float = Field(default=0)
created_time: str = Field(default_factory=lambda: datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
modified_time: str = Field(default_factory=lambda: datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
time_created: str = Field(default_factory=lambda: datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
time_modified: str = Field(default_factory=lambda: datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S"))
author: str = Field(default="")
metadata: dict = Field(default_factory=dict)
def update_modified_time(self):
self.modified_time = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def update_time_modified(self):
self.time_modified = datetime.datetime.now().strftime("%Y-%m-%d %H:%M:%S")
def to_vector_node(self) -> VectorNode:
raise NotImplementedError
@ -43,24 +43,25 @@ class TaskMemory(BaseMemory):
"memory_type": self.memory_type,
"content": self.content,
"score": self.score,
"created_time": self.created_time,
"modified_time": self.modified_time,
"time_created": self.time_created,
"time_modified": self.time_modified,
"author": self.author,
"metadata": self.metadata,
})
@classmethod
def from_vector_node(cls, node: VectorNode) -> "TaskMemory":
metadata = node.metadata.copy()
return cls(workspace_id=node.workspace_id,
memory_id=node.unique_id,
memory_type=node.metadata.get("memory_type"),
memory_type=metadata.pop("memory_type"),
when_to_use=node.content,
content=node.metadata.get("content"),
score=node.metadata.get("score"),
created_time=node.metadata.get("created_time"),
modified_time=node.metadata.get("modified_time"),
author=node.metadata.get("author"),
metadata=node.metadata.get("metadata"))
content=metadata.pop("content"),
score=metadata.pop("score"),
time_created=metadata.pop("time_created"),
time_modified=metadata.pop("time_modified"),
author=metadata.pop("author"),
metadata=metadata.pop("metadata", {}))
class PersonalMemory(BaseMemory):
@ -78,26 +79,27 @@ class PersonalMemory(BaseMemory):
"target": self.target,
"reflection_subject": self.reflection_subject,
"score": self.score,
"created_time": self.created_time,
"modified_time": self.modified_time,
"time_created": self.time_created,
"time_modified": self.time_modified,
"author": self.author,
"metadata": self.metadata,
})
@classmethod
def from_vector_node(cls, node: VectorNode) -> "PersonalMemory":
metadata = node.metadata.copy()
return cls(workspace_id=node.workspace_id,
memory_id=node.unique_id,
memory_type=node.metadata.get("memory_type"),
memory_type=metadata.pop("memory_type"),
when_to_use=node.content,
content=node.metadata.get("content"),
target=node.metadata.get("target", ""),
reflection_subject=node.metadata.get("reflection_subject", ""),
score=node.metadata.get("score"),
created_time=node.metadata.get("created_time"),
modified_time=node.metadata.get("modified_time"),
author=node.metadata.get("author"),
metadata=node.metadata.get("metadata"))
content=metadata.pop("content"),
target=metadata.pop("target", ""),
reflection_subject=metadata.pop("reflection_subject", ""),
score=metadata.pop("score"),
time_created=metadata.pop("time_created"),
time_modified=metadata.pop("time_modified"),
author=metadata.pop("author"),
metadata=metadata.pop("metadata", {}))
def vector_node_to_memory(node: VectorNode) -> BaseMemory:

View file

@ -34,9 +34,9 @@ class ContraRepeatOp(BaseLLMOp):
"""
# Get memory list from context - standardized key
memory_list: List[BaseMemory] = []
memory_list.extend(self.context.observation_memories)
memory_list.extend(self.context.observation_memories_with_time)
memory_list.extend(self.context.today_memories)
memory_list.extend(self.context.get("observation_memories", []))
memory_list.extend(self.context.get("observation_memories_with_time", []))
memory_list.extend(self.context.get("today_memories", []))
self.context.response.metadata["memory_list"] = memory_list
@ -55,7 +55,7 @@ class ContraRepeatOp(BaseLLMOp):
return
# Sort and limit memories by count
sorted_memories = sorted(memory_list, key=lambda x: x.created_time, reverse=True)[:contra_repeat_max_count]
sorted_memories = sorted(memory_list, key=lambda x: x.time_created, reverse=True)[:contra_repeat_max_count]
if len(sorted_memories) <= 1:
logger.info("sorted_memories.size<=1, stop.")

View file

@ -19,6 +19,7 @@ class InfoFilterOp(BaseLLMOp):
def execute(self):
"""Filter messages based on information content scores"""
# Get messages from context - guaranteed to exist by flow input
self.context.messages = [Message(**x) if isinstance(x, dict) else x for x in self.context.messages]
messages: List[Message] = self.context.messages
if not messages:
logger.warning("No messages found in context")
@ -33,7 +34,7 @@ class InfoFilterOp(BaseLLMOp):
info_messages = self._filter_and_process_messages(messages, user_name, info_filter_msg_max_size)
if not info_messages:
logger.warning("No messages left after filtering")
self.context.response.metadata["memory_list"] = []
self.context.messages = []
return
logger.info(f"Filtering {len(info_messages)} messages for information content")
@ -42,7 +43,7 @@ class InfoFilterOp(BaseLLMOp):
filtered_memories = self._filter_messages_with_llm(info_messages, user_name, preserved_scores)
# Store results in context using standardized key
self.context.response.metadata["memory_list"] = filtered_memories
self.context.messages = filtered_memories
logger.info(f"Filtered to {len(filtered_memories)} high-information messages")
@staticmethod
@ -117,28 +118,25 @@ class InfoFilterOp(BaseLLMOp):
# Check if score should be preserved
if score in preserved_scores:
message_obj = info_messages[msg_idx]
# Get original message metadata or create empty dict
original_metadata = getattr(message_obj, 'metadata', {}) or {}
message = info_messages[msg_idx]
# Create memory from filtered message with combined metadata
memory = PersonalMemory(
workspace_id=self.context.get("workspace_id", ""),
content=message_obj.content,
content=message.content,
target=user_name,
author=getattr(self.llm, "model_name", "system"),
metadata={
"info_score": score,
"filter_type": "info_content",
"original_message_time": getattr(message_obj, 'time_created', None),
"role_name": original_metadata.get('role_name', user_name),
"original_message_time": getattr(message, 'time_created', None),
"role_name": message.metadata.pop("role_name", user_name),
"memorized": True,
**original_metadata # Include all original metadata
**message.metadata # Include all original metadata
}
)
filtered_memories.append(memory)
logger.info(f"Info filter: kept message with score {score}: {message_obj.content[:50]}...")
logger.info(f"Info filter: kept message with score {score}: {message.content[:50]}...")
return filtered_memories

View file

@ -50,7 +50,7 @@ class LongContraRepeatOp(BaseLLMOp):
# Sort memories by creation time (most recent first) and limit count
sorted_memories = sorted(
updated_insights,
key=lambda x: getattr(x, 'created_time', ''),
key=lambda x: x.time_created,
reverse=True
)[:max_memories_to_process]
@ -151,7 +151,7 @@ class LongContraRepeatOp(BaseLLMOp):
author=memory.author,
metadata={**memory.metadata, 'modified_by': 'long_contra_repeat'}
)
modified_memory.update_modified_time()
modified_memory.update_time_modified()
filtered_memories.append(modified_memory)
logger.info(f"Modified contradictory memory {idx}: {modified_content.strip()[:50]}...")
else:

View file

@ -201,7 +201,7 @@ class UpdateInsightOp(BaseLLMOp):
"update_reason": "integrated_new_observations"
}
)
updated_insight.update_modified_time()
updated_insight.update_time_modified()
logger.info(f"Updated insight: {updated_content[:50]}...")
return updated_insight

View file

@ -17,9 +17,9 @@ class MemoryValidationOp(BaseLLMOp):
"""Validate quality of extracted task memories"""
task_memories: List[BaseMemory] = []
task_memories.extend(self.context.success_task_memories)
task_memories.extend(self.context.failure_task_memories)
task_memories.extend(self.context.comparative_task_memories)
task_memories.extend(self.context.get("success_task_memories", []))
task_memories.extend(self.context.get("failure_task_memories", []))
task_memories.extend(self.context.get("comparative_task_memories", []))
if not task_memories:
logger.info("No task memories found for validation")

View file

@ -1,10 +1,110 @@
import asyncio
import json
import aiohttp
base_url = "http://0.0.0.0:8002"
async def run1(session):
workspace_id = "default1"
async with session.post(
f"{base_url}/vector_store",
json={
"action": "delete",
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
trajectories = [
{
"task_id": "t1",
"messages": [
{"role": "user", "content": "搜索可以使用websearch工具"}
],
"score": 1,
},
{
"task_id": "t1",
"messages": [
{"role": "user", "content": "搜索可以使用code工具"}
],
"score": 0,
}
]
async with session.post(
# f"{base_url}/summary_task_memory",
f"{base_url}/summary_task_memory_simple",
json={
"trajectories": trajectories,
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
await asyncio.sleep(2)
async with session.post(
# f"{base_url}/retrieve_task_memory",
f"{base_url}/retrieve_task_memory_simple",
json={
"query": "茅台怎么样?",
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
async def run2(session):
workspace_id = "default2"
async with session.post(
f"{base_url}/vector_store",
json={
"action": "delete",
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
messages = [{"role": "user", "content": "我喜欢吃西瓜🍉"}]
async with session.post(
f"{base_url}/summary_personal_memory",
json={
"messages": messages,
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
await asyncio.sleep(2)
async with session.post(
f"{base_url}/retrieve_personal_memory",
json={
"query": "茅台怎么样?",
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(json.dumps(result, ensure_ascii=False))
async def main():
base_url = "http://0.0.0.0:8002"
async with aiohttp.ClientSession() as session:
# 获取工具列表
@ -14,45 +114,13 @@ async def main():
tools = await response.json()
print("可用工具:")
for tool in tools:
print(tool)
print(json.dumps(tool, ensure_ascii=False))
else:
print(f"获取工具列表失败: {response.status}")
return
workspace_id = "default1"
trajectories = [
{
"task_id": "t1",
"messages": [
{"role": "user", "content": "搜索可以使用websearch工具"}
],
"score": 0.9,
}
]
async with session.post(
f"{base_url}/summary_task_memory_simple",
json={
"trajectories": trajectories,
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(result)
async with session.post(
f"{base_url}/retrieve_task_memory_simple",
json={
"query": "茅台怎么样?",
"workspace_id": workspace_id,
},
headers={"Content-Type": "application/json"}
) as response:
result = await response.json()
print(result)
# await run1(session)
await run2(session)
if __name__ == "__main__":
asyncio.run(main())