feat(context): implement context offload with compaction and compression

This commit is contained in:
jinli.yl 2025-11-20 11:13:26 +08:00
parent 3156c3e0d3
commit 1886da9020
3 changed files with 102 additions and 22 deletions

View file

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

View file

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

View file

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