refactor(context): restructure file system operations and context compaction

This commit is contained in:
jinli.yl 2025-11-18 15:36:20 +08:00
parent 7fe953a8d5
commit a093a9428d
22 changed files with 307 additions and 1254 deletions

View file

@ -32,13 +32,26 @@ classifiers = [
keywords = ["llm", "memory", "experience", "memoryscope", "ai", "mcp", "http"]
dependencies = [
"flowllm[reme]>=0.2.0.0",
"flowllm[reme]>=0.2.0.3",
]
[project.optional-dependencies]
dev = ["jupyter-book", "ghp-import", "myst-nb", "sphinxcontrib-bibtex", "furo", "sphinxcontrib-mermaid"]
dev = [
"jupyter-book",
"ghp-import",
"myst-nb",
"sphinxcontrib-bibtex",
"furo",
"sphinxcontrib-mermaid"
]
full = ["reme_ai[dev]"]
token = [
"flowllm[token]>=0.2.0.3"
]
full = [
"reme_ai[dev,token]"
]
[tool.setuptools.packages.find]
where = ["."]

View file

@ -156,10 +156,6 @@ flow:
description: "user query"
required: true
compact:
flow_content: ContextCompactOp() >> WriteFile()
....
llm:
default:
backend: openai_compatible

View file

@ -1,31 +0,0 @@
"""File system tool package.
This package provides file-related operations that can be used in LLM-powered flows.
It includes ready-to-use operations for:
- EditOp: File editing operation for replacing text within files
- GlobOp: File search operation for finding files matching glob patterns
- GrepOp: Text search operation for finding patterns in file contents
- ReadFileOp: File reading operation for reading file contents
- RipGrepOp: Text search operation using ripgrep for efficient pattern matching
- WriteFileOp: File writing operation for writing content to files
- WriteTodosOp: To-do list management operation for tracking subtasks
"""
from .edit_op import EditOp
from .glob_op import GlobOp
from .grep_op import GrepOp
from .read_file_op import ReadFileOp
from .rip_grep_op import RipGrepOp
from .write_file_op import WriteFileOp
from .write_todos_op import WriteTodosOp
__all__ = [
"EditOp",
"GlobOp",
"GrepOp",
"ReadFileOp",
"RipGrepOp",
"WriteFileOp",
"WriteTodosOp",
]

View file

@ -1,153 +0,0 @@
"""File edit operation module.
This module provides a tool operation for editing files by replacing text.
It supports creating new files, editing existing files, and replacing multiple occurrences.
"""
from pathlib import Path
from typing import Optional
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
@C.register_op()
class EditOp(BaseAsyncToolOp):
"""File edit operation.
This operation replaces text within a file. By default, replaces a single
occurrence, but can replace multiple occurrences when expected_replacements
is specified. Supports creating new files when old_string is empty.
"""
file_path = __file__
def __init__(self, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "Edit",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"file_path": {
"type": "string",
"description": self.get_prompt("file_path"),
"required": True,
},
"old_string": {
"type": "string",
"description": self.get_prompt("old_string"),
"required": True,
},
"new_string": {
"type": "string",
"description": self.get_prompt("new_string"),
"required": True,
},
"expected_replacements": {
"type": "number",
"description": self.get_prompt("expected_replacements"),
"required": False,
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the file edit operation."""
file_path: str = self.input_dict.get("file_path", "").strip()
old_string: str = self.input_dict.get("old_string", "")
new_string: str = self.input_dict.get("new_string", "")
expected_replacements: Optional[int] = self.input_dict.get("expected_replacements")
# Validate inputs
if not file_path:
raise ValueError("The 'file_path' parameter cannot be empty.")
if expected_replacements is None:
expected_replacements = 1
if expected_replacements < 1:
raise ValueError("The 'expected_replacements' parameter must be at least 1.")
# Resolve file path
file_path_obj = Path(file_path).expanduser().resolve()
# Check if file exists
file_exists = file_path_obj.exists() and file_path_obj.is_file()
# Handle new file creation
if not old_string and not file_exists:
# Create new file
file_path_obj.parent.mkdir(parents=True, exist_ok=True)
file_path_obj.write_text(new_string, encoding="utf-8")
self.set_output(f"Created new file: {file_path_obj}")
return
# File must exist for editing
if not file_exists:
raise FileNotFoundError(
f"File not found: {file_path_obj}. Use an empty old_string to create a new file.",
)
# Cannot create file that already exists
if not old_string and file_exists:
raise ValueError(
f"Failed to edit. Attempted to create a file that already exists: {file_path_obj}",
)
# Read current content
current_content = file_path_obj.read_text(encoding="utf-8")
current_content = current_content.replace("\r\n", "\n")
# Count occurrences
occurrences = current_content.count(old_string)
# Validate occurrences
if occurrences == 0:
raise ValueError(
f"Failed to edit, could not find the string to replace. "
f"0 occurrences found for old_string in {file_path_obj}.",
)
if occurrences != expected_replacements:
occurrence_term = "occurrence" if expected_replacements == 1 else "occurrences"
raise ValueError(
f"Failed to edit, expected {expected_replacements} {occurrence_term} "
f"but found {occurrences} for old_string in {file_path_obj}.",
)
# Check if old_string and new_string are identical
if old_string == new_string:
raise ValueError(
f"No changes to apply. The old_string and new_string are identical in {file_path_obj}.",
)
# Perform replacement
new_content = current_content.replace(old_string, new_string, occurrences)
# Check if content actually changed
if current_content == new_content:
raise ValueError(
f"No changes to apply. The new content is identical to the current content in {file_path_obj}.",
)
# Write new content
file_path_obj.write_text(new_content, encoding="utf-8")
# Set output
result_msg = f"Successfully modified file: {file_path_obj} ({occurrences} replacements)."
self.set_output(result_msg)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
file_path: str = self.input_dict.get("file_path", "").strip()
error_msg = f'Failed to edit file "{file_path}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,23 +0,0 @@
tool_desc: |
Replaces text within a file. By default, replaces a single occurrence, but can replace multiple occurrences when \`expected_replacements\` is specified. This tool requires providing significant context around the change to ensure precise targeting. Always use the ${READ_FILE_TOOL_NAME} tool to examine the file's current content before attempting a text replacement.
The user has the ability to modify the \`new_string\` content. If modified, this will be stated in the response.
Expectation for required parameters:
1. \`file_path\` is the path to the file to modify.
2. \`old_string\` MUST be the exact literal text to replace (including all whitespace, indentation, newlines, and surrounding code etc.).
3. \`new_string\` MUST be the exact literal text to replace \`old_string\` with (also including all whitespace, indentation, newlines, and surrounding code etc.). Ensure the resulting code is correct and idiomatic.
4. NEVER escape \`old_string\` or \`new_string\`, that would break the exact literal text requirement.
**Important:** If ANY of the above are not satisfied, the tool will fail. CRITICAL for \`old_string\`: Must uniquely identify the single instance to change. Include at least 3 lines of context BEFORE and AFTER the target text, matching whitespace and indentation precisely. If this string matches multiple locations, or does not match exactly, the tool will fail.
**Multiple replacements:** Set \`expected_replacements\` to the number of occurrences you want to replace. The tool will replace ALL occurrences that match \`old_string\` exactly. Ensure the number of replacements matches your expectation.
file_path: |
The path to the file to modify.
old_string: |
The exact literal text to replace, preferably unescaped. For single replacements (default), include at least 3 lines of context BEFORE and AFTER the target text, matching whitespace and indentation precisely. For multiple replacements, specify expected_replacements parameter. If this string is not the exact literal text (i.e. you escaped it) or does not match exactly, the tool will fail.
new_string: |
The exact literal text to replace `old_string` with, preferably unescaped. Provide the EXACT text. Ensure the resulting code is correct and idiomatic.
expected_replacements: |
Optional: Number of replacements expected. Defaults to 1 if not specified. Use when you want to replace multiple occurrences.

View file

@ -1,295 +0,0 @@
"""Glob file search operation module.
This module provides a tool operation for finding files matching glob patterns.
It enables efficient file discovery based on name or path structure, especially
in large codebases. Files are sorted by modification time (newest first).
"""
import fnmatch
import time
from pathlib import Path
from typing import List, Optional
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
from loguru import logger
from pathspec import PathSpec
from pathspec.patterns.gitwildmatch import GitIgnorePattern
@C.register_op()
class GlobOp(BaseAsyncToolOp):
"""Glob file search operation.
This operation efficiently finds files matching specific glob patterns,
returning absolute paths sorted by modification time (newest first).
Supports gitignore patterns for filtering files.
"""
file_path = __file__
def __init__(self, gitignore_patterns: List[str] = None, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
self.gitignore_patterns = gitignore_patterns
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "FindFiles",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"pattern": {
"type": "string",
"description": self.get_prompt("pattern"),
"required": True,
},
"dir_path": {
"type": "string",
"description": self.get_prompt("dir_path"),
"required": False,
},
"case_sensitive": {
"type": "boolean",
"description": self.get_prompt("case_sensitive"),
"required": False,
},
},
}
return ToolCall(**tool_params)
def _should_ignore_file(
self,
file_path: Path,
root_dir: Path,
) -> bool:
"""Check if a file should be ignored based on ignore patterns.
Args:
file_path: Path to the file to check.
root_dir: Root directory for resolving relative paths.
Returns:
True if the file should be ignored, False otherwise.
"""
if not self.gitignore_patterns:
return False
ignore_spec = PathSpec.from_lines(GitIgnorePattern, self.gitignore_patterns)
# Get relative path from root_dir
relative_path = file_path.relative_to(root_dir)
# Check if file matches ignore patterns
return ignore_spec.match_file(str(relative_path))
async def async_execute(self):
"""Execute the glob search operation."""
pattern: str = self.input_dict.get("pattern", "").strip()
dir_path: Optional[str] = self.input_dict.get("dir_path")
case_sensitive: bool = self.input_dict.get("case_sensitive", False)
# Validate pattern
if not pattern:
error_msg = "The 'pattern' parameter cannot be empty."
logger.error(f"{self.name}: {error_msg}")
self.set_output(error_msg)
return
# Determine search directory
if dir_path:
search_dir = Path(dir_path).expanduser().resolve()
if not search_dir.exists():
error_msg = f"Search path does not exist: {search_dir}"
logger.error(f"{self.name}: {error_msg}")
self.set_output(error_msg)
return
if not search_dir.is_dir():
error_msg = f"Search path is not a directory: {search_dir}"
logger.error(f"{self.name}: {error_msg}")
self.set_output(error_msg)
return
else:
search_dir = Path.cwd()
# Collect matching files
all_entries: List[Path] = []
ignored_count = 0
# Check if pattern is an exact file path
full_path = search_dir / pattern
if full_path.exists() and full_path.is_file():
# Use exact match
if not self._should_ignore_file(
full_path,
search_dir,
):
all_entries.append(full_path)
else:
# Use glob pattern matching
matching_files = self._glob_match(
search_dir,
pattern,
case_sensitive=case_sensitive,
)
# Filter by ignore patterns
for file_path in matching_files:
if self._should_ignore_file(
file_path,
search_dir,
):
ignored_count += 1
else:
all_entries.append(file_path)
# Check if any files found
if not all_entries:
message = f'No files found matching pattern "{pattern}"'
if dir_path:
message += f" within {search_dir}"
if ignored_count > 0:
message += f" ({ignored_count} files were ignored)"
self.set_output(message)
return
# Sort files by modification time
now_timestamp = time.time()
# recency_threshold_days: Number of days to consider a file "recent" for sorting (default: 1).
recency_threshold_ms = self.op_params.get("recency_threshold_days", 1) * 24 * 60 * 60 * 1000
def get_sort_key(path: Path) -> tuple:
mtime = path.stat().st_mtime
mtime_ms = mtime * 1000
is_recent = (now_timestamp * 1000) - mtime_ms < recency_threshold_ms
if is_recent:
# Recent files: sort by mtime descending (newest first)
return 0, -mtime_ms
else:
# Old files: sort alphabetically
return 1, str(path)
sorted_entries = sorted(all_entries, key=get_sort_key)
# Format results
sorted_absolute_paths = [str(entry.resolve()) for entry in sorted_entries]
file_list_description = "\n".join(sorted_absolute_paths)
file_count = len(sorted_absolute_paths)
result_message = f'Found {file_count} file(s) matching "{pattern}"'
if dir_path:
result_message += f" within {search_dir}"
if ignored_count > 0:
result_message += f" ({ignored_count} additional files were ignored)"
result_message += ", sorted by modification time (newest first):\n"
result_message += file_list_description
self.set_output(result_message)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
pattern: str = self.input_dict.get("pattern", "").strip()
error_msg = f'Failed to search files matching pattern "{pattern}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)
def _glob_match(
self,
root_dir: Path,
pattern: str,
case_sensitive: bool = False,
) -> List[Path]:
"""Match files using glob pattern.
Args:
root_dir: Root directory to search in.
pattern: Glob pattern to match.
case_sensitive: Whether matching should be case-sensitive.
Returns:
List of matching file paths.
"""
matching_files: List[Path] = []
# Handle ** pattern (recursive match)
if "**" in pattern:
# Split pattern into parts
parts = pattern.split("**", 1)
prefix = parts[0].rstrip("/")
suffix = parts[1] if len(parts) > 1 else ""
# Walk through directory tree
for path in root_dir.rglob("*"):
if not path.is_file():
continue
rel_path = path.relative_to(root_dir)
rel_str = str(rel_path).replace("\\", "/")
# Check if path matches pattern
if self._matches_glob_pattern(
rel_str,
prefix,
suffix,
case_sensitive,
):
matching_files.append(path)
else:
# For patterns without **, use fnmatch for better compatibility
# Walk through directory tree and match manually
pattern_normalized = pattern.replace("\\", "/")
for path in root_dir.rglob("*"):
if not path.is_file():
continue
rel_path = path.relative_to(root_dir)
rel_str = str(rel_path).replace("\\", "/")
# Match using fnmatch
if case_sensitive:
if fnmatch.fnmatch(rel_str, pattern_normalized):
matching_files.append(path)
else:
if fnmatch.fnmatch(rel_str.lower(), pattern_normalized.lower()):
matching_files.append(path)
return matching_files
@staticmethod
def _matches_glob_pattern(
path_str: str,
prefix: str,
suffix: str,
case_sensitive: bool,
) -> bool:
"""Check if a path matches a glob pattern with **.
Args:
path_str: Path string to check (relative to root).
prefix: Prefix pattern before **.
suffix: Suffix pattern after **.
case_sensitive: Whether matching should be case-sensitive.
Returns:
True if path matches pattern, False otherwise.
"""
if not case_sensitive:
path_str = path_str.lower()
prefix = prefix.lower()
suffix = suffix.lower()
# Check prefix
if prefix:
if not path_str.startswith(prefix):
return False
# Check suffix
if suffix:
if not fnmatch.fnmatch(path_str, f"*{suffix}"):
return False
return True

View file

@ -1,11 +0,0 @@
tool_desc: |
Efficiently finds files matching specific glob patterns (e.g., `src/**/*.ts`, `**/*.md`), returning absolute paths sorted by modification time (newest first). Ideal for quickly locating files based on their name or path structure, especially in large codebases.
pattern: |
The glob pattern to match against (e.g., '**/*.py', 'docs/*.md').
dir_path: |
Optional: The absolute path to the directory to search within. If omitted, searches the current working directory.
case_sensitive: |
Optional: Whether the search should be case-sensitive. Defaults to false.

View file

@ -1,176 +0,0 @@
"""Grep text search operation module.
This module provides a tool operation for searching text patterns in files.
It enables efficient content-based search using regular expressions, with support
for glob pattern filtering and result limiting.
"""
import fnmatch
import re
from pathlib import Path
from typing import List, Optional
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
from loguru import logger
@C.register_op()
class GrepOp(BaseAsyncToolOp):
"""Grep text search operation.
This operation searches for text patterns in files using regular expressions.
Supports glob pattern filtering and result limiting.
"""
file_path = __file__
def __init__(self, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "Grep",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"pattern": {
"type": "string",
"description": self.get_prompt("pattern"),
"required": True,
},
"path": {
"type": "string",
"description": self.get_prompt("path"),
"required": False,
},
"glob": {
"type": "string",
"description": self.get_prompt("glob"),
"required": False,
},
"limit": {
"type": "number",
"description": self.get_prompt("limit"),
"required": False,
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the grep search operation."""
pattern: str = self.input_dict.get("pattern", "").strip()
path: Optional[str] = self.input_dict.get("path")
glob_pattern: Optional[str] = self.input_dict.get("glob")
limit: Optional[int] = self.input_dict.get("limit")
# Validate pattern
if not pattern:
raise ValueError("The 'pattern' parameter cannot be empty.")
# Validate regex pattern
try:
regex = re.compile(pattern, re.IGNORECASE)
except re.error as e:
raise ValueError(f"Invalid regular expression pattern: {pattern}. Error: {str(e)}")
# Determine search directory
if path:
search_dir = Path(path).expanduser().resolve()
if not search_dir.exists():
raise ValueError(f"Search path does not exist: {search_dir}")
if not search_dir.is_dir():
raise ValueError(f"Search path is not a directory: {search_dir}")
else:
search_dir = Path.cwd()
# Collect matching files based on glob pattern
files_to_search: List[Path] = []
if glob_pattern:
# Use glob pattern to filter files
glob_normalized = glob_pattern.replace("\\", "/")
for file_path in search_dir.rglob("*"):
if not file_path.is_file():
continue
rel_path = file_path.relative_to(search_dir)
rel_str = str(rel_path).replace("\\", "/")
if fnmatch.fnmatch(rel_str.lower(), glob_normalized.lower()):
files_to_search.append(file_path)
else:
# Search all files recursively
files_to_search = [f for f in search_dir.rglob("*") if f.is_file()]
# Search for matches
matches: List[dict] = []
for file_path in files_to_search:
try:
content = file_path.read_text(encoding="utf-8", errors="ignore")
lines = content.split("\n")
for line_num, line in enumerate(lines, start=1):
if regex.search(line):
try:
relative_path = file_path.relative_to(search_dir)
except ValueError:
relative_path = file_path.name
matches.append(
{
"file_path": str(relative_path),
"line_number": line_num,
"line": line,
}
)
if limit and len(matches) >= limit:
break
if limit and len(matches) >= limit:
break
except Exception as e:
logger.debug(f"Could not read {file_path}: {str(e)}")
continue
# Format results
if not matches:
search_location = f'in path "{path}"' if path else "in the workspace directory"
filter_desc = f' (filter: "{glob_pattern}")' if glob_pattern else ""
result_msg = f'No matches found for pattern "{pattern}" {search_location}{filter_desc}.'
self.set_output(result_msg)
return
# Group matches by file
matches_by_file = {}
for match in matches:
file_key = match["file_path"]
if file_key not in matches_by_file:
matches_by_file[file_key] = []
matches_by_file[file_key].append(match)
# Build output
total_matches = len(matches)
match_term = "match" if total_matches == 1 else "matches"
search_location = f'in path "{path}"' if path else "in the workspace directory"
filter_desc = f' (filter: "{glob_pattern}")' if glob_pattern else ""
output_lines = [
f'Found {total_matches} {match_term} for pattern "{pattern}" {search_location}{filter_desc}:\n---',
]
for file_path in sorted(matches_by_file.keys()):
output_lines.append(f"File: {file_path}")
for match in sorted(matches_by_file[file_path], key=lambda x: x["line_number"]):
trimmed_line = match["line"].strip()
output_lines.append(f"L{match['line_number']}: {trimmed_line}")
output_lines.append("---")
result_msg = "\n".join(output_lines)
self.set_output(result_msg)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
pattern: str = self.input_dict.get("pattern", "").strip()
error_msg = f'Failed to search for pattern "{pattern}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,120 +0,0 @@
"""Read file operation module.
This module provides a tool operation for reading file contents.
It supports reading entire files or specific line ranges for large files.
"""
from pathlib import Path
from typing import Optional
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
@C.register_op()
class ReadFileOp(BaseAsyncToolOp):
"""Read file operation.
This operation reads and returns the content of a specified file.
For text files, it can read specific line ranges using offset and limit.
"""
file_path = __file__
def __init__(self, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "ReadFile",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"absolute_path": {
"type": "string",
"description": self.get_prompt("absolute_path"),
"required": True,
},
"offset": {
"type": "number",
"description": self.get_prompt("offset"),
"required": False,
},
"limit": {
"type": "number",
"description": self.get_prompt("limit"),
"required": False,
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the read file operation."""
absolute_path: str = self.input_dict.get("absolute_path", "").strip()
offset: Optional[int] = self.input_dict.get("offset")
limit: Optional[int] = self.input_dict.get("limit")
# Validate absolute_path
if not absolute_path:
raise ValueError("The 'absolute_path' parameter cannot be empty.")
# Resolve file path
file_path_obj = Path(absolute_path).expanduser().resolve()
# Check if file exists
if not file_path_obj.exists():
raise FileNotFoundError(f"File not found: {file_path_obj}")
if not file_path_obj.is_file():
raise ValueError(f"Path is not a file: {file_path_obj}")
# Read file content
content = file_path_obj.read_text(encoding="utf-8")
lines = content.split("\n")
# Handle line range if specified
if offset is not None or limit is not None:
if offset is None:
offset = 0
if limit is None:
limit = len(lines)
# Validate offset and limit
if offset < 0:
raise ValueError("Offset must be a non-negative number")
if limit <= 0:
raise ValueError("Limit must be a positive number")
total_lines = len(lines)
start = offset
end = min(offset + limit, total_lines)
if start >= total_lines:
raise ValueError(
f"Offset {offset} is beyond file length ({total_lines} lines)",
)
selected_lines = lines[start:end]
result_content = "\n".join(selected_lines)
# Format output with range information
if end < total_lines:
result = f"Showing lines {start}-{end-1} of {total_lines} total lines.\n\n---\n\n{result_content}"
else:
result = result_content
else:
result = content
self.set_output(result)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
absolute_path: str = self.input_dict.get("absolute_path", "").strip()
error_msg = f'Failed to read file "{absolute_path}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,160 +0,0 @@
"""Ripgrep text search operation module.
This module provides a tool operation for searching text patterns in files
using ripgrep (rg). It enables efficient content-based search using regular
expressions with support for glob pattern filtering and result limiting.
"""
import asyncio
from pathlib import Path
from typing import Optional
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
from loguru import logger
@C.register_op()
class RipGrepOp(BaseAsyncToolOp):
"""Ripgrep text search operation.
This operation searches for text patterns in files using ripgrep (rg).
Supports glob pattern filtering and result limiting.
"""
file_path = __file__
def __init__(self, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "RipGrep",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"pattern": {
"type": "string",
"description": self.get_prompt("pattern"),
"required": True,
},
"path": {
"type": "string",
"description": self.get_prompt("path"),
"required": False,
},
"glob": {
"type": "string",
"description": self.get_prompt("glob"),
"required": False,
},
"limit": {
"type": "number",
"description": self.get_prompt("limit"),
"required": False,
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the ripgrep search operation."""
pattern: str = self.input_dict.get("pattern", "").strip()
path: Optional[str] = self.input_dict.get("path")
glob_pattern: Optional[str] = self.input_dict.get("glob")
limit: Optional[int] = self.input_dict.get("limit")
# Validate pattern
if not pattern:
raise ValueError("The 'pattern' parameter cannot be empty.")
# Determine search path
if path:
search_path = Path(path).expanduser().resolve()
if not search_path.exists():
raise ValueError(f"Search path does not exist: {search_path}")
else:
search_path = Path.cwd()
# Build ripgrep command
rg_args = [
"rg",
"--line-number",
"--no-heading",
"--with-filename",
"--ignore-case",
"--regexp",
pattern,
]
# Add glob pattern if provided
if glob_pattern:
rg_args.extend(["--glob", glob_pattern])
# Add search path
rg_args.append(str(search_path))
# Execute ripgrep
try:
process = await asyncio.create_subprocess_exec(
*rg_args,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
stdout, stderr = await process.communicate()
except FileNotFoundError:
raise ValueError("ripgrep (rg) is not installed or not in PATH")
# Handle ripgrep exit codes
if process.returncode == 0:
raw_output = stdout.decode("utf-8").strip()
elif process.returncode == 1:
# No matches found
raw_output = ""
else:
error_msg = stderr.decode("utf-8").strip()
raise ValueError(f"ripgrep exited with code {process.returncode}: {error_msg}")
# Build search description
search_location = f'in path "{path}"' if path else "in the workspace directory"
filter_desc = f' (filter: "{glob_pattern}")' if glob_pattern else ""
# Check if we have any matches
if not raw_output:
result_msg = f'No matches found for pattern "{pattern}" {search_location}{filter_desc}.'
self.set_output(result_msg)
return
# Split into lines and apply limit
all_lines = [line for line in raw_output.split("\n") if line.strip()]
total_matches = len(all_lines)
match_term = "match" if total_matches == 1 else "matches"
# Apply limit if specified
lines_to_include = all_lines
truncated = False
if limit and len(all_lines) > limit:
lines_to_include = all_lines[:limit]
truncated = True
# Build output
header = f'Found {total_matches} {match_term} for pattern "{pattern}" {search_location}{filter_desc}:\n---\n'
grep_output = "\n".join(lines_to_include)
result_msg = header + grep_output
if truncated:
omitted = total_matches - len(lines_to_include)
result_msg += f"\n---\n[{omitted} {'line' if omitted == 1 else 'lines'} truncated] ..."
self.set_output(result_msg)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
pattern: str = self.input_dict.get("pattern", "").strip()
error_msg = f'Failed to search for pattern "{pattern}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,15 +0,0 @@
tool_desc: |
A powerful search tool built on ripgrep for finding patterns in files using regular expressions. Supports full regex syntax (e.g., "log.*Error", "function\\s+\\w+"), glob pattern filtering, and result limiting. Ideal for efficient searching across large codebases.
pattern: |
The regular expression pattern to search for in file contents.
path: |
Optional: The file or directory to search in. Defaults to current working directory.
glob: |
Optional: Glob pattern to filter files (e.g., "*.js", "*.{ts,tsx}").
limit: |
Optional: Maximum number of matching lines to return. Shows all matches if not specified.

View file

@ -1,89 +0,0 @@
"""Write file operation module.
This module provides a tool operation for writing content to files.
It supports creating new files or overwriting existing files, and automatically
creates parent directories if they don't exist.
"""
from pathlib import Path
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
@C.register_op()
class WriteFileOp(BaseAsyncToolOp):
"""Write file operation.
This operation writes content to a specified file. If the file doesn't exist,
it will be created. If parent directories don't exist, they will be created automatically.
"""
file_path = __file__
def __init__(self, **kwargs):
kwargs.setdefault("raise_exception", False)
super().__init__(**kwargs)
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "WriteFile",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"file_path": {
"type": "string",
"description": self.get_prompt("file_path"),
"required": True,
},
"content": {
"type": "string",
"description": self.get_prompt("content"),
"required": True,
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the write file operation."""
file_path: str = self.input_dict.get("file_path", "").strip()
content: str = self.input_dict.get("content", "")
# Validate file_path
if not file_path:
raise ValueError("The 'file_path' parameter cannot be empty.")
# Resolve file path
file_path_obj = Path(file_path).expanduser().resolve()
# Check if path is a directory
if file_path_obj.exists() and file_path_obj.is_dir():
raise ValueError(f"Path is a directory, not a file: {file_path_obj}")
# Create parent directories if they don't exist
file_path_obj.parent.mkdir(parents=True, exist_ok=True)
# Check if file exists
file_exists = file_path_obj.exists() and file_path_obj.is_file()
# Write content to file
file_path_obj.write_text(content, encoding="utf-8")
# Format success message
if file_exists:
result = f"Successfully overwrote file: {file_path_obj}"
else:
result = f"Successfully created and wrote to new file: {file_path_obj}"
self.set_output(result)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
file_path: str = self.input_dict.get("file_path", "").strip()
error_msg = f'Failed to write file "{file_path}"'
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,9 +0,0 @@
tool_desc: |
Writes content to a specified file in the local filesystem. If the file doesn't exist, it will be created. If parent directories don't exist, they will be created automatically. If the file already exists, it will be overwritten with the new content.
file_path: |
The absolute path to the file to write to (e.g., '/home/user/project/file.txt'). Relative paths are not supported. You must provide an absolute path.
content: |
The content to write to the file.

View file

@ -1,102 +0,0 @@
"""Write todos operation module.
This module provides a tool operation for managing todo lists.
It enables tracking subtasks with status (pending, in_progress, completed, cancelled).
"""
from typing import List
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncToolOp
from flowllm.core.schema import ToolCall
TODO_STATUSES = ["pending", "in_progress", "completed", "cancelled"]
@C.register_op()
class WriteTodosOp(BaseAsyncToolOp):
"""Write todos operation.
This operation manages a todo list with subtasks that can be tracked
through different statuses: pending, in_progress, completed, cancelled.
"""
file_path = __file__
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "WriteTodos",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"todos": {
"type": "array",
"description": self.get_prompt("todos"),
"required": True,
"items": {
"type": "object",
"properties": {
"description": {
"type": "string",
"description": self.get_prompt("todo_description"),
},
"status": {
"type": "string",
"description": self.get_prompt("todo_status"),
"enum": TODO_STATUSES,
},
},
"required": ["description", "status"],
},
},
},
}
return ToolCall(**tool_params)
async def async_execute(self):
"""Execute the write todos operation."""
todos: List[dict] = self.input_dict.get("todos", [])
# Validate todos
if not isinstance(todos, list):
raise ValueError("The 'todos' parameter must be an array")
# Validate each todo item
for i, todo in enumerate(todos):
if not isinstance(todo, dict):
raise ValueError(f"Todo item at index {i} must be an object")
if "description" not in todo or not isinstance(todo["description"], str):
raise ValueError(f"Todo item at index {i} must have a non-empty description string")
if not todo["description"].strip():
raise ValueError(f"Todo item at index {i} must have a non-empty description string")
if "status" not in todo or todo["status"] not in TODO_STATUSES:
raise ValueError(
f"Todo item at index {i} must have a valid status ({', '.join(TODO_STATUSES)})",
)
# Validate only one in_progress task
in_progress_count = sum(1 for todo in todos if todo.get("status") == "in_progress")
if in_progress_count > 1:
raise ValueError("Only one task can be 'in_progress' at a time")
# Format todo list
if not todos:
result_message = "Successfully cleared the todo list."
else:
todo_list_string = "\n".join(
f"{i + 1}. [{todo['status']}] {todo['description']}" for i, todo in enumerate(todos)
)
result_message = f"Successfully updated the todo list. The current list is now:\n{todo_list_string}"
self.set_output(result_message)
async def async_default_execute(self, e: Exception = None, **kwargs):
"""Fill outputs with a default failure message when execution fails."""
error_msg = "Failed to update the todo list"
if e:
error_msg += f": {str(e)}"
self.set_output(error_msg)

View file

@ -1,35 +0,0 @@
tool_desc: |
This tool can help you list out the current subtasks that are required to be completed for a given user request. The list of subtasks helps you keep track of the current task, organize complex queries and help ensure that you don't miss any steps. With this list, the user can also see the current progress you are making in executing a given task.
Depending on the task complexity, you should first divide a given task into subtasks and then use this tool to list out the subtasks that are required to be completed for a given user request.
Each of the subtasks should be clear and distinct.
Use this tool for complex queries that require multiple steps. If you find that the request is actually complex after you have started executing the user task, create a todo list and use it. If execution of the user task requires multiple steps, planning and generally is higher complexity than a simple Q&A, use this tool.
DO NOT use this tool for simple tasks that can be completed in less than 2 steps. If the user query is simple and straightforward, do not use the tool. If you can respond with an answer in a single turn then this tool is not required.
## Task state definitions
- pending: Work has not begun on a given subtask.
- in_progress: Marked just prior to beginning work on a given subtask. You should only have one subtask as in_progress at a time.
- completed: Subtask was successfully completed with no errors or issues. If the subtask required more steps to complete, update the todo list with the subtasks. All steps should be identified as completed only when they are completed.
- cancelled: As you update the todo list, some tasks are not required anymore due to the dynamic nature of the task. In this case, mark the subtasks as cancelled.
## Methodology for using this tool
1. Use this todo list list as soon as you receive a user request based on the complexity of the task.
2. Keep track of every subtask that you update the list with.
3. Mark a subtask as in_progress before you begin working on it. You should only have one subtask as in_progress at a time.
4. Update the subtask list as you proceed in executing the task. The subtask list is not static and should reflect your progress and current plans, which may evolve as you acquire new information.
5. Mark a subtask as completed when you have completed it.
6. Mark a subtask as cancelled if the subtask is no longer needed.
7. You must update the todo list as soon as you start, stop or cancel a subtask. Don't batch or wait to update the todo list.
todos: |
The complete list of todo items. This will replace the existing list.
todo_description: |
The description of the task.
todo_status: |
The current status of the task. Must be one of: pending, in_progress, completed, cancelled.

View file

@ -0,0 +1,19 @@
"""File tool package for ReMe framework.
This package provides file-related operations that can be used in LLM-powered flows.
It includes ready-to-use operations for:
- BatchWriteFileOp: Batch write multiple files operation
- GrepOp: Text search operation for finding patterns in files
- ReadFileOp: Read single file operation
"""
from .batch_write_file_op import BatchWriteFileOp
from .grep_op import GrepOp
from .read_file_op import ReadFileOp
__all__ = [
"BatchWriteFileOp",
"GrepOp",
"ReadFileOp",
]

View file

@ -0,0 +1,45 @@
"""Batch write file operation module.
This module provides a tool operation for batch writing multiple files at once.
It processes a dictionary of file paths and contents, writing each file sequentially
and returning a combined result of all write operations.
"""
from flowllm.core.context import C
from flowllm.core.op import BaseAsyncOp
from flowllm.extensions.file_tool import WriteFileOp
from loguru import logger
@C.register_op()
class BatchWriteFileOp(BaseAsyncOp):
"""Batch write file operation.
This operation writes multiple files in a single batch. It takes a dictionary
of file paths and contents from the context, and writes each file using
WriteFileOp. Returns a combined result of all write operations.
"""
async def async_execute(self):
"""Execute the batch write file operation.
Reads write_file_dict from context, which should be a dictionary mapping
file paths to file contents. Writes each file sequentially and collects
the results.
"""
# Get write file dictionary from context
write_file_dict: dict = self.context.get("write_file_dict", {})
if not write_file_dict:
self.context.response.answer = "No write file task."
logger.info("No write file task.")
return
# Process each file in the dictionary
result = []
for file_path, content in write_file_dict.items():
write_op = WriteFileOp()
await write_op.async_call(file_path=file_path, content=content)
result.append(write_op.output)
# Combine all results into a single response
self.context.response.answer = "\n".join(result)

View file

@ -0,0 +1,52 @@
"""Grep text search operation module.
This module provides a tool operation for searching text patterns in files.
It enables efficient content-based search using regular expressions, with support
for glob pattern filtering and result limiting.
"""
from flowllm.core.context import C
from flowllm.core.schema import ToolCall
from flowllm.extensions.file_tool import GrepOp as FlowGrepOp
@C.register_op()
class GrepOp(FlowGrepOp):
"""Grep text search operation.
This operation searches for text patterns in files using regular expressions.
Supports glob pattern filtering and result limiting.
"""
file_path = __file__
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "Grep",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"pattern": {
"type": "string",
"description": self.get_prompt("pattern"),
"required": True,
},
"path": {
"type": "string",
"description": self.get_prompt("path"),
"required": False,
},
"glob": {
"type": "string",
"description": self.get_prompt("glob"),
"required": False,
},
"limit": {
"type": "number",
"description": self.get_prompt("limit"),
"required": False,
},
},
}
return ToolCall(**tool_params)

View file

@ -0,0 +1,46 @@
"""Read file operation module.
This module provides a tool operation for reading file contents.
It supports reading entire files or specific line ranges for large files.
"""
from flowllm.core.context import C
from flowllm.core.schema import ToolCall
from flowllm.extensions.file_tool import ReadFileOp as FlowReadFileOp
@C.register_op()
class ReadFileOp(FlowReadFileOp):
"""Read file operation.
This operation reads and returns the content of a specified file.
For text files, it can read specific line ranges using offset and limit.
"""
file_path = __file__
def build_tool_call(self) -> ToolCall:
"""Build and return the tool call schema for this operator."""
tool_params = {
"name": "ReadFile",
"description": self.get_prompt("tool_desc"),
"input_schema": {
"absolute_path": {
"type": "string",
"description": self.get_prompt("absolute_path"),
"required": True,
},
"offset": {
"type": "number",
"description": self.get_prompt("offset"),
"required": False,
},
"limit": {
"type": "number",
"description": self.get_prompt("limit"),
"required": False,
},
},
}
return ToolCall(**tool_params)

View file

@ -1,44 +1,145 @@
"""
Context compaction module for reducing token usage in conversation contexts.
This module provides functionality to compress large tool messages by storing
their full content in external files and keeping only previews in the context.
This helps manage context window limits while preserving important information.
"""
import json
from pathlib import Path
from typing import List
from uuid import uuid4
from flowllm.core.context import C
from flowllm.core.enumeration import Role
from flowllm.core.op import BaseAsyncOp
from flowllm.core.schema import ToolCall
from flowllm.core.schema import Message
from loguru import logger
@C.register_op()
class ContextCompactOp(BaseAsyncOp):
"""
Context compaction operation that reduces token usage by compressing tool messages.
def __init__(self, ratio: 0.3, **kwargs):
When the total token count exceeds the threshold, this operation compresses large tool
messages by truncating their content and storing the full content in external files.
This helps manage context window limits while preserving recent tool messages.
"""
def __init__(
self,
all_token_threshold: int = 20000,
tool_token_threshold: int = 2000,
tool_left_char_len: int = 100,
keep_recent: int = 1,
storage_path: str = "./",
exclude_tools: List[str] = None,
**kwargs,
):
"""
Initialize the context compaction operation.
Args:
all_token_threshold: Maximum total token count before compaction is triggered.
tool_token_threshold: Maximum token count for a single tool message before it's compressed.
tool_left_char_len: Number of characters to keep in the compressed tool message preview.
keep_recent: Number of recent tool messages to keep uncompressed.
storage_path: Directory path where compressed tool message contents will be stored.
exclude_tools: List of tool names to exclude from compaction (not currently used).
**kwargs: Additional arguments passed to the base class.
"""
super().__init__(**kwargs)
self.ratio: float = ratio
self.all_token_threshold: int = all_token_threshold
self.tool_token_threshold: int = tool_token_threshold
self.tool_left_char_len: int = tool_left_char_len
self.keep_recent: int = keep_recent
self.storage_path: Path = Path(storage_path)
self.exclude_tools: List[str] = exclude_tools
def async_execute(self):
messages = self.context.messages
async def async_execute(self):
"""
Execute the context compaction operation.
self.llm.achat(messages=messages, tools=None)
self.token_count()
The operation:
1. Calculates the total token count of all messages
2. If below threshold, returns messages unchanged
3. Otherwise, compresses large tool messages by:
- Keeping only a preview of the content
- Storing full content in external files
- Preserving recent tool messages
"""
# Convert context messages to Message objects
messages = [Message(**x) for x in self.context.messages]
self.context["messages"] = messages
self.context["offloaded_data"] = {"f": "ddd"}
# Calculate total token count
token_cnt: int = self.token_count(messages)
logger.info(f"Context compaction check: total token count={token_cnt}, threshold={self.all_token_threshold}")
# If token count is within threshold, no compaction needed
if token_cnt <= self.all_token_threshold:
self.context.response.answer = self.context.messages
logger.info(
f"Token count ({token_cnt}) is within threshold ({self.all_token_threshold}), no compaction needed",
)
return
"""
1. token counter 计数,最长上下文X20%=20K
1. openai tiktoken / hagggingface / modelscope / rule-based
self.token_counter.count()
2. tool:
1. dump_tool(write_file)
2. grep / rip_grep / read_file
3. compact:
if > threashold(20K):
tool_call -> write_file
写一个引用:/xx/xxx/xxx.txt。是否保留前面的token()
4. compress:
1. prompt
2. context -> write_file
# Filter tool messages for processing
tool_messages = [x for x in messages if x.role is Role.TOOL]
# If there are too few tool messages, no compaction needed
if len(tool_messages) <= self.keep_recent:
self.context.response.answer = self.context.messages
logger.info(
f"Tool message count ({len(tool_messages)}) is less than or "
f"equal to keep_recent ({self.keep_recent}), no compaction needed",
)
return
1. case study:
2. messages
1. agentscope: reme
2. http远程 reme
# Exclude recent tool messages from compaction (keep them intact)
tool_messages = tool_messages[: -self.keep_recent]
logger.info(
f"Processing {len(tool_messages)} tool messages for "
f"compaction (keeping {self.keep_recent} recent messages)",
)
"""
# Dictionary to store file paths and their compressed content (for potential batch writing)
write_file_dict = {}
# Process each tool message
for tool_message in tool_messages:
# Calculate token count for this specific tool message
tool_token_cnt = self.token_count([tool_message])
# Skip if token count is within threshold
if tool_token_cnt <= self.tool_token_threshold:
logger.info(
f"Skipping tool message (tool_call_id={tool_message.tool_call_id}): "
f"token count ({tool_token_cnt}) is within threshold ({self.tool_token_threshold})",
)
continue
# Create compressed preview of the tool message content
compact_result = tool_message.content[: self.tool_left_char_len] + "..."
# Generate file name from tool_call_id or create a unique identifier
file_name = tool_message.tool_call_id or uuid4().hex
path = self.storage_path / f"{file_name}.txt"
# Store the mapping for potential batch writing
write_file_dict[str(path)] = compact_result
# Log the compaction action
logger.info(
f"Compacting tool message (tool_call_id={tool_message.tool_call_id}): "
f"token count={tool_token_cnt}, saving full content to {path}",
)
# Update tool message content with preview and file reference
compact_result += f" (detailed result is stored in {path})"
tool_message.content = compact_result
# Return the compacted messages as JSON
self.context.response.answer = json.dumps([x.model_dump() for x in messages], ensure_ascii=False, indent=2)
logger.info(f"Context compaction completed: {len(write_file_dict)} tool messages were compacted")