diff --git a/frontend/terminal/src/App.tsx b/frontend/terminal/src/App.tsx index 702418a..33160a4 100644 --- a/frontend/terminal/src/App.tsx +++ b/frontend/terminal/src/App.tsx @@ -437,7 +437,7 @@ function AppInner({config}: {config: FrontendConfig}): React.JSX.Element { setInput={setInput} onSubmit={onSubmit} toolName={session.busy ? currentToolName : undefined} - statusLabel={session.busy ? (currentToolName ? `Running ${currentToolName}...` : 'Running agent loop...') : undefined} + statusLabel={session.busy ? (session.busyLabel ?? (currentToolName ? `Running ${currentToolName}...` : 'Running agent loop...')) : undefined} suppressSubmit={showPicker} /> )} diff --git a/frontend/terminal/src/hooks/useBackendSession.ts b/frontend/terminal/src/hooks/useBackendSession.ts index 6eede51..5e30fee 100644 --- a/frontend/terminal/src/hooks/useBackendSession.ts +++ b/frontend/terminal/src/hooks/useBackendSession.ts @@ -31,6 +31,7 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number const [modal, setModal] = useState | null>(null); const [selectRequest, setSelectRequest] = useState<{title: string; command: string; options: SelectOptionPayload[]} | null>(null); const [busy, setBusy] = useState(false); + const [busyLabel, setBusyLabel] = useState(undefined); const [ready, setReady] = useState(false); const [todoMarkdown, setTodoMarkdown] = useState(''); const [swarmTeammates, setSwarmTeammates] = useState([]); @@ -197,6 +198,43 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number return; } setTranscript((items) => [...items, {role: 'status', text: message}]); + if (busy) { + setBusyLabel(message); + } + return; + } + if (event.type === 'compact_progress') { + const phase = String(event.compact_phase ?? ''); + const trigger = String(event.compact_trigger ?? ''); + const attempt = event.attempt != null ? Number(event.attempt) : undefined; + if (phase === 'hooks_start') { + setBusyLabel( + trigger === 'reactive' + ? 'Preparing retry compaction…' + : 'Preparing conversation compaction…', + ); + } else if (phase === 'context_collapse_start') { + setBusyLabel('Collapsing oversized context…'); + } else if (phase === 'context_collapse_end') { + setBusyLabel('Context collapse complete…'); + } else if (phase === 'session_memory_start') { + setBusyLabel('Condensing earlier conversation…'); + } else if (phase === 'compact_start') { + setBusyLabel( + trigger === 'reactive' + ? 'Context is too large. Compacting and retrying…' + : 'Compacting conversation memory…', + ); + } else if (phase === 'compact_retry') { + setBusyLabel(attempt ? `Retrying compaction (${attempt})…` : 'Retrying compaction…'); + } else if (phase === 'compact_end') { + setBusyLabel('Compaction complete. Continuing…'); + } else if (phase === 'compact_failed') { + setBusyLabel('Compaction failed. Continuing without it…'); + } + if (event.message) { + setTranscript((items) => [...items, {role: 'status', text: event.message!}]); + } return; } if (event.type === 'assistant_delta') { @@ -227,6 +265,7 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number setTranscript((items) => [...items, {role: 'assistant', text}]); clearAssistantDelta(); setBusy(false); + setBusyLabel(undefined); return; } if (event.type === 'line_complete') { @@ -234,9 +273,13 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number // don't leave stale streaming text on screen. clearAssistantDelta(); setBusy(false); + setBusyLabel(undefined); return; } if ((event.type === 'tool_started' || event.type === 'tool_completed') && event.item) { + if (event.type === 'tool_started') { + setBusyLabel(event.tool_name ? `Running ${event.tool_name}...` : 'Running...'); + } const enrichedItem: TranscriptItem = { ...event.item, tool_name: event.item.tool_name ?? event.tool_name ?? undefined, @@ -249,6 +292,7 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number if (event.type === 'clear_transcript') { setTranscript([]); clearAssistantDelta(); + setBusyLabel(undefined); return; } if (event.type === 'select_request') { @@ -268,6 +312,7 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number setTranscript((items) => [...items, {role: 'system', text: `error: ${event.message ?? 'unknown error'}`}]); clearAssistantDelta(); setBusy(false); + setBusyLabel(undefined); return; } if (event.type === 'todo_update') { @@ -308,6 +353,7 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number modal, selectRequest, busy, + busyLabel, ready, todoMarkdown, swarmTeammates, @@ -317,6 +363,6 @@ export function useBackendSession(config: FrontendConfig, onExit: (code?: number setBusy, sendRequest, }), - [assistantBuffer, bridgeSessions, busy, commands, mcpServers, modal, ready, selectRequest, status, swarmNotifications, swarmTeammates, tasks, todoMarkdown, transcript] + [assistantBuffer, bridgeSessions, busy, busyLabel, commands, mcpServers, modal, ready, selectRequest, status, swarmNotifications, swarmTeammates, tasks, todoMarkdown, transcript] ); } diff --git a/frontend/terminal/src/types.ts b/frontend/terminal/src/types.ts index d72bf8f..2f23cff 100644 --- a/frontend/terminal/src/types.ts +++ b/frontend/terminal/src/types.ts @@ -78,6 +78,11 @@ export type BackendEvent = { tool_name?: string | null; output?: string | null; is_error?: boolean | null; + compact_phase?: string | null; + compact_trigger?: string | null; + attempt?: number | null; + compact_checkpoint?: string | null; + compact_metadata?: Record | null; // New event payloads todo_items?: TodoItemSnapshot[] | null; todo_markdown?: string | null; diff --git a/ohmo/gateway/runtime.py b/ohmo/gateway/runtime.py index 2354b26..167cf81 100644 --- a/ohmo/gateway/runtime.py +++ b/ohmo/gateway/runtime.py @@ -18,6 +18,7 @@ from openharness.engine.query import MaxTurnsExceeded from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, ErrorEvent, StatusEvent, ToolExecutionCompleted, @@ -326,6 +327,32 @@ class OhmoSessionRuntimePool: if isinstance(event, AssistantTextDelta): reply_parts.append(event.text) return + if isinstance(event, CompactProgressEvent): + logger.info( + "ohmo runtime compact progress session_key=%s session_id=%s phase=%s trigger=%s attempt=%s", + session_key, + bundle.session_id, + event.phase, + event.trigger, + event.attempt, + ) + rendered = _format_channel_progress( + channel=message.channel, + kind="compact_progress", + text=event.message or "", + session_key=session_key, + content=content, + compact_phase=event.phase, + compact_trigger=event.trigger, + attempt=event.attempt, + ) + if rendered: + yield GatewayStreamUpdate( + kind="progress", + text=rendered, + metadata={"_progress": True, "_session_key": session_key, "_compact": True}, + ) + return if isinstance(event, StatusEvent): logger.info( "ohmo runtime status session_key=%s session_id=%s message=%r", @@ -490,6 +517,9 @@ def _format_channel_progress( text: str, session_key: str, content: str, + compact_phase: str | None = None, + compact_trigger: str | None = None, + attempt: int | None = None, ) -> str: if channel not in { "feishu", @@ -525,13 +555,55 @@ def _format_channel_progress( if text.startswith(("🤔", "🧠", "✨", "🔎", "🪄", "🛠️", "🫧")): return text return f"🫧 {text}" + if kind == "compact_progress": + if compact_phase == "hooks_start": + if prefers_chinese: + if compact_trigger == "reactive": + return "🫧 上下文有点超长,我先准备压缩一下记忆,然后立刻继续重试~" + return "🫧 我先把上下文和记忆准备一下,马上开始压缩重点~" + if compact_trigger == "reactive": + return "🫧 The context got too large. I’m preparing a quick memory compaction before retrying." + return "🫧 Let me get the context ready before I compact the conversation." + if compact_phase == "context_collapse_start": + if prefers_chinese: + return "🫧 我先把太长的上下文折叠一下,让后面的压缩更快一点~" + return "🫧 I’m collapsing the oversized context first so compaction can move faster." + if compact_phase == "context_collapse_end": + if prefers_chinese: + return "🫧 上下文已经先收紧了一层,继续压缩重点~" + return "🫧 The context is trimmed down now. Continuing with the main compaction." + if compact_phase in {"session_memory_start", "compact_start"}: + if prefers_chinese: + if compact_phase == "session_memory_start": + return "🧠 我先把前面的聊天重点悄悄捋顺一下,马上继续~" + if compact_trigger == "reactive": + return "🧠 这轮上下文太长了,我先压缩一下记忆,然后马上继续重试~" + return "🧠 聊天有点长啦,我先帮你悄悄压缩一下记忆,马上继续~" + if compact_phase == "session_memory_start": + return "🧠 Let me quickly condense the earlier parts of this chat, then I’ll keep going." + if compact_trigger == "reactive": + return "🧠 The context is too large for this turn. I’ll compact the memory and retry." + return "🧠 This chat is getting long. I’ll compact the memory and keep going." + if compact_phase == "compact_retry": + suffix = f" (attempt {attempt})" if attempt is not None else "" + if prefers_chinese: + return f"🔁 压缩记忆这一步有点卡,我换个方式再试一次{suffix}。" + return f"🔁 Compaction got stuck, trying a lighter retry{suffix}." + if compact_phase == "compact_failed": + if prefers_chinese: + return "⚠️ 这次记忆压缩没成功,我先跳过它继续处理你的消息。" + return "⚠️ Memory compaction did not complete. I’m skipping it and continuing." + return "" return text def _build_inbound_user_message(message: InboundMessage) -> ConversationMessage: """Convert an inbound channel message into user content blocks.""" content: list[TextBlock | ImageBlock] = [] + speaker_context = _build_speaker_context(message) base = (message.content or "").strip() + if speaker_context: + content.append(TextBlock(text=speaker_context)) if base: content.append(TextBlock(text=base)) @@ -551,6 +623,26 @@ def _build_inbound_user_message(message: InboundMessage) -> ConversationMessage: return ConversationMessage.from_user_content(content) +def _build_speaker_context(message: InboundMessage) -> str: + """Return a lightweight speaker header for group-chat messages.""" + metadata = message.metadata or {} + chat_type = str(metadata.get("chat_type") or "").strip().lower() + sender_label = ( + str(metadata.get("sender_display_name") or "").strip() + or str(metadata.get("sender_label") or "").strip() + or str(message.sender_id).strip() + ) + if chat_type != "group": + return "" + if not sender_label: + sender_label = "unknown" + return ( + "[Channel speaker]\n" + f"This message was sent in a group chat by: {sender_label}\n" + f"Sender id: {message.sender_id}" + ) + + def _build_attachment_notes(media_paths: list[str]) -> str: """Build textual attachment notes for non-image context and persistence.""" if not media_paths: diff --git a/ohmo/runtime.py b/ohmo/runtime.py index 7494f56..5ade10e 100644 --- a/ohmo/runtime.py +++ b/ohmo/runtime.py @@ -9,7 +9,7 @@ import sys from pathlib import Path from openharness.api.client import SupportsStreamingMessages -from openharness.engine.stream_events import AssistantTextDelta, AssistantTurnComplete, ErrorEvent, StatusEvent +from openharness.engine.stream_events import AssistantTextDelta, AssistantTurnComplete, CompactProgressEvent, ErrorEvent, StatusEvent from openharness.ui.backend_host import run_backend_host from openharness.ui.runtime import build_runtime, close_runtime, handle_line, start_runtime from openharness.ui.react_launcher import _resolve_npm, _resolve_tsx, get_frontend_dir @@ -173,6 +173,9 @@ async def run_ohmo_print_mode( sys.stdout.flush() elif isinstance(event, ErrorEvent): print(event.message, file=sys.stderr) + elif isinstance(event, CompactProgressEvent): + if event.message: + print(event.message, file=sys.stderr) elif isinstance(event, StatusEvent): print(event.message, file=sys.stderr) diff --git a/src/openharness/channels/impl/feishu.py b/src/openharness/channels/impl/feishu.py index db378fd..c78751f 100644 --- a/src/openharness/channels/impl/feishu.py +++ b/src/openharness/channels/impl/feishu.py @@ -6,6 +6,7 @@ import os import re import threading from collections import OrderedDict +from dataclasses import dataclass from typing import Any @@ -31,6 +32,12 @@ MSG_TYPE_MAP = { } +@dataclass(frozen=True) +class _FeishuSenderInfo: + open_id: str + display_name: str + + def _extract_share_card_content(content_json: dict, msg_type: str) -> str: """Extract text representation from share cards and interactive messages.""" parts = [] @@ -253,6 +260,7 @@ class FeishuChannel(BaseChannel): self._ws_client: Any = None self._ws_thread: threading.Thread | None = None self._processed_message_ids: OrderedDict[str, None] = OrderedDict() # Ordered dedup cache + self._sender_cache: OrderedDict[str, _FeishuSenderInfo] = OrderedDict() self._loop: asyncio.AbstractEventLoop | None = None async def start(self) -> None: @@ -842,6 +850,44 @@ class FeishuChannel(BaseChannel): except Exception as e: logger.error("Error sending Feishu message: {}", e) + def _resolve_sender_display_name_sync(self, open_id: str) -> str: + """Resolve a human-friendly sender name from Feishu contact APIs.""" + cached = self._sender_cache.get(open_id) + if cached is not None: + self._sender_cache.move_to_end(open_id) + return cached.display_name + if not self._client or not open_id: + return open_id or "unknown" + + try: + import lark_oapi as lark + + request = ( + lark.api.contact.v3.GetUserRequest.builder() + .user_id(open_id) + .user_id_type("open_id") + .build() + ) + response = self._client.contact.v3.user.get(request) + if getattr(response, "success", lambda: False)(): + user = getattr(getattr(response, "data", None), "user", None) + display_name = ( + getattr(user, "name", None) + or getattr(user, "en_name", None) + or getattr(user, "nickname", None) + or open_id + ) + else: + display_name = open_id + except Exception: + logger.exception("Failed to resolve Feishu sender name open_id=%s", open_id) + display_name = open_id + + self._sender_cache[open_id] = _FeishuSenderInfo(open_id=open_id, display_name=display_name) + while len(self._sender_cache) > 512: + self._sender_cache.popitem(last=False) + return display_name + def _on_message_sync(self, data: "P2ImMessageReceiveV1") -> None: # noqa: F821 """ Sync handler for incoming messages (called from WebSocket thread). @@ -872,6 +918,12 @@ class FeishuChannel(BaseChannel): return sender_id = sender.sender_id.open_id if sender.sender_id else "unknown" + loop = asyncio.get_running_loop() + sender_display_name = await loop.run_in_executor( + None, + self._resolve_sender_display_name_sync, + sender_id, + ) chat_id = message.chat_id chat_type = message.chat_type msg_type = message.message_type @@ -937,6 +989,8 @@ class FeishuChannel(BaseChannel): "message_id": message_id, "chat_type": chat_type, "msg_type": msg_type, + "sender_display_name": sender_display_name, + "sender_label": sender_display_name or sender_id, } ) diff --git a/src/openharness/commands/registry.py b/src/openharness/commands/registry.py index d4d567b..c368d5f 100644 --- a/src/openharness/commands/registry.py +++ b/src/openharness/commands/registry.py @@ -41,7 +41,12 @@ from openharness.permissions import PermissionChecker, PermissionMode from openharness.plugins import load_plugins from openharness.prompts import build_runtime_system_prompt from openharness.plugins.installer import install_plugin_from_path, uninstall_plugin -from openharness.services import compact_messages, estimate_conversation_tokens, summarize_messages +from openharness.services import ( + compact_conversation, + compact_messages, + estimate_conversation_tokens, + summarize_messages, +) from openharness.services.session_backend import DEFAULT_SESSION_BACKEND, SessionBackend from openharness.skills import load_skill_registry from openharness.tasks import get_task_manager @@ -280,7 +285,17 @@ def create_default_command_registry( except ValueError: return CommandResult(message="Usage: /compact [PRESERVE_RECENT]") before = len(context.engine.messages) - compacted = compact_messages(context.engine.messages, preserve_recent=preserve_recent) + try: + compacted = await compact_conversation( + context.engine.messages, + api_client=context.engine.api_client, + model=context.engine.model, + system_prompt=context.engine.system_prompt, + preserve_recent=preserve_recent, + trigger="manual", + ) + except Exception: + compacted = compact_messages(context.engine.messages, preserve_recent=preserve_recent) context.engine.load_messages(compacted) return CommandResult( message=f"Compacted conversation from {before} messages to {len(compacted)}." diff --git a/src/openharness/engine/query.py b/src/openharness/engine/query.py index 84e106f..a77788f 100644 --- a/src/openharness/engine/query.py +++ b/src/openharness/engine/query.py @@ -7,7 +7,7 @@ import logging import time from dataclasses import dataclass from pathlib import Path -from typing import AsyncIterator, Awaitable, Callable +from typing import Any, AsyncIterator, Awaitable, Callable from openharness.api.client import ( ApiMessageCompleteEvent, @@ -21,6 +21,7 @@ from openharness.engine.messages import ConversationMessage, ToolResultBlock from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, ErrorEvent, StatusEvent, StreamEvent, @@ -33,6 +34,7 @@ from openharness.tools.base import ToolExecutionContext from openharness.tools.base import ToolRegistry AUTO_COMPACT_STATUS_MESSAGE = "Auto-compacting conversation memory to keep things fast and focused." +REACTIVE_COMPACT_STATUS_MESSAGE = "Prompt too long; compacting conversation memory and retrying." log = logging.getLogger(__name__) @@ -40,6 +42,26 @@ log = logging.getLogger(__name__) PermissionPrompt = Callable[[str, str], Awaitable[bool]] AskUserPrompt = Callable[[str], Awaitable[str]] +MAX_TRACKED_READ_FILES = 6 +MAX_TRACKED_SKILLS = 8 +MAX_TRACKED_ASYNC_AGENT_EVENTS = 8 + + +def _is_prompt_too_long_error(exc: Exception) -> bool: + text = str(exc).lower() + return any( + needle in text + for needle in ( + "prompt too long", + "context length", + "maximum context", + "context window", + "too many tokens", + "too large for the model", + "maximum context length", + ) + ) + class MaxTurnsExceeded(RuntimeError): """Raised when the agent exceeds the configured max_turns for one user prompt.""" @@ -67,6 +89,125 @@ class QueryContext: tool_metadata: dict[str, object] | None = None +def _tool_metadata_bucket( + tool_metadata: dict[str, object] | None, + key: str, +) -> list[Any]: + if tool_metadata is None: + return [] + value = tool_metadata.setdefault(key, []) + if isinstance(value, list): + return value + replacement: list[Any] = [] + tool_metadata[key] = replacement + return replacement + + +def _remember_read_file( + tool_metadata: dict[str, object] | None, + *, + path: str, + offset: int, + limit: int, + output: str, +) -> None: + bucket = _tool_metadata_bucket(tool_metadata, "read_file_state") + preview_lines = [line.strip() for line in output.splitlines()[:6] if line.strip()] + bucket.append( + { + "path": path, + "span": f"lines {offset + 1}-{offset + limit}", + "preview": " | ".join(preview_lines)[:320], + } + ) + if len(bucket) > MAX_TRACKED_READ_FILES: + del bucket[:-MAX_TRACKED_READ_FILES] + + +def _remember_skill_invocation( + tool_metadata: dict[str, object] | None, + *, + skill_name: str, +) -> None: + bucket = _tool_metadata_bucket(tool_metadata, "invoked_skills") + normalized = skill_name.strip() + if not normalized: + return + if normalized in bucket: + bucket.remove(normalized) + bucket.append(normalized) + if len(bucket) > MAX_TRACKED_SKILLS: + del bucket[:-MAX_TRACKED_SKILLS] + + +def _remember_async_agent_activity( + tool_metadata: dict[str, object] | None, + *, + tool_name: str, + tool_input: dict[str, object], + output: str, +) -> None: + bucket = _tool_metadata_bucket(tool_metadata, "async_agent_state") + if tool_name == "agent": + description = str(tool_input.get("description") or tool_input.get("prompt") or "").strip() + summary = f"Spawned async agent. {description}".strip() + if output.strip(): + summary = f"{summary} [{output.strip()[:180]}]".strip() + elif tool_name == "send_message": + target = str(tool_input.get("task_id") or "").strip() + summary = f"Sent follow-up message to async agent {target}".strip() + else: + summary = output.strip()[:220] or f"Async agent activity via {tool_name}" + bucket.append(summary) + if len(bucket) > MAX_TRACKED_ASYNC_AGENT_EVENTS: + del bucket[:-MAX_TRACKED_ASYNC_AGENT_EVENTS] + + +def _update_plan_mode(tool_metadata: dict[str, object] | None, mode: str) -> None: + if tool_metadata is None: + return + tool_metadata["permission_mode"] = mode + + +def _record_tool_carryover( + context: QueryContext, + *, + tool_name: str, + tool_input: dict[str, object], + tool_output: str, + is_error: bool, + resolved_file_path: str | None, +) -> None: + if is_error: + return + if tool_name == "read_file" and resolved_file_path is not None: + offset = int(tool_input.get("offset") or 0) + limit = int(tool_input.get("limit") or 200) + _remember_read_file( + context.tool_metadata, + path=resolved_file_path, + offset=offset, + limit=limit, + output=tool_output, + ) + elif tool_name == "skill": + _remember_skill_invocation( + context.tool_metadata, + skill_name=str(tool_input.get("name") or ""), + ) + elif tool_name in {"agent", "send_message"}: + _remember_async_agent_activity( + context.tool_metadata, + tool_name=tool_name, + tool_input=tool_input, + output=tool_output, + ) + elif tool_name == "enter_plan_mode": + _update_plan_mode(context.tool_metadata, "plan") + elif tool_name == "exit_plan_mode": + _update_plan_mode(context.tool_metadata, "default") + + async def run_query( context: QueryContext, messages: list[ConversationMessage], @@ -85,20 +226,54 @@ async def run_query( ) compact_state = AutoCompactState() + reactive_compact_attempted = False + last_compaction_result: tuple[list[ConversationMessage], bool] = (messages, False) + + async def _stream_compaction( + *, + trigger: str, + force: bool = False, + ) -> AsyncIterator[tuple[StreamEvent, UsageSnapshot | None]]: + nonlocal last_compaction_result + progress_queue: asyncio.Queue[CompactProgressEvent] = asyncio.Queue() + + async def _progress(event: CompactProgressEvent) -> None: + await progress_queue.put(event) + + task = asyncio.create_task( + auto_compact_if_needed( + messages, + api_client=context.api_client, + model=context.model, + system_prompt=context.system_prompt, + state=compact_state, + progress_callback=_progress, + force=force, + trigger=trigger, + hook_executor=context.hook_executor, + carryover_metadata=context.tool_metadata, + ) + ) + while True: + try: + event = await asyncio.wait_for(progress_queue.get(), timeout=0.05) + yield event, None + except asyncio.TimeoutError: + if task.done(): + break + continue + while not progress_queue.empty(): + yield progress_queue.get_nowait(), None + last_compaction_result = await task + return turn_count = 0 while context.max_turns is None or turn_count < context.max_turns: turn_count += 1 # --- auto-compact check before calling the model --------------- - messages, was_compacted = await auto_compact_if_needed( - messages, - api_client=context.api_client, - model=context.model, - system_prompt=context.system_prompt, - state=compact_state, - ) - if was_compacted: - yield StatusEvent(message=AUTO_COMPACT_STATUS_MESSAGE), None + async for event, usage in _stream_compaction(trigger="auto"): + yield event, usage + messages, was_compacted = last_compaction_result # --------------------------------------------------------------- final_message: ConversationMessage | None = None @@ -131,6 +306,14 @@ async def run_query( usage = event.usage except Exception as exc: error_msg = str(exc) + if not reactive_compact_attempted and _is_prompt_too_long_error(exc): + reactive_compact_attempted = True + yield StatusEvent(message=REACTIVE_COMPACT_STATUS_MESSAGE), None + async for event, usage in _stream_compaction(trigger="reactive", force=True): + yield event, usage + messages, was_compacted = last_compaction_result + if was_compacted: + continue if "connect" in error_msg.lower() or "timeout" in error_msg.lower() or "network" in error_msg.lower(): yield ErrorEvent(message=f"Network error: {error_msg}. Check your internet connection and try again."), None else: @@ -275,6 +458,14 @@ async def _execute_tool_call( content=result.output, is_error=result.is_error, ) + _record_tool_carryover( + context, + tool_name=tool_name, + tool_input=tool_input, + tool_output=tool_result.content, + is_error=tool_result.is_error, + resolved_file_path=_file_path, + ) if context.hook_executor is not None: await context.hook_executor.execute( HookEvent.POST_TOOL_USE, diff --git a/src/openharness/engine/query_engine.py b/src/openharness/engine/query_engine.py index 48e4ee7..2a6a990 100644 --- a/src/openharness/engine/query_engine.py +++ b/src/openharness/engine/query_engine.py @@ -59,6 +59,21 @@ class QueryEngine: """Return the maximum number of agentic turns per user input, if capped.""" return self._max_turns + @property + def api_client(self) -> SupportsStreamingMessages: + """Return the active API client.""" + return self._api_client + + @property + def model(self) -> str: + """Return the active model identifier.""" + return self._model + + @property + def system_prompt(self) -> str: + """Return the active system prompt.""" + return self._system_prompt + @property def total_usage(self): """Return the total usage across all turns.""" diff --git a/src/openharness/engine/stream_events.py b/src/openharness/engine/stream_events.py index bdbee81..27aef80 100644 --- a/src/openharness/engine/stream_events.py +++ b/src/openharness/engine/stream_events.py @@ -3,7 +3,7 @@ from __future__ import annotations from dataclasses import dataclass -from typing import Any +from typing import Any, Literal from openharness.api.usage import UsageSnapshot from openharness.engine.messages import ConversationMessage @@ -56,6 +56,28 @@ class StatusEvent: message: str +@dataclass(frozen=True) +class CompactProgressEvent: + """Structured progress event for conversation compaction.""" + + phase: Literal[ + "hooks_start", + "context_collapse_start", + "context_collapse_end", + "session_memory_start", + "session_memory_end", + "compact_start", + "compact_retry", + "compact_end", + "compact_failed", + ] + trigger: Literal["auto", "manual", "reactive"] + message: str | None = None + attempt: int | None = None + checkpoint: str | None = None + metadata: dict[str, Any] | None = None + + StreamEvent = ( AssistantTextDelta | AssistantTurnComplete @@ -63,4 +85,5 @@ StreamEvent = ( | ToolExecutionCompleted | ErrorEvent | StatusEvent + | CompactProgressEvent ) diff --git a/src/openharness/hooks/events.py b/src/openharness/hooks/events.py index c12369a..683a554 100644 --- a/src/openharness/hooks/events.py +++ b/src/openharness/hooks/events.py @@ -10,5 +10,7 @@ class HookEvent(str, Enum): SESSION_START = "session_start" SESSION_END = "session_end" + PRE_COMPACT = "pre_compact" + POST_COMPACT = "post_compact" PRE_TOOL_USE = "pre_tool_use" POST_TOOL_USE = "post_tool_use" diff --git a/src/openharness/services/__init__.py b/src/openharness/services/__init__.py index 385b421..c1fa9cf 100644 --- a/src/openharness/services/__init__.py +++ b/src/openharness/services/__init__.py @@ -1,6 +1,7 @@ """Service exports.""" from openharness.services.compact import ( + compact_conversation, compact_messages, estimate_conversation_tokens, summarize_messages, @@ -15,6 +16,7 @@ from openharness.services.token_estimation import estimate_message_tokens, estim __all__ = [ "compact_messages", + "compact_conversation", "estimate_conversation_tokens", "estimate_message_tokens", "estimate_tokens", diff --git a/src/openharness/services/compact/__init__.py b/src/openharness/services/compact/__init__.py index d3b9b70..092d27e 100644 --- a/src/openharness/services/compact/__init__.py +++ b/src/openharness/services/compact/__init__.py @@ -8,18 +8,25 @@ Faithfully translated from Claude Code's compaction system: from __future__ import annotations +import asyncio +import inspect import logging import re from dataclasses import dataclass -from typing import Any +from pathlib import Path +from typing import Any, Awaitable, Callable, Literal +from uuid import uuid4 from openharness.engine.messages import ( ConversationMessage, ContentBlock, + ImageBlock, TextBlock, ToolResultBlock, ToolUseBlock, ) +from openharness.engine.stream_events import CompactProgressEvent +from openharness.hooks import HookEvent, HookExecutor from openharness.services.token_estimation import estimate_tokens log = logging.getLogger(__name__) @@ -45,6 +52,17 @@ TIME_BASED_MC_CLEARED_MESSAGE = "[Old tool result content cleared]" AUTOCOMPACT_BUFFER_TOKENS = 13_000 MAX_OUTPUT_TOKENS_FOR_SUMMARY = 20_000 MAX_CONSECUTIVE_AUTOCOMPACT_FAILURES = 3 +COMPACT_TIMEOUT_SECONDS = 25 +MAX_COMPACT_STREAMING_RETRIES = 2 +MAX_PTL_RETRIES = 3 +SESSION_MEMORY_KEEP_RECENT = 12 +SESSION_MEMORY_MAX_LINES = 48 +SESSION_MEMORY_MAX_CHARS = 4_000 +CONTEXT_COLLAPSE_TEXT_CHAR_LIMIT = 2_400 +CONTEXT_COLLAPSE_HEAD_CHARS = 900 +CONTEXT_COLLAPSE_TAIL_CHARS = 500 +MAX_COMPACT_ATTACHMENTS = 6 +MAX_DISCOVERED_TOOLS = 12 # Microcompact defaults DEFAULT_KEEP_RECENT = 5 @@ -55,6 +73,11 @@ TOKEN_ESTIMATION_PADDING = 4 / 3 # Default context windows per model family _DEFAULT_CONTEXT_WINDOW = 200_000 +PTL_RETRY_MARKER = "[earlier conversation truncated for compaction retry]" +ERROR_MESSAGE_INCOMPLETE_RESPONSE = "Compaction interrupted before a complete summary was returned." + +CompactTrigger = Literal["auto", "manual", "reactive"] +CompactProgressCallback = Callable[[CompactProgressEvent], Awaitable[None]] # --------------------------------------------------------------------------- @@ -81,6 +104,290 @@ def estimate_conversation_tokens(messages: list[ConversationMessage]) -> int: return estimate_message_tokens(messages) +def _sanitize_metadata(value: Any) -> Any: + if isinstance(value, (str, int, float, bool)) or value is None: + return value + if isinstance(value, Path): + return str(value) + if isinstance(value, dict): + return {str(key): _sanitize_metadata(item) for key, item in value.items()} + if isinstance(value, (list, tuple, set)): + return [_sanitize_metadata(item) for item in value] + return str(value) + + +def _record_compact_checkpoint( + carryover_metadata: dict[str, Any] | None, + *, + checkpoint: str, + trigger: CompactTrigger, + message_count: int, + token_count: int, + attempt: int | None = None, + details: dict[str, Any] | None = None, +) -> dict[str, Any]: + payload: dict[str, Any] = { + "checkpoint": checkpoint, + "trigger": trigger, + "message_count": message_count, + "token_count": token_count, + } + if attempt is not None: + payload["attempt"] = attempt + if details: + payload.update(_sanitize_metadata(details)) + if carryover_metadata is not None: + checkpoints = carryover_metadata.setdefault("compact_checkpoints", []) + if isinstance(checkpoints, list): + checkpoints.append(payload) + carryover_metadata["compact_last"] = payload + return payload + + +async def _emit_progress( + callback: CompactProgressCallback | None, + *, + phase: Literal[ + "hooks_start", + "context_collapse_start", + "context_collapse_end", + "session_memory_start", + "session_memory_end", + "compact_start", + "compact_retry", + "compact_end", + "compact_failed", + ], + trigger: CompactTrigger, + message: str | None = None, + attempt: int | None = None, + checkpoint: str | None = None, + metadata: dict[str, Any] | None = None, +) -> None: + if callback is None: + return + await callback( + CompactProgressEvent( + phase=phase, + trigger=trigger, + message=message, + attempt=attempt, + checkpoint=checkpoint, + metadata=_sanitize_metadata(metadata) if metadata else None, + ) + ) + + +def _is_prompt_too_long_error(exc: Exception) -> bool: + text = str(exc).lower() + return any( + needle in text + for needle in ( + "prompt too long", + "context length", + "maximum context", + "context window", + "too many tokens", + "too large for the model", + ) + ) + + +def _group_messages_by_prompt_round( + messages: list[ConversationMessage], +) -> list[list[ConversationMessage]]: + groups: list[list[ConversationMessage]] = [] + current: list[ConversationMessage] = [] + for message in messages: + starts_new_round = ( + message.role == "user" + and not any(isinstance(block, ToolResultBlock) for block in message.content) + and bool(message.text.strip()) + ) + if starts_new_round and current: + groups.append(current) + current = [] + current.append(message) + if current: + groups.append(current) + return groups + + +def _collapse_text(text: str) -> str: + if len(text) <= CONTEXT_COLLAPSE_TEXT_CHAR_LIMIT: + return text + omitted = len(text) - CONTEXT_COLLAPSE_HEAD_CHARS - CONTEXT_COLLAPSE_TAIL_CHARS + head = text[:CONTEXT_COLLAPSE_HEAD_CHARS].rstrip() + tail = text[-CONTEXT_COLLAPSE_TAIL_CHARS:].lstrip() + return f"{head}\n...[collapsed {omitted} chars]...\n{tail}" + + +def try_context_collapse( + messages: list[ConversationMessage], + *, + preserve_recent: int, +) -> list[ConversationMessage] | None: + """Deterministically shrink oversized text blocks before full compact.""" + if len(messages) <= preserve_recent + 2: + return None + + older = messages[:-preserve_recent] + newer = messages[-preserve_recent:] + changed = False + collapsed_older: list[ConversationMessage] = [] + for message in older: + new_blocks: list[ContentBlock] = [] + for block in message.content: + if isinstance(block, TextBlock): + collapsed = _collapse_text(block.text) + if collapsed != block.text: + changed = True + new_blocks.append(TextBlock(text=collapsed)) + else: + new_blocks.append(block) + collapsed_older.append(ConversationMessage(role=message.role, content=new_blocks)) + + if not changed: + return None + + result = [*collapsed_older, *newer] + if estimate_message_tokens(result) >= estimate_message_tokens(messages): + return None + return result + + +def truncate_head_for_ptl_retry( + messages: list[ConversationMessage], +) -> list[ConversationMessage] | None: + """Drop the oldest prompt rounds when the compact request itself is too large.""" + groups = _group_messages_by_prompt_round(messages) + if len(groups) < 2: + return None + + drop_count = max(1, len(groups) // 5) + drop_count = min(drop_count, len(groups) - 1) + retained = [message for group in groups[drop_count:] for message in group] + if not retained: + return None + if retained[0].role == "assistant": + return [ConversationMessage.from_user_text(PTL_RETRY_MARKER), *retained] + return retained + + +def _extract_attachment_paths(messages: list[ConversationMessage]) -> list[str]: + found: list[str] = [] + seen: set[str] = set() + path_pattern = re.compile(r"path:\s*([^)\\n]+)") + attachment_pattern = re.compile(r"\[attachment:\s*([^\]]+)\]") + for message in messages: + for block in message.content: + if isinstance(block, ImageBlock) and block.source_path: + path = str(Path(block.source_path).expanduser()) + if path not in seen: + seen.add(path) + found.append(path) + elif isinstance(block, TextBlock): + for match in path_pattern.findall(block.text): + path = match.strip() + if path and path not in seen: + seen.add(path) + found.append(path) + for match in attachment_pattern.findall(block.text): + path = match.strip() + if path and "download failed" not in path and path not in seen: + seen.add(path) + found.append(path) + if len(found) >= MAX_COMPACT_ATTACHMENTS: + return found + return found + + +def _extract_discovered_tools(messages: list[ConversationMessage]) -> list[str]: + discovered: list[str] = [] + seen: set[str] = set() + for message in messages: + for tool_use in message.tool_uses: + if tool_use.name and tool_use.name not in seen: + seen.add(tool_use.name) + discovered.append(tool_use.name) + if len(discovered) >= MAX_DISCOVERED_TOOLS: + return discovered + return discovered + + +def build_compact_carryover_message( + messages: list[ConversationMessage], + *, + metadata: dict[str, Any] | None = None, + hook_note: str | None = None, +) -> ConversationMessage | None: + """Preserve lightweight runtime context that should survive compaction.""" + metadata = metadata or {} + attachment_paths = _extract_attachment_paths(messages) + discovered_tools = _extract_discovered_tools(messages) + permission_mode = str(metadata.get("permission_mode") or "").strip().lower() + read_file_state = metadata.get("read_file_state") + invoked_skills = metadata.get("invoked_skills") + async_agent_state = metadata.get("async_agent_state") + compact_last = metadata.get("compact_last") + + lines: list[str] = [] + if permission_mode == "plan": + lines.extend( + [ + "Plan mode is still active for this session.", + "Do not execute mutating tools until the user exits plan mode.", + ] + ) + if attachment_paths: + lines.append("Recent local attachments to keep in mind:") + lines.extend(f"- {path}" for path in attachment_paths) + if discovered_tools: + lines.append("Tools already discovered or used in this session:") + lines.append("- " + ", ".join(discovered_tools)) + if isinstance(read_file_state, list) and read_file_state: + lines.append("Recently read files to keep in working memory:") + for entry in read_file_state[-4:]: + if not isinstance(entry, dict): + continue + path = str(entry.get("path") or "").strip() + span = str(entry.get("span") or "").strip() + preview = str(entry.get("preview") or "").strip() + if not path: + continue + bullet = f"- {path}" + if span: + bullet += f" ({span})" + lines.append(bullet) + if preview: + lines.append(f" Preview: {preview}") + if isinstance(invoked_skills, list) and invoked_skills: + lines.append("Skills invoked earlier in the session:") + lines.append("- " + ", ".join(str(skill) for skill in invoked_skills[-8:])) + if isinstance(async_agent_state, list) and async_agent_state: + lines.append("Async agent / background task state:") + lines.extend(f"- {entry}" for entry in async_agent_state[-6:]) + if isinstance(compact_last, dict) and compact_last: + checkpoint = str(compact_last.get("checkpoint") or "").strip() + token_count = compact_last.get("token_count") + if checkpoint: + if token_count is not None: + lines.append( + f"Last compact checkpoint: {checkpoint} (token_count={token_count})" + ) + else: + lines.append(f"Last compact checkpoint: {checkpoint}") + if hook_note: + lines.append("Compact hook note:") + lines.append(hook_note) + + if not lines: + return None + return ConversationMessage.from_user_text( + "Carry-over context preserved after compaction:\n" + "\n".join(lines) + ) + + # --------------------------------------------------------------------------- # Microcompact — clear old tool results to reduce tokens cheaply # --------------------------------------------------------------------------- @@ -148,6 +455,62 @@ def microcompact_messages( return messages, tokens_saved +def _summarize_message_for_memory(message: ConversationMessage) -> str: + text = " ".join(message.text.split()) + if text: + text = text[:160] + return f"{message.role}: {text}" + tool_uses = [block.name for block in message.tool_uses] + if tool_uses: + return f"{message.role}: tool calls -> {', '.join(tool_uses[:4])}" + if any(isinstance(block, ToolResultBlock) for block in message.content): + return f"{message.role}: tool results returned" + return f"{message.role}: [non-text content]" + + +def _build_session_memory_message(messages: list[ConversationMessage]) -> ConversationMessage | None: + lines: list[str] = [] + total_chars = 0 + for message in messages: + line = _summarize_message_for_memory(message) + if not line: + continue + projected = total_chars + len(line) + 1 + if lines and (len(lines) >= SESSION_MEMORY_MAX_LINES or projected >= SESSION_MEMORY_MAX_CHARS): + lines.append("... earlier context condensed ...") + break + lines.append(line) + total_chars = projected + if not lines: + return None + body = "\n".join(lines) + return ConversationMessage.from_user_text( + "Session memory summary from earlier in this conversation:\n" + body + ) + + +def try_session_memory_compaction( + messages: list[ConversationMessage], + *, + preserve_recent: int = SESSION_MEMORY_KEEP_RECENT, +) -> list[ConversationMessage] | None: + """Cheap deterministic compaction for long chats before full LLM compaction.""" + if len(messages) <= preserve_recent + 4: + return None + older = messages[:-preserve_recent] + newer = messages[-preserve_recent:] + summary_message = _build_session_memory_message(older) + if summary_message is None: + return None + result = [summary_message, *newer] + if ( + estimate_message_tokens(result) >= estimate_message_tokens(messages) + and len(result) >= len(messages) + ): + return None + return result + + # --------------------------------------------------------------------------- # Full compact — LLM-based summarization # --------------------------------------------------------------------------- @@ -245,6 +608,7 @@ class AutoCompactState: compacted: bool = False turn_counter: int = 0 + turn_id: str = "" consecutive_failures: int = 0 @@ -299,6 +663,11 @@ async def compact_conversation( preserve_recent: int = 6, custom_instructions: str | None = None, suppress_follow_up: bool = True, + trigger: CompactTrigger = "manual", + progress_callback: CompactProgressCallback | None = None, + emit_hooks_start: bool = True, + hook_executor: HookExecutor | None = None, + carryover_metadata: dict[str, Any] | None = None, ) -> list[ConversationMessage]: """Compact messages by calling the LLM to produce a summary. @@ -337,21 +706,190 @@ async def compact_conversation( # Step 3: build compact request — send older messages + compact prompt compact_prompt = get_compact_prompt(custom_instructions) compact_messages = list(older) + [ConversationMessage.from_user_text(compact_prompt)] + attachment_paths = _extract_attachment_paths(older) + discovered_tools = _extract_discovered_tools(older) + hook_payload = { + "event": HookEvent.PRE_COMPACT.value, + "trigger": trigger, + "model": model, + "message_count": len(messages), + "token_count": pre_compact_tokens, + "preserve_recent": preserve_recent, + "attachments": attachment_paths, + "discovered_tools": discovered_tools, + **(carryover_metadata or {}), + } + start_checkpoint = _record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_prepare", + trigger=trigger, + message_count=len(messages), + token_count=pre_compact_tokens, + details={ + "preserve_recent": preserve_recent, + "attachments": attachment_paths, + "discovered_tools": discovered_tools, + }, + ) + + if emit_hooks_start: + await _emit_progress( + progress_callback, + phase="hooks_start", + trigger=trigger, + message="Preparing conversation compaction.", + checkpoint="compact_hooks_start", + metadata=start_checkpoint, + ) + if hook_executor is not None: + hook_result = await hook_executor.execute(HookEvent.PRE_COMPACT, hook_payload) + if hook_result.blocked: + reason = hook_result.reason or "pre-compact hook blocked compaction" + failed_checkpoint = _record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_failed", + trigger=trigger, + message_count=len(messages), + token_count=pre_compact_tokens, + details={"reason": reason}, + ) + await _emit_progress( + progress_callback, + phase="compact_failed", + trigger=trigger, + message=reason, + checkpoint="compact_failed", + metadata=failed_checkpoint, + ) + return messages + compact_start_checkpoint = _record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_start", + trigger=trigger, + message_count=len(messages), + token_count=pre_compact_tokens, + details={"preserve_recent": preserve_recent}, + ) + await _emit_progress( + progress_callback, + phase="compact_start", + trigger=trigger, + message="Compacting conversation memory.", + checkpoint="compact_start", + metadata=compact_start_checkpoint, + ) summary_text = "" - async for event in api_client.stream_message( - ApiMessageRequest( - model=model, - messages=compact_messages, - system_prompt=system_prompt or "You are a conversation summarizer.", - max_tokens=MAX_OUTPUT_TOKENS_FOR_SUMMARY, - tools=[], # no tools for compact call + messages_to_summarize = compact_messages + retry_messages = messages_to_summarize + ptl_retries = 0 + + async def _collect_summary(summary_request_messages: list[ConversationMessage]) -> str: + collected = "" + stream = api_client.stream_message( + ApiMessageRequest( + model=model, + messages=summary_request_messages, + system_prompt=system_prompt or "You are a conversation summarizer.", + max_tokens=MAX_OUTPUT_TOKENS_FOR_SUMMARY, + tools=[], # no tools for compact call + ) ) - ): - if isinstance(event, ApiMessageCompleteEvent): - summary_text = event.message.text + if inspect.isawaitable(stream): + stream = await stream + if not hasattr(stream, "__aiter__"): + raise RuntimeError("Compaction client did not provide a streaming response.") + async for event in stream: + if isinstance(event, ApiMessageCompleteEvent): + collected = event.message.text + if collected.strip(): + return collected + raise RuntimeError(ERROR_MESSAGE_INCOMPLETE_RESPONSE) + + for attempt in range(1, MAX_COMPACT_STREAMING_RETRIES + 2): + try: + summary_text = await asyncio.wait_for( + _collect_summary(retry_messages), + timeout=COMPACT_TIMEOUT_SECONDS, + ) + break + except Exception as exc: + if _is_prompt_too_long_error(exc) and ptl_retries < MAX_PTL_RETRIES: + truncated = truncate_head_for_ptl_retry(retry_messages[:-1]) + if truncated: + ptl_retries += 1 + retry_messages = [*truncated, retry_messages[-1]] + await _emit_progress( + progress_callback, + phase="compact_retry", + trigger=trigger, + message="Compaction prompt was too large; retrying with older context trimmed.", + attempt=ptl_retries, + checkpoint="compact_retry_prompt_too_long", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_retry_prompt_too_long", + trigger=trigger, + message_count=len(retry_messages), + token_count=estimate_message_tokens(retry_messages), + attempt=ptl_retries, + details={"ptl_retries": ptl_retries}, + ), + ) + continue + if attempt > MAX_COMPACT_STREAMING_RETRIES: + await _emit_progress( + progress_callback, + phase="compact_failed", + trigger=trigger, + message=str(exc), + attempt=attempt, + checkpoint="compact_failed", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_failed", + trigger=trigger, + message_count=len(retry_messages), + token_count=estimate_message_tokens(retry_messages), + attempt=attempt, + details={"reason": str(exc)}, + ), + ) + raise + await _emit_progress( + progress_callback, + phase="compact_retry", + trigger=trigger, + message=str(exc), + attempt=attempt, + checkpoint="compact_retry", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_retry", + trigger=trigger, + message_count=len(retry_messages), + token_count=estimate_message_tokens(retry_messages), + attempt=attempt, + details={"reason": str(exc)}, + ), + ) if not summary_text: + await _emit_progress( + progress_callback, + phase="compact_failed", + trigger=trigger, + message=ERROR_MESSAGE_INCOMPLETE_RESPONSE, + checkpoint="compact_failed", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_failed", + trigger=trigger, + message_count=len(messages), + token_count=pre_compact_tokens, + details={"reason": ERROR_MESSAGE_INCOMPLETE_RESPONSE}, + ), + ) log.warning("Compact summary was empty — returning original messages") return messages @@ -362,15 +900,78 @@ async def compact_conversation( recent_preserved=len(newer) > 0, ) summary_msg = ConversationMessage.from_user_text(summary_content) + carryover_msg = build_compact_carryover_message( + older, + metadata=carryover_metadata, + ) - result = [summary_msg, *newer] + result = [summary_msg] + if carryover_msg is not None: + result.append(carryover_msg) + result.extend(newer) post_compact_tokens = estimate_message_tokens(result) + post_hook_result = None + if hook_executor is not None: + post_hook_result = await hook_executor.execute( + HookEvent.POST_COMPACT, + { + "event": HookEvent.POST_COMPACT.value, + "trigger": trigger, + "model": model, + "pre_compact_message_count": len(messages), + "post_compact_message_count": len(result), + "pre_compact_tokens": pre_compact_tokens, + "post_compact_tokens": post_compact_tokens, + "attachments": attachment_paths, + "discovered_tools": discovered_tools, + **(carryover_metadata or {}), + }, + ) + hook_note = post_hook_result.reason or "\n".join( + result.output.strip() + for result in post_hook_result.results + if result.output.strip() + ) + if hook_note: + carryover_msg = build_compact_carryover_message( + older, + metadata=carryover_metadata, + hook_note=hook_note, + ) + result = [summary_msg] + if carryover_msg is not None: + result.append(carryover_msg) + result.extend(newer) + post_compact_tokens = estimate_message_tokens(result) log.info( "Compaction done: %d -> %d messages, ~%d -> ~%d tokens (saved ~%d)", len(messages), len(result), pre_compact_tokens, post_compact_tokens, pre_compact_tokens - post_compact_tokens, ) + await _emit_progress( + progress_callback, + phase="compact_end", + trigger=trigger, + message="Conversation compaction complete.", + checkpoint="compact_end", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="compact_end", + trigger=trigger, + message_count=len(result), + token_count=post_compact_tokens, + details={ + "pre_compact_message_count": len(messages), + "post_compact_message_count": len(result), + "pre_compact_tokens": pre_compact_tokens, + "post_compact_tokens": post_compact_tokens, + "tokens_saved": pre_compact_tokens - post_compact_tokens, + "attachments": attachment_paths, + "discovered_tools": discovered_tools, + }, + ), + ) return result @@ -386,6 +987,11 @@ async def auto_compact_if_needed( system_prompt: str = "", state: AutoCompactState, preserve_recent: int = 6, + progress_callback: CompactProgressCallback | None = None, + force: bool = False, + trigger: CompactTrigger = "auto", + hook_executor: HookExecutor | None = None, + carryover_metadata: dict[str, Any] | None = None, ) -> tuple[list[ConversationMessage], bool]: """Check if auto-compact should fire, and if so, compact. @@ -394,17 +1000,103 @@ async def auto_compact_if_needed( Returns: (messages, was_compacted) — if compacted, messages is the new list. """ - if not should_autocompact(messages, model, state): + if not force and not should_autocompact(messages, model, state): return messages, False log.info("Auto-compact triggered (failures=%d)", state.consecutive_failures) + _record_compact_checkpoint( + carryover_metadata, + checkpoint=f"query_{trigger}_triggered", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + details={"consecutive_failures": state.consecutive_failures}, + ) # Try microcompact first — may be enough messages, tokens_freed = microcompact_messages(messages) + _record_compact_checkpoint( + carryover_metadata, + checkpoint="query_microcompact_end", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + details={"tokens_freed": tokens_freed}, + ) if tokens_freed > 0 and not should_autocompact(messages, model, state): log.info("Microcompact freed ~%d tokens, auto-compact no longer needed", tokens_freed) return messages, True + context_collapsed = try_context_collapse(messages, preserve_recent=preserve_recent) + if context_collapsed is not None: + await _emit_progress( + progress_callback, + phase="context_collapse_start", + trigger=trigger, + message="Collapsing oversized context before full compaction.", + checkpoint="query_context_collapse_start", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="query_context_collapse_start", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + ), + ) + messages = context_collapsed + await _emit_progress( + progress_callback, + phase="context_collapse_end", + trigger=trigger, + message="Context collapse complete.", + checkpoint="query_context_collapse_end", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="query_context_collapse_end", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + ), + ) + if not force and not should_autocompact(messages, model, state): + return messages, True + + session_memory = try_session_memory_compaction(messages, preserve_recent=max(preserve_recent, SESSION_MEMORY_KEEP_RECENT)) + if session_memory is not None: + await _emit_progress( + progress_callback, + phase="session_memory_start", + trigger=trigger, + message="Condensing earlier conversation into session memory.", + checkpoint="query_session_memory_start", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="query_session_memory_start", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + ), + ) + await _emit_progress( + progress_callback, + phase="session_memory_end", + trigger=trigger, + message="Session memory condensation complete.", + checkpoint="query_session_memory_end", + metadata=_record_compact_checkpoint( + carryover_metadata, + checkpoint="query_session_memory_end", + trigger=trigger, + message_count=len(session_memory), + token_count=estimate_message_tokens(session_memory), + ), + ) + state.compacted = True + state.turn_counter += 1 + state.turn_id = uuid4().hex + state.consecutive_failures = 0 + return session_memory, True + # Full compact needed try: result = await compact_conversation( @@ -414,13 +1106,26 @@ async def auto_compact_if_needed( system_prompt=system_prompt, preserve_recent=preserve_recent, suppress_follow_up=True, + trigger=trigger, + progress_callback=progress_callback, + hook_executor=hook_executor, + carryover_metadata=carryover_metadata, ) state.compacted = True state.turn_counter += 1 + state.turn_id = uuid4().hex state.consecutive_failures = 0 return result, True except Exception as exc: state.consecutive_failures += 1 + _record_compact_checkpoint( + carryover_metadata, + checkpoint=f"query_{trigger}_failed", + trigger=trigger, + message_count=len(messages), + token_count=estimate_message_tokens(messages), + details={"reason": str(exc), "consecutive_failures": state.consecutive_failures}, + ) log.error( "Auto-compact failed (attempt %d/%d): %s", state.consecutive_failures, diff --git a/src/openharness/ui/app.py b/src/openharness/ui/app.py index f991bdb..eef5a18 100644 --- a/src/openharness/ui/app.py +++ b/src/openharness/ui/app.py @@ -78,6 +78,7 @@ async def run_print_mode( from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, ErrorEvent, StatusEvent, ToolExecutionCompleted, @@ -155,6 +156,19 @@ async def run_print_mode( obj = {"type": "error", "message": event.message, "recoverable": event.recoverable} print(json.dumps(obj), flush=True) events_list.append(obj) + elif isinstance(event, CompactProgressEvent): + if output_format == "text" and event.message: + print(event.message, file=sys.stderr) + elif output_format == "stream-json": + obj = { + "type": "compact_progress", + "phase": event.phase, + "trigger": event.trigger, + "attempt": event.attempt, + "message": event.message, + } + print(json.dumps(obj), flush=True) + events_list.append(obj) elif isinstance(event, StatusEvent): if output_format == "text": print(event.message, file=sys.stderr) diff --git a/src/openharness/ui/backend_host.py b/src/openharness/ui/backend_host.py index 873bdb1..b63f6b4 100644 --- a/src/openharness/ui/backend_host.py +++ b/src/openharness/ui/backend_host.py @@ -20,6 +20,7 @@ from openharness.themes import list_themes from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, ErrorEvent, StatusEvent, StreamEvent, @@ -203,6 +204,19 @@ class ReactBackendHost: if isinstance(event, AssistantTextDelta): await self._emit(BackendEvent(type="assistant_delta", message=event.text)) return + if isinstance(event, CompactProgressEvent): + await self._emit( + BackendEvent( + type="compact_progress", + compact_phase=event.phase, + compact_trigger=event.trigger, + attempt=event.attempt, + compact_checkpoint=event.checkpoint, + compact_metadata=event.metadata, + message=event.message, + ) + ) + return if isinstance(event, AssistantTurnComplete): await self._emit( BackendEvent( diff --git a/src/openharness/ui/output.py b/src/openharness/ui/output.py index 50a1c2b..e79192c 100644 --- a/src/openharness/ui/output.py +++ b/src/openharness/ui/output.py @@ -10,6 +10,7 @@ from rich.syntax import Syntax from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, StreamEvent, ToolExecutionCompleted, ToolExecutionStarted, @@ -70,6 +71,41 @@ class OutputRenderer: self._assistant_buffer = "" return + if isinstance(event, CompactProgressEvent): + self._stop_spinner() + if event.message: + label = event.message + elif event.phase == "hooks_start": + label = ( + "Preparing retry compaction..." + if event.trigger == "reactive" + else "Preparing conversation compaction..." + ) + elif event.phase == "session_memory_start": + label = "Condensing earlier conversation..." + elif event.phase == "session_memory_end": + label = "Conversation condensed." + elif event.phase == "context_collapse_start": + label = "Collapsing oversized context..." + elif event.phase == "context_collapse_end": + label = "Context collapse complete." + elif event.phase == "compact_start": + label = ( + "Context is too large. Compacting and retrying..." + if event.trigger == "reactive" + else "Compacting conversation memory..." + ) + elif event.phase == "compact_retry": + label = "Retrying compaction..." + elif event.phase == "compact_end": + label = "Compaction complete." + elif event.phase == "compact_failed": + label = "Compaction failed." + else: + label = "Compacting..." + self.console.print(f"[yellow]\u2139 {label}[/yellow]") + return + if isinstance(event, ToolExecutionStarted): self._stop_spinner() if self._assistant_line_open: diff --git a/src/openharness/ui/protocol.py b/src/openharness/ui/protocol.py index 4a3b9d0..3f46aad 100644 --- a/src/openharness/ui/protocol.py +++ b/src/openharness/ui/protocol.py @@ -70,6 +70,7 @@ class BackendEvent(BaseModel): "state_snapshot", "tasks_snapshot", "transcript_item", + "compact_progress", "assistant_delta", "assistant_complete", "line_complete", @@ -97,6 +98,11 @@ class BackendEvent(BaseModel): tool_input: dict[str, Any] | None = None output: str | None = None is_error: bool | None = None + compact_phase: str | None = None + compact_trigger: str | None = None + attempt: int | None = None + compact_checkpoint: str | None = None + compact_metadata: dict[str, Any] | None = None # New fields for enhanced events todo_markdown: str | None = None plan_mode: str | None = None diff --git a/src/openharness/ui/runtime.py b/src/openharness/ui/runtime.py index 167840b..6aee250 100644 --- a/src/openharness/ui/runtime.py +++ b/src/openharness/ui/runtime.py @@ -247,6 +247,10 @@ async def build_runtime( extra_skill_dirs=normalized_skill_dirs, extra_plugin_roots=normalized_plugin_roots, ) + from uuid import uuid4 + + session_id = uuid4().hex[:12] + engine = QueryEngine( api_client=resolved_api_client, tool_registry=tool_registry, @@ -264,6 +268,12 @@ async def build_runtime( "bridge_manager": bridge_manager, "extra_skill_dirs": normalized_skill_dirs, "extra_plugin_roots": normalized_plugin_roots, + "permission_mode": settings.permission.mode.value, + "session_id": session_id, + "read_file_state": [], + "invoked_skills": [], + "async_agent_state": [], + "compact_checkpoints": [], }, ) # Restore messages from a saved session if provided @@ -273,8 +283,6 @@ async def build_runtime( ] engine.load_messages(restored) - from uuid import uuid4 - return RuntimeBundle( api_client=resolved_api_client, cwd=cwd, @@ -293,7 +301,7 @@ async def build_runtime( ), external_api_client=api_client is not None, enforce_max_turns=enforce_max_turns or max_turns is not None, - session_id=uuid4().hex[:12], + session_id=session_id, settings_overrides=settings_overrides, session_backend=session_backend or DEFAULT_SESSION_BACKEND, extra_skill_dirs=normalized_skill_dirs, diff --git a/src/openharness/ui/textual_app.py b/src/openharness/ui/textual_app.py index cbd2ae0..2b0760c 100644 --- a/src/openharness/ui/textual_app.py +++ b/src/openharness/ui/textual_app.py @@ -19,6 +19,7 @@ from openharness.config.settings import load_settings, save_settings from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, ErrorEvent, StatusEvent, StreamEvent, @@ -317,6 +318,35 @@ class OpenHarnessTerminalApp(App[None]): self._set_current_response(f"[bold]assistant>[/bold] {self._assistant_buffer}") return + if isinstance(event, CompactProgressEvent): + if event.phase == "hooks_start": + if event.trigger == "reactive": + self._set_current_response("[dim]Preparing retry compaction...[/dim]") + else: + self._set_current_response("[dim]Preparing conversation compaction...[/dim]") + elif event.phase == "compact_start": + if event.trigger == "reactive": + self._set_current_response("[dim]Context too large. Compacting and retrying...[/dim]") + else: + self._set_current_response("[dim]Compacting conversation memory...[/dim]") + elif event.phase == "compact_retry": + attempt = f" (attempt {event.attempt})" if event.attempt is not None else "" + self._set_current_response(f"[dim]Retrying compaction{attempt}...[/dim]") + elif event.phase == "compact_failed": + self._append_line(f"system> Compaction failed: {event.message or 'unknown error'}") + self._set_current_response("Ready.") + elif event.phase == "compact_end": + self._set_current_response("[dim]Compaction complete.[/dim]") + elif event.phase == "session_memory_start": + self._set_current_response("[dim]Condensing earlier conversation...[/dim]") + elif event.phase == "session_memory_end": + self._set_current_response("[dim]Condensed earlier conversation.[/dim]") + elif event.phase == "context_collapse_start": + self._set_current_response("[dim]Collapsing oversized context...[/dim]") + elif event.phase == "context_collapse_end": + self._set_current_response("[dim]Context collapse complete.[/dim]") + return + if isinstance(event, AssistantTurnComplete): text = self._assistant_buffer or event.message.text or "(empty response)" self._append_line(f"assistant> {text}") diff --git a/tests/test_engine/test_query_engine.py b/tests/test_engine/test_query_engine.py index b3ebcaf..87e3cee 100644 --- a/tests/test_engine/test_query_engine.py +++ b/tests/test_engine/test_query_engine.py @@ -8,6 +8,7 @@ from pathlib import Path import pytest from openharness.api.client import ApiMessageCompleteEvent, ApiRetryEvent, ApiTextDeltaEvent +from openharness.api.errors import RequestFailure from openharness.api.usage import UsageSnapshot from openharness.config.settings import PermissionSettings from openharness.engine.messages import ConversationMessage, TextBlock, ToolUseBlock @@ -15,12 +16,14 @@ from openharness.engine.query_engine import QueryEngine from openharness.engine.stream_events import ( AssistantTextDelta, AssistantTurnComplete, + CompactProgressEvent, StatusEvent, ToolExecutionCompleted, ToolExecutionStarted, ) from openharness.permissions import PermissionChecker, PermissionMode from openharness.tools import create_default_tool_registry +from openharness.tools.base import ToolResult from openharness.hooks import HookExecutionContext, HookExecutor, HookEvent from openharness.hooks.loader import HookRegistry from openharness.hooks.schemas import PromptHookDefinition @@ -77,6 +80,28 @@ class RetryThenSuccessApiClient: ) +class PromptTooLongThenSuccessApiClient: + def __init__(self) -> None: + self._calls = 0 + + async def stream_message(self, request): + self._calls += 1 + if self._calls == 1: + raise RequestFailure("prompt too long") + if self._calls == 2: + yield ApiMessageCompleteEvent( + message=ConversationMessage(role="assistant", content=[TextBlock(text="compressed")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + stop_reason=None, + ) + return + yield ApiMessageCompleteEvent( + message=ConversationMessage(role="assistant", content=[TextBlock(text="after reactive compact")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + stop_reason=None, + ) + + @pytest.mark.asyncio async def test_query_engine_plain_text_reply(tmp_path: Path): engine = QueryEngine( @@ -220,6 +245,195 @@ async def test_query_engine_surfaces_retry_status_events(tmp_path: Path): assert isinstance(events[-1], AssistantTurnComplete) +@pytest.mark.asyncio +async def test_query_engine_emits_compact_progress_before_reply(tmp_path: Path, monkeypatch): + long_text = "alpha " * 50000 + monkeypatch.setattr("openharness.services.compact.try_session_memory_compaction", lambda *args, **kwargs: None) + monkeypatch.setattr("openharness.services.compact.should_autocompact", lambda *args, **kwargs: True) + engine = QueryEngine( + api_client=FakeApiClient( + [ + _FakeResponse( + message=ConversationMessage(role="assistant", content=[TextBlock(text="trimmed")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + _FakeResponse( + message=ConversationMessage(role="assistant", content=[TextBlock(text="after compact")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + ] + ), + tool_registry=create_default_tool_registry(), + permission_checker=PermissionChecker(PermissionSettings()), + cwd=tmp_path, + model="claude-sonnet-4-6", + system_prompt="system", + ) + engine.load_messages( + [ + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ] + ) + + events = [event async for event in engine.submit_message("hello")] + + hooks_start_index = next(i for i, event in enumerate(events) if isinstance(event, CompactProgressEvent) and event.phase == "hooks_start") + compact_start_index = next(i for i, event in enumerate(events) if isinstance(event, CompactProgressEvent) and event.phase == "compact_start") + final_index = next(i for i, event in enumerate(events) if isinstance(event, AssistantTurnComplete)) + assert hooks_start_index < compact_start_index + assert compact_start_index < final_index + assert any(isinstance(event, CompactProgressEvent) and event.phase == "compact_end" for event in events) + + +@pytest.mark.asyncio +async def test_query_engine_reactive_compacts_after_prompt_too_long(tmp_path: Path, monkeypatch): + monkeypatch.setattr("openharness.services.compact.try_session_memory_compaction", lambda *args, **kwargs: None) + monkeypatch.setattr("openharness.services.compact.should_autocompact", lambda *args, **kwargs: False) + engine = QueryEngine( + api_client=PromptTooLongThenSuccessApiClient(), + tool_registry=create_default_tool_registry(), + permission_checker=PermissionChecker(PermissionSettings()), + cwd=tmp_path, + model="claude-test", + system_prompt="system", + ) + engine.load_messages( + [ + ConversationMessage(role="user", content=[TextBlock(text="one")]), + ConversationMessage(role="assistant", content=[TextBlock(text="two")]), + ConversationMessage(role="user", content=[TextBlock(text="three")]), + ConversationMessage(role="assistant", content=[TextBlock(text="four")]), + ConversationMessage(role="user", content=[TextBlock(text="five")]), + ConversationMessage(role="assistant", content=[TextBlock(text="six")]), + ConversationMessage(role="user", content=[TextBlock(text="seven")]), + ConversationMessage(role="assistant", content=[TextBlock(text="eight")]), + ] + ) + + events = [event async for event in engine.submit_message("nine")] + + assert any( + isinstance(event, CompactProgressEvent) + and event.trigger == "reactive" + and event.phase == "compact_start" + for event in events + ) + assert isinstance(events[-1], AssistantTurnComplete) + assert events[-1].message.text == "after reactive compact" + + +@pytest.mark.asyncio +async def test_query_engine_tracks_recent_read_files_and_skills(tmp_path: Path): + sample = tmp_path / "hello.txt" + sample.write_text("alpha\nbeta\n", encoding="utf-8") + registry = create_default_tool_registry() + skill_tool = registry.get("skill") + assert skill_tool is not None + + async def _fake_skill_execute(arguments, context): + del context + return ToolResult(output=f"Loaded skill: {arguments.name}") + + monkeypatch = pytest.MonkeyPatch() + monkeypatch.setattr(skill_tool, "execute", _fake_skill_execute) + + engine = QueryEngine( + api_client=FakeApiClient( + [ + _FakeResponse( + message=ConversationMessage( + role="assistant", + content=[ + ToolUseBlock(name="read_file", input={"path": str(sample)}), + ToolUseBlock(name="skill", input={"name": "demo-skill"}), + ], + ), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + _FakeResponse( + message=ConversationMessage(role="assistant", content=[TextBlock(text="done")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + ] + ), + tool_registry=registry, + permission_checker=PermissionChecker(PermissionSettings()), + cwd=tmp_path, + model="claude-test", + system_prompt="system", + tool_metadata={}, + ) + + try: + events = [event async for event in engine.submit_message("track context")] + finally: + monkeypatch.undo() + + assert isinstance(events[-1], AssistantTurnComplete) + read_state = engine._tool_metadata.get("read_file_state") + assert isinstance(read_state, list) and read_state + assert read_state[-1]["path"] == str(sample.resolve()) + assert "alpha" in read_state[-1]["preview"] + invoked_skills = engine._tool_metadata.get("invoked_skills") + assert isinstance(invoked_skills, list) + assert invoked_skills[-1] == "demo-skill" + + +@pytest.mark.asyncio +async def test_query_engine_tracks_async_agent_activity(tmp_path: Path, monkeypatch): + registry = create_default_tool_registry() + agent_tool = registry.get("agent") + assert agent_tool is not None + + async def _fake_execute(arguments, context): + del arguments, context + return ToolResult(output="Spawned agent worker@team (task_id=task_123, backend=subprocess)") + + monkeypatch.setattr(agent_tool, "execute", _fake_execute) + engine = QueryEngine( + api_client=FakeApiClient( + [ + _FakeResponse( + message=ConversationMessage( + role="assistant", + content=[ + ToolUseBlock( + name="agent", + input={"description": "Inspect CI", "prompt": "Inspect CI"}, + ) + ], + ), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + _FakeResponse( + message=ConversationMessage(role="assistant", content=[TextBlock(text="spawned")]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + ), + ] + ), + tool_registry=registry, + permission_checker=PermissionChecker(PermissionSettings(mode=PermissionMode.FULL_AUTO)), + cwd=tmp_path, + model="claude-test", + system_prompt="system", + tool_metadata={}, + ) + + events = [event async for event in engine.submit_message("spawn helper")] + + assert isinstance(events[-1], AssistantTurnComplete) + async_state = engine._tool_metadata.get("async_agent_state") + assert isinstance(async_state, list) + assert async_state[-1].startswith("Spawned async agent") + + @pytest.mark.asyncio async def test_query_engine_respects_pre_tool_hook_blocks(tmp_path: Path): sample = tmp_path / "hello.txt" diff --git a/tests/test_ohmo/test_gateway.py b/tests/test_ohmo/test_gateway.py index 4ddd6a4..50b1852 100644 --- a/tests/test_ohmo/test_gateway.py +++ b/tests/test_ohmo/test_gateway.py @@ -13,11 +13,11 @@ from openharness.channels.bus.events import InboundMessage from openharness.channels.bus.queue import MessageBus from openharness.commands import CommandResult from openharness.engine.messages import ConversationMessage, ImageBlock, TextBlock -from openharness.engine.stream_events import AssistantTextDelta, StatusEvent, ToolExecutionStarted +from openharness.engine.stream_events import AssistantTextDelta, CompactProgressEvent, ToolExecutionStarted from ohmo.gateway.bridge import OhmoGatewayBridge, _format_gateway_error from ohmo.gateway.models import GatewayState -from ohmo.gateway.runtime import OhmoSessionRuntimePool +from ohmo.gateway.runtime import OhmoSessionRuntimePool, _build_inbound_user_message, _format_channel_progress from ohmo.gateway.service import OhmoGatewayService, gateway_status, stop_gateway_process from ohmo.gateway.router import session_key_for_message from ohmo.session_storage import save_session_snapshot @@ -58,6 +58,20 @@ def test_gateway_error_formats_generic_auth_failure(): assert "Authentication failed" in _format_gateway_error(exc) +def test_compact_progress_formats_reactive_channel_hint_in_chinese(): + text = _format_channel_progress( + channel="feishu", + kind="compact_progress", + text="", + session_key="feishu:c1", + content="帮我继续处理", + compact_phase="compact_start", + compact_trigger="reactive", + attempt=None, + ) + assert "重试" in text + + def test_gateway_status_prefers_live_config_over_stale_state(tmp_path): workspace = tmp_path / ".ohmo-home" workspace.mkdir() @@ -199,7 +213,7 @@ async def test_runtime_pool_stream_message_formats_auto_compact_status_for_feish return None async def submit_message(self, content): - yield StatusEvent(message="Auto-compacting conversation memory to keep things fast and focused.") + yield CompactProgressEvent(phase="compact_start", trigger="auto") yield AssistantTextDelta(text="done") return SimpleNamespace( @@ -220,11 +234,87 @@ async def test_runtime_pool_stream_message_formats_auto_compact_status_for_feish updates = [u async for u in pool.stream_message(message, "feishu:c1")] assert updates[1].kind == "progress" - assert updates[1].text == "🧠 聊天有点长啦,我先帮你蹦蹦跳跳压缩一下记忆,马上带着重点回来~" + assert updates[1].text == "🧠 聊天有点长啦,我先帮你悄悄压缩一下记忆,马上继续~" assert updates[-1].kind == "final" assert updates[-1].text == "done" +@pytest.mark.asyncio +async def test_runtime_pool_stream_message_formats_compact_retry_for_feishu(tmp_path, monkeypatch): + workspace = tmp_path / ".ohmo-home" + initialize_workspace(workspace) + + async def fake_build_runtime(**kwargs): + class FakeEngine: + messages = [] + total_usage = UsageSnapshot() + + def set_system_prompt(self, prompt): + return None + + async def submit_message(self, content): + yield CompactProgressEvent(phase="compact_retry", trigger="auto", attempt=2, message="retrying") + yield AssistantTextDelta(text="done") + + return SimpleNamespace( + engine=FakeEngine(), + session_id="sess123", + current_settings=lambda: SimpleNamespace(model="gpt-5.4"), + commands=SimpleNamespace(lookup=lambda raw: None), + ) + + async def fake_start_runtime(bundle): + return None + + monkeypatch.setattr("ohmo.gateway.runtime.build_runtime", fake_build_runtime) + monkeypatch.setattr("ohmo.gateway.runtime.start_runtime", fake_start_runtime) + + pool = OhmoSessionRuntimePool(cwd=tmp_path, workspace=workspace, provider_profile="codex") + message = InboundMessage(channel="feishu", sender_id="u1", chat_id="c1", content="继续") + updates = [u async for u in pool.stream_message(message, "feishu:c1")] + + assert updates[1].kind == "progress" + assert "再试一次" in updates[1].text + + +@pytest.mark.asyncio +async def test_runtime_pool_stream_message_formats_compact_hooks_start_for_feishu(tmp_path, monkeypatch): + workspace = tmp_path / ".ohmo-home" + initialize_workspace(workspace) + + async def fake_build_runtime(**kwargs): + class FakeEngine: + messages = [] + total_usage = UsageSnapshot() + + def set_system_prompt(self, prompt): + return None + + async def submit_message(self, content): + yield CompactProgressEvent(phase="hooks_start", trigger="auto") + yield AssistantTextDelta(text="done") + + return SimpleNamespace( + engine=FakeEngine(), + session_id="sess123", + current_settings=lambda: SimpleNamespace(model="gpt-5.4"), + commands=SimpleNamespace(lookup=lambda raw: None), + ) + + async def fake_start_runtime(bundle): + return None + + monkeypatch.setattr("ohmo.gateway.runtime.build_runtime", fake_build_runtime) + monkeypatch.setattr("ohmo.gateway.runtime.start_runtime", fake_start_runtime) + + pool = OhmoSessionRuntimePool(cwd=tmp_path, workspace=workspace, provider_profile="codex") + message = InboundMessage(channel="feishu", sender_id="u1", chat_id="c1", content="继续") + updates = [u async for u in pool.stream_message(message, "feishu:c1")] + + assert updates[1].kind == "progress" + assert "准备" in updates[1].text + + @pytest.mark.asyncio async def test_runtime_pool_stream_message_uses_english_progress_for_english_input(tmp_path, monkeypatch): workspace = tmp_path / ".ohmo-home" @@ -328,6 +418,23 @@ async def test_runtime_pool_includes_media_paths_in_prompt(tmp_path, monkeypatch assert "text preview: Quarterly summary Revenue up 12%" in text +def test_runtime_pool_includes_group_speaker_context(): + built = _build_inbound_user_message( + InboundMessage( + channel="feishu", + sender_id="ou_123", + chat_id="oc_group", + content="请帮我看一下", + metadata={"chat_type": "group", "sender_display_name": "Tang Jiabin"}, + ) + ) + text = "".join(block.text for block in built.content if isinstance(block, TextBlock)) + assert "[Channel speaker]" in text + assert "Tang Jiabin" in text + assert "Sender id: ou_123" in text + assert "请帮我看一下" in text + + @pytest.mark.asyncio async def test_gateway_bridge_publishes_progress_updates(): bus = MessageBus() diff --git a/tests/test_services/test_compact.py b/tests/test_services/test_compact.py index 135a84b..8397954 100644 --- a/tests/test_services/test_compact.py +++ b/tests/test_services/test_compact.py @@ -2,14 +2,28 @@ from __future__ import annotations -from openharness.engine.messages import ConversationMessage, TextBlock +import asyncio + +import pytest + +from openharness.api.client import ApiMessageCompleteEvent +from openharness.api.usage import UsageSnapshot +from openharness.engine.messages import ConversationMessage, ImageBlock, TextBlock, ToolUseBlock +from openharness.hooks import HookEvent from openharness.services import ( + compact_conversation, compact_messages, estimate_conversation_tokens, estimate_message_tokens, estimate_tokens, summarize_messages, ) +from openharness.services.compact import ( + AutoCompactState, + auto_compact_if_needed, + try_context_collapse, + try_session_memory_compaction, +) def test_token_estimation_helpers(): @@ -34,3 +48,210 @@ def test_compact_and_summarize_messages(): assert len(compacted) == 3 assert "[conversation summary]" in compacted[0].text assert estimate_conversation_tokens(compacted) >= 1 + + +class _CompactApiClient: + def __init__(self, responses): + self._responses = list(responses) + + async def stream_message(self, request): + del request + response = self._responses.pop(0) + if isinstance(response, Exception): + raise response + if asyncio.iscoroutinefunction(response): + await response() + return + yield ApiMessageCompleteEvent( + message=ConversationMessage(role="assistant", content=[TextBlock(text=response)]), + usage=UsageSnapshot(input_tokens=1, output_tokens=1), + stop_reason=None, + ) + + +class _HookExecutorStub: + def __init__(self) -> None: + self.events: list[tuple[HookEvent, dict[str, object]]] = [] + + async def execute(self, event: HookEvent, payload: dict[str, object]): + self.events.append((event, payload)) + from openharness.hooks.types import AggregatedHookResult + + return AggregatedHookResult() + + +def test_try_session_memory_compaction_reduces_long_history(): + messages = [ + ConversationMessage(role="user", content=[TextBlock(text=(f"user {index} " * 200).strip())]) + if index % 2 == 0 + else ConversationMessage(role="assistant", content=[TextBlock(text=(f"assistant {index} " * 200).strip())]) + for index in range(20) + ] + + result = try_session_memory_compaction(messages) + + assert result is not None + assert len(result) < len(messages) + assert "Session memory summary" in result[0].text + + +def test_try_context_collapse_trims_oversized_messages(): + giant = ("alpha " * 1200).strip() + messages = [ + ConversationMessage(role="user", content=[TextBlock(text=giant)]), + ConversationMessage(role="assistant", content=[TextBlock(text=giant)]), + ConversationMessage(role="user", content=[TextBlock(text=giant)]), + ConversationMessage(role="assistant", content=[TextBlock(text=giant)]), + ConversationMessage(role="user", content=[TextBlock(text=giant)]), + ConversationMessage(role="assistant", content=[TextBlock(text="keep recent")]), + ConversationMessage(role="user", content=[TextBlock(text="latest")]), + ] + + result = try_context_collapse(messages, preserve_recent=2) + + assert result is not None + assert "[collapsed" in result[0].text + + +@pytest.mark.asyncio +async def test_compact_conversation_retries_after_incomplete_response(): + messages = [ + ConversationMessage(role="user", content=[TextBlock(text="alpha")]), + ConversationMessage(role="assistant", content=[TextBlock(text="beta")]), + ConversationMessage(role="user", content=[TextBlock(text="gamma")]), + ConversationMessage(role="assistant", content=[TextBlock(text="delta")]), + ConversationMessage(role="user", content=[TextBlock(text="epsilon")]), + ConversationMessage(role="assistant", content=[TextBlock(text="zeta")]), + ConversationMessage(role="user", content=[TextBlock(text="eta")]), + ] + + compacted = await compact_conversation( + messages, + api_client=_CompactApiClient(["", "condensed"]), + model="claude-test", + ) + + assert compacted[0].text.startswith("This session is being continued") + + +@pytest.mark.asyncio +async def test_compact_conversation_runs_hooks_and_preserves_carryover_state(tmp_path): + image_path = tmp_path / "sample.png" + image_path.write_bytes( + b"\x89PNG\r\n\x1a\n" + b"\x00\x00\x00\rIHDR\x00\x00\x00\x01\x00\x00\x00\x01\x08\x02\x00\x00\x00\x90wS\xde" + b"\x00\x00\x00\x0cIDAT\x08\x99c``\x00\x00\x00\x04\x00\x01\xf6\x178U" + b"\x00\x00\x00\x00IEND\xaeB`\x82" + ) + hook_executor = _HookExecutorStub() + messages = [ + ConversationMessage(role="user", content=[ImageBlock.from_path(image_path)]), + ConversationMessage(role="assistant", content=[TextBlock(text="Looking at the attachment")]), + ConversationMessage( + role="assistant", + content=[ToolUseBlock(name="read_file", input={"path": str(image_path)})], + ), + ConversationMessage(role="user", content=[TextBlock(text="Please keep going")]), + ConversationMessage(role="assistant", content=[TextBlock(text="Working through it")]), + ConversationMessage(role="user", content=[TextBlock(text="And preserve context")]), + ConversationMessage(role="assistant", content=[TextBlock(text="Sure")]), + ] + + compacted = await compact_conversation( + messages, + api_client=_CompactApiClient(["condensed"]), + model="claude-test", + preserve_recent=2, + hook_executor=hook_executor, + carryover_metadata={ + "permission_mode": "plan", + "session_id": "sess123", + "read_file_state": [ + { + "path": str(image_path), + "span": "lines 1-20", + "preview": "1\tPNG header", + } + ], + "invoked_skills": ["pikastream-video-meeting"], + "async_agent_state": ["Spawned async agent [task_id=task_123]"], + "compact_last": {"checkpoint": "query_auto_triggered", "token_count": 12345}, + }, + ) + + assert [event for event, _payload in hook_executor.events] == [HookEvent.PRE_COMPACT, HookEvent.POST_COMPACT] + assert compacted[0].text.startswith("This session is being continued") + assert "Carry-over context preserved after compaction" in compacted[1].text + assert "Plan mode is still active" in compacted[1].text + assert str(image_path) in compacted[1].text + assert "read_file" in compacted[1].text + assert "Recently read files" in compacted[1].text + assert "Skills invoked earlier" in compacted[1].text + assert "Async agent / background task state" in compacted[1].text + assert "Last compact checkpoint" in compacted[1].text + + +@pytest.mark.asyncio +async def test_auto_compact_records_richer_checkpoint_metadata(monkeypatch): + monkeypatch.setattr("openharness.services.compact.try_session_memory_compaction", lambda *args, **kwargs: None) + monkeypatch.setattr("openharness.services.compact.should_autocompact", lambda *args, **kwargs: True) + long_text = "alpha " * 50000 + messages = [ + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ] + metadata: dict[str, object] = {} + + result, was_compacted = await auto_compact_if_needed( + messages, + api_client=_CompactApiClient(["condensed"]), + model="claude-sonnet-4-6", + state=AutoCompactState(), + carryover_metadata=metadata, + ) + + assert was_compacted is True + assert result[0].text.startswith("This session is being continued") + checkpoints = metadata.get("compact_checkpoints") + assert isinstance(checkpoints, list) + checkpoint_names = [entry["checkpoint"] for entry in checkpoints] + assert "query_auto_triggered" in checkpoint_names + assert "query_microcompact_end" in checkpoint_names + assert "compact_end" in checkpoint_names + assert isinstance(metadata.get("compact_last"), dict) + assert metadata["compact_last"]["checkpoint"] == "compact_end" + + +@pytest.mark.asyncio +async def test_auto_compact_if_needed_returns_original_messages_after_timeout(monkeypatch): + async def _stall(): + await asyncio.sleep(0.05) + + monkeypatch.setattr("openharness.services.compact.COMPACT_TIMEOUT_SECONDS", 0.01) + monkeypatch.setattr("openharness.services.compact.try_session_memory_compaction", lambda *args, **kwargs: None) + monkeypatch.setattr("openharness.services.compact.should_autocompact", lambda *args, **kwargs: True) + long_text = "alpha " * 50000 + messages = [ + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ConversationMessage(role="assistant", content=[TextBlock(text=long_text)]), + ConversationMessage(role="user", content=[TextBlock(text=long_text)]), + ] + + result, was_compacted = await auto_compact_if_needed( + messages, + api_client=_CompactApiClient([_stall]), + model="claude-sonnet-4-6", + state=AutoCompactState(), + ) + + assert was_compacted is False + assert result == messages diff --git a/tests/test_ui/test_react_backend.py b/tests/test_ui/test_react_backend.py index 76d5896..e22130a 100644 --- a/tests/test_ui/test_react_backend.py +++ b/tests/test_ui/test_react_backend.py @@ -10,6 +10,7 @@ import pytest from openharness.api.client import ApiMessageCompleteEvent from openharness.api.usage import UsageSnapshot +from openharness.engine.stream_events import CompactProgressEvent from openharness.engine.messages import ConversationMessage, TextBlock from openharness.ui.backend_host import BackendHostConfig, ReactBackendHost, run_backend_host from openharness.ui.protocol import BackendEvent @@ -171,6 +172,50 @@ async def test_backend_host_processes_model_turn(tmp_path, monkeypatch): ) +@pytest.mark.asyncio +async def test_backend_host_emits_compact_progress_event(tmp_path, monkeypatch): + monkeypatch.chdir(tmp_path) + monkeypatch.setenv("OPENHARNESS_CONFIG_DIR", str(tmp_path / "config")) + monkeypatch.setenv("OPENHARNESS_DATA_DIR", str(tmp_path / "data")) + + host = ReactBackendHost(BackendHostConfig(api_client=StaticApiClient("unused"))) + host._bundle = await build_runtime(api_client=StaticApiClient("unused")) + events = [] + + async def _emit(event): + events.append(event) + + async def _fake_handle_line(bundle, line, print_system, render_event, clear_output): + del bundle, line, print_system, clear_output + await render_event( + CompactProgressEvent( + phase="compact_start", + trigger="auto", + message="Compacting conversation memory.", + checkpoint="compact_start", + metadata={"token_count": 12345}, + ) + ) + return True + + monkeypatch.setattr("openharness.ui.backend_host.handle_line", _fake_handle_line) + host._emit = _emit # type: ignore[method-assign] + await start_runtime(host._bundle) + try: + should_continue = await host._process_line("hi") + finally: + await close_runtime(host._bundle) + + assert should_continue is True + assert any( + event.type == "compact_progress" + and event.compact_phase == "compact_start" + and event.compact_checkpoint == "compact_start" + and event.compact_metadata == {"token_count": 12345} + for event in events + ) + + @pytest.mark.asyncio async def test_backend_host_surfaces_query_errors(tmp_path, monkeypatch): monkeypatch.chdir(tmp_path)