restructure summarizer & retriever

This commit is contained in:
鸣山 2025-07-16 19:51:29 +08:00
parent 93e88d4a11
commit 3088cfb4ac
25 changed files with 2355 additions and 4 deletions

View file

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

View file

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

View file

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

View file

@ -1 +1,163 @@
# at jiaji
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 []

View file

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

View file

@ -1 +1,154 @@
# @jiaji
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 ""

View file

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

View file

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

View file

@ -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"
}}
]
```

View file

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

View file

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

View file

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

View file

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

View file

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

View file

@ -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"]
}}
]
```

View file

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

View file

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

View file

@ -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"]
}}
]
```

View file

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

View file

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

View file

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

View file

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