refactor: enhance summarization widget with streaming support and text accumulation
This commit is contained in:
@@ -712,9 +712,12 @@ def run_textual_interactive(
|
||||
if summarization_w is None:
|
||||
summarization_w = SummarizationWidget()
|
||||
await container.mount(summarization_w)
|
||||
summarization_w.set_content(content)
|
||||
summarization_w.append_text(content)
|
||||
|
||||
elif event_type == "text":
|
||||
# Finalize summarization widget when regular text resumes
|
||||
if summarization_w is not None and summarization_w._is_active:
|
||||
summarization_w.finalize()
|
||||
if thinking_w is not None and thinking_w._is_active:
|
||||
thinking_w.finalize()
|
||||
# Clear processing indicator
|
||||
@@ -982,6 +985,9 @@ def run_textual_interactive(
|
||||
)
|
||||
|
||||
elif event_type == "done":
|
||||
# Finalize summarization if still active
|
||||
if summarization_w is not None and summarization_w._is_active:
|
||||
summarization_w.finalize()
|
||||
# Clean up transient indicators
|
||||
await _remove_w(narration_w)
|
||||
narration_w = None
|
||||
|
||||
@@ -4,6 +4,10 @@ Renders a Rich Panel showing that LangGraph's summarization middleware
|
||||
has compressed older conversation history. Yellow/amber border to
|
||||
distinguish from the blue thinking panel. Default collapsed; click to
|
||||
expand/collapse.
|
||||
|
||||
Supports streaming: text arrives incrementally via ``append_text()``
|
||||
while the widget shows a live "Summarizing..." indicator, then switches
|
||||
to a collapsed preview once ``finalize()`` is called.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -21,14 +25,16 @@ _MAX_EXPANDED_CHARS = 3000
|
||||
class SummarizationWidget(Static):
|
||||
"""Collapsible panel showing context summarization.
|
||||
|
||||
Unlike ThinkingWidget this is not streamed — the full text arrives
|
||||
in a single event. Defaults to collapsed with a one-line preview.
|
||||
Streams text via ``append_text()`` (shows live spinner while active).
|
||||
Defaults to collapsed after streaming ends; click to expand/collapse.
|
||||
|
||||
Usage::
|
||||
|
||||
w = SummarizationWidget()
|
||||
await container.mount(w)
|
||||
w.set_content("The conversation covered ...")
|
||||
w.append_text("The conversation ")
|
||||
w.append_text("covered ...")
|
||||
w.finalize() # stop spinner, collapse
|
||||
"""
|
||||
|
||||
DEFAULT_CSS = """
|
||||
@@ -42,6 +48,7 @@ class SummarizationWidget(Static):
|
||||
super().__init__("")
|
||||
self._content = ""
|
||||
self._collapsed = True
|
||||
self._is_active = True # still receiving chunks
|
||||
|
||||
def _char_count_label(self) -> str:
|
||||
n = len(self._content)
|
||||
@@ -51,10 +58,27 @@ class SummarizationWidget(Static):
|
||||
|
||||
def _refresh_display(self) -> None:
|
||||
if not self._content:
|
||||
self.update("")
|
||||
if self._is_active:
|
||||
self.update(
|
||||
Panel(
|
||||
Text("Summarizing...", style="dim italic"),
|
||||
title="Context Summarizing",
|
||||
border_style="#f59e0b",
|
||||
padding=(0, 1),
|
||||
)
|
||||
)
|
||||
else:
|
||||
self.update("")
|
||||
return
|
||||
|
||||
if self._collapsed:
|
||||
if self._is_active:
|
||||
# While streaming: show latest content tail (like thinking widget)
|
||||
title = "Context Summarizing..."
|
||||
tail = self._content.rstrip()
|
||||
if len(tail) > 200:
|
||||
tail = tail[-200:]
|
||||
body = Text(tail, style="dim italic")
|
||||
elif self._collapsed:
|
||||
title = f"Context Summarized ({self._char_count_label()})"
|
||||
first_line = self._content.strip().split("\n")[0].strip()
|
||||
if len(first_line) > _MAX_COLLAPSED_CHARS:
|
||||
@@ -80,13 +104,25 @@ class SummarizationWidget(Static):
|
||||
Panel(body, title=title, border_style="#f59e0b", padding=(0, 1))
|
||||
)
|
||||
|
||||
def append_text(self, text: str) -> None:
|
||||
"""Append a chunk of summarization text (streaming)."""
|
||||
self._content += text
|
||||
self._refresh_display()
|
||||
|
||||
def finalize(self) -> None:
|
||||
"""Mark streaming as complete — switch to collapsed preview."""
|
||||
self._is_active = False
|
||||
self._collapsed = True
|
||||
self._refresh_display()
|
||||
|
||||
def set_content(self, text: str) -> None:
|
||||
"""Set the summarization text and refresh display."""
|
||||
"""Set the full summarization text at once (non-streaming fallback)."""
|
||||
self._content = text
|
||||
self._is_active = False
|
||||
self._refresh_display()
|
||||
|
||||
def on_click(self, event: Click) -> None:
|
||||
"""Toggle collapsed/expanded state."""
|
||||
if self._content:
|
||||
if self._content and not self._is_active:
|
||||
self._collapsed = not self._collapsed
|
||||
self._refresh_display()
|
||||
|
||||
@@ -416,11 +416,13 @@ def create_streaming_display(
|
||||
# Summarization panel (context was compressed by LangGraph middleware)
|
||||
if summarization_text:
|
||||
summary_display = summarization_text.rstrip()
|
||||
if len(summary_display) > 300:
|
||||
n = len(summary_display)
|
||||
char_label = f"{n / 1000:.1f}k chars" if n >= 1000 else f"{n:,} chars"
|
||||
if n > 300:
|
||||
summary_display = summary_display[:300] + " ..."
|
||||
elements.append(Panel(
|
||||
Text(summary_display, style="dim italic"),
|
||||
title="Context Summarized",
|
||||
title=f"Context Summarized ({char_label})",
|
||||
border_style="#f59e0b",
|
||||
padding=(0, 1),
|
||||
))
|
||||
|
||||
@@ -63,6 +63,30 @@ def _extract_tool_content(msg) -> tuple[str, bool]:
|
||||
return str(content), False
|
||||
|
||||
|
||||
def _extract_summarization_text(msg: Any) -> str:
|
||||
"""Extract plain text from a summarization chunk.
|
||||
|
||||
The summarization LLM streams ``AIMessageChunk`` objects whose
|
||||
``content`` may be a plain string **or** a list of content blocks
|
||||
(e.g. ``[{'type': 'text', 'text': '...', 'index': 1}]``) depending
|
||||
on the provider. This helper normalises both forms to a plain string.
|
||||
"""
|
||||
if not hasattr(msg, "content"):
|
||||
return ""
|
||||
content = msg.content
|
||||
if isinstance(content, str):
|
||||
return content
|
||||
if isinstance(content, list):
|
||||
parts: list[str] = []
|
||||
for block in content:
|
||||
if isinstance(block, dict) and block.get("type") == "text":
|
||||
parts.append(block.get("text", ""))
|
||||
elif isinstance(block, str):
|
||||
parts.append(block)
|
||||
return "".join(parts)
|
||||
return ""
|
||||
|
||||
|
||||
async def stream_agent_events(
|
||||
agent: Any,
|
||||
message: Any,
|
||||
@@ -295,6 +319,8 @@ async def stream_agent_events(
|
||||
# HITL resume: Command object passed directly to agent
|
||||
astream_input = message
|
||||
|
||||
_summarization_in_progress = False
|
||||
|
||||
try:
|
||||
async for chunk in agent.astream(
|
||||
astream_input,
|
||||
@@ -381,13 +407,15 @@ async def stream_agent_events(
|
||||
else:
|
||||
msg = data
|
||||
|
||||
# Emit summarization middleware messages as a dedicated event
|
||||
# Accumulate summarization middleware chunks and emit text incrementally.
|
||||
# The summarization LLM streams AIMessageChunks; content may be a
|
||||
# plain string or a list of content blocks (provider-dependent).
|
||||
if isinstance(metadata, dict) and metadata.get("lc_source") == "summarization":
|
||||
content = ""
|
||||
if hasattr(msg, "content"):
|
||||
content = msg.content if isinstance(msg.content, str) else str(msg.content)
|
||||
if content:
|
||||
yield emitter.summarization(content).data
|
||||
if not _summarization_in_progress:
|
||||
_summarization_in_progress = True
|
||||
chunk_text = _extract_summarization_text(msg)
|
||||
if chunk_text:
|
||||
yield emitter.summarization(chunk_text).data
|
||||
continue
|
||||
|
||||
subagent = _get_subagent_name(namespace, metadata)
|
||||
|
||||
@@ -266,7 +266,7 @@ class StreamState:
|
||||
self.pending_ask_user = event
|
||||
|
||||
elif event_type == "summarization":
|
||||
self.summarization_text = event.get("content", "")
|
||||
self.summarization_text += event.get("content", "")
|
||||
|
||||
elif event_type == "usage_stats":
|
||||
self.total_input_tokens += event.get("input_tokens", 0)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""Tests for the summarization event pipeline and display widgets."""
|
||||
|
||||
from EvoScientist.stream.emitter import StreamEventEmitter
|
||||
from EvoScientist.stream.events import _extract_summarization_text
|
||||
from EvoScientist.stream.state import StreamState
|
||||
|
||||
|
||||
@@ -44,12 +45,12 @@ class TestSummarizationState:
|
||||
assert etype == "summarization"
|
||||
assert state.summarization_text == "summary"
|
||||
|
||||
def test_overwrites_previous(self):
|
||||
"""Each summarization replaces (not appends) the previous text."""
|
||||
def test_accumulates_chunks(self):
|
||||
"""Summarization chunks are accumulated (streaming)."""
|
||||
state = StreamState()
|
||||
state.handle_event({"type": "summarization", "content": "first"})
|
||||
state.handle_event({"type": "summarization", "content": "second"})
|
||||
assert state.summarization_text == "second"
|
||||
assert state.summarization_text == "firstsecond"
|
||||
|
||||
def test_get_display_args_includes_field(self):
|
||||
state = StreamState()
|
||||
@@ -158,6 +159,24 @@ class TestSummarizationWidget:
|
||||
w._content = "x" * 2500
|
||||
assert w._char_count_label() == "2.5k chars"
|
||||
|
||||
def test_append_text(self):
|
||||
from EvoScientist.cli.widgets.summarization_widget import SummarizationWidget
|
||||
|
||||
w = SummarizationWidget()
|
||||
w.append_text("hello ")
|
||||
w.append_text("world")
|
||||
assert w._content == "hello world"
|
||||
assert w._is_active is True
|
||||
|
||||
def test_finalize(self):
|
||||
from EvoScientist.cli.widgets.summarization_widget import SummarizationWidget
|
||||
|
||||
w = SummarizationWidget()
|
||||
w.append_text("some summary")
|
||||
w.finalize()
|
||||
assert w._is_active is False
|
||||
assert w._collapsed is True
|
||||
|
||||
def test_toggle_collapsed(self):
|
||||
from EvoScientist.cli.widgets.summarization_widget import SummarizationWidget
|
||||
|
||||
@@ -167,3 +186,48 @@ class TestSummarizationWidget:
|
||||
assert w._collapsed is False
|
||||
w._collapsed = not w._collapsed
|
||||
assert w._collapsed is True
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _extract_summarization_text helper
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
class TestExtractSummarizationText:
|
||||
"""Content extraction from summarization chunks."""
|
||||
|
||||
def test_string_content(self):
|
||||
class Msg:
|
||||
content = "hello world"
|
||||
assert _extract_summarization_text(Msg()) == "hello world"
|
||||
|
||||
def test_content_blocks(self):
|
||||
class Msg:
|
||||
content = [{"type": "text", "text": "part1"}, {"type": "text", "text": "part2"}]
|
||||
assert _extract_summarization_text(Msg()) == "part1part2"
|
||||
|
||||
def test_content_blocks_with_index(self):
|
||||
"""Content blocks may include 'index' field — should still extract text."""
|
||||
class Msg:
|
||||
content = [{"type": "text", "text": " vs", "index": 1}]
|
||||
assert _extract_summarization_text(Msg()) == " vs"
|
||||
|
||||
def test_empty_list(self):
|
||||
class Msg:
|
||||
content = []
|
||||
assert _extract_summarization_text(Msg()) == ""
|
||||
|
||||
def test_no_content_attr(self):
|
||||
class Msg:
|
||||
pass
|
||||
assert _extract_summarization_text(Msg()) == ""
|
||||
|
||||
def test_mixed_block_types(self):
|
||||
class Msg:
|
||||
content = [{"type": "text", "text": "hello"}, {"type": "image", "url": "..."}]
|
||||
assert _extract_summarization_text(Msg()) == "hello"
|
||||
|
||||
def test_string_blocks_in_list(self):
|
||||
class Msg:
|
||||
content = ["hello", "world"]
|
||||
assert _extract_summarization_text(Msg()) == "helloworld"
|
||||
|
||||
Reference in New Issue
Block a user