From 77ce2d22d341f977443d4919cf3fc8d4f9ae37e0 Mon Sep 17 00:00:00 2001 From: "jinli.yl" Date: Tue, 26 Aug 2025 23:05:49 +0800 Subject: [PATCH] refactor(summary): rename experiences to task memories and optimize related ops - Rename ExperienceDeduplicationOp to TaskMemoryDeduplicationOp - Rename ExperienceValidationOp to TaskMemoryValidationOp - Update comparative extraction, failure extraction, and success extraction ops to use task memories - Refactor PDFPreprocessOp and ReactV1Op for better code structure - Remove unused simple_config.yaml - Update default.yaml with new task memory related flows --- reme_ai/config/default.yaml | 19 ++- reme_ai/config/simple_config.yaml | 61 ------- reme_ai/react/react_v1_op.py | 19 +-- reme_ai/schema/__init__.py | 2 +- .../summary/task/comparative_extraction_op.py | 91 +++++----- .../task/comparative_extraction_prompt.yaml | 4 +- .../task/experience_deduplication_op.py | 160 ------------------ .../summary/task/experience_validation_op.py | 57 ++++--- .../task/experience_validation_prompt.yaml | 22 +-- reme_ai/summary/task/failure_extraction_op.py | 61 ++++--- .../task/failure_extraction_prompt.yaml | 4 +- .../summary/task/pdf_preprocess_op_wrapper.py | 7 +- .../task/simple_comparative_summary_op.py | 36 ++-- reme_ai/summary/task/success_extraction_op.py | 61 ++++--- .../task/success_extraction_prompt.yaml | 4 +- .../task/task_memory_deduplication_op.py | 160 ++++++++++++++++++ .../summary/task/trajectory_preprocess_op.py | 9 +- .../task/trajectory_segmentation_op.py | 6 +- 18 files changed, 369 insertions(+), 414 deletions(-) delete mode 100644 reme_ai/config/simple_config.yaml delete mode 100644 reme_ai/summary/task/experience_deduplication_op.py create mode 100644 reme_ai/summary/task/task_memory_deduplication_op.py diff --git a/reme_ai/config/default.yaml b/reme_ai/config/default.yaml index 26eafc6e..bd6c213c 100644 --- a/reme_ai/config/default.yaml +++ b/reme_ai/config/default.yaml @@ -34,7 +34,24 @@ flow: description: "A list of conversation trajectory information, including message content and score. This field does not need to be filled in, the system will complete it automatically." required: false -# task_summarizer: simple_summary_op->update_vector_store_op + 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" + input_schema: + query: + type: "str" + description: "current query" + required: true + + summary_task_memory: + flow_content: trajectory_preprocess_op->[success_extraction_op|failure_extraction_op|comparative_extraction_op]->experience_validation_op->experience_deduplication_op->update_vector_store_op + description: "Summarize trajectories or messages into memories" + input_schema: + trajectories: + type: "list" + description: "A list of conversation trajectory information, including message content and score. This field does not need to be filled in, the system will complete it automatically." + required: false + # vector_store: vector_store_action_op # agent: react_op diff --git a/reme_ai/config/simple_config.yaml b/reme_ai/config/simple_config.yaml deleted file mode 100644 index 6a0deb20..00000000 --- a/reme_ai/config/simple_config.yaml +++ /dev/null @@ -1,61 +0,0 @@ -# demo config.yaml - -http_service: - host: "0.0.0.0" - port: 8001 - timeout_keep_alive: 600 - limit_concurrency: 64 - -thread_pool: - max_workers: 64 - -api: - retriever: build_query_op->recall_vector_store_op->merge_experience_op - summarizer: simple_summary_op->update_vector_store_op - vector_store: vector_store_action_op - agent: react_op - -op: - build_query_op: - backend: build_query_op - llm: default - params: - enable_llm_build: false - recall_vector_store_op: - backend: recall_vector_store_op - vector_store: default - merge_experience_op: - backend: merge_experience_op - simple_summary_op: - backend: simple_summary_op - llm: default - params: - success_score_threshold: 0.9 - update_vector_store_op: - backend: update_vector_store_op - vector_store: default - vector_store_action_op: - backend: vector_store_action_op - vector_store: default - react_op: - backend: react_v1_op - llm: default - -llm: - default: - backend: openai_compatible - model_name: qwen-max-2025-01-25 - params: - temperature: 0.6 - -embedding_model: - default: - backend: openai_compatible - model_name: text-embedding-v4 - params: - dimensions: 1024 - -vector_store: - default: - backend: elasticsearch - embedding_model: default diff --git a/reme_ai/react/react_v1_op.py b/reme_ai/react/react_v1_op.py index 489b5d68..41b612c3 100644 --- a/reme_ai/react/react_v1_op.py +++ b/reme_ai/react/react_v1_op.py @@ -2,25 +2,24 @@ import datetime import time from typing import List, Dict +from flowllm import C, BaseLLMOp +from flowllm.flow.base_tool_flow import BaseToolFlow +from flowllm.flow.gallery import DashscopeSearchToolFlow, CodeToolFlow, TerminateToolFlow from loguru import logger -from flowllm import C, BaseLLMOp -from reme_ai.schema.message import Message, Role -from experiencemaker.tool import TOOL_REGISTRY -from experiencemaker.tool.base_tool import BaseTool +from reme_ai.schema import Message, Role @C.register_op() class ReactV1Op(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): query: str = self.context.query max_steps: int = int(self.op_params.get("max_steps", 10)) - tool_names = self.op_params.get("tool_names", "code_tool,tavily_search_tool,terminate_tool") - tools: List[BaseTool] = [TOOL_REGISTRY[x.strip()]() for x in tool_names.split(",") if x] - tool_dict: Dict[str, BaseTool] = {x.name: x for x in tools} + tools: List[BaseToolFlow] = [DashscopeSearchToolFlow(), CodeToolFlow(), TerminateToolFlow()] + tool_dict: Dict[str, BaseToolFlow] = {x.name: x for x in tools} now_time = datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S') has_terminate_tool = False @@ -35,7 +34,7 @@ class ReactV1Op(BaseLLMOp): if has_terminate_tool: assistant_message: Message = self.llm.chat(messages) else: - assistant_message: Message = self.llm.chat(messages, tools=tools) + assistant_message: Message = self.llm.chat(messages, tools=[x.tool_call for x in tools]) messages.append(assistant_message) logger.info(f"assistant.{i}.reasoning_content={assistant_message.reasoning_content}\n" @@ -61,7 +60,7 @@ class ReactV1Op(BaseLLMOp): if tool_call.name not in tool_dict: continue - self.submit_task(tool_dict[tool_call.name].execute, **tool_call.argument_dict) + self.submit_task(tool_dict[tool_call.name].__call__, **tool_call.argument_dict) time.sleep(1) if not has_terminate_tool: diff --git a/reme_ai/schema/__init__.py b/reme_ai/schema/__init__.py index 7089e89b..adb4a0af 100644 --- a/reme_ai/schema/__init__.py +++ b/reme_ai/schema/__init__.py @@ -1 +1 @@ -from flowllm.schema.message import Message, Trajectory, ToolCall \ No newline at end of file +from flowllm.schema.message import Message, Trajectory, ToolCall, Role \ No newline at end of file diff --git a/reme_ai/summary/task/comparative_extraction_op.py b/reme_ai/summary/task/comparative_extraction_op.py index 0a7309c8..14d52550 100644 --- a/reme_ai/summary/task/comparative_extraction_op.py +++ b/reme_ai/summary/task/comparative_extraction_op.py @@ -1,33 +1,33 @@ -import json -import re from typing import List, Tuple, Optional -from loguru import logger from flowllm import C, BaseLLMOp -from reme_ai.schema.memory import BaseMemory, TaskMemory +from loguru import logger + from reme_ai.schema import Message, Trajectory -from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience_response, get_trajectory_context +from reme_ai.schema.memory import BaseMemory, TaskMemory +from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience_response @C.register_op() class ComparativeExtractionOp(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): - """Extract comparative experiences by comparing different scoring trajectories""" + """Extract comparative task memories by comparing different scoring trajectories""" all_trajectories: List[Trajectory] = self.context.get("all_trajectories", []) success_trajectories: List[Trajectory] = self.context.get("success_trajectories", []) failure_trajectories: List[Trajectory] = self.context.get("failure_trajectories", []) - comparative_experiences = [] + comparative_task_memories = [] # Soft comparison: highest score vs lowest score if len(all_trajectories) >= 2 and self.op_params.get("enable_soft_comparison", True): highest_traj, lowest_traj = self._find_highest_lowest_scoring_trajectories(all_trajectories) if highest_traj and lowest_traj and highest_traj.score > lowest_traj.score: - logger.info(f"Extracting soft comparative experiences: highest ({highest_traj.score:.2f}) vs lowest ({lowest_traj.score:.2f})") - soft_experiences = self._extract_soft_comparative_experience(highest_traj, lowest_traj) - comparative_experiences.extend(soft_experiences) + logger.info( + f"Extracting soft comparative task memories: highest ({highest_traj.score:.2f}) vs lowest ({lowest_traj.score:.2f})") + soft_task_memories = self._extract_soft_comparative_task_memory(highest_traj, lowest_traj) + comparative_task_memories.extend(soft_task_memories) # Hard comparison: success vs failure (if similarity search is enabled) if (success_trajectories and failure_trajectories and @@ -37,13 +37,14 @@ class ComparativeExtractionOp(BaseLLMOp): logger.info(f"Found {len(similar_pairs)} similar pairs for hard comparison") for success_steps, failure_steps, similarity_score in similar_pairs: - hard_experiences = self._extract_hard_comparative_experience(success_steps, failure_steps, similarity_score) - comparative_experiences.extend(hard_experiences) + hard_task_memories = self._extract_hard_comparative_task_memory(success_steps, failure_steps, + similarity_score) + comparative_task_memories.extend(hard_task_memories) - logger.info(f"Extracted {len(comparative_experiences)} comparative experiences") + logger.info(f"Extracted {len(comparative_task_memories)} comparative task memories") - # Add experiences to context - self.context.comparative_experiences = comparative_experiences + # Add task memories to context + self.context.comparative_task_memories = comparative_task_memories def _find_highest_lowest_scoring_trajectories(self, trajectories: List[Trajectory]) -> Tuple[Optional[Trajectory], Optional[Trajectory]]: """Find the highest and lowest scoring trajectories""" @@ -69,67 +70,67 @@ class ComparativeExtractionOp(BaseLLMOp): """Get trajectory score""" return trajectory.score - def _extract_soft_comparative_experience(self, higher_traj: Trajectory, lower_traj: Trajectory) -> List[BaseMemory]: - """Extract soft comparative experience (high score vs low score)""" + def _extract_soft_comparative_task_memory(self, higher_traj: Trajectory, lower_traj: Trajectory) -> List[ + BaseMemory]: + """Extract soft comparative task memory (high score vs low score)""" higher_steps = self._get_trajectory_steps(higher_traj) lower_steps = self._get_trajectory_steps(lower_traj) higher_score = self._get_trajectory_score(higher_traj) lower_score = self._get_trajectory_score(lower_traj) prompt = self.prompt_format( - prompt_name="soft_comparative_step_experience_prompt", + prompt_name="soft_comparative_step_task_memory_prompt", higher_steps=merge_messages_content(higher_steps), lower_steps=merge_messages_content(lower_steps), higher_score=f"{higher_score:.2f}", lower_score=f"{lower_score:.2f}" ) - def parse_experiences(message: Message) -> List[BaseMemory]: - experiences_data = parse_json_experience_response(message.content) - experiences = [] + def parse_task_memories(message: Message) -> List[BaseMemory]: + task_memories_data = parse_json_experience_response(message.content) + task_memories = [] - for exp_data in experiences_data: - experience = TaskMemory( + for tm_data in task_memories_data: + task_memory = TaskMemory( workspace_id=self.context.get("workspace_id", ""), - when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), - content=exp_data.get("experience", ""), + when_to_use=tm_data.get("when_to_use", tm_data.get("condition", "")), + content=tm_data.get("experience", ""), author=getattr(self.llm, 'model_name', 'system'), - metadata=exp_data + metadata=tm_data ) - experiences.append(experience) + task_memories.append(task_memory) - return experiences + return task_memories - return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_task_memories) - - def _extract_hard_comparative_experience(self, success_steps: List[Message], + def _extract_hard_comparative_task_memory(self, success_steps: List[Message], failure_steps: List[Message], similarity_score: float) -> List[BaseMemory]: - """Extract hard comparative experience (success vs failure)""" + """Extract hard comparative task memory (success vs failure)""" prompt = self.prompt_format( - prompt_name="comparative_step_experience_prompt", + prompt_name="comparative_step_task_memory_prompt", success_steps=merge_messages_content(success_steps), failure_steps=merge_messages_content(failure_steps), similarity_score=similarity_score ) - def parse_experiences(message: Message) -> List[BaseMemory]: - experiences_data = parse_json_experience_response(message.content) - experiences = [] + def parse_task_memories(message: Message) -> List[BaseMemory]: + task_memories_data = parse_json_experience_response(message.content) + task_memories = [] - for exp_data in experiences_data: - experience = TaskMemory( + for tm_data in task_memories_data: + task_memory = TaskMemory( workspace_id=self.context.get("workspace_id", ""), - when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), - content=exp_data.get("experience", ""), + when_to_use=tm_data.get("when_to_use", tm_data.get("condition", "")), + content=tm_data.get("experience", ""), author=getattr(self.llm, 'model_name', 'system'), - metadata=exp_data + metadata=tm_data ) - experiences.append(experience) + task_memories.append(task_memory) - return experiences + return task_memories - return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_task_memories) def _get_trajectory_steps(self, trajectory: Trajectory) -> List[Message]: diff --git a/reme_ai/summary/task/comparative_extraction_prompt.yaml b/reme_ai/summary/task/comparative_extraction_prompt.yaml index 8b7724b7..c817ac67 100644 --- a/reme_ai/summary/task/comparative_extraction_prompt.yaml +++ b/reme_ai/summary/task/comparative_extraction_prompt.yaml @@ -1,4 +1,4 @@ -soft_comparative_step_experience_prompt: | +soft_comparative_step_task_memory_prompt: | You are an expert AI analyst comparing higher-scoring and lower-scoring step sequences to extract performance insights. Your task is to identify the key differences between higher and lower performing approaches at the step level. @@ -38,7 +38,7 @@ soft_comparative_step_experience_prompt: | ] ``` -comparative_step_experience_prompt: | +comparative_step_task_memory_prompt: | You are an expert AI analyst comparing successful and failed step sequences to extract differential insights. Your task is to identify the key differences between success and failure patterns at the step level. diff --git a/reme_ai/summary/task/experience_deduplication_op.py b/reme_ai/summary/task/experience_deduplication_op.py deleted file mode 100644 index 5b57d8be..00000000 --- a/reme_ai/summary/task/experience_deduplication_op.py +++ /dev/null @@ -1,160 +0,0 @@ -from typing import List -from loguru import logger - -from flowllm import C, BaseOp -from reme_ai.schema.memory import BaseMemory - - -@C.register_op() -class ExperienceDeduplicationOp(BaseOp): - current_path: str = __file__ - - def execute(self): - """Remove duplicate experiences""" - # Get experiences to deduplicate - experiences: List[BaseMemory] = self.context.get("experiences", []) - - if not experiences: - logger.info("No experiences found for deduplication") - return - - logger.info(f"Starting deduplication for {len(experiences)} experiences") - - # Perform deduplication - deduplicated_experiences = self._deduplicate_experiences(experiences) - - logger.info(f"Deduplication complete: {len(deduplicated_experiences)} deduplicated experiences out of {len(experiences)}") - - # Update context - self.context.deduplicated_experiences = deduplicated_experiences - - def _deduplicate_experiences(self, experiences: List[BaseMemory]) -> List[BaseMemory]: - """Remove duplicate experiences""" - if not experiences: - return experiences - - similarity_threshold = self.op_params.get("similarity_threshold", 0.5) - workspace_id = self.context.get("workspace_id") - - unique_experiences = [] - - # Get existing experience embeddings - existing_embeddings = self._get_existing_experience_embeddings(workspace_id) - - for experience in experiences: - # Generate embedding for current experience - current_embedding = self._get_experience_embedding(experience) - - if current_embedding is None: - logger.warning(f"Failed to generate embedding for experience: {str(experience.when_to_use)[:50]}...") - continue - - # Check similarity with existing experiences - if self._is_similar_to_existing_experiences(current_embedding, existing_embeddings, similarity_threshold): - logger.debug(f"Skipping similar experience: {str(experience.when_to_use)[:50]}...") - continue - - # Check similarity with current batch experiences - if self._is_similar_to_current_experiences(current_embedding, unique_experiences, similarity_threshold): - logger.debug(f"Skipping duplicate in current batch: {str(experience.when_to_use)[:50]}...") - continue - - # Add to unique experiences list - unique_experiences.append(experience) - logger.debug(f"Added unique experience: {str(experience.when_to_use)[:50]}...") - - return unique_experiences - - def _get_existing_experience_embeddings(self, workspace_id: str) -> List[List[float]]: - """Get embeddings of existing experiences""" - try: - if not hasattr(self.context, 'vector_store') or not self.context.vector_store or not workspace_id: - return [] - - # Query existing experience nodes - existing_nodes = self.context.vector_store.search( - query="...", # Empty query to get all - workspace_id=workspace_id, - top_k=self.op_params.get("max_existing_experiences", 1000) - ) - - # Extract embeddings - existing_embeddings = [] - for node in existing_nodes: - if hasattr(node, 'embedding') and node.embedding: - existing_embeddings.append(node.embedding) - - logger.debug(f"Retrieved {len(existing_embeddings)} existing experience embeddings from workspace {workspace_id}") - return existing_embeddings - - except Exception as e: - logger.warning(f"Failed to retrieve existing experience embeddings: {e}") - return [] - - def _get_experience_embedding(self, experience: BaseMemory) -> List[float]: - """Generate embedding for experience""" - try: - if not hasattr(self.context, 'vector_store') or not self.context.vector_store: - return None - - # Combine experience description and content for embedding - text_for_embedding = f"{experience.when_to_use} {experience.content}" - embeddings = self.context.vector_store.embedding_model.get_embeddings([text_for_embedding]) - - if embeddings and len(embeddings) > 0: - return embeddings[0] - else: - logger.warning("Empty embedding generated for experience") - return None - - except Exception as e: - logger.error(f"Error generating embedding for experience: {e}") - return None - - - def _is_similar_to_existing_experiences(self, current_embedding: List[float], - existing_embeddings: List[List[float]], - threshold: float) -> bool: - """Check if current embedding is similar to existing embeddings""" - for existing_embedding in existing_embeddings: - similarity = self._calculate_cosine_similarity(current_embedding, existing_embedding) - if similarity > threshold: - logger.debug(f"Found similar existing experience with similarity: {similarity:.3f}") - return True - return False - - def _is_similar_to_current_experiences(self, current_embedding: List[float], - current_experiences: List[BaseMemory], - threshold: float) -> bool: - for existing_experience in current_experiences: - existing_embedding = self._get_experience_embedding(existing_experience) - if existing_embedding is None: - continue - - similarity = self._calculate_cosine_similarity(current_embedding, existing_embedding) - if similarity > threshold: - logger.debug(f"Found similar experience in current batch with similarity: {similarity:.3f}") - return True - return False - - def _calculate_cosine_similarity(self, embedding1: List[float], embedding2: List[float]) -> float: - """Calculate cosine similarity""" - try: - import numpy as np - - vec1 = np.array(embedding1) - vec2 = np.array(embedding2) - - # Calculate cosine similarity - dot_product = np.dot(vec1, vec2) - norm1 = np.linalg.norm(vec1) - norm2 = np.linalg.norm(vec2) - - if norm1 == 0 or norm2 == 0: - return 0.0 - - return dot_product / (norm1 * norm2) - - except Exception as e: - logger.error(f"Error calculating cosine similarity: {e}") - return 0.0 \ No newline at end of file diff --git a/reme_ai/summary/task/experience_validation_op.py b/reme_ai/summary/task/experience_validation_op.py index 22590d2e..bf959a85 100644 --- a/reme_ai/summary/task/experience_validation_op.py +++ b/reme_ai/summary/task/experience_validation_op.py @@ -1,56 +1,57 @@ +import json import re from typing import List, Dict, Any -from loguru import logger -import json from flowllm import C, BaseLLMOp +from loguru import logger + +from reme_ai.schema import Message from reme_ai.schema.memory import BaseMemory -from reme_ai.schema.message import Message @C.register_op() -class ExperienceValidationOp(BaseLLMOp): - current_path: str = __file__ +class TaskMemoryValidationOp(BaseLLMOp): + file_path: str = __file__ def execute(self): - """Validate quality of extracted experiences""" - experiences: List[BaseMemory] = self.context.get("experiences", []) - - if not experiences: - logger.info("No experiences found for validation") + """Validate quality of extracted task memories""" + task_memories: List[BaseMemory] = self.context.get("task_memories", []) + + if not task_memories: + logger.info("No task memories found for validation") return - logger.info(f"Validating {len(experiences)} extracted experiences") + logger.info(f"Validating {len(task_memories)} extracted task memories") - # Validate experiences - validated_experiences = [] - - for experience in experiences: - validation_result = self._validate_single_experience(experience) + # Validate task memories + validated_task_memories = [] + + for task_memory in task_memories: + validation_result = self._validate_single_task_memory(task_memory) if validation_result and validation_result.get("is_valid", False): - validated_experiences.append(experience) + validated_task_memories.append(task_memory) else: reason = validation_result.get("reason", "Unknown reason") if validation_result else "Validation failed" - logger.warning(f"Experience validation failed: {reason}") + logger.warning(f"Task memory validation failed: {reason}") - logger.info(f"Validated {len(validated_experiences)} out of {len(experiences)} experiences") + logger.info(f"Validated {len(validated_task_memories)} out of {len(task_memories)} task memories") # Update context - self.context.validated_experiences = validated_experiences + self.context.validated_task_memories = validated_task_memories - def _validate_single_experience(self, experience: BaseMemory) -> Dict[str, Any]: - """Validate single experience""" - validation_info = self._llm_validate_experience(experience) + def _validate_single_task_memory(self, task_memory: BaseMemory) -> Dict[str, Any]: + """Validate single task memory""" + validation_info = self._llm_validate_task_memory(task_memory) logger.info(f"Validating: {validation_info}") return validation_info - def _llm_validate_experience(self, experience: BaseMemory) -> Dict[str, Any]: - """Validate experience using LLM""" + def _llm_validate_task_memory(self, task_memory: BaseMemory) -> Dict[str, Any]: + """Validate task memory using LLM""" try: prompt = self.prompt_format( - prompt_name="experience_validation_prompt", - condition=experience.when_to_use, - experience_content=experience.content + prompt_name="task_memory_validation_prompt", + condition=task_memory.when_to_use, + task_memory_content=task_memory.content ) def parse_validation(message: Message) -> Dict[str, Any]: diff --git a/reme_ai/summary/task/experience_validation_prompt.yaml b/reme_ai/summary/task/experience_validation_prompt.yaml index 408bc834..681c78ac 100644 --- a/reme_ai/summary/task/experience_validation_prompt.yaml +++ b/reme_ai/summary/task/experience_validation_prompt.yaml @@ -1,18 +1,18 @@ -experience_validation_prompt: | - You are an expert AI analyst tasked with validating the quality and usefulness of extracted step-level experiences. +task_memory_validation_prompt: | + You are an expert AI analyst tasked with validating the quality and usefulness of extracted step-level task memories. - Your task is to assess whether the extracted experience is actionable, accurate, and valuable for future agent executions. + Your task is to assess whether the extracted task memory is actionable, accurate, and valuable for future agent executions. VALIDATION CRITERIA: - ● ACTIONABILITY: Is the experience specific enough to guide future actions? - ● ACCURACY: Does the experience correctly reflect the patterns observed? - ● RELEVANCE: Is the experience applicable to similar future scenarios? - ● CLARITY: Is the experience clearly articulated and understandable? - ● UNIQUENESS: Does the experience provide novel insights or common knowledge? + ● ACTIONABILITY: Is the task memory specific enough to guide future actions? + ● ACCURACY: Does the task memory correctly reflect the patterns observed? + ● RELEVANCE: Is the task memory applicable to similar future scenarios? + ● CLARITY: Is the task memory clearly articulated and understandable? + ● UNIQUENESS: Does the task memory provide novel insights or common knowledge? - # Experience to Validate + # Task Memory to Validate Condition: {condition} - Experience Content: {experience_content} + Task Memory Content: {task_memory_content} OUTPUT FORMAT: Provide validation assessment: @@ -26,4 +26,4 @@ experience_validation_prompt: | ``` Score should be between 0.0 (poor quality) and 1.0 (excellent quality). - Mark as invalid if score is below 0.3 or if there are fundamental issues with the experience. + Mark as invalid if score is below 0.3 or if there are fundamental issues with the task memory. diff --git a/reme_ai/summary/task/failure_extraction_op.py b/reme_ai/summary/task/failure_extraction_op.py index 2220ec7d..e64f0baf 100644 --- a/reme_ai/summary/task/failure_extraction_op.py +++ b/reme_ai/summary/task/failure_extraction_op.py @@ -1,74 +1,73 @@ -import json -import re from typing import List -from loguru import logger from flowllm import C, BaseLLMOp +from loguru import logger + +from reme_ai.schema import Message, Trajectory from reme_ai.schema.memory import BaseMemory, TaskMemory -from reme_ai.schema.message import Message, Trajectory -from reme_ai.utils.memory_utils import merge_messages_content, parse_json_experience_response, get_trajectory_context +from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience_response, get_trajectory_context @C.register_op() class FailureExtractionOp(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): - """Extract experiences from failed trajectories""" + """Extract task memories from failed trajectories""" failure_trajectories: List[Trajectory] = self.context.get("failure_trajectories", []) if not failure_trajectories: logger.info("No failure trajectories found for extraction") return - logger.info(f"Extracting experiences from {len(failure_trajectories)} failed trajectories") + logger.info(f"Extracting task memories from {len(failure_trajectories)} failed trajectories") - failure_experiences = [] + failure_task_memories = [] # Process trajectories for trajectory in failure_trajectories: if hasattr(trajectory, 'segments') and trajectory.segments: # Process segmented step sequences for segment in trajectory.segments: - experiences = self._extract_failure_experience_from_steps(segment, trajectory) - failure_experiences.extend(experiences) + task_memories = self._extract_failure_task_memory_from_steps(segment, trajectory) + failure_task_memories.extend(task_memories) else: # Process entire trajectory - experiences = self._extract_failure_experience_from_steps(trajectory.messages, trajectory) - failure_experiences.extend(experiences) + task_memories = self._extract_failure_task_memory_from_steps(trajectory.messages, trajectory) + failure_task_memories.extend(task_memories) - logger.info(f"Extracted {len(failure_experiences)} failure experiences") - - # Add experiences to context - self.context.failure_experiences = failure_experiences + logger.info(f"Extracted {len(failure_task_memories)} failure task memories") - def _extract_failure_experience_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[BaseMemory]: - """Extract experience from failed step sequences""" + # Add task memories to context + self.context.failure_task_memories = failure_task_memories + + def _extract_failure_task_memory_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[BaseMemory]: + """Extract task memory from failed step sequences""" step_content = merge_messages_content(steps) context = get_trajectory_context(trajectory, steps) prompt = self.prompt_format( - prompt_name="failure_step_experience_prompt", + prompt_name="failure_step_task_memory_prompt", query=trajectory.metadata.get('query', ''), step_sequence=step_content, context=context, outcome="failed" ) - def parse_experiences(message: Message) -> List[BaseMemory]: - experiences_data = parse_json_experience_response(message.content) - experiences = [] + def parse_task_memories(message: Message) -> List[BaseMemory]: + task_memories_data = parse_json_experience_response(message.content) + task_memories = [] - for exp_data in experiences_data: - experience = TaskMemory( + for tm_data in task_memories_data: + task_memory = TaskMemory( workspace_id=self.context.get("workspace_id", ""), - when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), - content=exp_data.get("experience", ""), + when_to_use=tm_data.get("when_to_use", tm_data.get("condition", "")), + content=tm_data.get("experience", ""), author=getattr(self.llm, 'model_name', 'system'), - metadata=exp_data + metadata=tm_data ) - experiences.append(experience) + task_memories.append(task_memory) - return experiences + return task_memories - return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) \ No newline at end of file + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_task_memories) diff --git a/reme_ai/summary/task/failure_extraction_prompt.yaml b/reme_ai/summary/task/failure_extraction_prompt.yaml index 8318f023..5b55cf57 100644 --- a/reme_ai/summary/task/failure_extraction_prompt.yaml +++ b/reme_ai/summary/task/failure_extraction_prompt.yaml @@ -1,7 +1,7 @@ -failure_step_experience_prompt: | +failure_step_task_memory_prompt: | You are an expert AI analyst reviewing failed step sequences from an AI agent execution. - Your task is to extract learning experiences from failures to prevent similar mistakes in future executions. + Your task is to extract learning task memories from failures to prevent similar mistakes in future executions. Focus on identifying error patterns, missed opportunities, and alternative approaches. ANALYSIS FRAMEWORK: diff --git a/reme_ai/summary/task/pdf_preprocess_op_wrapper.py b/reme_ai/summary/task/pdf_preprocess_op_wrapper.py index 8baa70de..38556433 100644 --- a/reme_ai/summary/task/pdf_preprocess_op_wrapper.py +++ b/reme_ai/summary/task/pdf_preprocess_op_wrapper.py @@ -1,13 +1,12 @@ -from typing import List, Dict, Any, Tuple +from flowllm import C, BaseOp from loguru import logger -from flowllm import C, BaseOp -from .pdf_preprocess_op import MinerUPDFProcessor, chunk_pdf_content +from reme_ai.utils.miner_u_pdf_processor import MinerUPDFProcessor, chunk_pdf_content @C.register_op() class PDFPreprocessOp(BaseOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): """Process PDF files using MinerU and chunk content""" diff --git a/reme_ai/summary/task/simple_comparative_summary_op.py b/reme_ai/summary/task/simple_comparative_summary_op.py index a9919d0a..a8569020 100644 --- a/reme_ai/summary/task/simple_comparative_summary_op.py +++ b/reme_ai/summary/task/simple_comparative_summary_op.py @@ -1,17 +1,17 @@ import json from typing import List, Dict +from flowllm import C, BaseLLMOp from loguru import logger -from flowllm import C, BaseLLMOp +from reme_ai.schema import Message, Trajectory from reme_ai.schema.memory import BaseMemory, TaskMemory -from reme_ai.schema.message import Message, Trajectory -from reme_ai.utils.memory_utils import merge_messages_content +from reme_ai.utils.op_utils import merge_messages_content @C.register_op() class SimpleComparativeSummaryOp(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def compare_summary_trajectory(self, trajectory_a: Trajectory, trajectory_b: Trajectory) -> List[BaseMemory]: summary_prompt = self.prompt_format(prompt_name="summary_prompt", @@ -21,22 +21,22 @@ class SimpleComparativeSummaryOp(BaseLLMOp): def parse_content(message: Message): content = message.content - experience_list = [] + task_memory_list = [] try: content = content.split("```")[1].strip() if content.startswith("json"): content = content.strip("json") - for exp_dict in json.loads(content): - when_to_use = exp_dict.get("when_to_use", "").strip() - experience = exp_dict.get("experience", "").strip() - if when_to_use and experience: - experience_list.append(TaskMemory(workspace_id=self.context.get("workspace_id", ""), + for tm_dict in json.loads(content): + when_to_use = tm_dict.get("when_to_use", "").strip() + task_memory_content = tm_dict.get("experience", "").strip() + if when_to_use and task_memory_content: + task_memory_list.append(TaskMemory(workspace_id=self.context.get("workspace_id", ""), when_to_use=when_to_use, - content=experience, + content=task_memory_content, author=getattr(self.llm, 'model_name', 'system'))) - return experience_list + return task_memory_list except Exception as e: logger.exception(f"parse content failed!\n{content}") @@ -53,17 +53,17 @@ class SimpleComparativeSummaryOp(BaseLLMOp): task_id_dict[trajectory.task_id] = [] task_id_dict[trajectory.task_id].append(trajectory) - experience_list = [] + task_memory_list = [] for task_id, task_trajectories in task_id_dict.items(): task_trajectories: List[Trajectory] = sorted(task_trajectories, key=lambda x: x.score, reverse=True) if len(task_trajectories) < 2: continue if task_trajectories[0].score > task_trajectories[-1].score: - experiences = self.compare_summary_trajectory(trajectory_a=task_trajectories[0], + task_memories = self.compare_summary_trajectory(trajectory_a=task_trajectories[0], trajectory_b=task_trajectories[-1]) - experience_list.extend(experiences) + task_memory_list.extend(task_memories) - self.context.comparative_summary_experiences = experience_list - for e in experience_list: - logger.info(f"add experience when_to_use={e.when_to_use}\ncontent={e.content}") + self.context.comparative_summary_task_memories = task_memory_list + for tm in task_memory_list: + logger.info(f"add task memory when_to_use={tm.when_to_use}\ncontent={tm.content}") diff --git a/reme_ai/summary/task/success_extraction_op.py b/reme_ai/summary/task/success_extraction_op.py index 6c2bda7b..5c73a436 100644 --- a/reme_ai/summary/task/success_extraction_op.py +++ b/reme_ai/summary/task/success_extraction_op.py @@ -1,74 +1,73 @@ -import json -import re from typing import List +from flowllm import C, BaseLLMOp from loguru import logger -from flowllm import C, BaseLLMOp +from reme_ai.schema import Message, Trajectory from reme_ai.schema.memory import BaseMemory, TaskMemory -from reme_ai.schema.message import Message, Trajectory -from reme_ai.utils.memory_utils import merge_messages_content, parse_json_experience_response, get_trajectory_context +from reme_ai.utils.op_utils import merge_messages_content, parse_json_experience_response, get_trajectory_context + @C.register_op() class SuccessExtractionOp(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): - """Extract experiences from successful trajectories""" + """Extract task memories from successful trajectories""" success_trajectories: List[Trajectory] = self.context.get("success_trajectories", []) if not success_trajectories: logger.info("No success trajectories found for extraction") return - logger.info(f"Extracting experiences from {len(success_trajectories)} successful trajectories") + logger.info(f"Extracting task memories from {len(success_trajectories)} successful trajectories") - success_experiences = [] + success_task_memories = [] # Process trajectories for trajectory in success_trajectories: if "segments" in trajectory.metadata: # Process segmented step sequences for segment in trajectory.metadata["segments"]: - experiences = self._extract_success_experience_from_steps(segment, trajectory) - success_experiences.extend(experiences) + task_memories = self._extract_success_task_memory_from_steps(segment, trajectory) + success_task_memories.extend(task_memories) else: # Process entire trajectory - experiences = self._extract_success_experience_from_steps(trajectory.messages, trajectory) - success_experiences.extend(experiences) + task_memories = self._extract_success_task_memory_from_steps(trajectory.messages, trajectory) + success_task_memories.extend(task_memories) - logger.info(f"Extracted {len(success_experiences)} success experiences") - - # Add experiences to context - self.context.success_experiences = success_experiences + logger.info(f"Extracted {len(success_task_memories)} success task memories") - def _extract_success_experience_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[BaseMemory]: - """Extract experience from successful step sequences""" + # Add task memories to context + self.context.success_task_memories = success_task_memories + + def _extract_success_task_memory_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[BaseMemory]: + """Extract task memory from successful step sequences""" step_content = merge_messages_content(steps) context = get_trajectory_context(trajectory, steps) prompt = self.prompt_format( - prompt_name="success_step_experience_prompt", + prompt_name="success_step_task_memory_prompt", query=trajectory.metadata.get('query', ''), step_sequence=step_content, context=context, outcome="successful" ) - def parse_experiences(message: Message) -> List[BaseMemory]: - experiences_data = parse_json_experience_response(message.content) - experiences = [] + def parse_task_memories(message: Message) -> List[BaseMemory]: + task_memories_data = parse_json_experience_response(message.content) + task_memories = [] - for exp_data in experiences_data: - experience = TaskMemory( + for tm_data in task_memories_data: + task_memory = TaskMemory( workspace_id=self.context.get("workspace_id", ""), - when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), - content=exp_data.get("experience", ""), + when_to_use=tm_data.get("when_to_use", tm_data.get("condition", "")), + content=tm_data.get("experience", ""), author=getattr(self.llm, 'model_name', 'system'), - metadata=exp_data + metadata=tm_data ) - experiences.append(experience) + task_memories.append(task_memory) - return experiences + return task_memories - return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) \ No newline at end of file + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_task_memories) diff --git a/reme_ai/summary/task/success_extraction_prompt.yaml b/reme_ai/summary/task/success_extraction_prompt.yaml index 5b7e742a..d63365bb 100644 --- a/reme_ai/summary/task/success_extraction_prompt.yaml +++ b/reme_ai/summary/task/success_extraction_prompt.yaml @@ -1,7 +1,7 @@ -success_step_experience_prompt: | +success_step_task_memory_prompt: | You are an expert AI analyst reviewing successful step sequences from an AI agent execution. - Your task is to extract reusable, actionable step-level experiences that can guide future agent executions. + Your task is to extract reusable, actionable step-level task memories that can guide future agent executions. Focus on identifying specific patterns, techniques, and decision points that contributed to success. ANALYSIS FRAMEWORK: diff --git a/reme_ai/summary/task/task_memory_deduplication_op.py b/reme_ai/summary/task/task_memory_deduplication_op.py new file mode 100644 index 00000000..c57af218 --- /dev/null +++ b/reme_ai/summary/task/task_memory_deduplication_op.py @@ -0,0 +1,160 @@ +from typing import List +from loguru import logger + +from flowllm import C, BaseOp +from reme_ai.schema.memory import BaseMemory + + +@C.register_op() +class TaskMemoryDeduplicationOp(BaseOp): + file_path: str = __file__ + + def execute(self): + """Remove duplicate task memories""" + # Get task memories to deduplicate + task_memories: List[BaseMemory] = self.context.get("task_memories", []) + + if not task_memories: + logger.info("No task memories found for deduplication") + return + + logger.info(f"Starting deduplication for {len(task_memories)} task memories") + + # Perform deduplication + deduplicated_task_memories = self._deduplicate_task_memories(task_memories) + + logger.info(f"Deduplication complete: {len(deduplicated_task_memories)} deduplicated task memories out of {len(task_memories)}") + + # Update context + self.context.deduplicated_task_memories = deduplicated_task_memories + + def _deduplicate_task_memories(self, task_memories: List[BaseMemory]) -> List[BaseMemory]: + """Remove duplicate task memories""" + if not task_memories: + return task_memories + + similarity_threshold = self.op_params.get("similarity_threshold", 0.5) + workspace_id = self.context.get("workspace_id") + + unique_task_memories = [] + + # Get existing task memory embeddings + existing_embeddings = self._get_existing_task_memory_embeddings(workspace_id) + + for task_memory in task_memories: + # Generate embedding for current task memory + current_embedding = self._get_task_memory_embedding(task_memory) + + if current_embedding is None: + logger.warning(f"Failed to generate embedding for task memory: {str(task_memory.when_to_use)[:50]}...") + continue + + # Check similarity with existing task memories + if self._is_similar_to_existing_task_memories(current_embedding, existing_embeddings, similarity_threshold): + logger.debug(f"Skipping similar task memory: {str(task_memory.when_to_use)[:50]}...") + continue + + # Check similarity with current batch task memories + if self._is_similar_to_current_task_memories(current_embedding, unique_task_memories, similarity_threshold): + logger.debug(f"Skipping duplicate in current batch: {str(task_memory.when_to_use)[:50]}...") + continue + + # Add to unique task memories list + unique_task_memories.append(task_memory) + logger.debug(f"Added unique task memory: {str(task_memory.when_to_use)[:50]}...") + + return unique_task_memories + + def _get_existing_task_memory_embeddings(self, workspace_id: str) -> List[List[float]]: + """Get embeddings of existing task memories""" + try: + if not hasattr(self.context, 'vector_store') or not self.context.vector_store or not workspace_id: + return [] + + # Query existing task memory nodes + existing_nodes = self.context.vector_store.search( + query="...", # Empty query to get all + workspace_id=workspace_id, + top_k=self.op_params.get("max_existing_task_memories", 1000) + ) + + # Extract embeddings + existing_embeddings = [] + for node in existing_nodes: + if hasattr(node, 'embedding') and node.embedding: + existing_embeddings.append(node.embedding) + + logger.debug(f"Retrieved {len(existing_embeddings)} existing task memory embeddings from workspace {workspace_id}") + return existing_embeddings + + except Exception as e: + logger.warning(f"Failed to retrieve existing task memory embeddings: {e}") + return [] + + def _get_task_memory_embedding(self, task_memory: BaseMemory) -> List[float]: + """Generate embedding for task memory""" + try: + if not hasattr(self.context, 'vector_store') or not self.context.vector_store: + return None + + # Combine task memory description and content for embedding + text_for_embedding = f"{task_memory.when_to_use} {task_memory.content}" + embeddings = self.context.vector_store.embedding_model.get_embeddings([text_for_embedding]) + + if embeddings and len(embeddings) > 0: + return embeddings[0] + else: + logger.warning("Empty embedding generated for task memory") + return None + + except Exception as e: + logger.error(f"Error generating embedding for task memory: {e}") + return None + + + def _is_similar_to_existing_task_memories(self, current_embedding: List[float], + existing_embeddings: List[List[float]], + threshold: float) -> bool: + """Check if current embedding is similar to existing embeddings""" + for existing_embedding in existing_embeddings: + similarity = self._calculate_cosine_similarity(current_embedding, existing_embedding) + if similarity > threshold: + logger.debug(f"Found similar existing task memory with similarity: {similarity:.3f}") + return True + return False + + def _is_similar_to_current_task_memories(self, current_embedding: List[float], + current_task_memories: List[BaseMemory], + threshold: float) -> bool: + for existing_task_memory in current_task_memories: + existing_embedding = self._get_task_memory_embedding(existing_task_memory) + if existing_embedding is None: + continue + + similarity = self._calculate_cosine_similarity(current_embedding, existing_embedding) + if similarity > threshold: + logger.debug(f"Found similar task memory in current batch with similarity: {similarity:.3f}") + return True + return False + + def _calculate_cosine_similarity(self, embedding1: List[float], embedding2: List[float]) -> float: + """Calculate cosine similarity""" + try: + import numpy as np + + vec1 = np.array(embedding1) + vec2 = np.array(embedding2) + + # Calculate cosine similarity + dot_product = np.dot(vec1, vec2) + norm1 = np.linalg.norm(vec1) + norm2 = np.linalg.norm(vec2) + + if norm1 == 0 or norm2 == 0: + return 0.0 + + return dot_product / (norm1 * norm2) + + except Exception as e: + logger.error(f"Error calculating cosine similarity: {e}") + return 0.0 \ No newline at end of file diff --git a/reme_ai/summary/task/trajectory_preprocess_op.py b/reme_ai/summary/task/trajectory_preprocess_op.py index 60cb6e8a..501abbd1 100644 --- a/reme_ai/summary/task/trajectory_preprocess_op.py +++ b/reme_ai/summary/task/trajectory_preprocess_op.py @@ -1,18 +1,19 @@ from typing import List, Dict +from flowllm import C, BaseOp from loguru import logger -from flowllm import C, BaseOp -from reme_ai.schema.message import Trajectory +from reme_ai.schema import Trajectory @C.register_op() class TrajectoryPreprocessOp(BaseOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): """Preprocess trajectories: validate and classify""" - trajectories: List[Trajectory] = self.context.get("trajectories", []) + trajectories: list = self.context.get("trajectories", []) + trajectories: List[Trajectory] = [Trajectory(**x) if isinstance(x, dict) else x for x in trajectories] # Classify trajectories classified = self._classify_trajectories(trajectories) diff --git a/reme_ai/summary/task/trajectory_segmentation_op.py b/reme_ai/summary/task/trajectory_segmentation_op.py index 6b209a67..f81d7d98 100644 --- a/reme_ai/summary/task/trajectory_segmentation_op.py +++ b/reme_ai/summary/task/trajectory_segmentation_op.py @@ -2,15 +2,15 @@ import json import re from typing import List +from flowllm import C, BaseLLMOp from loguru import logger -from flowllm import C, BaseLLMOp -from reme_ai.schema.message import Message, Trajectory +from reme_ai.schema import Message, Trajectory @C.register_op() class TrajectorySegmentationOp(BaseLLMOp): - current_path: str = __file__ + file_path: str = __file__ def execute(self): """Segment trajectories into meaningful steps"""