Migrate memory middleware to profile files (#253)
* feat(memory): migrate to profile memory files * chore(stream): read profile headings from templates * fix(display): keep assistant responses if response_text has started * fix(memory): do not treat failed bootstraps as profile creation * chore(memory): unlink blank legacy memory * fix(memory): resolve project_id once * fix(memory): preserve unreadable profile files * chore(tui): render streamed narration inline with tool timeline Update the TUI streaming timeline so assistant text emitted before or between tool calls is rendered inline where it occurs, rather than being kept as a single answer bubble above or below the tools. If the model begins an assistant response and then emits another tool call, the provisional response is converted into inline narration before that tool. The final assistant message then renders only the remaining response suffix, avoiding duplicate text in the completed transcript. Stop/cancel handling now preserves any active inline narration, appends the visible stopped marker only to the remaining displayed segment, and still returns the full normalized stopped response for channel callers. Completed tools continue to collapse while long runs are active, but expand again when the turn reaches a final state so the completed transcript shows the full tool timeline. * fix(stream): preserve narration around tool timelines Keep assistant narration attached to the tool call that follows it instead of folding all streamed text into the final answer block. Track narrated response segments in stream state, render them before their corresponding regular or task tool entries, and keep final answers limited to the response suffix that has not already been shown inline. Preserve narration across normal completion, stop/error final frames, sub-agent task calls, and collapsed live tool summaries. Add regression coverage for pending tools, completed tools, sub-agent task delegations, collapsed completed/running tool summaries, and final stop frames. * fix(tui): finalize inline narration transitions * test(memory): use canonical project id helper
This commit is contained in:
@@ -449,7 +449,11 @@ def _get_default_backend():
|
||||
)
|
||||
|
||||
|
||||
def _get_default_middleware(*, for_async_subagent: bool = False):
|
||||
def _get_default_middleware(
|
||||
*,
|
||||
for_async_subagent: bool = False,
|
||||
workspace_dir: str | Path | None = None,
|
||||
):
|
||||
"""Build the default middleware list.
|
||||
|
||||
Args:
|
||||
@@ -492,7 +496,7 @@ def _get_default_middleware(*, for_async_subagent: bool = False):
|
||||
ContextOverflowMapperMiddleware(),
|
||||
ToolErrorHandlerMiddleware(),
|
||||
*create_tool_selector_middleware(model=model),
|
||||
create_memory_middleware(memory_dir, extraction_model=model),
|
||||
create_memory_middleware(memory_dir, workspace_dir=workspace_dir),
|
||||
]
|
||||
|
||||
if cfg.enable_ask_user and not cfg.auto_mode and not for_async_subagent:
|
||||
@@ -669,7 +673,7 @@ def create_cli_agent(
|
||||
# Delegate middleware construction to the single source of truth so the
|
||||
# CLI agent never drifts from the default chain. Anything CLI-specific
|
||||
# (e.g. ``HumanInTheLoopMiddleware``) is appended below.
|
||||
mw: list[AgentMiddleware] = _get_default_middleware()
|
||||
mw: list[AgentMiddleware] = _get_default_middleware(workspace_dir=workspace_dir)
|
||||
|
||||
# HITL on main agent only — passing `interrupt_on=` to create_deep_agent
|
||||
# would propagate it to every subagent, breaking parallel execute calls
|
||||
|
||||
@@ -12,6 +12,7 @@ import queue
|
||||
import random
|
||||
import sys
|
||||
from collections.abc import Callable
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime
|
||||
from typing import Any, ClassVar
|
||||
|
||||
@@ -33,7 +34,7 @@ from ..sessions import (
|
||||
thread_exists,
|
||||
)
|
||||
from ..stream.events import stream_agent_events
|
||||
from ..stream.state import _INTERNAL_TOOLS, ResearchPhase, StreamState
|
||||
from ..stream.state import ResearchPhase, StreamState
|
||||
from ._agent_loader import BackgroundAgentLoader, MCPProgressTracker
|
||||
from ._constants import LOGO_GRADIENT, LOGO_LINES, WELCOME_SLOGANS, build_metadata
|
||||
from .channel import (
|
||||
@@ -220,6 +221,32 @@ def _should_finalize_active_summarization(event_type: str) -> bool:
|
||||
return bool(event_type) and event_type not in _SUMMARY_CONTINUATION_EVENTS
|
||||
|
||||
|
||||
def _response_after_narration(response_text: str, narrated_response_end: int) -> str:
|
||||
"""Return the response suffix that has not already been shown inline."""
|
||||
return response_text[max(0, narrated_response_end) :]
|
||||
|
||||
|
||||
def _strip_trailing_placeholder_ellipsis(text: str) -> str:
|
||||
"""Remove the standalone streaming placeholder suffix from final TUI text."""
|
||||
clean = text.strip()
|
||||
while clean.endswith("\n...") or clean.rstrip() == "...":
|
||||
clean = clean.rstrip().removesuffix("...").rstrip()
|
||||
return clean
|
||||
|
||||
|
||||
def _stopped_response_after_narration(
|
||||
previous_text: str,
|
||||
narrated_response_end: int,
|
||||
) -> tuple[str, str, str]:
|
||||
"""Return display-current, display-stopped, and full stopped response text."""
|
||||
from ..stream.display import build_stopped_response_text
|
||||
|
||||
display_previous = _response_after_narration(previous_text, narrated_response_end)
|
||||
display_current, display_stopped = build_stopped_response_text(display_previous)
|
||||
_, full_stopped = build_stopped_response_text(previous_text)
|
||||
return display_current, display_stopped, full_stopped
|
||||
|
||||
|
||||
def run_textual_interactive(
|
||||
*,
|
||||
show_thinking: bool,
|
||||
@@ -240,6 +267,7 @@ def run_textual_interactive(
|
||||
from textual.binding import Binding
|
||||
from textual.containers import Container, Horizontal, VerticalScroll
|
||||
from textual.events import MouseUp
|
||||
from textual.widget import Widget
|
||||
from textual.widgets import Static
|
||||
|
||||
from .clipboard import copy_selection_to_clipboard, get_clipboard_text
|
||||
@@ -1220,7 +1248,6 @@ def run_textual_interactive(
|
||||
of mounting the AskUserWidget.
|
||||
"""
|
||||
from ..stream.display import (
|
||||
build_stopped_response_text,
|
||||
is_stream_cancel_requested,
|
||||
)
|
||||
|
||||
@@ -1242,20 +1269,26 @@ def run_textual_interactive(
|
||||
loading_removed = False
|
||||
thinking_w: ThinkingWidget | None = None
|
||||
summarization_w: SummarizationWidget | None = None
|
||||
assistant_w: AssistantMessage | None = None
|
||||
todo_w: TodoWidget | None = None
|
||||
tool_widgets: dict[str, ToolCallWidget] = {}
|
||||
subagent_widgets: dict[str, SubAgentWidget] = {}
|
||||
|
||||
@dataclass
|
||||
class _ResponseDisplayState:
|
||||
assistant: AssistantMessage | None = None
|
||||
narration: AssistantMessage | None = None
|
||||
assistant_start: int = 0
|
||||
narrated_end: int = 0
|
||||
|
||||
response_display = _ResponseDisplayState()
|
||||
|
||||
# Transient indicator widgets (auto-removed on state transitions)
|
||||
narration_w: Static | None = None # dim italic intermediate text
|
||||
processing_w: Static | None = None # "Analyzing results..."
|
||||
|
||||
# Tool collapsing (matches CLI MAX_VISIBLE_TOOLS)
|
||||
_MAX_VISIBLE_TOOLS = 4
|
||||
completed_tool_order: list[str] = [] # tool_ids in completion order
|
||||
collapse_summary_w: Static | None = None
|
||||
has_used_tools = False
|
||||
|
||||
_thinking_sent = False
|
||||
_todo_sent = False
|
||||
@@ -1285,7 +1318,7 @@ def run_textual_interactive(
|
||||
metadata = build_metadata(self._workspace_dir, self._current_model)
|
||||
response = ""
|
||||
|
||||
async def _remove_w(w: Static | None) -> None:
|
||||
async def _remove_w(w: Widget | None) -> None:
|
||||
"""Safely remove a transient indicator widget."""
|
||||
if w is not None:
|
||||
try:
|
||||
@@ -1293,26 +1326,78 @@ def run_textual_interactive(
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def _mount_or_update_narration(text: str) -> None:
|
||||
"""Show the current between-tool text segment inline."""
|
||||
content = text.strip()
|
||||
if not content:
|
||||
return
|
||||
if response_display.narration is None:
|
||||
response_display.narration = AssistantMessage(content)
|
||||
await container.mount(response_display.narration)
|
||||
else:
|
||||
response_display.narration._content = content
|
||||
await response_display.narration.stop_stream()
|
||||
response_display.narrated_end = len(state.response_text)
|
||||
|
||||
async def _convert_assistant_to_narration() -> None:
|
||||
"""Turn a provisional answer into inline narration before a tool."""
|
||||
if response_display.assistant is None:
|
||||
return
|
||||
content = response_display.assistant._content
|
||||
if content.strip():
|
||||
narration = AssistantMessage(content.strip())
|
||||
try:
|
||||
await container.mount(
|
||||
narration, before=response_display.assistant
|
||||
)
|
||||
except Exception:
|
||||
await container.mount(narration)
|
||||
response_display.narrated_end = max(
|
||||
response_display.narrated_end,
|
||||
response_display.assistant_start + len(content),
|
||||
)
|
||||
try:
|
||||
await response_display.assistant.remove()
|
||||
except Exception:
|
||||
pass
|
||||
response_display.assistant = None
|
||||
response_display.assistant_start = len(state.response_text)
|
||||
response_display.narration = None
|
||||
|
||||
async def _preserve_active_narration() -> None:
|
||||
"""Finalize inline narration without removing it from the transcript."""
|
||||
if response_display.narration is not None:
|
||||
await response_display.narration.stop_stream()
|
||||
response_display.narration = None
|
||||
|
||||
async def _mark_cancelled_response() -> str:
|
||||
nonlocal assistant_w
|
||||
previous_text = state.response_text or ""
|
||||
current, final_text = build_stopped_response_text(previous_text)
|
||||
current, final_segment, final_text = _stopped_response_after_narration(
|
||||
previous_text,
|
||||
response_display.narrated_end,
|
||||
)
|
||||
|
||||
state.response_text = final_text
|
||||
self._set_status_streaming_text(final_text)
|
||||
self._set_status_streaming_text(final_segment)
|
||||
|
||||
if assistant_w is None:
|
||||
if final_text:
|
||||
assistant_w = AssistantMessage(final_text)
|
||||
await container.mount(assistant_w)
|
||||
await _preserve_active_narration()
|
||||
_expand_completed_tools()
|
||||
|
||||
if response_display.assistant is None:
|
||||
if final_segment:
|
||||
response_display.assistant_start = len(final_text) - len(
|
||||
final_segment
|
||||
)
|
||||
response_display.assistant = AssistantMessage(final_segment)
|
||||
await container.mount(response_display.assistant)
|
||||
else:
|
||||
if previous_text != current:
|
||||
assistant_w._content = final_text
|
||||
await assistant_w.stop_stream()
|
||||
if response_display.assistant._content != current:
|
||||
response_display.assistant._content = final_segment
|
||||
await response_display.assistant.stop_stream()
|
||||
else:
|
||||
suffix = final_text[len(current) :]
|
||||
suffix = final_segment[len(current) :]
|
||||
if suffix:
|
||||
await assistant_w.append_content(suffix)
|
||||
await response_display.assistant.append_content(suffix)
|
||||
|
||||
_schedule_scroll()
|
||||
return final_text
|
||||
@@ -1370,6 +1455,13 @@ def run_textual_interactive(
|
||||
collapse_summary_w.update(summary)
|
||||
collapse_summary_w.display = True
|
||||
|
||||
def _expand_completed_tools() -> None:
|
||||
"""Show every completed tool once the turn reaches a final state."""
|
||||
if collapse_summary_w is not None:
|
||||
collapse_summary_w.display = False
|
||||
for tw in tool_widgets.values():
|
||||
tw.display = True
|
||||
|
||||
def _find_or_rename_sa_widget(
|
||||
resolved_name: str,
|
||||
description: str = "",
|
||||
@@ -1520,39 +1612,38 @@ def run_textual_interactive(
|
||||
_schedule_scroll()
|
||||
|
||||
elif event_type == "text":
|
||||
chunk = event.get("content", "")
|
||||
if thinking_w is not None and thinking_w._is_active:
|
||||
thinking_w.finalize()
|
||||
# Clear processing indicator
|
||||
await _remove_w(processing_w)
|
||||
processing_w = None
|
||||
|
||||
if has_used_tools and not _is_final_response(state):
|
||||
# Tools still running — show intermediate narration
|
||||
await _remove_w(narration_w)
|
||||
narration_w = None
|
||||
last_line = (
|
||||
state.latest_text.strip().split("\n")[-1].strip()
|
||||
)
|
||||
if last_line:
|
||||
if len(last_line) > 60:
|
||||
last_line = last_line[:57] + "\u2026"
|
||||
narration_w = Static(
|
||||
Text(f" {last_line}", style="dim italic"),
|
||||
)
|
||||
await container.mount(narration_w)
|
||||
if not _is_final_response(state):
|
||||
await _mount_or_update_narration(state.latest_text)
|
||||
self._set_status_streaming_text(state.latest_text)
|
||||
else:
|
||||
# Stream final response incrementally (both
|
||||
# text-only replies and post-tool responses).
|
||||
await _remove_w(narration_w)
|
||||
narration_w = None
|
||||
if assistant_w is None:
|
||||
assistant_w = AssistantMessage(state.response_text)
|
||||
await container.mount(assistant_w)
|
||||
else:
|
||||
await assistant_w.append_content(
|
||||
event.get("content", ""),
|
||||
await _preserve_active_narration()
|
||||
if response_display.assistant is None:
|
||||
response_display.assistant_start = max(
|
||||
response_display.narrated_end,
|
||||
len(state.response_text) - len(chunk),
|
||||
)
|
||||
self._set_status_streaming_text(state.response_text)
|
||||
response_display.assistant = AssistantMessage(
|
||||
state.response_text[
|
||||
response_display.assistant_start :
|
||||
]
|
||||
)
|
||||
await container.mount(response_display.assistant)
|
||||
else:
|
||||
await response_display.assistant.append_content(
|
||||
chunk
|
||||
)
|
||||
self._set_status_streaming_text(
|
||||
response_display.assistant._content
|
||||
)
|
||||
|
||||
elif event_type == "tool_call":
|
||||
tool_name = event.get("name", "unknown")
|
||||
@@ -1562,21 +1653,17 @@ def run_textual_interactive(
|
||||
if thinking_w is not None and thinking_w._is_active:
|
||||
thinking_w.finalize()
|
||||
# Clear transient indicators
|
||||
await _remove_w(narration_w)
|
||||
narration_w = None
|
||||
await _preserve_active_narration()
|
||||
await _remove_w(processing_w)
|
||||
processing_w = None
|
||||
# Remove early AssistantMessage (text arrived before tools)
|
||||
if assistant_w is not None:
|
||||
try:
|
||||
await assistant_w.remove()
|
||||
except Exception:
|
||||
pass
|
||||
assistant_w = None
|
||||
# Skip internal tools and task (handled by SubAgentWidget)
|
||||
if tool_name not in _INTERNAL_TOOLS and tool_name != "task":
|
||||
has_used_tools = True
|
||||
if tool_id and tool_id in tool_widgets:
|
||||
existing_tool = bool(tool_id and tool_id in tool_widgets)
|
||||
if response_display.assistant is not None and (
|
||||
tool_name == "task" or not existing_tool
|
||||
):
|
||||
await _convert_assistant_to_narration()
|
||||
# Task tools are handled by SubAgentWidget.
|
||||
if tool_name != "task":
|
||||
if existing_tool:
|
||||
# Re-emitted with updated args — update in place
|
||||
existing = tool_widgets[tool_id]
|
||||
existing._tool_name = tool_name
|
||||
@@ -1839,27 +1926,39 @@ def run_textual_interactive(
|
||||
|
||||
elif event_type == "done":
|
||||
# Clean up transient indicators
|
||||
await _remove_w(narration_w)
|
||||
narration_w = None
|
||||
await _preserve_active_narration()
|
||||
await _remove_w(processing_w)
|
||||
processing_w = None
|
||||
_expand_completed_tools()
|
||||
# Mount final response
|
||||
if assistant_w is None and state.response_text:
|
||||
# Strip trailing standalone "..."
|
||||
clean = state.response_text.strip()
|
||||
while (
|
||||
clean.endswith("\n...") or clean.rstrip() == "..."
|
||||
):
|
||||
clean = clean.rstrip().removesuffix("...").rstrip()
|
||||
assistant_w = AssistantMessage(
|
||||
clean or state.response_text
|
||||
)
|
||||
await container.mount(assistant_w)
|
||||
self._schedule_scroll_to_bottom(
|
||||
container,
|
||||
delays=(0.15, 0.4, 0.8, 1.5),
|
||||
immediate=False,
|
||||
)
|
||||
final_response_text = _response_after_narration(
|
||||
state.response_text,
|
||||
response_display.narrated_end,
|
||||
)
|
||||
clean = _strip_trailing_placeholder_ellipsis(
|
||||
final_response_text
|
||||
)
|
||||
if clean:
|
||||
response_display.assistant_start = len(
|
||||
state.response_text
|
||||
) - len(final_response_text)
|
||||
if response_display.assistant is None:
|
||||
response_display.assistant = AssistantMessage(clean)
|
||||
await container.mount(response_display.assistant)
|
||||
self._schedule_scroll_to_bottom(
|
||||
container,
|
||||
delays=(0.15, 0.4, 0.8, 1.5),
|
||||
immediate=False,
|
||||
)
|
||||
elif response_display.assistant._content != clean:
|
||||
response_display.assistant._content = clean
|
||||
await response_display.assistant.stop_stream()
|
||||
elif response_display.assistant is not None:
|
||||
try:
|
||||
await response_display.assistant.remove()
|
||||
except Exception:
|
||||
pass
|
||||
response_display.assistant = None
|
||||
# Mount token usage stats with elapsed time
|
||||
if state.total_input_tokens or state.total_output_tokens:
|
||||
elapsed = None
|
||||
@@ -1885,7 +1984,9 @@ def run_textual_interactive(
|
||||
response = (state.response_text or "").strip()
|
||||
|
||||
except asyncio.CancelledError:
|
||||
# Ctrl+C cancellation — re-raise so _run_turn can handle it
|
||||
# Ctrl+C cancellation: preserve streamed text before the outer
|
||||
# turn handler appends its interruption notice.
|
||||
await _mark_cancelled_response()
|
||||
raise
|
||||
except Exception as exc:
|
||||
error_msg = str(exc)
|
||||
@@ -1911,12 +2012,14 @@ def run_textual_interactive(
|
||||
await loading.cleanup()
|
||||
except Exception:
|
||||
pass
|
||||
# Clean up transient indicators
|
||||
for w in (narration_w, processing_w):
|
||||
await _remove_w(w)
|
||||
# Clean up transient indicators. Inline narration is transcript
|
||||
# content, so finalize it without removing the widget.
|
||||
await _preserve_active_narration()
|
||||
await _remove_w(processing_w)
|
||||
# Mark any still-running tool widgets as interrupted
|
||||
# (skip if HITL approved — tools will continue next round)
|
||||
if not _hitl_resuming:
|
||||
_expand_completed_tools()
|
||||
for tw in tool_widgets.values():
|
||||
if tw._status == "running":
|
||||
try:
|
||||
@@ -1937,8 +2040,8 @@ def run_textual_interactive(
|
||||
except Exception:
|
||||
pass
|
||||
# Finalize assistant message stream
|
||||
if assistant_w is not None:
|
||||
await assistant_w.stop_stream()
|
||||
if response_display.assistant is not None:
|
||||
await response_display.assistant.stop_stream()
|
||||
# Flush remaining thinking callback
|
||||
if (
|
||||
on_thinking_cb
|
||||
|
||||
@@ -20,8 +20,6 @@ from .context_editing import (
|
||||
from .context_overflow import ContextOverflowMapperMiddleware
|
||||
from .memory import (
|
||||
EvoMemoryMiddleware,
|
||||
EvoMemoryState,
|
||||
ExtractedMemory,
|
||||
create_memory_middleware,
|
||||
)
|
||||
from .model_fallback import ModelFallbackMiddleware, load_fallback_chain
|
||||
@@ -37,8 +35,6 @@ __all__ = [
|
||||
"ConfigurableModelMiddleware",
|
||||
"ContextOverflowMapperMiddleware",
|
||||
"EvoMemoryMiddleware",
|
||||
"EvoMemoryState",
|
||||
"ExtractedMemory",
|
||||
"ModelFallbackMiddleware",
|
||||
"Question",
|
||||
"ToolErrorHandlerMiddleware",
|
||||
|
||||
+347
-741
File diff suppressed because it is too large
Load Diff
+168
-98
@@ -27,7 +27,6 @@ from .diff_format import build_edit_diff
|
||||
from .events import stream_agent_events
|
||||
from .formatter import ToolResultFormatter
|
||||
from .state import (
|
||||
_INTERNAL_TOOLS,
|
||||
StreamState,
|
||||
SubAgentState,
|
||||
_build_todo_stats,
|
||||
@@ -65,6 +64,38 @@ def _fix_markdown_heading_spacing(text: str) -> str:
|
||||
return _HEADING_FIX_RE.sub(r"\1 ", text)
|
||||
|
||||
|
||||
def _split_response_for_display(
|
||||
response_text: str,
|
||||
narrated_response_end: int,
|
||||
) -> tuple[str, str]:
|
||||
"""Split cumulative response text into narrated prefix and answer suffix."""
|
||||
boundary = max(0, min(len(response_text), narrated_response_end))
|
||||
return response_text[:boundary], response_text[boundary:]
|
||||
|
||||
|
||||
def _clean_response_text(text: str) -> str:
|
||||
"""Trim a streamed response copy for display."""
|
||||
clean = text.strip()
|
||||
while clean.endswith("\n...") or clean.rstrip() == "...":
|
||||
clean = clean.rstrip().removesuffix("...").rstrip()
|
||||
return clean
|
||||
|
||||
|
||||
def _response_markdown_for_display(
|
||||
text: str,
|
||||
*,
|
||||
response_markdown: Any = None,
|
||||
full_response_text: str = "",
|
||||
) -> Any | None:
|
||||
"""Build Markdown for the answer text, reusing the full-response cache if valid."""
|
||||
clean = _clean_response_text(text)
|
||||
if not clean:
|
||||
return None
|
||||
if response_markdown is not None and text == full_response_text:
|
||||
return response_markdown
|
||||
return Markdown(_fix_markdown_heading_spacing(clean))
|
||||
|
||||
|
||||
formatter = ToolResultFormatter()
|
||||
|
||||
|
||||
@@ -472,6 +503,8 @@ def create_streaming_display(
|
||||
final_show_thinking: bool = False,
|
||||
final_thinking_max_length: int = DisplayLimits.THINKING_FINAL,
|
||||
response_markdown: Any = None,
|
||||
narrated_response_end: int = 0,
|
||||
narration_segments: list[tuple[int, str]] | None = None,
|
||||
total_input_tokens: int = 0,
|
||||
total_output_tokens: int = 0,
|
||||
summarization_text: str = "",
|
||||
@@ -568,11 +601,72 @@ def create_streaming_display(
|
||||
)
|
||||
)
|
||||
|
||||
# Response text handling: keep the final answer behind pending tool calls.
|
||||
_n_tools = len(tool_calls)
|
||||
_n_done = min(len(tool_results), _n_tools)
|
||||
has_pending_tools = _n_tools > _n_done
|
||||
any_active_subagent = any(sa.is_active for sa in subagents)
|
||||
is_processing_blocking = is_processing
|
||||
all_done = (
|
||||
not has_pending_tools and not any_active_subagent and not is_processing_blocking
|
||||
)
|
||||
_, answer_text = _split_response_for_display(
|
||||
response_text,
|
||||
narrated_response_end,
|
||||
)
|
||||
narration_by_tool: dict[int, list[str]] = {}
|
||||
for tool_index, text in narration_segments or []:
|
||||
if text.strip():
|
||||
narration_by_tool.setdefault(tool_index, []).append(text)
|
||||
|
||||
def _append_narration_before_tool(tool_index: int) -> None:
|
||||
for text in narration_by_tool.get(tool_index, []):
|
||||
narration_markdown = _response_markdown_for_display(text)
|
||||
if narration_markdown is not None:
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(narration_markdown)
|
||||
|
||||
def _find_task_subagent(tc: dict, shown_sa_names: set[str]) -> SubAgentState | None:
|
||||
sa_name = tc.get("args", {}).get("subagent_type", "")
|
||||
task_desc = tc.get("args", {}).get("description", "")
|
||||
for sa in subagents:
|
||||
if sa.name in shown_sa_names:
|
||||
continue
|
||||
if sa.name == sa_name or (
|
||||
task_desc and task_desc in (sa.description or "")
|
||||
):
|
||||
return sa
|
||||
|
||||
candidates = [
|
||||
sa
|
||||
for sa in subagents
|
||||
if sa.name not in shown_sa_names and (sa.tool_calls or sa.is_active)
|
||||
]
|
||||
if len(candidates) == 1:
|
||||
return candidates[0]
|
||||
return None
|
||||
|
||||
def _append_task_entry(
|
||||
tool_index: int,
|
||||
tc: dict,
|
||||
tr: dict | None,
|
||||
*,
|
||||
shown_sa_names: set[str],
|
||||
compact: bool,
|
||||
) -> None:
|
||||
_append_narration_before_tool(tool_index)
|
||||
elements.append(_render_tool_call_line(tc, tr))
|
||||
matched_sa = _find_task_subagent(tc, shown_sa_names)
|
||||
if matched_sa is not None:
|
||||
shown_sa_names.add(matched_sa.name)
|
||||
elements.extend(_render_subagent_section(matched_sa, compact=compact))
|
||||
|
||||
# Tool calls and results paired display
|
||||
# Collapse older completed tools to prevent overflow in Live mode
|
||||
# Task tool calls are ALWAYS visible (they represent sub-agent delegations)
|
||||
MAX_VISIBLE_TOOLS = 4
|
||||
MAX_VISIBLE_RUNNING = 3
|
||||
shown_sa_names: set[str] = set()
|
||||
|
||||
if tool_calls:
|
||||
# Split into categories
|
||||
@@ -585,24 +679,30 @@ def create_streaming_display(
|
||||
tr = tool_results[i] if has_result else None
|
||||
is_task = tc.get("name") == "task"
|
||||
|
||||
# Skip internal middleware tools
|
||||
if tc.get("name") in _INTERNAL_TOOLS:
|
||||
continue
|
||||
|
||||
if is_task:
|
||||
# Skip task calls with empty args (still streaming)
|
||||
if tc.get("args"):
|
||||
task_tools.append((tc, tr))
|
||||
task_tools.append((i, tc, tr))
|
||||
elif has_result:
|
||||
completed_regular.append((tc, tr))
|
||||
completed_regular.append((i, tc, tr))
|
||||
else:
|
||||
running_regular.append((tc, None))
|
||||
running_regular.append((i, tc, None))
|
||||
|
||||
if is_final:
|
||||
# Final frame: show ALL tools expanded, no spinners, no collapsing
|
||||
shown_sa_names: set[str] = set()
|
||||
for tool_index, tc, tr in sorted(
|
||||
completed_regular + running_regular + task_tools,
|
||||
key=lambda item: item[0],
|
||||
):
|
||||
if tc.get("name") == "task":
|
||||
_append_task_entry(
|
||||
tool_index,
|
||||
tc,
|
||||
tr,
|
||||
shown_sa_names=shown_sa_names,
|
||||
compact=True,
|
||||
)
|
||||
continue
|
||||
|
||||
for tc, tr in completed_regular:
|
||||
_append_narration_before_tool(tool_index)
|
||||
elements.append(_render_tool_call_line(tc, tr))
|
||||
content = tr.get("content", "") if tr else ""
|
||||
if tr and (not is_success(content) or tc.get("name") == "edit_file"):
|
||||
@@ -614,22 +714,6 @@ def create_streaming_display(
|
||||
)
|
||||
elements.extend(result_elements)
|
||||
|
||||
# Task tools with compact sub-agent summaries
|
||||
for tc, tr in task_tools:
|
||||
elements.append(_render_tool_call_line(tc, tr))
|
||||
sa_name = tc.get("args", {}).get("subagent_type", "")
|
||||
task_desc = tc.get("args", {}).get("description", "")
|
||||
matched_sa = None
|
||||
for sa in subagents:
|
||||
if sa.name == sa_name or (
|
||||
task_desc and task_desc in (sa.description or "")
|
||||
):
|
||||
matched_sa = sa
|
||||
break
|
||||
if matched_sa:
|
||||
shown_sa_names.add(matched_sa.name)
|
||||
elements.extend(_render_subagent_section(matched_sa, compact=True))
|
||||
|
||||
# Render any sub-agents not already shown via task tool calls
|
||||
for sa in subagents:
|
||||
if sa.name not in shown_sa_names and (sa.tool_calls or sa.is_active):
|
||||
@@ -647,7 +731,9 @@ def create_streaming_display(
|
||||
visible = completed_regular[-slots:] if slots else []
|
||||
|
||||
if hidden:
|
||||
ok = sum(1 for _, tr in hidden if is_success(tr.get("content", "")))
|
||||
for tool_index, _, _ in hidden:
|
||||
_append_narration_before_tool(tool_index)
|
||||
ok = sum(1 for _, _, tr in hidden if is_success(tr.get("content", "")))
|
||||
fail = len(hidden) - ok
|
||||
summary = Text()
|
||||
summary.append(f"\u2713 {ok} completed", style="dim green")
|
||||
@@ -655,10 +741,42 @@ def create_streaming_display(
|
||||
summary.append(f" | {fail} failed", style="dim red")
|
||||
elements.append(summary)
|
||||
|
||||
for tc, tr in visible:
|
||||
# --- Running regular tools (limit visible) ---
|
||||
hidden_running = len(running_regular) - MAX_VISIBLE_RUNNING
|
||||
if hidden_running > 0:
|
||||
hidden_running_tools = running_regular[:-MAX_VISIBLE_RUNNING]
|
||||
for tool_index, _, _ in hidden_running_tools:
|
||||
_append_narration_before_tool(tool_index)
|
||||
summary = Text()
|
||||
summary.append(
|
||||
f"\u25cf {hidden_running} more running...", style="dim yellow"
|
||||
)
|
||||
elements.append(summary)
|
||||
running_regular = running_regular[-MAX_VISIBLE_RUNNING:]
|
||||
|
||||
for tool_index, tc, tr in sorted(
|
||||
visible + running_regular + task_tools,
|
||||
key=lambda item: item[0],
|
||||
):
|
||||
if tc.get("name") == "task":
|
||||
matched_sa = _find_task_subagent(tc, shown_sa_names)
|
||||
_append_task_entry(
|
||||
tool_index,
|
||||
tc,
|
||||
tr,
|
||||
shown_sa_names=shown_sa_names,
|
||||
compact=not matched_sa.is_active if matched_sa else True,
|
||||
)
|
||||
continue
|
||||
|
||||
_append_narration_before_tool(tool_index)
|
||||
elements.append(_render_tool_call_line(tc, tr))
|
||||
content = tr.get("content", "") if tr else ""
|
||||
if tr and (not is_success(content) or tc.get("name") == "edit_file"):
|
||||
if tr is None:
|
||||
elements.append(Spinner("dots", text=" Running...", style="yellow"))
|
||||
continue
|
||||
|
||||
content = tr.get("content", "")
|
||||
if not is_success(content) or tc.get("name") == "edit_file":
|
||||
result_elements = format_tool_result_compact(
|
||||
tr["name"],
|
||||
content,
|
||||
@@ -667,36 +785,7 @@ def create_streaming_display(
|
||||
)
|
||||
elements.extend(result_elements)
|
||||
|
||||
# --- Running regular tools (limit visible) ---
|
||||
hidden_running = len(running_regular) - MAX_VISIBLE_RUNNING
|
||||
if hidden_running > 0:
|
||||
summary = Text()
|
||||
summary.append(
|
||||
f"\u25cf {hidden_running} more running...", style="dim yellow"
|
||||
)
|
||||
elements.append(summary)
|
||||
running_regular = running_regular[-MAX_VISIBLE_RUNNING:]
|
||||
|
||||
for tc, tr in running_regular:
|
||||
elements.append(_render_tool_call_line(tc, tr))
|
||||
elements.append(Spinner("dots", text=" Running...", style="yellow"))
|
||||
|
||||
# Task tool calls are rendered as part of sub-agent sections below
|
||||
|
||||
# Response text handling — exclude internal tools (e.g. ExtractedMemory)
|
||||
# from the "done" calculation so they don't block final Markdown rendering.
|
||||
_n_visible = 0
|
||||
_n_visible_done = 0
|
||||
for i, tc in enumerate(tool_calls):
|
||||
if tc.get("name") in _INTERNAL_TOOLS:
|
||||
continue
|
||||
_n_visible += 1
|
||||
if i < len(tool_results):
|
||||
_n_visible_done += 1
|
||||
has_pending_tools = _n_visible > _n_visible_done
|
||||
any_active_subagent = any(sa.is_active for sa in subagents)
|
||||
has_used_tools = _n_visible > 0
|
||||
all_done = not has_pending_tools and not any_active_subagent and not is_processing
|
||||
# Remaining sub-agent sections are rendered below.
|
||||
|
||||
if is_final:
|
||||
# Final frame: render todo panel + response (tools/subagents handled above).
|
||||
@@ -707,17 +796,14 @@ def create_streaming_display(
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(_render_todo_panel(todo_items))
|
||||
|
||||
# Include response in final frame so it stays visible after Live exits
|
||||
if response_text:
|
||||
clean_response = response_text.strip()
|
||||
while clean_response.endswith("\n...") or clean_response.rstrip() == "...":
|
||||
clean_response = clean_response.rstrip().removesuffix("...").rstrip()
|
||||
if clean_response:
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(
|
||||
response_markdown
|
||||
or Markdown(_fix_markdown_heading_spacing(clean_response))
|
||||
)
|
||||
answer_markdown = _response_markdown_for_display(
|
||||
answer_text,
|
||||
response_markdown=response_markdown,
|
||||
full_response_text=response_text,
|
||||
)
|
||||
if answer_markdown is not None:
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(answer_markdown)
|
||||
|
||||
# Token usage stats (right-aligned)
|
||||
if total_input_tokens or total_output_tokens:
|
||||
@@ -731,16 +817,6 @@ def create_streaming_display(
|
||||
stats.append("]", style="dim italic")
|
||||
elements.append(stats)
|
||||
else:
|
||||
# Intermediate narration (tools still running) -- dim italic above Task List
|
||||
if latest_text and has_used_tools and not all_done:
|
||||
preview = latest_text.strip()
|
||||
if preview:
|
||||
last_line = preview.split("\n")[-1].strip()
|
||||
if last_line:
|
||||
if len(last_line) > 60:
|
||||
last_line = last_line[:57] + "\u2026"
|
||||
elements.append(Text(f" {last_line}", style="dim italic"))
|
||||
|
||||
# Task List panel (persistent, updates on write_todos / read_todos)
|
||||
todo_items = todo_items or []
|
||||
if todo_items:
|
||||
@@ -750,16 +826,11 @@ def create_streaming_display(
|
||||
# Sub-agent activity sections
|
||||
# Active: full bordered view; Completed: compact 1-line summary
|
||||
for sa in subagents:
|
||||
if sa.tool_calls or sa.is_active:
|
||||
if sa.name not in shown_sa_names and (sa.tool_calls or sa.is_active):
|
||||
elements.extend(_render_subagent_section(sa, compact=not sa.is_active))
|
||||
|
||||
# Processing state after tool execution
|
||||
if (
|
||||
is_processing
|
||||
and not is_thinking
|
||||
and not is_responding
|
||||
and not response_text
|
||||
):
|
||||
if is_processing and not is_thinking and not is_responding:
|
||||
# Check if any sub-agent is active
|
||||
any_active = any(sa.is_active for sa in subagents)
|
||||
if not any_active:
|
||||
@@ -769,11 +840,14 @@ def create_streaming_display(
|
||||
|
||||
# Stream response in real-time as tokens arrive (all tools done)
|
||||
if response_text and all_done:
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(
|
||||
response_markdown
|
||||
or Markdown(_fix_markdown_heading_spacing(response_text))
|
||||
answer_markdown = _response_markdown_for_display(
|
||||
answer_text,
|
||||
response_markdown=response_markdown,
|
||||
full_response_text=response_text,
|
||||
)
|
||||
if answer_markdown is not None:
|
||||
elements.append(Text("")) # blank separator
|
||||
elements.append(answer_markdown)
|
||||
|
||||
if not elements:
|
||||
elements.append(Spinner("dots", text=" Processing...", style="cyan"))
|
||||
@@ -849,10 +923,6 @@ def display_final_results(
|
||||
tool_name = tc.get("name", "")
|
||||
is_task = tool_name.lower() == "task"
|
||||
|
||||
# Skip internal middleware tools
|
||||
if tool_name in _INTERNAL_TOOLS:
|
||||
continue
|
||||
|
||||
# Task tools: show delegation line + compact sub-agent summary
|
||||
if is_task:
|
||||
console.print(_render_tool_call_line(tc, tr))
|
||||
|
||||
@@ -9,10 +9,6 @@ import ast
|
||||
import json
|
||||
from enum import StrEnum
|
||||
|
||||
# Tool names that are internal middleware artifacts (not user-visible actions).
|
||||
# These should be excluded from display rendering and "all_done" calculations.
|
||||
_INTERNAL_TOOLS = {"ExtractedMemory"}
|
||||
|
||||
|
||||
class ResearchPhase(StrEnum):
|
||||
"""Research phase constants used by the TUI status bar."""
|
||||
@@ -122,6 +118,8 @@ class StreamState:
|
||||
self.todo_items: list[dict] = []
|
||||
# Latest text segment (reset on each tool_call)
|
||||
self.latest_text = ""
|
||||
self.narrated_response_end = 0
|
||||
self.narration_segments = []
|
||||
# Token usage tracking
|
||||
self.total_input_tokens = 0
|
||||
self.total_output_tokens = 0
|
||||
@@ -139,7 +137,7 @@ class StreamState:
|
||||
|
||||
def get_response_markdown(self):
|
||||
"""Return cached Markdown object, only re-parsing when text changes."""
|
||||
from rich.markdown import Markdown # type: ignore[import-untyped]
|
||||
from rich.markdown import Markdown
|
||||
|
||||
from .display import _fix_markdown_heading_spacing
|
||||
|
||||
@@ -221,7 +219,6 @@ class StreamState:
|
||||
self.is_thinking = False
|
||||
self.is_responding = False
|
||||
self.is_processing = False
|
||||
self.latest_text = "" # Reset -- next text segment is a new message
|
||||
|
||||
tool_id = event.get("id", "")
|
||||
tool_name = event.get("name", "unknown")
|
||||
@@ -240,10 +237,22 @@ class StreamState:
|
||||
updated = True
|
||||
break
|
||||
if not updated:
|
||||
if self.latest_text.strip():
|
||||
self.narration_segments.append(
|
||||
(len(self.tool_calls), self.latest_text)
|
||||
)
|
||||
self.narrated_response_end = len(self.response_text)
|
||||
self.tool_calls.append(tc_data)
|
||||
else:
|
||||
if self.latest_text.strip():
|
||||
self.narration_segments.append(
|
||||
(len(self.tool_calls), self.latest_text)
|
||||
)
|
||||
self.narrated_response_end = len(self.response_text)
|
||||
self.tool_calls.append(tc_data)
|
||||
|
||||
self.latest_text = "" # Reset -- next text segment is a new message
|
||||
|
||||
# Capture todo items from write_todos args (most reliable source)
|
||||
if tool_name == "write_todos":
|
||||
todos = tool_args.get("todos", [])
|
||||
@@ -252,8 +261,7 @@ class StreamState:
|
||||
|
||||
elif event_type == "tool_result":
|
||||
result_name = event.get("name", "unknown")
|
||||
if result_name not in _INTERNAL_TOOLS:
|
||||
self.is_processing = True
|
||||
self.is_processing = True
|
||||
result_content = event.get("content", "")
|
||||
self.tool_results.append(
|
||||
{
|
||||
@@ -352,16 +360,9 @@ class StreamState:
|
||||
return event_type
|
||||
|
||||
def visible_tool_counts(self) -> tuple[int, int]:
|
||||
"""Return (completed, total) counts for visible (non-internal) tools."""
|
||||
n_visible = 0
|
||||
n_done = 0
|
||||
for i, tc in enumerate(self.tool_calls):
|
||||
if tc.get("name") in _INTERNAL_TOOLS:
|
||||
continue
|
||||
n_visible += 1
|
||||
if i < len(self.tool_results):
|
||||
n_done += 1
|
||||
return n_done, n_visible
|
||||
"""Return (completed, total) counts for tool calls."""
|
||||
n_total = len(self.tool_calls)
|
||||
return min(len(self.tool_results), n_total), n_total
|
||||
|
||||
def has_pending_work(self) -> bool:
|
||||
"""Return True if tools or sub-agents are still running."""
|
||||
@@ -396,6 +397,8 @@ class StreamState:
|
||||
"is_summarizing": self.is_summarizing,
|
||||
"response_text": self.response_text,
|
||||
"latest_text": self.latest_text,
|
||||
"narrated_response_end": self.narrated_response_end,
|
||||
"narration_segments": self.narration_segments,
|
||||
"tool_calls": self.tool_calls,
|
||||
"tool_results": self.tool_results,
|
||||
"is_thinking": self.is_thinking,
|
||||
|
||||
@@ -7,6 +7,7 @@ adapted for deepagents tool names.
|
||||
|
||||
import sys
|
||||
from enum import StrEnum
|
||||
from functools import lru_cache
|
||||
from pathlib import PurePath
|
||||
|
||||
# === Status marker constants ===
|
||||
@@ -119,6 +120,23 @@ def _is_memory_path(path: str) -> bool:
|
||||
return normalized == "/memories" or normalized.startswith("/memories/")
|
||||
|
||||
|
||||
@lru_cache(maxsize=1)
|
||||
def _profile_memory_headings() -> tuple[str, ...]:
|
||||
"""Return profile headings from the canonical profile templates."""
|
||||
from EvoScientist.middleware.memory import PROFILE_TEMPLATES
|
||||
|
||||
return tuple(
|
||||
template.strip().splitlines()[0].strip()
|
||||
for template in PROFILE_TEMPLATES.values()
|
||||
if template.strip()
|
||||
)
|
||||
|
||||
|
||||
def _looks_like_profile_memory(content: str) -> bool:
|
||||
"""Recognize profile-memory content when streamed without file args."""
|
||||
return any(heading in content for heading in _profile_memory_headings())
|
||||
|
||||
|
||||
def format_tool_compact(name: str, args: dict | None) -> str:
|
||||
"""Format as compact tool call string: ToolName(key_arg).
|
||||
|
||||
@@ -141,19 +159,19 @@ def format_tool_compact(name: str, args: dict | None) -> str:
|
||||
# File operations (with special case for memory files)
|
||||
if name_lower == "read_file":
|
||||
path = _tool_path_arg(args)
|
||||
if _is_memory_path(path) or path.endswith("/MEMORY.md") or path == "/MEMORY.md":
|
||||
if _is_memory_path(path):
|
||||
return "Reading memory"
|
||||
return f"read_file({_shorten_path(path)})"
|
||||
|
||||
if name_lower == "write_file":
|
||||
path = _tool_path_arg(args)
|
||||
if _is_memory_path(path) or path.endswith("/MEMORY.md") or path == "/MEMORY.md":
|
||||
if _is_memory_path(path):
|
||||
return "Updating memory"
|
||||
return f"write_file({_shorten_path(path)})"
|
||||
|
||||
if name_lower == "edit_file":
|
||||
path = _tool_path_arg(args)
|
||||
if _is_memory_path(path) or path.endswith("/MEMORY.md") or path == "/MEMORY.md":
|
||||
if _is_memory_path(path):
|
||||
return "Updating memory"
|
||||
return f"edit_file({_shorten_path(path)})"
|
||||
|
||||
@@ -249,15 +267,11 @@ def format_tool_compact_with_result(
|
||||
result_content = result_content or ""
|
||||
|
||||
if name_lower in ("write_file", "edit_file"):
|
||||
if (
|
||||
"/memories/" in result_content
|
||||
or "/MEMORY.md" in result_content
|
||||
or "MEMORY.md" in result_content
|
||||
):
|
||||
if "/memories/" in result_content:
|
||||
return "Updating memory"
|
||||
elif name_lower == "read_file":
|
||||
path = _tool_path_arg(args)
|
||||
if not path and "# EvoScientist Memory" in result_content:
|
||||
if not path and _looks_like_profile_memory(result_content):
|
||||
return "Reading memory"
|
||||
|
||||
return compact
|
||||
|
||||
@@ -10,7 +10,7 @@ runnable graph.
|
||||
Reuses the main EvoScientist agent's chat model, backend, and middleware so
|
||||
the deployed sub-agent has full capability parity with its in-process
|
||||
synchronous counterpart: same workspace files, same ``/skills/`` and
|
||||
``/memory/`` routes, same error-handling and context-overflow middleware.
|
||||
``/memories/`` routes, same error-handling and context-overflow middleware.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -96,12 +96,8 @@ def build_async_subagent_graph(name: str) -> Any:
|
||||
# which uses ``interrupt()`` for the same purpose (waiting on a user
|
||||
# reply) and would deadlock an async sub-agent for the same reason.
|
||||
#
|
||||
# Memory middleware is included so async sub-agents can READ
|
||||
# /memory/MEMORY.md, but the extraction trigger (20+ human messages,
|
||||
# see middleware/memory.py) never fires here — sub-agents only receive
|
||||
# the parent's task delegation as a system prompt, not human messages.
|
||||
# Net effect: sub-agents have read-only memory access. Memory writes
|
||||
# happen exclusively from the main agent's user-facing conversation.
|
||||
# Memory middleware is included so async sub-agents get the same profile
|
||||
# context and `/memories/profile/...` file guidance as the main agent.
|
||||
return create_deep_agent(
|
||||
name=name,
|
||||
model=_ensure_chat_model(),
|
||||
|
||||
@@ -1,73 +0,0 @@
|
||||
"""Tests for _merge_memory — backslash-safe regex replacement."""
|
||||
|
||||
import pytest
|
||||
|
||||
from EvoScientist.middleware.memory import DEFAULT_MEMORY_TEMPLATE, _merge_memory
|
||||
|
||||
|
||||
class TestMergeMemoryBackslashSafety:
|
||||
"""Ensure values containing regex-special sequences survive _merge_memory."""
|
||||
|
||||
def test_backslash_n_preserved(self):
|
||||
"""A value containing literal '\\n' must not become a newline."""
|
||||
extracted = {
|
||||
"user_profile": {"name": r"C:\new_user"},
|
||||
}
|
||||
result = _merge_memory(DEFAULT_MEMORY_TEMPLATE, extracted)
|
||||
assert r"C:\new_user" in result
|
||||
# The replacement must not introduce an actual newline inside the Name line
|
||||
for line in result.splitlines():
|
||||
if "**Name**" in line:
|
||||
assert r"C:\new_user" in line
|
||||
break
|
||||
else:
|
||||
pytest.fail("Name line not found")
|
||||
|
||||
def test_backreference_preserved(self):
|
||||
r"""A value containing '\\1' must not be treated as a backreference."""
|
||||
extracted = {
|
||||
"user_profile": {"role": r"A\1B"},
|
||||
}
|
||||
result = _merge_memory(DEFAULT_MEMORY_TEMPLATE, extracted)
|
||||
assert r"A\1B" in result
|
||||
for line in result.splitlines():
|
||||
if "**Role**" in line:
|
||||
assert r"A\1B" in line
|
||||
break
|
||||
else:
|
||||
pytest.fail("Role line not found")
|
||||
|
||||
def test_windows_path_preserved(self):
|
||||
r"""A Windows-style path must survive without corruption."""
|
||||
extracted = {
|
||||
"research_preferences": {
|
||||
"preferred_frameworks": r"C:\path\to\file",
|
||||
},
|
||||
}
|
||||
result = _merge_memory(DEFAULT_MEMORY_TEMPLATE, extracted)
|
||||
assert r"C:\path\to\file" in result
|
||||
|
||||
def test_multiple_backslash_fields(self):
|
||||
"""Multiple fields with backslashes all survive."""
|
||||
extracted = {
|
||||
"user_profile": {
|
||||
"name": r"user\name",
|
||||
"institution": r"MIT\Lab\42",
|
||||
},
|
||||
"research_preferences": {
|
||||
"hardware": r"GPU\0",
|
||||
},
|
||||
}
|
||||
result = _merge_memory(DEFAULT_MEMORY_TEMPLATE, extracted)
|
||||
assert r"user\name" in result
|
||||
assert r"MIT\Lab\42" in result
|
||||
assert r"GPU\0" in result
|
||||
|
||||
def test_plain_value_still_works(self):
|
||||
"""Sanity check: normal values without backslashes work fine."""
|
||||
extracted = {
|
||||
"user_profile": {"name": "Alice", "role": "Researcher"},
|
||||
}
|
||||
result = _merge_memory(DEFAULT_MEMORY_TEMPLATE, extracted)
|
||||
assert "- **Name**: Alice" in result
|
||||
assert "- **Role**: Researcher" in result
|
||||
@@ -0,0 +1,338 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from types import SimpleNamespace
|
||||
|
||||
from langchain_core.messages import SystemMessage
|
||||
|
||||
import EvoScientist.middleware.memory as memory_module
|
||||
from EvoScientist import paths
|
||||
|
||||
|
||||
def _request():
|
||||
request = SimpleNamespace(
|
||||
state={},
|
||||
runtime=object(),
|
||||
system_message=SystemMessage(content="base system"),
|
||||
)
|
||||
request.override = lambda **kwargs: SimpleNamespace(
|
||||
**{
|
||||
"state": request.state,
|
||||
"runtime": request.runtime,
|
||||
"system_message": kwargs.get("system_message", request.system_message),
|
||||
}
|
||||
)
|
||||
return request
|
||||
|
||||
|
||||
def _system_text(modified) -> str:
|
||||
system_message = modified.system_message
|
||||
assert system_message is not None
|
||||
return str(system_message.content)
|
||||
|
||||
|
||||
def _path_project_id(workspace) -> str:
|
||||
return memory_module._resolve_project_id(workspace)
|
||||
|
||||
|
||||
def _profile_texts(memories):
|
||||
return [
|
||||
path.read_text(encoding="utf-8")
|
||||
for path in (memories / "profile").rglob("*.md")
|
||||
]
|
||||
|
||||
|
||||
def test_profile_memory_bootstraps_and_injects_profile_files(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
modified = middleware.modify_request(_request())
|
||||
system_text = _system_text(modified)
|
||||
|
||||
assert "Today's date" not in system_text
|
||||
assert "<profile_memory>" in system_text
|
||||
assert "# User profile" in system_text
|
||||
assert "/memories/profile/USER_PROFILE.md" in system_text
|
||||
assert (memories / "profile" / "SOUL.md").exists()
|
||||
assert (memories / "profile" / "USER_PROFILE.md").exists()
|
||||
assert (memories / "profile" / "RESEARCH_TASTE.md").exists()
|
||||
assert list((memories / "profile" / "projects").glob("*/PROJECT_PROFILE.md"))
|
||||
|
||||
|
||||
def test_profile_memory_uses_path_pointers_when_profiles_exceed_budget(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(
|
||||
str(memories), max_inline_profile_chars=10
|
||||
)
|
||||
modified = middleware.modify_request(_request())
|
||||
system_text = _system_text(modified)
|
||||
|
||||
assert "Profile files are available at:" in system_text
|
||||
assert "File: /memories/profile/SOUL.md" not in system_text
|
||||
assert "/memories/profile/USER_PROFILE.md" in system_text
|
||||
|
||||
|
||||
def test_profile_memory_async_path_bootstraps_and_injects(
|
||||
tmp_path, monkeypatch, run_async
|
||||
):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
|
||||
async def _handler(request):
|
||||
return request
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
modified = run_async(middleware.awrap_model_call(_request(), _handler))
|
||||
system_text = _system_text(modified)
|
||||
|
||||
assert "<profile_memory>" in system_text
|
||||
assert "/memories/profile/USER_PROFILE.md" in system_text
|
||||
assert (memories / "profile" / "USER_PROFILE.md").exists()
|
||||
|
||||
|
||||
def test_profile_memory_write_failure_uses_path_pointers(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
monkeypatch.setattr(middleware, "_write_text", lambda _path, _content: False)
|
||||
|
||||
modified = middleware.modify_request(_request())
|
||||
system_text = _system_text(modified)
|
||||
|
||||
assert "Profile files are available at:" in system_text
|
||||
assert "File: /memories/profile/SOUL.md" not in system_text
|
||||
assert "# User profile" not in system_text
|
||||
assert not (memories / "profile" / "USER_PROFILE.md").exists()
|
||||
|
||||
|
||||
def test_profile_memory_read_failure_uses_path_pointers_without_overwriting(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
|
||||
profile_dir = memories / "profile"
|
||||
profile_dir.mkdir(parents=True)
|
||||
soul_path = profile_dir / "SOUL.md"
|
||||
original_bytes = b"\xff\xfe\xfa existing profile bytes"
|
||||
soul_path.write_bytes(original_bytes)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
modified = middleware.modify_request(_request())
|
||||
system_text = _system_text(modified)
|
||||
|
||||
assert "Profile files are available at:" in system_text
|
||||
assert soul_path.read_bytes() == original_bytes
|
||||
|
||||
|
||||
def test_profile_memory_migrates_legacy_memory_once(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
memories.mkdir()
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
(memories / "MEMORY.md").write_text(
|
||||
"\n".join(
|
||||
[
|
||||
"# EvoScientist Memory",
|
||||
"",
|
||||
"## User Profile",
|
||||
"- **Name**: Alice",
|
||||
"",
|
||||
"## Research Preferences",
|
||||
"- **Primary Domain**: RL",
|
||||
"",
|
||||
"## Experiment History",
|
||||
"### [2026-01-01] Baseline",
|
||||
"- **Conclusion**: Worked",
|
||||
"",
|
||||
"## Learned Preferences",
|
||||
"- Prefers concise plans.",
|
||||
]
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
middleware.modify_request(_request())
|
||||
middleware.modify_request(_request())
|
||||
|
||||
user_profile = (memories / "profile" / "USER_PROFILE.md").read_text(
|
||||
encoding="utf-8"
|
||||
)
|
||||
research_taste = (memories / "profile" / "RESEARCH_TASTE.md").read_text(
|
||||
encoding="utf-8"
|
||||
)
|
||||
|
||||
assert user_profile.count("- **Name**: Alice") == 1
|
||||
assert user_profile.count("Prefers concise plans.") == 1
|
||||
assert user_profile.count("### Experiment History") == 1
|
||||
assert user_profile.count("- **Conclusion**: Worked") == 1
|
||||
assert research_taste.count("- **Primary Domain**: RL") == 1
|
||||
assert "Migrated from /memories/MEMORY.md" not in user_profile
|
||||
assert "Migrated from /memories/MEMORY.md" not in research_taste
|
||||
assert not (memories / "MEMORY.md").exists()
|
||||
|
||||
|
||||
def test_profile_memory_deletes_blank_legacy_memory(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
memories.mkdir()
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
legacy_path = memories / "MEMORY.md"
|
||||
legacy_path.write_text(" \n\n", encoding="utf-8")
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
middleware.modify_request(_request())
|
||||
|
||||
assert not legacy_path.exists()
|
||||
|
||||
|
||||
def test_profile_memory_uses_explicit_workspace_for_project_profile(
|
||||
tmp_path, monkeypatch
|
||||
):
|
||||
memories = tmp_path / "memories"
|
||||
global_workspace = tmp_path / "global-workspace"
|
||||
active_workspace = tmp_path / "active-workspace"
|
||||
global_workspace.mkdir()
|
||||
active_workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", global_workspace)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(
|
||||
str(memories), workspace_dir=str(active_workspace)
|
||||
)
|
||||
modified = middleware.modify_request(_request())
|
||||
system_text = _system_text(modified)
|
||||
|
||||
expected_project_id = _path_project_id(active_workspace)
|
||||
wrong_project_id = _path_project_id(global_workspace)
|
||||
|
||||
assert (
|
||||
f"/memories/profile/projects/{expected_project_id}/PROJECT_PROFILE.md"
|
||||
in system_text
|
||||
)
|
||||
assert (
|
||||
memories / "profile" / "projects" / expected_project_id / "PROJECT_PROFILE.md"
|
||||
).exists()
|
||||
assert wrong_project_id not in system_text
|
||||
assert not (
|
||||
memories / "profile" / "projects" / wrong_project_id / "PROJECT_PROFILE.md"
|
||||
).exists()
|
||||
|
||||
|
||||
def test_profile_memory_resolves_project_id_once_per_middleware(
|
||||
tmp_path, monkeypatch, run_async
|
||||
):
|
||||
memories = tmp_path / "memories"
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
calls = []
|
||||
|
||||
def _resolve_project_id(workspace_dir):
|
||||
calls.append(workspace_dir)
|
||||
return "P-cached-project"
|
||||
|
||||
monkeypatch.setattr(memory_module, "_resolve_project_id", _resolve_project_id)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(
|
||||
str(memories), workspace_dir=str(workspace), max_inline_profile_chars=10
|
||||
)
|
||||
sync_modified = middleware.modify_request(_request())
|
||||
async_modified = run_async(middleware.amodify_request(_request()))
|
||||
|
||||
assert calls == [workspace]
|
||||
assert (
|
||||
"/memories/profile/projects/P-cached-project/PROJECT_PROFILE.md"
|
||||
in _system_text(sync_modified)
|
||||
)
|
||||
assert (
|
||||
"/memories/profile/projects/P-cached-project/PROJECT_PROFILE.md"
|
||||
in _system_text(async_modified)
|
||||
)
|
||||
|
||||
|
||||
def test_profile_memory_preserves_unmapped_legacy_memory(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
memories.mkdir()
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
legacy_path = memories / "MEMORY.md"
|
||||
custom_note = "Keep this custom deployment note."
|
||||
legacy_path.write_text(
|
||||
"\n".join(
|
||||
[
|
||||
"# EvoScientist Memory",
|
||||
"",
|
||||
"## User Profile",
|
||||
"- **Name**: Alice",
|
||||
"",
|
||||
"## Custom Notes",
|
||||
custom_note,
|
||||
]
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
middleware.modify_request(_request())
|
||||
|
||||
user_profile = (memories / "profile" / "USER_PROFILE.md").read_text(
|
||||
encoding="utf-8"
|
||||
)
|
||||
assert custom_note in user_profile
|
||||
assert not legacy_path.exists()
|
||||
|
||||
|
||||
def test_profile_memory_skips_legacy_unknown_placeholders(tmp_path, monkeypatch):
|
||||
memories = tmp_path / "memories"
|
||||
memories.mkdir()
|
||||
workspace = tmp_path / "workspace"
|
||||
workspace.mkdir()
|
||||
monkeypatch.setattr(paths, "WORKSPACE_ROOT", workspace)
|
||||
(memories / "MEMORY.md").write_text(
|
||||
"\n".join(
|
||||
[
|
||||
"# EvoScientist Memory",
|
||||
"",
|
||||
"## User Profile",
|
||||
"- **Name**: (unknown)",
|
||||
"- **Role**: (unknown)",
|
||||
"",
|
||||
"## Research Preferences",
|
||||
"- **Primary Domain**: (unknown)",
|
||||
"- **Preferred Methods**: (unknown)",
|
||||
"",
|
||||
"## Experiment History",
|
||||
"(No experiments yet)",
|
||||
"",
|
||||
"## Learned Preferences",
|
||||
"- (none yet)",
|
||||
]
|
||||
),
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
middleware = memory_module.create_memory_middleware(str(memories))
|
||||
middleware.modify_request(_request())
|
||||
|
||||
migrated_profile_text = "\n".join(_profile_texts(memories))
|
||||
assert "(unknown)" not in migrated_profile_text
|
||||
assert "Imported from legacy MEMORY.md" not in migrated_profile_text
|
||||
assert not (memories / "MEMORY.md").exists()
|
||||
@@ -1,9 +1,22 @@
|
||||
"""Tests for Rich streaming display helpers."""
|
||||
|
||||
from typing import Any, cast
|
||||
|
||||
from rich.console import Console
|
||||
from rich.markdown import Markdown
|
||||
|
||||
from EvoScientist.stream.display import (
|
||||
_fix_markdown_heading_spacing,
|
||||
create_streaming_display,
|
||||
resolve_final_status_footer,
|
||||
)
|
||||
from EvoScientist.stream.state import SubAgentState
|
||||
|
||||
|
||||
def _render_text(renderable) -> str:
|
||||
console = Console(record=True, width=100, color_system=None)
|
||||
console.print(renderable)
|
||||
return console.export_text()
|
||||
|
||||
|
||||
def test_resolve_final_status_footer_hides_footer_for_interactive_cli():
|
||||
@@ -16,6 +29,416 @@ def test_resolve_final_status_footer_keeps_footer_for_noninteractive():
|
||||
assert resolve_final_status_footer(False, lambda: "footer") == "footer"
|
||||
|
||||
|
||||
def test_streaming_display_keeps_narration_visible_with_pending_memory_tool():
|
||||
"""Profile-memory reads still block, while lead-in text remains visible."""
|
||||
narration = "Here is the answer."
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "read_file",
|
||||
"args": {"path": "/memories/profile/USER_PROFILE.md"},
|
||||
}
|
||||
],
|
||||
tool_results=[],
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "Here is the answer." in rendered
|
||||
assert "Reading memory" in rendered
|
||||
assert "Running" in rendered
|
||||
assert rendered.index("Here is the answer.") < rendered.index("Reading memory")
|
||||
|
||||
|
||||
def test_streaming_display_keeps_narration_visible_while_processing_tool_result():
|
||||
"""Completed tools still block while their result is being processed."""
|
||||
narration = "Here is the answer."
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "read_file",
|
||||
"args": {"path": "/memories/profile/USER_PROFILE.md"},
|
||||
}
|
||||
],
|
||||
tool_results=[
|
||||
{
|
||||
"name": "read_file",
|
||||
"content": "# User profile\n\n- Likes concise updates.",
|
||||
}
|
||||
],
|
||||
is_processing=True,
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "Here is the answer." in rendered
|
||||
assert "Reading memory" in rendered
|
||||
assert "Analyzing results" in rendered
|
||||
assert rendered.index("Here is the answer.") < rendered.index("Reading memory")
|
||||
|
||||
|
||||
def test_streaming_display_keeps_narration_visible_with_pending_normal_tool():
|
||||
"""Ordinary tools use the same pending-tool behavior as memory reads."""
|
||||
narration = "I will inspect the files first."
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "execute",
|
||||
"args": {"command": "rg -n TODO ."},
|
||||
}
|
||||
],
|
||||
tool_results=[],
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "execute(rg -n TODO .)" in rendered
|
||||
assert "Running" in rendered
|
||||
assert "I will inspect the files first." in rendered
|
||||
assert rendered.index("I will inspect the files first.") < rendered.index(
|
||||
"execute(rg -n TODO .)"
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_keeps_narration_separate_when_answer_streams():
|
||||
"""Post-tool answers should not re-render pre-tool narration as Markdown."""
|
||||
narration = "I will inspect the files first.\n"
|
||||
answer = "The delayed check completed."
|
||||
renderable = create_streaming_display(
|
||||
response_text=f"{narration}{answer}",
|
||||
latest_text=answer,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "execute",
|
||||
"args": {"command": "python check.py"},
|
||||
}
|
||||
],
|
||||
tool_results=[
|
||||
{
|
||||
"name": "execute",
|
||||
"content": "check complete",
|
||||
}
|
||||
],
|
||||
response_markdown=Markdown("SHOULD NOT RENDER"),
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I will inspect the files first." in rendered
|
||||
assert "The delayed check completed." in rendered
|
||||
assert "SHOULD NOT RENDER" not in rendered
|
||||
assert rendered.index("I will inspect the files first.") < rendered.index(
|
||||
"execute(python check.py)"
|
||||
)
|
||||
assert rendered.index("execute(python check.py)") < rendered.index(
|
||||
"The delayed check completed."
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_keeps_narration_separate_in_final_frame():
|
||||
"""The final Rich frame should preserve narration without one concat block."""
|
||||
narration = "I will inspect the files first.\n"
|
||||
answer = "The delayed check completed."
|
||||
renderable = create_streaming_display(
|
||||
response_text=f"{narration}{answer}",
|
||||
latest_text=answer,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "execute",
|
||||
"args": {"command": "python check.py"},
|
||||
}
|
||||
],
|
||||
tool_results=[
|
||||
{
|
||||
"name": "execute",
|
||||
"content": "check complete",
|
||||
}
|
||||
],
|
||||
is_final=True,
|
||||
response_markdown=Markdown("SHOULD NOT RENDER"),
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I will inspect the files first." in rendered
|
||||
assert "The delayed check completed." in rendered
|
||||
assert "SHOULD NOT RENDER" not in rendered
|
||||
assert rendered.index("I will inspect the files first.") < rendered.index(
|
||||
"execute(python check.py)"
|
||||
)
|
||||
assert rendered.index("execute(python check.py)") < rendered.index(
|
||||
"The delayed check completed."
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_preserves_narration_for_pending_final_tool():
|
||||
"""Stopped/error final frames keep narration attached to pending tools."""
|
||||
narration = "I will inspect the files first.\n"
|
||||
answer = "[Stopped.]"
|
||||
renderable = create_streaming_display(
|
||||
response_text=f"{narration}{answer}",
|
||||
latest_text=answer,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "execute",
|
||||
"args": {"command": "sleep 30"},
|
||||
}
|
||||
],
|
||||
tool_results=[],
|
||||
is_final=True,
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I will inspect the files first." in rendered
|
||||
assert "execute(sleep 30)" in rendered
|
||||
assert "[Stopped.]" in rendered
|
||||
assert rendered.index("I will inspect the files first.") < rendered.index(
|
||||
"execute(sleep 30)"
|
||||
)
|
||||
assert rendered.index("execute(sleep 30)") < rendered.index("[Stopped.]")
|
||||
|
||||
|
||||
def test_streaming_display_interleaves_multiple_narration_segments():
|
||||
"""Multiple narrated segments should stay attached to their following tools."""
|
||||
first = "I will inspect the files first.\n"
|
||||
second = "I found one file, now I will run it.\n"
|
||||
answer = "The delayed check completed."
|
||||
renderable = create_streaming_display(
|
||||
response_text=f"{first}{second}{answer}",
|
||||
latest_text=answer,
|
||||
narrated_response_end=len(first) + len(second),
|
||||
narration_segments=[
|
||||
(0, first),
|
||||
(1, second),
|
||||
],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "tc1",
|
||||
"name": "execute",
|
||||
"args": {"command": "rg -n delayed ."},
|
||||
},
|
||||
{
|
||||
"id": "tc2",
|
||||
"name": "execute",
|
||||
"args": {"command": "python check.py"},
|
||||
},
|
||||
],
|
||||
tool_results=[
|
||||
{
|
||||
"name": "execute",
|
||||
"content": "check.py",
|
||||
},
|
||||
{
|
||||
"name": "execute",
|
||||
"content": "check complete",
|
||||
},
|
||||
],
|
||||
is_final=True,
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert rendered.index("I will inspect the files first.") < rendered.index(
|
||||
"execute(rg -n delayed .)"
|
||||
)
|
||||
assert rendered.index("execute(rg -n delayed .)") < rendered.index(
|
||||
"I found one file, now I will run it."
|
||||
)
|
||||
assert rendered.index("I found one file, now I will run it.") < rendered.index(
|
||||
"execute(python check.py)"
|
||||
)
|
||||
assert rendered.index("execute(python check.py)") < rendered.index(
|
||||
"The delayed check completed."
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_preserves_narration_for_collapsed_completed_tool():
|
||||
"""Live collapsed completed summaries keep narration from hidden tools."""
|
||||
narration = "I will check the early result.\n"
|
||||
tool_calls = [
|
||||
{
|
||||
"id": f"tc{i}",
|
||||
"name": "execute",
|
||||
"args": {"command": f"python step_{i}.py"},
|
||||
}
|
||||
for i in range(5)
|
||||
]
|
||||
tool_results = [
|
||||
{
|
||||
"name": "execute",
|
||||
"content": f"step {i} complete",
|
||||
}
|
||||
for i in range(5)
|
||||
]
|
||||
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=tool_calls,
|
||||
tool_results=tool_results,
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I will check the early result." in rendered
|
||||
assert "1 completed" in rendered
|
||||
assert "execute(python step_0.py)" not in rendered
|
||||
assert rendered.index("I will check the early result.") < rendered.index(
|
||||
"1 completed"
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_preserves_narration_for_collapsed_running_tool():
|
||||
"""Live collapsed running summaries keep narration from hidden tools."""
|
||||
narration = "I will start the long-running check.\n"
|
||||
tool_calls = [
|
||||
{
|
||||
"id": f"tc{i}",
|
||||
"name": "execute",
|
||||
"args": {"command": f"sleep {i + 1}"},
|
||||
}
|
||||
for i in range(4)
|
||||
]
|
||||
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=tool_calls,
|
||||
tool_results=[],
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I will start the long-running check." in rendered
|
||||
assert "1 more running" in rendered
|
||||
assert "execute(sleep 1)" not in rendered
|
||||
assert rendered.index("I will start the long-running check.") < rendered.index(
|
||||
"1 more running"
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_preserves_task_narration_while_subagent_runs():
|
||||
"""Narration before a task call should stay attached to the task section."""
|
||||
narration = "I'll ask a specialist to inspect this.\n"
|
||||
subagent = SubAgentState("code-agent", "inspect this")
|
||||
subagent.is_active = True
|
||||
subagent.add_tool_call("execute", {"command": "rg -n TODO ."}, "sa1")
|
||||
|
||||
renderable = create_streaming_display(
|
||||
response_text=narration,
|
||||
narrated_response_end=len(narration),
|
||||
narration_segments=[(0, narration)],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "task1",
|
||||
"name": "task",
|
||||
"args": {
|
||||
"subagent_type": "code-agent",
|
||||
"description": "inspect this",
|
||||
},
|
||||
}
|
||||
],
|
||||
tool_results=[],
|
||||
subagents=[subagent],
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert "I'll ask a specialist to inspect this." in rendered
|
||||
assert "Cooking with code-agent" in rendered
|
||||
assert "execute(rg -n TODO .)" in rendered
|
||||
assert rendered.index("I'll ask a specialist to inspect this.") < rendered.index(
|
||||
"Cooking with code-agent"
|
||||
)
|
||||
|
||||
|
||||
def test_streaming_display_orders_final_task_narration_by_tool_index():
|
||||
"""Task narration should not move after regular-tool narration in final frames."""
|
||||
first = "I'll ask a specialist to inspect this.\n"
|
||||
second = "Now I will run the result locally.\n"
|
||||
answer = "The local run passed."
|
||||
subagent = SubAgentState("code-agent", "inspect this")
|
||||
subagent.is_active = False
|
||||
subagent.add_tool_call("execute", {"command": "rg -n TODO ."}, "sa1")
|
||||
subagent.add_tool_result("execute", "todo.py", True, "sa1")
|
||||
|
||||
renderable = create_streaming_display(
|
||||
response_text=f"{first}{second}{answer}",
|
||||
latest_text=answer,
|
||||
narrated_response_end=len(first) + len(second),
|
||||
narration_segments=[
|
||||
(0, first),
|
||||
(1, second),
|
||||
],
|
||||
tool_calls=[
|
||||
{
|
||||
"id": "task1",
|
||||
"name": "task",
|
||||
"args": {
|
||||
"subagent_type": "code-agent",
|
||||
"description": "inspect this",
|
||||
},
|
||||
},
|
||||
{
|
||||
"id": "tc2",
|
||||
"name": "execute",
|
||||
"args": {"command": "python todo.py"},
|
||||
},
|
||||
],
|
||||
tool_results=[
|
||||
{
|
||||
"name": "task",
|
||||
"content": "todo.py",
|
||||
},
|
||||
{
|
||||
"name": "execute",
|
||||
"content": "passed",
|
||||
},
|
||||
],
|
||||
subagents=[subagent],
|
||||
is_final=True,
|
||||
)
|
||||
|
||||
rendered = _render_text(renderable)
|
||||
|
||||
assert rendered.index("I'll ask a specialist to inspect this.") < rendered.index(
|
||||
"Cooking with code-agent"
|
||||
)
|
||||
assert rendered.index("Cooking with code-agent") < rendered.index(
|
||||
"Now I will run the result locally."
|
||||
)
|
||||
assert rendered.index("Now I will run the result locally.") < rendered.index(
|
||||
"execute(python todo.py)"
|
||||
)
|
||||
assert rendered.index("execute(python todo.py)") < rendered.index(
|
||||
"The local run passed."
|
||||
)
|
||||
|
||||
|
||||
class TestFixMarkdownHeadingSpacing:
|
||||
"""Pure-helper tests: heading levels, idempotence, EOS / CRLF / fenced
|
||||
code. The display-copy-only contract at call sites is covered by
|
||||
@@ -105,7 +528,7 @@ class TestAssistantMessageBufferContract:
|
||||
|
||||
msg = AssistantMessage(initial_content=initial)
|
||||
fake_md = MagicMock()
|
||||
msg.query_one = MagicMock(return_value=fake_md)
|
||||
cast(Any, msg).query_one = MagicMock(return_value=fake_md)
|
||||
return msg, fake_md
|
||||
|
||||
def test_flush_markdown_does_not_mutate_buffer(self):
|
||||
|
||||
+32
-25
@@ -579,6 +579,37 @@ class TestLatestTextReset:
|
||||
# response_text still has everything
|
||||
assert state.response_text == "first segmentsecond segment"
|
||||
|
||||
def test_tool_call_marks_existing_text_as_narration(self):
|
||||
state = StreamState()
|
||||
state.handle_event({"type": "text", "content": "first segment"})
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc1", "name": "execute", "args": {}}
|
||||
)
|
||||
|
||||
assert state.narrated_response_end == len("first segment")
|
||||
assert state.narration_segments == [(0, "first segment")]
|
||||
|
||||
state.handle_event({"type": "text", "content": "second segment"})
|
||||
assert state.narrated_response_end == len("first segment")
|
||||
assert state.narration_segments == [(0, "first segment")]
|
||||
|
||||
def test_later_tool_call_extends_narrated_boundary(self):
|
||||
state = StreamState()
|
||||
state.handle_event({"type": "text", "content": "first segment"})
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc1", "name": "execute", "args": {}}
|
||||
)
|
||||
state.handle_event({"type": "text", "content": "second segment"})
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc2", "name": "execute", "args": {}}
|
||||
)
|
||||
|
||||
assert state.narrated_response_end == len("first segmentsecond segment")
|
||||
assert state.narration_segments == [
|
||||
(0, "first segment"),
|
||||
(1, "second segment"),
|
||||
]
|
||||
|
||||
|
||||
# =============================================================================
|
||||
# Name merging edge cases
|
||||
@@ -859,17 +890,6 @@ class TestHasPendingWork:
|
||||
state.handle_event({"type": "text", "content": "done"})
|
||||
assert state.has_pending_work() is False
|
||||
|
||||
def test_internal_tool_ignored(self):
|
||||
state = StreamState()
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc1", "name": "ExtractedMemory", "args": {}}
|
||||
)
|
||||
state.handle_event(
|
||||
{"type": "tool_result", "name": "ExtractedMemory", "content": "ok"}
|
||||
)
|
||||
state.is_processing = False
|
||||
assert state.has_pending_work() is False
|
||||
|
||||
|
||||
class TestVisibleToolCounts:
|
||||
"""Tests for StreamState.visible_tool_counts()."""
|
||||
@@ -893,29 +913,16 @@ class TestVisibleToolCounts:
|
||||
state.handle_event({"type": "tool_result", "name": "execute", "content": "ok"})
|
||||
assert state.visible_tool_counts() == (1, 1)
|
||||
|
||||
def test_internal_tool_excluded(self):
|
||||
state = StreamState()
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc1", "name": "ExtractedMemory", "args": {}}
|
||||
)
|
||||
assert state.visible_tool_counts() == (0, 0)
|
||||
|
||||
def test_mixed(self):
|
||||
state = StreamState()
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc1", "name": "execute", "args": {}}
|
||||
)
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc2", "name": "ExtractedMemory", "args": {}}
|
||||
)
|
||||
state.handle_event(
|
||||
{"type": "tool_call", "id": "tc3", "name": "search", "args": {}}
|
||||
)
|
||||
state.handle_event({"type": "tool_result", "name": "execute", "content": "ok"})
|
||||
state.handle_event(
|
||||
{"type": "tool_result", "name": "ExtractedMemory", "content": "ok"}
|
||||
)
|
||||
# execute done, ExtractedMemory done but invisible, search pending
|
||||
# execute done, search pending
|
||||
assert state.visible_tool_counts() == (1, 2)
|
||||
|
||||
|
||||
|
||||
+37
-10
@@ -95,13 +95,17 @@ class TestFormatToolCompact:
|
||||
result = format_tool_compact("edit_file", {"path": "f.py"})
|
||||
assert result == "edit_file(f.py)"
|
||||
|
||||
# Global memory file special display (/memories/ = global)
|
||||
# Global profile memory display (/memories/ = global)
|
||||
def test_read_file_global_memory(self):
|
||||
result = format_tool_compact("read_file", {"path": "/memories/MEMORY.md"})
|
||||
result = format_tool_compact(
|
||||
"read_file", {"path": "/memories/profile/USER_PROFILE.md"}
|
||||
)
|
||||
assert result == "Reading memory"
|
||||
|
||||
def test_read_file_global_memory_file_path_alias(self):
|
||||
result = format_tool_compact("read_file", {"file_path": "/memories/MEMORY.md"})
|
||||
result = format_tool_compact(
|
||||
"read_file", {"file_path": "/memories/profile/USER_PROFILE.md"}
|
||||
)
|
||||
assert result == "Reading memory"
|
||||
|
||||
def test_read_file_any_global_memory_file(self):
|
||||
@@ -109,14 +113,15 @@ class TestFormatToolCompact:
|
||||
assert result == "Reading memory"
|
||||
|
||||
def test_write_file_global_memory(self):
|
||||
result = format_tool_compact("write_file", {"path": "/MEMORY.md"})
|
||||
result = format_tool_compact(
|
||||
"write_file", {"path": "/memories/profile/USER_PROFILE.md"}
|
||||
)
|
||||
assert result == "Updating memory"
|
||||
# Also covers paths with /memories/ prefix
|
||||
result2 = format_tool_compact("write_file", {"path": "/memories/MEMORY.md"})
|
||||
assert result2 == "Updating memory"
|
||||
|
||||
def test_edit_file_global_memory(self):
|
||||
result = format_tool_compact("edit_file", {"path": "/memories/MEMORY.md"})
|
||||
result = format_tool_compact(
|
||||
"edit_file", {"path": "/memories/profile/USER_PROFILE.md"}
|
||||
)
|
||||
assert result == "Updating memory"
|
||||
|
||||
def test_write_edit_any_global_memory_file(self):
|
||||
@@ -144,14 +149,14 @@ class TestFormatToolCompact:
|
||||
read_result = format_tool_compact_with_result(
|
||||
"read_file",
|
||||
{},
|
||||
"# EvoScientist Memory\n\nFounder: Zachary",
|
||||
"# User profile\n\nFounder: Zachary",
|
||||
)
|
||||
assert read_result == "Reading memory"
|
||||
|
||||
edit_result = format_tool_compact_with_result(
|
||||
"edit_file",
|
||||
{},
|
||||
"Successfully replaced 1 instance(s) of the string in '/memories/MEMORY.md'",
|
||||
"Successfully replaced 1 instance(s) of the string in '/memories/profile/USER_PROFILE.md'",
|
||||
)
|
||||
assert edit_result == "Updating memory"
|
||||
|
||||
@@ -162,6 +167,28 @@ class TestFormatToolCompact:
|
||||
)
|
||||
assert write_result == "Updating memory"
|
||||
|
||||
def test_profile_memory_inference_uses_profile_template_headings(self, monkeypatch):
|
||||
from EvoScientist.middleware import memory
|
||||
from EvoScientist.stream import utils
|
||||
|
||||
monkeypatch.setitem(
|
||||
memory.PROFILE_TEMPLATES,
|
||||
"/profile/CUSTOM.md",
|
||||
"# Custom profile\n\n## Notes\n",
|
||||
)
|
||||
utils._profile_memory_headings.cache_clear()
|
||||
|
||||
try:
|
||||
result = format_tool_compact_with_result(
|
||||
"read_file",
|
||||
{},
|
||||
"# Custom profile\n\n- remembered",
|
||||
)
|
||||
finally:
|
||||
utils._profile_memory_headings.cache_clear()
|
||||
|
||||
assert result == "Reading memory"
|
||||
|
||||
def test_project_memory_result_not_special(self):
|
||||
result = format_tool_compact_with_result(
|
||||
"write_file",
|
||||
|
||||
+15
-12
@@ -170,6 +170,20 @@ class TestStoppedResponseText(unittest.TestCase):
|
||||
assert current == "partial\n[Stopped.]"
|
||||
assert final_text == "partial\n[Stopped.]"
|
||||
|
||||
def test_strips_trailing_placeholder_ellipsis(self):
|
||||
from EvoScientist.cli.tui_interactive import (
|
||||
_strip_trailing_placeholder_ellipsis,
|
||||
)
|
||||
|
||||
assert (
|
||||
_strip_trailing_placeholder_ellipsis("final answer\n...") == "final answer"
|
||||
)
|
||||
assert _strip_trailing_placeholder_ellipsis("...") == ""
|
||||
assert (
|
||||
_strip_trailing_placeholder_ellipsis("final answer\n...\n...")
|
||||
== "final answer"
|
||||
)
|
||||
|
||||
|
||||
@unittest.skipUnless(_has_textual, "textual not installed")
|
||||
class TestAssistantMessage(unittest.TestCase):
|
||||
@@ -225,9 +239,7 @@ class TestToolCallWidget(unittest.TestCase):
|
||||
from EvoScientist.cli.widgets.tool_call_widget import ToolCallWidget
|
||||
|
||||
w = ToolCallWidget("edit_file", {}, "mem-1")
|
||||
w._result_content = (
|
||||
"Successfully replaced 1 instance(s) of the string in '/memories/MEMORY.md'"
|
||||
)
|
||||
w._result_content = "Successfully replaced 1 instance(s) of the string in '/memories/profile/USER_PROFILE.md'"
|
||||
|
||||
class _Header:
|
||||
def __init__(self) -> None:
|
||||
@@ -510,15 +522,6 @@ class TestIsFinalResponse(unittest.TestCase):
|
||||
state.subagents = [sa]
|
||||
assert _is_final_response(state) is True
|
||||
|
||||
def test_internal_tools_ignored(self):
|
||||
from EvoScientist.cli.tui_interactive import _is_final_response
|
||||
from EvoScientist.stream.state import StreamState
|
||||
|
||||
state = StreamState()
|
||||
state.tool_calls = [{"name": "ExtractedMemory", "args": {}}]
|
||||
# No result for internal tool -- should still be considered final
|
||||
assert _is_final_response(state) is True
|
||||
|
||||
def test_processing_not_final(self):
|
||||
from EvoScientist.cli.tui_interactive import _is_final_response
|
||||
from EvoScientist.stream.state import StreamState
|
||||
|
||||
Reference in New Issue
Block a user