diff --git a/experiencemaker/config/full_pipeline_config.yaml b/experiencemaker/config/full_pipeline_config.yaml new file mode 100644 index 00000000..5d04ed31 --- /dev/null +++ b/experiencemaker/config/full_pipeline_config.yaml @@ -0,0 +1,96 @@ +# 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: recall_experience_op->rerank_experience_op->rewrite_experience_op + summarizer: trajectory_preprocess_op->[success_extraction_op|failure_extraction_op|comparative_extraction_op]->experience_validation_op->experience_deduplication_op->experience_storage_op + vector_store: vector_store_action_op + +op: + # retriever ops + recall_experience_op: + backend: recall_experience_op + vector_store: default + params: + retrieve_top_k: 15 + + rerank_experience_op: + backend: rerank_experience_op + llm: default + params: + enable_llm_rerank: true + enable_score_filter: false + top_k: 5 + + rewrite_experience_op: + backend: rewrite_experience_op + llm: default + params: + enable_llm_rewrite: true + + #summarizer ops + trajectory_preprocess_op: + backend: trajectory_preprocess_op + params: + success_threshold: 1.0 + + success_extraction_op: + backend: success_extraction_op + llm: default + + failure_extraction_op: + backend: failure_extraction_op + llm: default + + comparative_extraction_op: + backend: comparative_extraction_op + llm: default + params: + enable_soft_comparison: true + + experience_validation_op: + backend: experience_validation_op + llm: default + params: + validation_threshold: 0.5 + + experience_deduplication_op: + backend: experience_deduplication_op + vector_store: default + params: + similarity_threshold: 0.5 + + experience_storage_op: + backend: experience_storage_op + vector_store: default + + vector_store_action_op: + backend: vector_store_action_op + vector_store: default + +llm: + default: + backend: openai_compatible + model_name: qwen3-32b + params: + temperature: 0.6 + +embedding_model: + default: + backend: openai_compatible + model_name: text-embedding-v4 + params: + dimensions: 1024 + +vector_store: + default: + backend: local_file + embedding_model: default diff --git a/experiencemaker/op/__init__.py b/experiencemaker/op/__init__.py index c8c85319..52b2d055 100644 --- a/experiencemaker/op/__init__.py +++ b/experiencemaker/op/__init__.py @@ -6,6 +6,21 @@ from experiencemaker.op.mock_op import Mock1Op, Mock2Op, Mock3Op, Mock4Op, Mock5 from experiencemaker.op.retriever.build_query_op import BuildQueryOp from experiencemaker.op.retriever.merge_experience_op import MergeExperienceOp from experiencemaker.op.summarizer.simple_summary_op import SimpleSummaryOp + +from experiencemaker.op.summarizer.trajectory_preprocess_op import TrajectoryPreprocessOp +from experiencemaker.op.summarizer.comparative_extraction_op import ComparativeExtractionOp +from experiencemaker.op.summarizer.success_extraction_op import SuccessExtractionOp +from experiencemaker.op.summarizer.failure_extraction_op import FailureExtractionOp +from experiencemaker.op.summarizer.experience_validation_op import ExperienceValidationOp +from experiencemaker.op.summarizer.experience_deduplication_op import ExperienceDeduplicationOp +from experiencemaker.op.summarizer.experience_validation_op import ExperienceValidationOp +from experiencemaker.op.summarizer.trajectory_segmentation_op import TrajectorySegmentationOp +from experiencemaker.op.summarizer.experience_storage_op import ExperienceStorageOp + +from experiencemaker.op.retriever.recall_experience_op import RecallExperienceOp +from experiencemaker.op.retriever.rerank_experience_op import RerankExperienceOp +from experiencemaker.op.retriever.rewrite_experience_op import RewriteExperienceOp + from experiencemaker.op.vector_store.update_vector_store_op import UpdateVectorStoreOp from experiencemaker.op.vector_store.recall_vector_store_op import RecallVectorStoreOp from experiencemaker.op.vector_store.vector_store_action_op import VectorStoreActionOp diff --git a/experiencemaker/op/retriever/recall_experience_op.py b/experiencemaker/op/retriever/recall_experience_op.py new file mode 100644 index 00000000..51a03329 --- /dev/null +++ b/experiencemaker/op/retriever/recall_experience_op.py @@ -0,0 +1,104 @@ +from typing import List +from loguru import logger +from pydantic import Field + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.message import Message +from experiencemaker.schema.vector_node import VectorNode +from experiencemaker.schema.request import RetrieverRequest + + +@OP_REGISTRY.register() +class RecallExperienceOp(BaseOp): + """ + Recall relevant experiences from vector store based on query + """ + current_path: str = __file__ + + # Configuration parameters + + def execute(self): + """Execute recall operation""" + request: RetrieverRequest = self.context.request + retrieve_top_k = self.op_params.get("retrieve_top_k",15) + query_enhancement = self.op_params.get("query_enhancement",False) + logger.info(request) + + # Extract query and messages from request + query = getattr(request, 'query', None) + messages = getattr(request, 'messages', None) + + # Store in context for downstream ops + self.context.set_context("query", request.query) + self.context.set_context("messages", request.messages) + + try: + # Build retrieval query + retrieval_query = self._build_retrieve_query(query, messages, query_enhancement) + logger.info(f"Built retrieval query: {retrieval_query}") + + # Retrieve experiences from vector store + recalled_experiences = self._retrieve_experiences(request.workspace_id, retrieval_query, retrieve_top_k) + + logger.info(f"Recalled {len(recalled_experiences)} experiences") + + # Store results in context for downstream ops + self.context.set_context("recalled_experiences", recalled_experiences) + self.context.set_context("retrieval_query", retrieval_query) + + except Exception as e: + logger.error(f"Error in recall operation: {e}") + self.context.set_context("recalled_experiences", []) + + def _build_retrieve_query(self, query: str, messages: List[Message] = None, query_enhancement = False) -> str: + """Build retrieval query from query and messages""" + # Use the original query as base + base_query = query + + # Optionally enhance with current step context if enabled + if query_enhancement and messages: + current_context = self._extract_context(messages) + if current_context: + base_query = f"{base_query} {current_context}" + + return base_query + + def _retrieve_experiences(self, workspace_id: str, query: str, retrieve_top_k: int) -> List[VectorNode]: + """Retrieve experiences from vector store""" + if not query: + logger.warning("Empty query provided for vector retrieval") + return [] + + try: + retrieved_nodes = self.vector_store.search( + workspace_id=workspace_id, + query=query, + top_k=retrieve_top_k + ) + + logger.info(f"Vector retrieval found {len(retrieved_nodes)} candidates") + return retrieved_nodes + + except Exception as e: + logger.exception(f"Error in vector retrieval: {e}") + return [] + + def _extract_context(self, messages: List[Message]) -> str: + """Extract relevant context from messages for query enhancement""" + if not messages: + return "" + + context_parts = [] + + # Add recent steps if available + recent_messages = messages[-3:] # Last 3 messages + message_summaries = [] + for message in recent_messages: + content = message.content[:200] + "..." if len(message.content) > 200 else message.content + message_summaries.append(f"- {message.role.value}: {content}") + + if message_summaries: + context_parts.append("Recent messages:\n" + "\n".join(message_summaries)) + + return "\n\n".join(context_parts) \ No newline at end of file diff --git a/experiencemaker/op/retriever/rerank_experience_op.py b/experiencemaker/op/retriever/rerank_experience_op.py index c98c96bc..8d8999ca 100644 --- a/experiencemaker/op/retriever/rerank_experience_op.py +++ b/experiencemaker/op/retriever/rerank_experience_op.py @@ -1 +1,163 @@ -# at jiaji \ No newline at end of file +import json +import re +from typing import List +from loguru import logger +from pydantic import Field + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.message import Message +from experiencemaker.schema.vector_node import VectorNode +from experiencemaker.enumeration.role import Role + + +@OP_REGISTRY.register() +class RerankExperienceOp(BaseOp): + """ + Rerank and filter recalled experiences using LLM and score-based filtering + """ + current_path: str = __file__ + + def execute(self): + + """Execute rerank operation""" + recalled_experiences: List[VectorNode] = self.context.get_context("recalled_experiences", []) + retrieval_query: str = self.context.get_context("retrieval_query", "") + enable_llm_rerank = self.op_params.get("enable_llm_rerank", True) + enable_score_filter = self.op_params.get("enable_score_filter", False) + min_score_threshold = self.op_params.get("min_score_threshold", 0.3) + top_k = self.op_params.get("top_k", 5) + logger.info(f"top_k: {top_k}") + + if not recalled_experiences: + logger.info("No recalled experiences to rerank") + self.context.set_context("reranked_experiences", []) + return + + try: + logger.info(f"Reranking {len(recalled_experiences)} experiences") + + # Step 1: LLM reranking (optional) + if enable_llm_rerank: + recalled_experiences = self._llm_rerank(retrieval_query, recalled_experiences) + logger.info(f"After LLM reranking: {len(recalled_experiences)} experiences") + + # Step 2: Score-based filtering (optional) + if enable_score_filter: + recalled_experiences = self._score_based_filter(recalled_experiences, min_score_threshold) + logger.info(f"After score filtering: {len(recalled_experiences)} experiences") + + # Step 3: Return top-k results + final_results = recalled_experiences[:top_k] + logger.info(f"Final reranked results: {len(final_results)} experiences") + + # Store results in context + self.context.set_context("reranked_experiences", final_results) + + except Exception as e: + logger.error(f"Error in rerank operation: {e}") + self.context.set_context("reranked_experiences", recalled_experiences[:top_k]) + + def _llm_rerank(self, query: str, candidates: List[VectorNode]) -> List[VectorNode]: + """LLM-based reranking of candidate experiences""" + if not candidates: + return candidates + + try: + # Format candidates for LLM evaluation + candidates_text = self._format_candidates_for_rerank(candidates) + + prompt = self.prompt_format( + prompt_name="experience_rerank_prompt", + query=query, + candidates=candidates_text, + num_candidates=len(candidates) + ) + + response = self.llm.chat([Message(role=Role.USER, content=prompt)]) + + # Parse reranking results + reranked_indices = self._parse_rerank_response(response.content) + + # Reorder candidates based on LLM ranking + if reranked_indices: + reranked_candidates = [] + for idx in reranked_indices: + if 0 <= idx < len(candidates): + reranked_candidates.append(candidates[idx]) + + # Add any remaining candidates that weren't explicitly ranked + ranked_indices_set = set(reranked_indices) + for i, candidate in enumerate(candidates): + if i not in ranked_indices_set: + reranked_candidates.append(candidate) + + return reranked_candidates + + return candidates + + except Exception as e: + logger.error(f"Error in LLM reranking: {e}") + return candidates + + def _score_based_filter(self, experiences: List[VectorNode], min_score: float) -> List[VectorNode]: + """Filter experiences based on quality scores""" + filtered_experiences = [] + + for exp in experiences: + # Get confidence score from metadata + confidence = exp.metadata.get("confidence", 0.5) + validation_score = exp.metadata.get("validation_score", 0.5) + + # Calculate combined score + combined_score = (confidence + validation_score) / 2 + + if combined_score >= min_score: + filtered_experiences.append(exp) + else: + logger.debug(f"Filtered out experience with score {combined_score:.2f}") + + logger.info(f"Score filtering: {len(filtered_experiences)}/{len(experiences)} experiences retained") + return filtered_experiences + + def _format_candidates_for_rerank(self, candidates: List[VectorNode]) -> str: + """Format candidates for LLM reranking""" + formatted_candidates = [] + + for i, candidate in enumerate(candidates): + condition = candidate.content + experience = candidate.metadata.get("experience", "") + tags = candidate.metadata.get("tags", []) + confidence = candidate.metadata.get("confidence", 0.5) + + candidate_text = f"Candidate {i}:\n" + candidate_text += f"Condition: {condition}\n" + candidate_text += f"Experience: {experience}\n" + candidate_text += f"Tags: {', '.join(tags) if tags else 'None'}\n" + candidate_text += f"Confidence: {confidence}\n" + + formatted_candidates.append(candidate_text) + + return "\n---\n".join(formatted_candidates) + + def _parse_rerank_response(self, response: str) -> List[int]: + """Parse LLM reranking response to extract ranked indices""" + try: + # Try to extract JSON format + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + parsed = json.loads(json_blocks[0]) + if isinstance(parsed, dict) and "ranked_indices" in parsed: + return parsed["ranked_indices"] + elif isinstance(parsed, list): + return parsed + + # Try to extract numbers from text + numbers = re.findall(r'\b\d+\b', response) + return [int(num) for num in numbers if int(num) < 100] # Reasonable upper bound + + except Exception as e: + logger.error(f"Error parsing rerank response: {e}") + return [] \ No newline at end of file diff --git a/experiencemaker/op/retriever/rerank_experience_prompt.yaml b/experiencemaker/op/retriever/rerank_experience_prompt.yaml new file mode 100644 index 00000000..2d344d6e --- /dev/null +++ b/experiencemaker/op/retriever/rerank_experience_prompt.yaml @@ -0,0 +1,25 @@ +experience_rerank_prompt: | + You are an expert AI analyst tasked with reranking retrieved experiences based on their relevance to a specific query. + + Your task is to analyze the candidates and rank them by relevance, considering: + ● DIRECT RELEVANCE: How directly applicable the experience is to the current query + ● SITUATION SIMILARITY: How similar the experience context is to the current situation + ● ACTIONABILITY: How actionable and specific the experience is + ● QUALITY: The overall quality and clarity of the experience + + # Current Query + {query} + + # Candidate Experiences (Total: {num_candidates}) + {candidates} + + OUTPUT FORMAT: + Provide a ranked list of candidate indices (0-based) from most relevant to least relevant: + ```json + {{ + "ranked_indices": [2, 0, 4, 1, 3], + "reasoning": "Brief explanation of ranking rationale" + }} + ``` + + Note: Include ALL candidate indices in the ranking, even if some are less relevant. \ No newline at end of file diff --git a/experiencemaker/op/retriever/rewrite_experience_op.py b/experiencemaker/op/retriever/rewrite_experience_op.py index 122c231b..f7fabb4b 100644 --- a/experiencemaker/op/retriever/rewrite_experience_op.py +++ b/experiencemaker/op/retriever/rewrite_experience_op.py @@ -1 +1,154 @@ -# @jiaji \ No newline at end of file +import json +import re +from typing import List +from loguru import logger +from pydantic import Field + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import TextExperience +from experiencemaker.schema.message import Message +from experiencemaker.schema.vector_node import VectorNode +from experiencemaker.enumeration.role import Role +from experiencemaker.schema.response import RetrieverResponse + + +@OP_REGISTRY.register() +class RewriteExperienceOp(BaseOp): + """ + Generate and rewrite context messages from reranked experiences + """ + current_path: str = __file__ + + def execute(self): + """Execute rewrite operation""" + reranked_experiences: List[VectorNode] = self.context.get_context("reranked_experiences", []) + query: str = self.context.get_context("query", "") + messages: List[Message] = self.context.get_context("messages", []) + retrieval_query: str = self.context.get_context("retrieval_query", "") + + if not reranked_experiences: + logger.info("No reranked experiences to rewrite") + self.context.set_context("context_message", Message(content="")) + return + + logger.info(f"Generating context from {len(reranked_experiences)} experiences") + + # Generate initial context message + context_message = self._generate_context_message(query, messages, reranked_experiences, retrieval_query) + + # Store results in context + self.context.set_context("context_message", context_message) + + response: RetrieverResponse = self.context.response + response.experience_list = [TextExperience.from_vector_node(node) for node in reranked_experiences] + response.experience_merged = context_message + + + def _generate_context_message(self, query: str, messages: List[Message], nodes: List[VectorNode], + retrieval_query: str) -> Message: + """Generate context message from retrieved experiences""" + if not nodes: + return "" + + try: + # Format retrieved experiences + formatted_experiences = self._format_experiences_for_context(nodes) + + if self.op_params.get("enable_llm_rewrite", True): + context_content = self._rewrite_context(query, formatted_experiences, messages) + else: + context_content = formatted_experiences + + return context_content + + except Exception as e: + logger.error(f"Error generating context message: {e}") + return self._format_experiences_for_context(nodes) + + def _rewrite_context(self, query: str, context_content: str, messages: List[Message]) -> str: + """LLM-based context rewriting to make experiences more relevant and actionable""" + if not context_content: + return context_content + + try: + # Extract current context + current_context = self._extract_context(messages) + + prompt = self.prompt_format( + prompt_name="experience_rewrite_prompt", + current_query=query, + current_context=current_context, + original_context=context_content + ) + + response = self.llm.chat([Message(role=Role.USER, content=prompt)]) + + # Extract rewritten context + rewritten_context = self._parse_json_response(response.content, "rewritten_context") + + if rewritten_context and rewritten_context.strip(): + logger.info("Context successfully rewritten for current task") + return rewritten_context.strip() + + return context_content + + except Exception as e: + logger.error(f"Error in context rewriting: {e}") + return context_content + + def _format_experiences_for_context(self, experiences: List[VectorNode]) -> str: + """Format experiences for context generation""" + formatted_experiences = [] + + for i, exp in enumerate(experiences, 1): + condition = exp.content + experience_content = exp.metadata.get("experience_content", "") + exp_text = f"Experience {i} :\n When to use: {condition}\n Content: {experience_content}\n" + + formatted_experiences.append(exp_text) + + return "\n".join(formatted_experiences) + + def _extract_context(self, messages: List[Message]) -> str: + """Extract relevant context from messages""" + if not messages: + return "" + + context_parts = [] + + # Add recent messages if available + recent_messages = messages[-3:] # Last 3 messages + message_summaries = [] + for message in recent_messages: + content = message.content[:300] + "..." if len(message.content) > 300 else message.content + message_summaries.append(f"- {message.role.value}: {content}") + + if message_summaries: + context_parts.append("Recent conversation:\n" + "\n".join(message_summaries)) + + return "\n\n".join(context_parts) + + def _parse_json_response(self, response: str, key: str) -> str: + """Parse JSON response to extract specific key""" + try: + # Try to extract JSON blocks + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + parsed = json.loads(json_blocks[0]) + if isinstance(parsed, dict) and key in parsed: + return parsed[key] + + # Fallback: try to parse the entire response as JSON + parsed = json.loads(response) + if isinstance(parsed, dict) and key in parsed: + return parsed[key] + + except json.JSONDecodeError: + logger.warning(f"Failed to parse JSON response for key '{key}', using raw response") + # If JSON parsing fails, return the response as-is for fallback + return response.strip() + + return "" \ No newline at end of file diff --git a/experiencemaker/op/retriever/rewrite_experience_prompt.yaml b/experiencemaker/op/retriever/rewrite_experience_prompt.yaml new file mode 100644 index 00000000..e3459c04 --- /dev/null +++ b/experiencemaker/op/retriever/rewrite_experience_prompt.yaml @@ -0,0 +1,77 @@ +experience_rewrite_prompt: | + You are an expert AI assistant tasked with rewriting and reorganizing context content to make it more relevant and actionable for the current task. + + Your task is to take the original context (containing multiple experiences) and rewrite it as a cohesive, task-specific guidance that directly addresses the current situation. + + REWRITING GUIDELINES: + ● RELEVANCE FOCUS: Emphasize the most relevant aspects of each experience. Prioritize the most relevant experiences. Use clear, direct language. + ● ACTIONABLE INSIGHTS: Extract specific, actionable guidance. Make the context immediately actionable + ● COHERENT NARRATIVE: Create a flowing narrative rather than disconnected tips + ● SITUATIONAL AWARENESS: Adapt the guidance to the current situation + + # Current Task/Query + {current_query} + + # Current Trajectory + {current_context} + + # Original Context Content (Multiple Experiences) + {original_context} + + OUTPUT FORMAT: + Provide the rewritten context: + ```json + {{ + "rewritten_context": "A cohesive, task-specific context message that reorganizes and adapts the original experiences for the current task. This should be written as a unified guidance rather than separate experience items.", + }} + ``` + + Guidelines: + - Rewrite as a unified, flowing guidance + - Adapt terminology and examples to match the current task domain + - Consolidate overlapping insights into coherent recommendations + - Prioritize experiences most relevant to the current situation + - Make the guidance feel custom-written for this specific task + +context_generation_prompt: | + You are an expert AI assistant tasked with synthesizing retrieved experiences into actionable context for an AI agent. + + Your task is to create a coherent, actionable context message that helps the agent leverage relevant past experiences. + + SYNTHESIS GUIDELINES: + ● RELEVANCE FOCUS: Emphasize the most relevant aspects of each experience + ● ACTIONABLE INSIGHTS: Extract specific, actionable guidance + ● COHERENT NARRATIVE: Create a flowing narrative rather than disconnected tips + ● SITUATIONAL AWARENESS: Adapt the guidance to the current situation + + # Current Query/Task + {query} + + # Current Step Info + {context} + + # Retrieved Experiences ({num_experiences} total) + {experiences} + + OUTPUT FORMAT: + Create a synthesized context message: + ```json + {{ + "context": "A coherent, actionable context message that synthesizes the relevant experiences and provides specific guidance for the current task", + "key_insights": [ + "Key pattern from successful approaches", + "Common pitfall to avoid", + "Specific technique that worked well" + ], + "recommended_actions": [ + "Specific action recommendation based on experiences", + "Decision point with recommended choice" + ] + }} + ``` + + Guidelines: + - Make the context immediately actionable + - Prioritize the most relevant experiences + - Use clear, direct language + - Focus on practical guidance rather than abstract principles \ No newline at end of file diff --git a/experiencemaker/op/summarizer/comparative_extraction_op.py b/experiencemaker/op/summarizer/comparative_extraction_op.py new file mode 100644 index 00000000..24094e3c --- /dev/null +++ b/experiencemaker/op/summarizer/comparative_extraction_op.py @@ -0,0 +1,313 @@ +import json +import re +from typing import List, Tuple, Optional +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import TextExperience, ExperienceMeta +from experiencemaker.schema.message import Message, Trajectory + + +@OP_REGISTRY.register() +class ComparativeExtractionOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Extract comparative experiences by comparing different scoring trajectories""" + all_trajectories: List[Trajectory] = self.context.get_context("all_trajectories", []) + success_trajectories: List[Trajectory] = self.context.get_context("success_trajectories", []) + failure_trajectories: List[Trajectory] = self.context.get_context("failure_trajectories", []) + + all_experiences = [] + + # 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})") + self.submit_task(self._extract_soft_comparative_experience, + higher_traj=highest_traj, lower_traj=lowest_traj) + + # Hard comparison: success vs failure (if similarity search is enabled) + if (success_trajectories and failure_trajectories and + self.op_params.get("enable_similarity_comparison", False)): + + similar_pairs = self._find_similar_step_sequences(success_trajectories, failure_trajectories) + logger.info(f"Found {len(similar_pairs)} similar pairs for hard comparison") + + for success_steps, failure_steps, similarity_score in similar_pairs: + self.submit_task(self._extract_hard_comparative_experience, + success_steps=success_steps, failure_steps=failure_steps, + similarity_score=similarity_score) + + # Collect all experiences + for task_result in self.join_task(): + if task_result: + all_experiences.extend(task_result) + + logger.info(f"Extracted {len(all_experiences)} comparative experiences") + + # Add experiences to context + existing_experiences = self.context.get_context("extracted_experiences", []) + existing_experiences.extend(all_experiences) + self.context.set_context("extracted_experiences", existing_experiences) + + def _find_highest_lowest_scoring_trajectories(self, trajectories: List[Trajectory]) -> Tuple[Optional[Trajectory], Optional[Trajectory]]: + """Find the highest and lowest scoring trajectories""" + if len(trajectories) < 2: + return None, None + + # Filter trajectories with valid scores + valid_trajectories = [traj for traj in trajectories if traj.score is not None] + + if len(valid_trajectories) < 2: + logger.warning("Not enough trajectories with valid scores for comparison") + return None, None + + # Sort by score + sorted_trajectories = sorted(valid_trajectories, key=lambda x: x.score, reverse=True) + + highest_traj = sorted_trajectories[0] + lowest_traj = sorted_trajectories[-1] + + return highest_traj, lowest_traj + + def _get_trajectory_score(self, trajectory: Trajectory) -> Optional[float]: + """Get trajectory score""" + return trajectory.score + + def _extract_soft_comparative_experience(self, higher_traj: Trajectory, lower_traj: Trajectory) -> List[TextExperience]: + """Extract soft comparative experience (high score vs low score)""" + try: + 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", + higher_steps=self._format_step_sequence(higher_steps), + lower_steps=self._format_step_sequence(lower_steps), + higher_score=f"{higher_score:.2f}", + lower_score=f"{lower_score:.2f}" + ) + + def parse_experiences(message: Message) -> List[TextExperience]: + try: + experiences_data = self._parse_json_experience_response(message.content) + experiences = [] + + for exp_data in experiences_data: + experience = TextExperience( + workspace_id=self.context.request.workspace_id, + when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), + content=exp_data.get("experience", ""), + metadata=exp_data + ) + experiences.append(experience) + + return experiences + + except Exception as e: + logger.error(f"Error parsing soft comparative experiences: {e}") + return [] + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + + except Exception as e: + logger.error(f"Error extracting soft comparative experience: {e}") + return [] + + def _extract_hard_comparative_experience(self, success_steps: List[Message], + failure_steps: List[Message], similarity_score: float) -> List[TextExperience]: + """Extract hard comparative experience (success vs failure)""" + try: + prompt = self.prompt_format( + prompt_name="comparative_step_experience_prompt", + success_steps=self._format_step_sequence(success_steps), + failure_steps=self._format_step_sequence(failure_steps), + similarity_score=similarity_score + ) + + def parse_experiences(message: Message) -> List[TextExperience]: + try: + experiences_data = self._parse_json_experience_response(message.content) + experiences = [] + + for exp_data in experiences_data: + experience = TextExperience( + workspace_id=self.context.request.workspace_id, + when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), + content=exp_data.get("experience", ""), + metadata=exp_data + ) + experiences.append(experience) + + return experiences + + except Exception as e: + logger.error(f"Error parsing hard comparative experiences: {e}") + return [] + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + + except Exception as e: + logger.error(f"Error extracting hard comparative experience: {e}") + return [] + + def _get_trajectory_steps(self, trajectory: Trajectory) -> List[Message]: + """Get trajectory steps, prioritizing segmented steps""" + if hasattr(trajectory, 'segments') and trajectory.segments: + # If there are segments, merge all segments + all_steps = [] + for segment in trajectory.segments: + all_steps.extend(segment) + return all_steps + else: + return trajectory.messages + + def _find_similar_step_sequences(self, success_trajectories: List[Trajectory], + failure_trajectories: List[Trajectory]) -> List[Tuple[List[Message], List[Message], float]]: + """Find similar step sequences for comparison""" + if not self.op_params.get("enable_similarity_comparison", False): + return [] + + try: + similar_pairs = [] + + # Get step sequences + success_step_sequences = [] + for traj in success_trajectories: + if hasattr(traj, 'segments') and traj.segments: + success_step_sequences.extend(traj.segments) + else: + success_step_sequences.append(traj.steps) + + failure_step_sequences = [] + for traj in failure_trajectories: + if hasattr(traj, 'segments') and traj.segments: + failure_step_sequences.extend(traj.segments) + else: + failure_step_sequences.append(traj.steps) + + # Limit comparison count to avoid computational overload + max_sequences = self.op_params.get("max_similarity_sequences", 5) + success_step_sequences = success_step_sequences[:max_sequences] + failure_step_sequences = failure_step_sequences[:max_sequences] + + if not success_step_sequences or not failure_step_sequences: + return [] + + # Generate text representation for embedding + success_texts = [self._format_step_sequence(seq) for seq in success_step_sequences] + failure_texts = [self._format_step_sequence(seq) for seq in failure_step_sequences] + + # Get embedding vectors + if hasattr(self, 'vector_store') and self.vector_store and hasattr(self.vector_store, 'embedding_model'): + success_embeddings = self.vector_store.embedding_model.get_embeddings(success_texts) + failure_embeddings = self.vector_store.embedding_model.get_embeddings(failure_texts) + + # Calculate similarity and find most similar pairs + similarity_threshold = self.op_params.get("similarity_threshold", 0.3) + + for i, s_emb in enumerate(success_embeddings): + for j, f_emb in enumerate(failure_embeddings): + similarity = self._calculate_cosine_similarity(s_emb, f_emb) + + if similarity > similarity_threshold: + similar_pairs.append(( + success_step_sequences[i], + failure_step_sequences[j], + similarity + )) + + # Return top most similar pairs + max_pairs = self.op_params.get("max_similarity_pairs", 3) + return sorted(similar_pairs, key=lambda x: x[2], reverse=True)[:max_pairs] + + except Exception as e: + logger.error(f"Error finding similar step sequences: {e}") + + return [] + + 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 + + def _format_step_sequence(self, steps: List[Message]) -> str: + """Format step sequence""" + formatted_steps = [] + for i, step in enumerate(steps): + step_info = f"Step {i + 1} [{step.role.value}]:" + + if hasattr(step, 'reasoning_content') and step.reasoning_content: + step_info += f"\nReasoning: {step.reasoning_content}" + + step_info += f"\nContent: {step.content}" + + if hasattr(step, 'tool_calls') and step.tool_calls: + for tool_call in step.tool_calls: + step_info += f"\nTool: {tool_call.name}({tool_call.arguments})" + + formatted_steps.append(step_info) + + return "\n\n".join(formatted_steps) + + def _parse_json_experience_response(self, response: str) -> List[dict]: + """Parse JSON format experience response""" + try: + # Extract JSON blocks + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + parsed = json.loads(json_blocks[0]) + + # Handle array format + if isinstance(parsed, list): + valid_experiences = [] + for exp_data in parsed: + if isinstance(exp_data, dict) and ( + ("when_to_use" in exp_data and "experience" in exp_data) or + ("condition" in exp_data and "experience" in exp_data) + ): + valid_experiences.append(exp_data) + return valid_experiences + + # Handle single object + elif isinstance(parsed, dict) and ( + ("when_to_use" in parsed and "experience" in parsed) or + ("condition" in parsed and "experience" in parsed) + ): + return [parsed] + + # Fallback: try to parse entire response + parsed = json.loads(response) + if isinstance(parsed, list): + return parsed + elif isinstance(parsed, dict): + return [parsed] + + except json.JSONDecodeError as e: + logger.warning(f"Failed to parse JSON experience response: {e}") + + return [] \ No newline at end of file diff --git a/experiencemaker/op/summarizer/comparative_extraction_prompt.yaml b/experiencemaker/op/summarizer/comparative_extraction_prompt.yaml new file mode 100644 index 00000000..8b7724b7 --- /dev/null +++ b/experiencemaker/op/summarizer/comparative_extraction_prompt.yaml @@ -0,0 +1,79 @@ +soft_comparative_step_experience_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. + Focus on what made the higher-scoring approach more effective, even when both approaches may have had partial success. + + SOFT COMPARATIVE ANALYSIS FRAMEWORK: + ● PERFORMANCE FACTORS: Identify what specifically contributed to the higher score + ● APPROACH DIFFERENCES: Compare methodologies and execution strategies + ● EFFICIENCY ANALYSIS: Analyze why one approach was more efficient or effective + ● OPTIMIZATION INSIGHTS: Extract lessons for improving performance + + EXTRACTION PRINCIPLES: + ● Focus on INCREMENTAL IMPROVEMENTS and performance optimization + ● Extract QUALITY INDICATORS that differentiate better vs good approaches + ● Identify REFINEMENT STRATEGIES that lead to higher scores + ● Frame insights as PERFORMANCE ENHANCEMENT guidelines + + # Higher-Scoring Step Sequence (Score: {higher_score}) + {higher_steps} + + # Lower-Scoring Step Sequence (Score: {lower_score}) + {lower_steps} + + + OUTPUT FORMAT: + Generate 1-2 performance improvement insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific scenarios where this performance insight applies", + "experience": "Detailed analysis of what made the higher-scoring approach more effective", + "tags": ["performance_optimization", "score_improvement", "relevant_keywords"], + "confidence": 0.7, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` + +comparative_step_experience_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. + Focus on critical decision points, technique variations, and approach differences. + + COMPARATIVE ANALYSIS FRAMEWORK: + ● DECISION CONTRAST: Compare critical decisions made in success vs failure cases + ● TECHNIQUE VARIATIONS: Identify different approaches and their outcomes + ● TIMING DIFFERENCES: Analyze when certain actions were taken and their impact + ● SUCCESS FACTORS: Extract what specifically made the difference + + EXTRACTION PRINCIPLES: + ● Frame comparisons as PRINCIPLES as well as case-specific SOLUTIONS + ● Identify PATTERNS that differentiate effective vs ineffective approaches + ● Extract RULES that can guide future similar situations + ● Focus on UNDERLYING MECHANISMS rather than surface-level differences + + # Successful Step Sequence + {success_steps} + + # Failed Step Sequence + {failure_steps} + + # Similarity Score: {similarity_score} + + OUTPUT FORMAT: + Generate 1-2 comparative insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific scenarios where this comparative insight applies", + "experience": "Detailed comparison highlighting why success approach works better", + "tags": ["comparative_analysis", "success_factors", "relevant_keywords"], + "confidence": 0.8, + "step_type": "reasoning|action|observation|decision" + }} + ] + ``` \ No newline at end of file diff --git a/experiencemaker/op/summarizer/compare_summary_op.py b/experiencemaker/op/summarizer/compare_summary_op.py deleted file mode 100644 index e69de29b..00000000 diff --git a/experiencemaker/op/summarizer/experience_deduplication_op.py b/experiencemaker/op/summarizer/experience_deduplication_op.py new file mode 100644 index 00000000..b0b85a09 --- /dev/null +++ b/experiencemaker/op/summarizer/experience_deduplication_op.py @@ -0,0 +1,163 @@ +from typing import List +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import BaseExperience + + +@OP_REGISTRY.register() +class ExperienceDeduplicationOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Remove duplicate experiences""" + # Get experiences to deduplicate + experiences: List[BaseExperience] = self.context.get_context("validated_experiences", []) + if not experiences: + experiences = self.context.get_context("extracted_experiences", []) + + if not experiences: + logger.info("No experiences found for deduplication") + return + + logger.info(f"Starting deduplication for {len(experiences)} experiences") + + # Perform deduplication + unique_experiences = self._deduplicate_experiences(experiences) + + logger.info(f"Deduplication complete: {len(unique_experiences)} unique experiences out of {len(experiences)}") + + # Update context + self.context.set_context("deduplicated_experiences", unique_experiences) + + def _deduplicate_experiences(self, experiences: List[BaseExperience]) -> List[BaseExperience]: + """Remove duplicate experiences""" + if not experiences: + return experiences + + similarity_threshold = self.op_params.get("similarity_threshold", 0.5) + workspace_id = self.context.request.workspace_id if hasattr(self.context, 'request') else None + + 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, 'vector_store') or not self.vector_store or not workspace_id: + return [] + + # Query existing experience nodes + existing_nodes = self.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: BaseExperience) -> List[float]: + """Generate embedding for experience""" + try: + if not hasattr(self, 'vector_store') or not self.vector_store: + return None + + # Combine experience description and content for embedding + text_for_embedding = f"{experience.when_to_use} {experience.content}" + embeddings = self.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[BaseExperience], + 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/experiencemaker/op/summarizer/experience_storage_op.py b/experiencemaker/op/summarizer/experience_storage_op.py new file mode 100644 index 00000000..cd906f43 --- /dev/null +++ b/experiencemaker/op/summarizer/experience_storage_op.py @@ -0,0 +1,59 @@ +from typing import List +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import BaseExperience +from experiencemaker.schema.vector_node import VectorNode + + +@OP_REGISTRY.register() +class ExperienceStorageOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Store experiences to vector database""" + # Get experiences to store + experiences: List[BaseExperience] = self.context.get_context("deduplicated_experiences", []) + if not experiences: + experiences = self.context.get_context("validated_experiences", []) + if not experiences: + experiences = self.context.get_context("extracted_experiences", []) + + if not experiences: + logger.info("No experiences found for storage") + return + + logger.info(f"Storing {len(experiences)} experiences to vector database") + + try: + # Convert to vector storage nodes + nodes: List[VectorNode] = [experience.to_vector_node() for experience in experiences] + + # Get workspace_id + workspace_id = self.context.request.workspace_id if hasattr(self.context, 'request') else None + if not workspace_id: + workspace_id = self.op_params.get("default_workspace_id", "default") + + # Store to vector database + if hasattr(self, 'vector_store') and self.vector_store: + self.vector_store.insert(nodes, workspace_id=workspace_id) + logger.info(f"Successfully stored {len(experiences)} experiences to workspace: {workspace_id}") + + # Set storage result to context + self.context.set_context("storage_success", True) + self.context.set_context("stored_count", len(experiences)) + + else: + logger.error("Vector store not available for storage") + self.context.set_context("storage_success", False) + self.context.set_context("storage_error", "Vector store not available") + + # Log stored experiences + for experience in experiences: + logger.info(f"Stored experience - Description: {str(experience.when_to_use)[:100]}...") + + except Exception as e: + logger.error(f"Error storing experiences: {e}") + self.context.set_context("storage_success", False) + self.context.set_context("storage_error", str(e)) \ No newline at end of file diff --git a/experiencemaker/op/summarizer/experience_validation_op.py b/experiencemaker/op/summarizer/experience_validation_op.py new file mode 100644 index 00000000..5524e92e --- /dev/null +++ b/experiencemaker/op/summarizer/experience_validation_op.py @@ -0,0 +1,101 @@ +import re +from typing import List, Dict, Any +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import BaseExperience +from experiencemaker.schema.message import Message +from experiencemaker.enumeration.role import Role + + +@OP_REGISTRY.register() +class ExperienceValidationOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Validate quality of extracted experiences""" + experiences: List[BaseExperience] = self.context.get_context("extracted_experiences", []) + + if not experiences: + logger.info("No experiences found for validation") + return + + logger.info(f"Validating {len(experiences)} extracted experiences") + + # Use thread pool for parallel validation + for experience in experiences: + self.submit_task(self._validate_single_experience, experience=experience) + + # Collect validation results + validated_experiences = [] + validation_results = list(self.join_task()) + + for i, result in enumerate(validation_results): + if result and result.get("is_valid", False): + validated_experiences.append(experiences[i]) + else: + reason = result.get("reason", "Unknown reason") if result else "Validation failed" + logger.warning(f"Experience validation failed: {reason}") + + logger.info(f"Validated {len(validated_experiences)} out of {len(experiences)} experiences") + + # Update context + self.context.set_context("validated_experiences", validated_experiences) + + def _validate_single_experience(self, experience: BaseExperience) -> Dict[str, Any]: + """Validate single experience""" + return self._llm_validate_experience(experience) + + def _llm_validate_experience(self, experience: BaseExperience) -> Dict[str, Any]: + """Validate experience using LLM""" + try: + prompt = self.prompt_format( + prompt_name="experience_validation_prompt", + condition=experience.when_to_use, + experience_content=experience.content + ) + + def parse_validation(message: Message) -> Dict[str, Any]: + try: + response_content = message.content + + # Parse validation result + is_valid = "valid" in response_content.lower() and "invalid" not in response_content.lower() + + # Extract score + score_match = re.search(r'score[:\s]*([0-9.]+)', response_content.lower()) + try: + score = float(score_match.group(1)) if score_match else 0.5 + except (ValueError, AttributeError): + score = 0.5 + + # Set validation threshold + validation_threshold = self.op_params.get("validation_threshold", 0.3) + + return { + "is_valid": is_valid and score > validation_threshold, + "score": score, + "feedback": response_content, + "reason": "" if (is_valid and score > validation_threshold) else f"Low validation score ({score:.2f}) or marked as invalid" + } + + except Exception as e: + logger.error(f"Error parsing validation response: {e}") + return { + "is_valid": False, + "score": 0.0, + "feedback": "", + "reason": f"Parse error: {str(e)}" + } + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_validation) + + except Exception as e: + logger.error(f"LLM validation failed: {e}") + return { + "is_valid": False, + "score": 0.0, + "feedback": "", + "reason": f"LLM validation error: {str(e)}" + } \ No newline at end of file diff --git a/experiencemaker/op/summarizer/experience_validation_prompt.yaml b/experiencemaker/op/summarizer/experience_validation_prompt.yaml new file mode 100644 index 00000000..408bc834 --- /dev/null +++ b/experiencemaker/op/summarizer/experience_validation_prompt.yaml @@ -0,0 +1,29 @@ +experience_validation_prompt: | + You are an expert AI analyst tasked with validating the quality and usefulness of extracted step-level experiences. + + Your task is to assess whether the extracted experience 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? + + # Experience to Validate + Condition: {condition} + Experience Content: {experience_content} + + OUTPUT FORMAT: + Provide validation assessment: + ```json + {{ + "is_valid": true/false, + "score": 0.8, + "feedback": "Detailed explanation of validation decision", + "recommendations": "Suggestions for improvement if applicable" + }} + ``` + + 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. diff --git a/experiencemaker/op/summarizer/failure_extraction_op.py b/experiencemaker/op/summarizer/failure_extraction_op.py new file mode 100644 index 00000000..67f2b865 --- /dev/null +++ b/experiencemaker/op/summarizer/failure_extraction_op.py @@ -0,0 +1,184 @@ +import json +import re +from typing import List +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import TextExperience, ExperienceMeta +from experiencemaker.schema.message import Message, Trajectory +from experiencemaker.enumeration.role import Role + + +@OP_REGISTRY.register() +class FailureExtractionOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Extract experiences from failed trajectories""" + failure_trajectories: List[Trajectory] = self.context.get_context("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") + + # Use thread pool for parallel processing + for trajectory in failure_trajectories: + if hasattr(trajectory, 'segments') and trajectory.segments: + # Process segmented step sequences + for segment in trajectory.segments: + self.submit_task(self._extract_failure_experience_from_steps, + steps=segment, trajectory=trajectory) + else: + # Process entire trajectory + self.submit_task(self._extract_failure_experience_from_steps, + steps=trajectory.messages, trajectory=trajectory) + + # Collect all experiences + all_experiences = [] + for task_result in self.join_task(): + if task_result: + all_experiences.extend(task_result) + + logger.info(f"Extracted {len(all_experiences)} failure experiences") + + # Add experiences to context + existing_experiences = self.context.get_context("extracted_experiences", []) + existing_experiences.extend(all_experiences) + self.context.set_context("extracted_experiences", existing_experiences) + + def _extract_failure_experience_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[TextExperience]: + """Extract experience from failed step sequences""" + try: + step_content = self._format_step_sequence(steps) + context = self._get_trajectory_context(trajectory, steps) + + prompt = self.prompt_format( + prompt_name="failure_step_experience_prompt", + query=trajectory.metadata.get('query', ''), + step_sequence=step_content, + context=context, + outcome="failed" + ) + + def parse_experiences(message: Message) -> List[TextExperience]: + try: + experiences_data = self._parse_json_experience_response(message.content) + experiences = [] + + for exp_data in experiences_data: + experience = TextExperience( + workspace_id=self.context.request.workspace_id, + when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), + content=exp_data.get("experience", ""), + metadata=ExperienceMeta(author=self.llm.model_name if hasattr(self, 'llm') else "system") + ) + experiences.append(experience) + + return experiences + + except Exception as e: + logger.error(f"Error parsing failure experiences: {e}") + return [] + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + + except Exception as e: + logger.error(f"Error extracting failure experience: {e}") + return [] + + def _format_step_sequence(self, steps: List[Message]) -> str: + """Format step sequence to string""" + step_content_collector = [] + + for step in steps: + step_index = len(step_content_collector) + + if step.role == Role.ASSISTANT: + line = f"### step.{step_index} role={step.role.value} content=\n{step.content}\n" + if hasattr(step, 'reasoning_content') and step.reasoning_content: + line += f"{step.reasoning_content}\n" + if hasattr(step, 'tool_calls') and step.tool_calls: + for tool_call in step.tool_calls: + line += f" - tool call={tool_call.name}\n params={tool_call.arguments}\n" + step_content_collector.append(line) + + elif step.role == Role.USER: + line = f"### step.{step_index} role={step.role.value} content=\n{step.content}\n" + step_content_collector.append(line) + + elif step.role == Role.TOOL: + line = f"### step.{step_index} role={step.role.value} tool call result=\n{step.content}\n" + step_content_collector.append(line) + + return "\n".join(step_content_collector).strip() + + def _get_trajectory_context(self, trajectory: Trajectory, step_sequence: List[Message]) -> str: + """Get context of step sequence within trajectory""" + try: + # Find position of step sequence in trajectory + start_idx = 0 + for i, step in enumerate(trajectory.messages): + if step == step_sequence[0]: + start_idx = i + break + + # Extract before and after context + context_before = trajectory.messages[max(0, start_idx - 2):start_idx] + context_after = trajectory.messages[start_idx + len(step_sequence):start_idx + len(step_sequence) + 2] + + context = f"Query: {trajectory.metadata.get('query', 'N/A')}\n" + + if context_before: + context += "Previous steps:\n" + "\n".join([f"- {step.content[:100]}..." for step in context_before]) + "\n" + + if context_after: + context += "Following steps:\n" + "\n".join([f"- {step.content[:100]}..." for step in context_after]) + + return context + + except Exception as e: + logger.error(f"Error getting trajectory context: {e}") + return f"Query: {trajectory.metadata.get('query', 'N/A')}" + + def _parse_json_experience_response(self, response: str) -> List[dict]: + """Parse JSON formatted experience response""" + try: + # Extract JSON blocks + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + parsed = json.loads(json_blocks[0]) + + # Handle array format + if isinstance(parsed, list): + valid_experiences = [] + for exp_data in parsed: + if isinstance(exp_data, dict) and ( + ("when_to_use" in exp_data and "experience" in exp_data) or + ("condition" in exp_data and "experience" in exp_data) + ): + valid_experiences.append(exp_data) + return valid_experiences + + # Handle single object + elif isinstance(parsed, dict) and ( + ("when_to_use" in parsed and "experience" in parsed) or + ("condition" in parsed and "experience" in parsed) + ): + return [parsed] + + # Fallback: try to parse entire response + parsed = json.loads(response) + if isinstance(parsed, list): + return parsed + elif isinstance(parsed, dict): + return [parsed] + + except json.JSONDecodeError as e: + logger.warning(f"Failed to parse JSON experience response: {e}") + + return [] \ No newline at end of file diff --git a/experiencemaker/op/summarizer/failure_extraction_prompt.yaml b/experiencemaker/op/summarizer/failure_extraction_prompt.yaml new file mode 100644 index 00000000..8318f023 --- /dev/null +++ b/experiencemaker/op/summarizer/failure_extraction_prompt.yaml @@ -0,0 +1,42 @@ +failure_step_experience_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. + Focus on identifying error patterns, missed opportunities, and alternative approaches. + + ANALYSIS FRAMEWORK: + ● FAILURE POINT IDENTIFICATION: Pinpoint where and why the steps went wrong + ● ERROR PATTERN ANALYSIS: Identify recurring mistakes or problematic approaches + ● ALTERNATIVE APPROACHES: Suggest what could have been done differently + ● PREVENTION STRATEGIES: Extract actionable insights to avoid similar failures + + EXTRACTION PRINCIPLES: + ● Extract GENERAL PRINCIPLES as well as SPECIFIC INSTRUCTIONS + ● Focus on PATTERNS and RULES as well as particular instances + + # Original Query + {query} + + # Step Sequence Analysis + {step_sequence} + + # Context Information + {context} + + # Outcome + This step sequence was part of a {outcome} trajectory. + + OUTPUT FORMAT: + Generate 1-3 step-level failure prevention insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific situations where this lesson should be remembered", + "experience": "Universal principle or rule extracted from the failure pattern ", + "tags": ["error_prevention", "failure_analysis", "relevant_keywords"], + "confidence": 0.7, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` \ No newline at end of file diff --git a/experiencemaker/op/summarizer/step_summarizer_prompts.yaml b/experiencemaker/op/summarizer/step_summarizer_prompts.yaml new file mode 100644 index 00000000..8c7441b4 --- /dev/null +++ b/experiencemaker/op/summarizer/step_summarizer_prompts.yaml @@ -0,0 +1,274 @@ +success_step_experience_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. + Focus on identifying specific patterns, techniques, and decision points that contributed to success. + + ANALYSIS FRAMEWORK: + ● STEP PATTERN ANALYSIS: Identify the specific sequence of actions that led to success + ● DECISION POINTS: Highlight critical decisions made during these steps + ● TECHNIQUE EFFECTIVENESS: Analyze why specific approaches worked well + ● REUSABILITY: Extract patterns that can be applied to similar scenarios + + EXTRACTION PRINCIPLES: + ● Focus on TRANSFERABLE TECHNIQUES and decision frameworks + ● Frame insights as actionable guidelines and best practices + + # Original Query + {query} + + # Step Sequence Analysis + {step_sequence} + + # Context Information + {context} + + # Outcome + This step sequence was part of a {outcome} trajectory. + + OUTPUT FORMAT: + Generate 1-3 step-level success insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific conditions when this step pattern should be applied", + "experience": "Detailed description of the successful step pattern and why it works", + "tags": ["relevant", "keywords", "for", "categorization"], + "confidence": 0.8, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` + +failure_step_experience_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. + Focus on identifying error patterns, missed opportunities, and alternative approaches. + + ANALYSIS FRAMEWORK: + ● FAILURE POINT IDENTIFICATION: Pinpoint where and why the steps went wrong + ● ERROR PATTERN ANALYSIS: Identify recurring mistakes or problematic approaches + ● ALTERNATIVE APPROACHES: Suggest what could have been done differently + ● PREVENTION STRATEGIES: Extract actionable insights to avoid similar failures + + EXTRACTION PRINCIPLES: + ● Extract GENERAL PRINCIPLES as well as SPECIFIC INSTRUCTIONS + ● Focus on PATTERNS and RULES as well as particular instances + + # Original Query + {query} + + # Step Sequence Analysis + {step_sequence} + + # Context Information + {context} + + # Outcome + This step sequence was part of a {outcome} trajectory. + + OUTPUT FORMAT: + Generate 1-3 step-level failure prevention insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific situations where this lesson should be remembered", + "experience": "Universal principle or rule extracted from the failure pattern ", + "tags": ["error_prevention", "failure_analysis", "relevant_keywords"], + "confidence": 0.7, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` + +comparative_step_experience_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. + Focus on critical decision points, technique variations, and approach differences. + + COMPARATIVE ANALYSIS FRAMEWORK: + ● DECISION CONTRAST: Compare critical decisions made in success vs failure cases + ● TECHNIQUE VARIATIONS: Identify different approaches and their outcomes + ● TIMING DIFFERENCES: Analyze when certain actions were taken and their impact + ● SUCCESS FACTORS: Extract what specifically made the difference + + EXTRACTION PRINCIPLES: + ● Frame comparisons as PRINCIPLES as well as case-specific SOLUTIONS + ● Identify PATTERNS that differentiate effective vs ineffective approaches + ● Extract RULES that can guide future similar situations + ● Focus on UNDERLYING MECHANISMS rather than surface-level differences + + # Successful Step Sequence + {success_steps} + + # Failed Step Sequence + {failure_steps} + + # Similarity Score: {similarity_score} + + OUTPUT FORMAT: + Generate 1-2 comparative insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific scenarios where this comparative insight applies", + "experience": "Detailed comparison highlighting why success approach works better", + "tags": ["comparative_analysis", "success_factors", "relevant_keywords"], + "confidence": 0.8, + "step_type": "reasoning|action|observation|decision" + }} + ] + ``` + +general_step_experience_prompt: | + You are an expert AI analyst reviewing step sequences to extract general patterns and insights. + + Your task is to identify valuable step-level patterns without explicit success/failure labels. + Focus on effective techniques, common patterns, and general best practices. + + ANALYSIS FRAMEWORK: + ● PATTERN RECOGNITION: Identify recurring effective patterns in the steps + ● TECHNIQUE ANALYSIS: Analyze the effectiveness of different approaches + ● BEST PRACTICES: Extract general principles that appear beneficial + ● APPLICABILITY: Determine when these patterns would be most useful + + GENERALIZATION PRINCIPLES: + ● Extract UNIVERSAL PATTERNS that transcend specific contexts + ● Identify TRANSFERABLE METHODOLOGIES and approaches + ● Focus on PRINCIPLE-LEVEL insights as well as tactical details + ● Formulate insights as REUSABLE FRAMEWORKS and guidelines + + GENERALIZATION PRINCIPLES: + ● Extract UNIVERSAL PATTERNS that transcend specific contexts + ● Identify TRANSFERABLE METHODOLOGIES and approaches + ● Focus on PRINCIPLE-LEVEL insights + ● Formulate insights as REUSABLE FRAMEWORKS and guidelines + + # Original Query + {query} + + # Step Sequence Analysis + {step_sequence} + + # Context Information + {context} + + OUTPUT FORMAT: + Generate 1-2 general step insights as JSON objects: + ```json + [ + {{ + "when_to_use": "General conditions where this pattern is applicable", + "experience": "Detailed description of the effective step pattern", + "tags": ["general_pattern", "best_practice", "relevant_keywords"], + "confidence": 0.6, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` + +step_segmentation_prompt: | + You are an expert AI analyst tasked with segmenting a trajectory into meaningful step sequences. + + Your task is to identify natural breakpoints in the execution where one logical unit of work ends and another begins. + Consider factors like: task completion, context switches, tool changes, reasoning phases, and logical groupings. + + SEGMENTATION CRITERIA: + ● LOGICAL COMPLETION: Steps that complete a specific sub-task or reasoning phase + ● CONTEXT SWITCHES: Points where the agent shifts focus or approach + ● TOOL BOUNDARIES: Natural breaks around tool usage patterns + ● REASONING PHASES: Distinct phases of analysis, planning, or execution + + # Original Query + {query} + + # Full Trajectory (Total steps: {total_steps}) + {trajectory_content} + + OUTPUT FORMAT: + Provide segmentation points as a JSON array of step indices where splits should occur: + ```json + {{ + "segment_points": [3, 7, 12, 18], + "reasoning": "Brief explanation of segmentation logic" + }} + ``` + + Note: Segment points indicate the END of each segment. For example, [3, 7] means: + - Segment 1: steps 0-3 + - Segment 2: steps 4-7 + - Segment 3: steps 8-end + +experience_validation_prompt: | + You are an expert AI analyst tasked with validating the quality and usefulness of extracted step-level experiences. + + Your task is to assess whether the extracted experience 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? + + # Experience to Validate + Condition: {condition} + Experience Content: {experience_content} + + OUTPUT FORMAT: + Provide validation assessment: + ```json + {{ + "is_valid": true/false, + "score": 0.8, + "feedback": "Detailed explanation of validation decision", + "recommendations": "Suggestions for improvement if applicable" + }} + ``` + + 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. + +soft_comparative_step_experience_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. + Focus on what made the higher-scoring approach more effective, even when both approaches may have had partial success. + + SOFT COMPARATIVE ANALYSIS FRAMEWORK: + ● PERFORMANCE FACTORS: Identify what specifically contributed to the higher score + ● APPROACH DIFFERENCES: Compare methodologies and execution strategies + ● EFFICIENCY ANALYSIS: Analyze why one approach was more efficient or effective + ● OPTIMIZATION INSIGHTS: Extract lessons for improving performance + + EXTRACTION PRINCIPLES: + ● Focus on INCREMENTAL IMPROVEMENTS and performance optimization + ● Extract QUALITY INDICATORS that differentiate better vs good approaches + ● Identify REFINEMENT STRATEGIES that lead to higher scores + ● Frame insights as PERFORMANCE ENHANCEMENT guidelines + + # Higher-Scoring Step Sequence (Score: {higher_score}) + {higher_steps} + + # Lower-Scoring Step Sequence (Score: {lower_score}) + {lower_steps} + + + OUTPUT FORMAT: + Generate 1-2 performance improvement insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific scenarios where this performance insight applies", + "experience": "Detailed analysis of what made the higher-scoring approach more effective", + "tags": ["performance_optimization", "score_improvement", "relevant_keywords"], + "confidence": 0.7, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] \ No newline at end of file diff --git a/experiencemaker/op/summarizer/success_extraction_op.py b/experiencemaker/op/summarizer/success_extraction_op.py new file mode 100644 index 00000000..9932fcc8 --- /dev/null +++ b/experiencemaker/op/summarizer/success_extraction_op.py @@ -0,0 +1,184 @@ +import json +import re +from typing import List +from loguru import logger + +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.experience import TextExperience, ExperienceMeta +from experiencemaker.schema.message import Message, Trajectory +from experiencemaker.enumeration.role import Role + + +@OP_REGISTRY.register() +class SuccessExtractionOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Extract experiences from successful trajectories""" + success_trajectories: List[Trajectory] = self.context.get_context("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") + + # Use thread pool for parallel processing + for trajectory in success_trajectories: + if hasattr(trajectory, 'segments') and trajectory.segments: + # Process segmented step sequences + for segment in trajectory.segments: + self.submit_task(self._extract_success_experience_from_steps, + steps=segment, trajectory=trajectory) + else: + # Process entire trajectory + self.submit_task(self._extract_success_experience_from_steps, + steps=trajectory.messages, trajectory=trajectory) + + # Collect all experiences + all_experiences = [] + for task_result in self.join_task(): + if task_result: + all_experiences.extend(task_result) + + logger.info(f"Extracted {len(all_experiences)} success experiences") + + # Add experiences to context + existing_experiences = self.context.get_context("extracted_experiences", []) + existing_experiences.extend(all_experiences) + self.context.set_context("extracted_experiences", existing_experiences) + + def _extract_success_experience_from_steps(self, steps: List[Message], trajectory: Trajectory) -> List[TextExperience]: + """Extract experience from successful step sequences""" + try: + step_content = self._format_step_sequence(steps) + context = self._get_trajectory_context(trajectory, steps) + + prompt = self.prompt_format( + prompt_name="success_step_experience_prompt", + query=trajectory.metadata.get('query', ''), + step_sequence=step_content, + context=context, + outcome="successful" + ) + + def parse_experiences(message: Message) -> List[TextExperience]: + try: + experiences_data = self._parse_json_experience_response(message.content) + experiences = [] + + for exp_data in experiences_data: + experience = TextExperience( + workspace_id=self.context.request.workspace_id, + when_to_use=exp_data.get("when_to_use", exp_data.get("condition", "")), + content=exp_data.get("experience", ""), + metadata=ExperienceMeta(author=self.llm.model_name if hasattr(self, 'llm') else "system") + ) + experiences.append(experience) + + return experiences + + except Exception as e: + logger.error(f"Error parsing success experiences: {e}") + return [] + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_experiences) + + except Exception as e: + logger.error(f"Error extracting success experience: {e}") + return [] + + def _format_step_sequence(self, steps: List[Message]) -> str: + """Format step sequence to string""" + step_content_collector = [] + + for step in steps: + step_index = len(step_content_collector) + + if step.role == Role.ASSISTANT: + line = f"### step.{step_index} role={step.role.value} content=\n{step.content}\n" + if hasattr(step, 'reasoning_content') and step.reasoning_content: + line += f"{step.reasoning_content}\n" + if hasattr(step, 'tool_calls') and step.tool_calls: + for tool_call in step.tool_calls: + line += f" - tool call={tool_call.name}\n params={tool_call.arguments}\n" + step_content_collector.append(line) + + elif step.role == Role.USER: + line = f"### step.{step_index} role={step.role.value} content=\n{step.content}\n" + step_content_collector.append(line) + + elif step.role == Role.TOOL: + line = f"### step.{step_index} role={step.role.value} tool call result=\n{step.content}\n" + step_content_collector.append(line) + + return "\n".join(step_content_collector).strip() + + def _get_trajectory_context(self, trajectory: Trajectory, step_sequence: List[Message]) -> str: + """Get context of step sequence within trajectory""" + try: + # Find position of step sequence in trajectory + start_idx = 0 + for i, step in enumerate(trajectory.messages): + if step == step_sequence[0]: + start_idx = i + break + + # Extract before and after context + context_before = trajectory.messages[max(0, start_idx - 2):start_idx] + context_after = trajectory.messages[start_idx + len(step_sequence):start_idx + len(step_sequence) + 2] + + context = f"Query: {trajectory.metadata.get('query', 'N/A')}\n" + + if context_before: + context += "Previous steps:\n" + "\n".join([f"- {step.content[:100]}..." for step in context_before]) + "\n" + + if context_after: + context += "Following steps:\n" + "\n".join([f"- {step.content[:100]}..." for step in context_after]) + + return context + + except Exception as e: + logger.error(f"Error getting trajectory context: {e}") + return f"Query: {trajectory.metadata.get('query', 'N/A')}" + + def _parse_json_experience_response(self, response: str) -> List[dict]: + """Parse JSON formatted experience response""" + try: + # Extract JSON blocks + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + parsed = json.loads(json_blocks[0]) + + # Handle array format + if isinstance(parsed, list): + valid_experiences = [] + for exp_data in parsed: + if isinstance(exp_data, dict) and ( + ("when_to_use" in exp_data and "experience" in exp_data) or + ("condition" in exp_data and "experience" in exp_data) + ): + valid_experiences.append(exp_data) + return valid_experiences + + # Handle single object + elif isinstance(parsed, dict) and ( + ("when_to_use" in parsed and "experience" in parsed) or + ("condition" in parsed and "experience" in parsed) + ): + return [parsed] + + # Fallback: try to parse entire response + parsed = json.loads(response) + if isinstance(parsed, list): + return parsed + elif isinstance(parsed, dict): + return [parsed] + + except json.JSONDecodeError as e: + logger.warning(f"Failed to parse JSON experience response: {e}") + + return [] \ No newline at end of file diff --git a/experiencemaker/op/summarizer/success_extraction_prompt.yaml b/experiencemaker/op/summarizer/success_extraction_prompt.yaml new file mode 100644 index 00000000..5b7e742a --- /dev/null +++ b/experiencemaker/op/summarizer/success_extraction_prompt.yaml @@ -0,0 +1,42 @@ +success_step_experience_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. + Focus on identifying specific patterns, techniques, and decision points that contributed to success. + + ANALYSIS FRAMEWORK: + ● STEP PATTERN ANALYSIS: Identify the specific sequence of actions that led to success + ● DECISION POINTS: Highlight critical decisions made during these steps + ● TECHNIQUE EFFECTIVENESS: Analyze why specific approaches worked well + ● REUSABILITY: Extract patterns that can be applied to similar scenarios + + EXTRACTION PRINCIPLES: + ● Focus on TRANSFERABLE TECHNIQUES and decision frameworks + ● Frame insights as actionable guidelines and best practices + + # Original Query + {query} + + # Step Sequence Analysis + {step_sequence} + + # Context Information + {context} + + # Outcome + This step sequence was part of a {outcome} trajectory. + + OUTPUT FORMAT: + Generate 1-3 step-level success insights as JSON objects: + ```json + [ + {{ + "when_to_use": "Specific conditions when this step pattern should be applied", + "experience": "Detailed description of the successful step pattern and why it works", + "tags": ["relevant", "keywords", "for", "categorization"], + "confidence": 0.8, + "step_type": "reasoning|action|observation|decision", + "tools_used": ["list", "of", "tools"] + }} + ] + ``` \ No newline at end of file diff --git a/experiencemaker/op/summarizer/trajectory_preprocess_op.py b/experiencemaker/op/summarizer/trajectory_preprocess_op.py new file mode 100644 index 00000000..ba8f47d8 --- /dev/null +++ b/experiencemaker/op/summarizer/trajectory_preprocess_op.py @@ -0,0 +1,84 @@ +from typing import List, Dict +from loguru import logger +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.message import Trajectory +from experiencemaker.schema.request import SummarizerRequest + + +@OP_REGISTRY.register() +class TrajectoryPreprocessOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Preprocess trajectories: validate and classify""" + request: SummarizerRequest = self.context.request + if request.traj_list: + self.context.set_context("trajectories", request.traj_list) + elif request.trajectories: + self.context.set_context("trajectories", request.trajectories) + else: + logger.error("No trajectories is recognized. Please send requests containing traj_list.") + + trajectories: List[Trajectory] = self.context.get_context("trajectories", []) + + if not trajectories: + logger.warning("No trajectories found in context") + return + + # Validate trajectories + valid_trajectories = self._validate_trajectories(trajectories) + logger.info(f"Validated {len(valid_trajectories)} out of {len(trajectories)} trajectories") + + # Classify trajectories + classified = self._classify_trajectories(valid_trajectories) + logger.info(f"Classified trajectories - Success: {len(classified['success'])}, " + f"Failure: {len(classified['failure'])}, All: {len(classified['all'])}") + + # Set context for downstream operators + self.context.set_context("success_trajectories", classified['success']) + self.context.set_context("failure_trajectories", classified['failure']) + self.context.set_context("all_trajectories", classified['all']) + + def _validate_trajectories(self, trajectories: List[Trajectory]) -> List[Trajectory]: + """Validate trajectory validity""" + valid_trajectories = [] + + for traj in trajectories: + if self._is_valid_trajectory(traj): + valid_trajectories.append(traj) + else: + logger.debug("Invalid trajectory filtered out") + + return valid_trajectories + + def _is_valid_trajectory(self, traj: Trajectory) -> bool: + """Check if trajectory is valid""" + if traj is None: + return False + if not hasattr(traj, 'score') or traj.score is None: + return False + if not hasattr(traj, 'messages') or not traj.messages or len(traj.messages) == 0: + return False + return True + + def _classify_trajectories(self, trajectories: List[Trajectory]) -> Dict[str, List[Trajectory]]: + """Classify trajectories based on score threshold""" + success_trajectories = [] + failure_trajectories = [] + + success_threshold = self.op_params.get("success_threshold", 1.0) + + for traj in trajectories: + is_success = traj.score >= success_threshold + + if is_success: + success_trajectories.append(traj) + else: + failure_trajectories.append(traj) + + return { + 'success': success_trajectories, + 'failure': failure_trajectories, + 'all': trajectories + } \ No newline at end of file diff --git a/experiencemaker/op/summarizer/trajectory_segmentation_op.py b/experiencemaker/op/summarizer/trajectory_segmentation_op.py new file mode 100644 index 00000000..5a72a499 --- /dev/null +++ b/experiencemaker/op/summarizer/trajectory_segmentation_op.py @@ -0,0 +1,135 @@ +import re +import json +from typing import List, Dict, Any +from loguru import logger +from experiencemaker.op import OP_REGISTRY +from experiencemaker.op.base_op import BaseOp +from experiencemaker.schema.message import Message, Trajectory +from experiencemaker.enumeration.role import Role + + +@OP_REGISTRY.register() +class TrajectorySegmentationOp(BaseOp): + current_path: str = __file__ + + def execute(self): + """Segment trajectories into meaningful steps""" + # Get trajectories from context + all_trajectories: List[Trajectory] = self.context.get_context("all_trajectories", []) + success_trajectories: List[Trajectory] = self.context.get_context("success_trajectories", []) + failure_trajectories: List[Trajectory] = self.context.get_context("failure_trajectories", []) + + if not all_trajectories: + logger.warning("No trajectories found in context") + return + + # Determine which trajectories to segment + target_trajectories = self._get_target_trajectories(all_trajectories, success_trajectories, + failure_trajectories) + + # Add segmentation info to trajectories + segmented_count = 0 + for trajectory in target_trajectories: + segments = self._segment_trajectory(trajectory) + trajectory.segments = segments + segmented_count += 1 + + logger.info(f"Segmented {segmented_count} trajectories") + + # Update context with segmented trajectories + self.context.set_context("segmented_trajectories", target_trajectories) + + def _get_target_trajectories(self, all_trajectories: List[Trajectory], + success_trajectories: List[Trajectory], + failure_trajectories: List[Trajectory]) -> List[Trajectory]: + """Determine which trajectories to segment based on configuration""" + segment_target = self.op_params.get("segment_target", "all") + + if segment_target == "success": + return success_trajectories + elif segment_target == "failure": + return failure_trajectories + else: + return all_trajectories + + def _segment_trajectory(self, trajectory: Trajectory) -> List[List[Message]]: + """Segment trajectory into step sequences using LLM""" + try: + return self._llm_segment_trajectory(trajectory) + except Exception as e: + logger.error(f"Error segmenting trajectory: {e}") + return [trajectory.messages] + + def _llm_segment_trajectory(self, trajectory: Trajectory) -> List[List[Message]]: + """Use LLM for trajectory segmentation""" + try: + trajectory_content = self._format_trajectory_content(trajectory) + + prompt = self.prompt_format( + prompt_name="step_segmentation_prompt", + query=trajectory.metadata.get('query', ''), + trajectory_content=trajectory_content, + total_steps=len(trajectory.messages) + ) + + def parse_segmentation(message: Message) -> List[List[Message]]: + try: + content = message.content + segment_points = self._parse_segmentation_response(content) + + # Segment trajectory based on segmentation points + segments = [] + start_idx = 0 + + for end_idx in segment_points: + if start_idx < end_idx <= len(trajectory.messages): + segments.append(trajectory.messages[start_idx:end_idx]) + start_idx = end_idx + + # Add remaining steps + if start_idx < len(trajectory.messages): + segments.append(trajectory.messages[start_idx:]) + + return segments if segments else [trajectory.messages] + + except Exception as e: + logger.error(f"Error parsing segmentation: {e}") + return [trajectory.messages] + + return self.llm.chat(messages=[Message(content=prompt)], callback_fn=parse_segmentation) + + except Exception as e: + logger.error(f"LLM segmentation failed: {e}") + return [trajectory.messages] + + def _format_trajectory_content(self, trajectory: Trajectory) -> str: + """Format trajectory content for LLM processing""" + content = "" + for i, step in enumerate(trajectory.messages): + content += f"Step {i + 1} ({step.role.value}):\n{step.content}\n\n" + return content + + def _parse_segmentation_response(self, response: str) -> List[int]: + """Parse segmentation response from LLM""" + segment_points = [] + + # Try to extract JSON format + json_pattern = r'```json\s*([\s\S]*?)\s*```' + json_blocks = re.findall(json_pattern, response) + + if json_blocks: + try: + parsed = json.loads(json_blocks[0]) + if isinstance(parsed, dict) and "segment_points" in parsed: + segment_points = parsed["segment_points"] + elif isinstance(parsed, list): + segment_points = parsed + except json.JSONDecodeError: + pass + + # Fallback: extract numbers + if not segment_points: + numbers = re.findall(r'\b\d+\b', response) + segment_points = [int(num) for num in numbers if int(num) > 0] + + return sorted(list(set(segment_points))) \ No newline at end of file diff --git a/experiencemaker/op/summarizer/trajectory_segmentation_prompt.yaml b/experiencemaker/op/summarizer/trajectory_segmentation_prompt.yaml new file mode 100644 index 00000000..50d325a5 --- /dev/null +++ b/experiencemaker/op/summarizer/trajectory_segmentation_prompt.yaml @@ -0,0 +1,31 @@ +step_segmentation_prompt: | + You are an expert AI analyst tasked with segmenting a trajectory into meaningful step sequences. + + Your task is to identify natural breakpoints in the execution where one logical unit of work ends and another begins. + Consider factors like: task completion, context switches, tool changes, reasoning phases, and logical groupings. + + SEGMENTATION CRITERIA: + ● LOGICAL COMPLETION: Steps that complete a specific sub-task or reasoning phase + ● CONTEXT SWITCHES: Points where the agent shifts focus or approach + ● TOOL BOUNDARIES: Natural breaks around tool usage patterns + ● REASONING PHASES: Distinct phases of analysis, planning, or execution + + # Original Query + {query} + + # Full Trajectory (Total steps: {total_steps}) + {trajectory_content} + + OUTPUT FORMAT: + Provide segmentation points as a JSON array of step indices where splits should occur: + ```json + {{ + "segment_points": [3, 7, 12, 18], + "reasoning": "Brief explanation of segmentation logic" + }} + ``` + + Note: Segment points indicate the END of each segment. For example, [3, 7] means: + - Segment 1: steps 0-3 + - Segment 2: steps 4-7 + - Segment 3: steps 8-end \ No newline at end of file diff --git a/experiencemaker/op/summarizer/update_experience_op.py b/experiencemaker/op/summarizer/update_experience_op.py deleted file mode 100644 index e69de29b..00000000 diff --git a/experiencemaker/op/summarizer/validate_experience_op.py b/experiencemaker/op/summarizer/validate_experience_op.py deleted file mode 100644 index e69de29b..00000000 diff --git a/experiencemaker/pipeline/pipeline.py b/experiencemaker/pipeline/pipeline.py index 7d56c0f8..67ef8284 100644 --- a/experiencemaker/pipeline/pipeline.py +++ b/experiencemaker/pipeline/pipeline.py @@ -68,10 +68,9 @@ class Pipeline: elif isinstance(pipeline, list): parallel_pipeline = [self._parse_sub_pipeline(x) for x in pipeline] - for op_list in zip_longest(parallel_pipeline, fillvalue="-"): + for op_list in zip_longest(*parallel_pipeline, fillvalue="-"): i += 1 logger.info(f"stage{i}: {' | '.join(op_list)}") - else: raise ValueError(f"unknown pipeline.type={type(pipeline)}")