diff --git a/reme_ai/config/default.yaml b/reme_ai/config/default.yaml index 8e8c233f..4c8fdc81 100644 --- a/reme_ai/config/default.yaml +++ b/reme_ai/config/default.yaml @@ -156,6 +156,83 @@ flow: description: "user query" required: true + context_offload: + flow_content: ContextOffloadOp() >> BatchWriteFileOp() + description: "Manages context window limits by compacting tool messages and compressing conversation history. First compacts large tool messages by storing full content in external files, then applies LLM-based compression if compaction ratio exceeds threshold. This helps reduce token usage while preserving important information." + input_schema: + messages: + type: array + description: "List of conversation messages to process for context offloading" + required: true + context_manage_mode: + type: string + description: "Context management mode: 'compact' only applies compaction to tool messages, 'compress' only applies LLM-based compression, 'auto' applies compaction first then compression if compaction ratio exceeds threshold. Defaults to 'auto'." + required: false + enum: ["compact", "compress", "auto"] + max_total_tokens: + type: integer + description: "Maximum token count threshold for triggering compression/compaction. For compaction, this is the total token count threshold. For compression, this excludes keep_recent_count messages and system messages. Defaults to 20000." + required: false + max_tool_message_tokens: + type: integer + description: "Maximum token count per tool message before compaction is applied. Tool messages exceeding this threshold will have their full content stored in external files with only a preview kept in context. Defaults to 2000." + required: false + group_token_threshold: + type: integer + description: "Maximum token count per compression group when using LLM-based compression. If None or 0, all messages are compressed in a single group. Messages exceeding this threshold individually will form their own group. Only used in 'compress' or 'auto' mode." + required: false + keep_recent_count: + type: integer + description: "Number of recent messages to preserve without compression or compaction. These messages remain unchanged to maintain conversation context. Defaults to 1 for compaction and 2 for compression." + required: false + store_dir: + type: string + description: "Directory path for storing offloaded message content. Full tool message content and compressed message groups are saved as files in this directory. Required for compaction and compression operations." + required: false + chat_id: + type: string + description: "Unique identifier for the chat session, used for file naming when storing compressed message groups. If not provided, a UUID will be generated automatically." + required: false + + context_offload_for_agentscope: + flow_content: ContextOffloadOp() + description: "Context offload operation for AgentScope integration. Manages context window limits by compacting tool messages and compressing conversation history without batch file writing. Same functionality as context_offload but without the BatchWriteFileOp step." + input_schema: + messages: + type: array + description: "List of conversation messages to process for context offloading" + required: true + context_manage_mode: + type: string + description: "Context management mode: 'compact' only applies compaction to tool messages, 'compress' only applies LLM-based compression, 'auto' applies compaction first then compression if compaction ratio exceeds threshold. Defaults to 'auto'." + required: false + enum: ["compact", "compress", "auto"] + max_total_tokens: + type: integer + description: "Maximum token count threshold for triggering compression/compaction. For compaction, this is the total token count threshold. For compression, this excludes keep_recent_count messages and system messages. Defaults to 20000." + required: false + max_tool_message_tokens: + type: integer + description: "Maximum token count per tool message before compaction is applied. Tool messages exceeding this threshold will have their full content stored in external files with only a preview kept in context. Defaults to 2000." + required: false + group_token_threshold: + type: integer + description: "Maximum token count per compression group when using LLM-based compression. If None or 0, all messages are compressed in a single group. Messages exceeding this threshold individually will form their own group. Only used in 'compress' or 'auto' mode." + required: false + keep_recent_count: + type: integer + description: "Number of recent messages to preserve without compression or compaction. These messages remain unchanged to maintain conversation context. Defaults to 1 for compaction and 2 for compression." + required: false + store_dir: + type: string + description: "Directory path for storing offloaded message content. Full tool message content and compressed message groups are saved as files in this directory. Required for compaction and compression operations." + required: false + chat_id: + type: string + description: "Unique identifier for the chat session, used for file naming when storing compressed message groups. If not provided, a UUID will be generated automatically." + required: false + + llm: default: backend: openai_compatible @@ -164,9 +241,7 @@ llm: temperature: 0.6 token_count: # Optional model_name: Qwen/Qwen3-30B-A3B-Instruct-2507 - backend: hf - params: - use_mirror: true + backend: base qwen3_30b_instruct: backend: openai_compatible diff --git a/reme_ai/context/offload/context_compact_op.py b/reme_ai/context/offload/context_compact_op.py index 96a5b1c1..7bd878a1 100644 --- a/reme_ai/context/offload/context_compact_op.py +++ b/reme_ai/context/offload/context_compact_op.py @@ -43,7 +43,11 @@ class ContextCompactOp(BaseAsyncOp): max_tool_message_tokens: int = self.context.get("max_tool_message_tokens", 2000) preview_char_length: int = self.context.get("preview_char_length", 100) keep_recent_count: int = self.context.get("keep_recent_count", 1) - storage_path: Path = Path(self.context.get("storage_path", "")) + store_dir: Path = Path(self.context.get("store_dir", "")) + + assert max_total_tokens > 0, "max_total_tokens must be greater than 0" + assert max_tool_message_tokens > 0, "max_tool_message_tokens must be greater than 0" + assert preview_char_length >= 0, "preview_char_length must be greater than 0" assert keep_recent_count > 0, "keep_recent_count must be greater than 0" # Convert context messages to Message objects @@ -53,7 +57,6 @@ class ContextCompactOp(BaseAsyncOp): # If nothing to compress after filtering, return original messages if not messages_to_compress: self.context.response.answer = self.context.messages - self.context.response.success = True logger.info("No messages to compress after filtering, returning original messages") return @@ -93,10 +96,10 @@ class ContextCompactOp(BaseAsyncOp): # Generate file name from tool_call_id or create a unique identifier file_name = tool_message.tool_call_id or uuid4().hex - path = storage_path / f"{file_name}.txt" + store_path = store_dir / f"{file_name}.txt" # Store the full content for batch writing - write_file_dict[path.as_posix()] = original_content + write_file_dict[store_path.as_posix()] = original_content # Create compressed preview of the tool message content compact_result = original_content[:preview_char_length] + "..." @@ -104,11 +107,11 @@ class ContextCompactOp(BaseAsyncOp): # 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}", + f"token count={tool_token_cnt}, saving full content to {store_path}", ) # Update tool message content with preview and file reference - compact_result += f" (detailed result is stored in {path})" + compact_result += f" (detailed result is stored in {store_path})" tool_message.content = compact_result # Store write_file_dict in context for potential batch writing @@ -117,6 +120,8 @@ class ContextCompactOp(BaseAsyncOp): # Return the compacted messages as JSON self.context.response.answer = [x.simple_dump() for x in messages] + self.context.response.metadata["write_file_dict"] = write_file_dict + logger.info(f"Context compaction completed: {len(write_file_dict)} tool messages were compacted") async def async_default_execute(self, e: Exception = None, **_kwargs): diff --git a/reme_ai/context/offload/context_compress_op.py b/reme_ai/context/offload/context_compress_op.py index cf71a9b3..218f6e21 100644 --- a/reme_ai/context/offload/context_compress_op.py +++ b/reme_ai/context/offload/context_compress_op.py @@ -144,6 +144,7 @@ class ContextCompressOp(BaseAsyncOp): if state_snapshot is None: logger.warning("Failed to extract state_snapshot from LLM response, using full content as fallback") return content + return state_snapshot # Call LLM to generate compressed summary @@ -193,7 +194,7 @@ class ContextCompressOp(BaseAsyncOp): for g_idx, messages in enumerate(message_groups): group_original_tokens = self.token_count(messages) messages_str = json.dumps([x.simple_dump() for x in messages], ensure_ascii=False, indent=2) - store_path = Path(self.context.store_dir) / f"{chat_id}_{g_idx}" + store_path = Path(self.context.get("store_dir", "")) / f"{chat_id}_{g_idx}.json" logger.info(f"Compress {g_idx}/{len(message_groups)} ({len(messages)}, {group_original_tokens} tokens)") group_summary = await self._compress_messages_with_llm(messages) @@ -216,7 +217,7 @@ class ContextCompressOp(BaseAsyncOp): ) return_messages.extend(messages) else: - system_message_copy.content += compress_content + system_message_copy.content += compress_content + "\n\n" write_file_dict[store_path.as_posix()] = messages_str logger.info( f"Group {g_idx} compression successful: " @@ -248,18 +249,15 @@ class ContextCompressOp(BaseAsyncOp): group_token_threshold: int = self.context.get("group_token_threshold", None) keep_recent_count: int = self.context.get("keep_recent_count", 2) - # Validate keep_recent_count - if keep_recent_count < 0: - logger.warning(f"keep_recent_count ({keep_recent_count}) is negative, setting to 0") - keep_recent_count = 0 + assert max_total_tokens > 0, "max_total_tokens must be positive" + assert keep_recent_count >= 0, "keep_recent_count must be non-negative" # Convert context messages to Message objects messages = [Message(**x) for x in self.context.messages] # Extract system message (should be exactly one) system_message = [x for x in messages if x.role is Role.SYSTEM] - if len(system_message) > 1: - raise ValueError("Expected exactly one system message") + assert len(system_message) <= 1, f"Expected at most one system message, got {len(system_message)}" if len(system_message) == 0: system_message = Message(role=Role.SYSTEM, content="") @@ -267,10 +265,8 @@ class ContextCompressOp(BaseAsyncOp): system_message = system_message[0] messages_without_system = [x for x in messages if x.role is not Role.SYSTEM] - messages_to_compress = ( - messages_without_system[:-keep_recent_count] if keep_recent_count > 0 else messages_without_system - ) - recent_messages = messages_without_system[-keep_recent_count:] if keep_recent_count > 0 else [] + messages_to_compress = messages_without_system[:-keep_recent_count] + recent_messages = messages_without_system[-keep_recent_count:] # If nothing to compress after filtering, return original messages if not messages_to_compress: @@ -297,8 +293,12 @@ class ContextCompressOp(BaseAsyncOp): write_file_dict, return_messages = await self._compress_with_groups(system_message, message_groups) - self.context.write_file_dict = write_file_dict + # Store write_file_dict in context for potential batch writing + if write_file_dict: + self.context.write_file_dict = write_file_dict + self.context.response.answer = [x.simple_dump() for x in (return_messages + recent_messages)] + self.context.response.metadata["write_file_dict"] = write_file_dict async def async_default_execute(self, e: Exception = None, **_kwargs): """Handle execution errors by returning original messages.