Skip to content

Commit 8c1cdcc

Browse files
author
Ferrum
committed
refactor(console): decouple chat title refresh from memory abstraction layer (agentscope-ai#7032)
Remove title_refresh_callback from BaseMemoryManager and its concrete subclasses (ADBPGMemoryManager, ReMeLightMemoryManager). MemoryMiddleware no longer holds or invokes the callback. Introduce ChatTitleRefreshMiddleware backed by ChatTitleRefreshService. The middleware observes conversation lifecycle via on_reply hook, reads the _auto_memory_flushed context variable set by MemoryMiddleware after a successful auto-memory flush, and delegates title generation + compare-and-set persistence to the chat service. Registered from runtime/builder.py where both chat services and middleware composition are available, keeping memory backends independent of chat UI behavior.
1 parent 2db7cd5 commit 8c1cdcc

7 files changed

Lines changed: 320 additions & 80 deletions

File tree

src/qwenpaw/agents/memory/adbpg_memory_manager.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,12 +43,10 @@ def __init__(
4343
self,
4444
working_dir: str,
4545
agent_id: str,
46-
title_refresh_callback=None,
4746
) -> None:
4847
super().__init__(
4948
working_dir=working_dir,
5049
agent_id=agent_id,
51-
title_refresh_callback=title_refresh_callback,
5250
)
5351
self._adbpg_config = None
5452
self._client: ADBPGMemoryClient | None = None

src/qwenpaw/agents/memory/base_memory_manager.py

Lines changed: 0 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -53,11 +53,9 @@ def __init__(
5353
self,
5454
working_dir: str,
5555
agent_id: str,
56-
title_refresh_callback: Callable[..., Awaitable[None]] | None = None,
5756
):
5857
self.working_dir: str = working_dir
5958
self.agent_id: str = agent_id
60-
self.title_refresh_callback = title_refresh_callback
6159
self._summary_task_info: dict[str, dict[str, Any]] = {}
6260
self._task_counter: int = 0
6361
self._task_queue: asyncio.Queue[
@@ -110,7 +108,6 @@ def build_middlewares(self) -> list[MiddlewareBase]:
110108
return [
111109
MemoryMiddleware(
112110
memory_manager=self,
113-
title_refresh_callback=self.title_refresh_callback,
114111
),
115112
]
116113

src/qwenpaw/agents/memory/reme_light_memory_manager.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -135,12 +135,10 @@ def __init__(
135135
self,
136136
working_dir: str,
137137
agent_id: str,
138-
title_refresh_callback=None,
139138
):
140139
super().__init__(
141140
working_dir=working_dir,
142141
agent_id=agent_id,
143-
title_refresh_callback=title_refresh_callback,
144142
)
145143
self._reme: "ReMe | None" = None
146144
self._reindex_lock = asyncio.Lock()

src/qwenpaw/agents/middlewares.py

Lines changed: 81 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,14 @@
4848

4949
logger = logging.getLogger(__name__)
5050
MAX_AUTO_MEMORY_TURN_MARKERS = 1000
51+
52+
# Set by MemoryMiddleware._flush_auto_memory() to signal that an auto-memory
53+
# flush just completed successfully. ChatTitleRefreshMiddleware reads this
54+
# in its on_reply hook to decide whether a title refresh is warranted.
55+
_auto_memory_flushed_var: ContextVar[bool] = ContextVar(
56+
"_auto_memory_flushed",
57+
default=False,
58+
)
5159
AUTO_MEMORY_TURN_STATE_KEY = "qwenpaw_auto_memory_turn_state"
5260
_AUTOMATION_MEMORY_SKIP_SOURCES = frozenset({"cron", "heartbeat"})
5361
_TOOL_RESULT_METADATA_KEY = "qwenpaw_tool_result_metadata"
@@ -115,10 +123,8 @@ def __init__(
115123
self,
116124
*,
117125
memory_manager: Any,
118-
title_refresh_callback: Callable[..., Awaitable[None]] | None = None,
119126
) -> None:
120127
self._memory_manager = memory_manager
121-
self._title_refresh_callback = title_refresh_callback
122128

123129
async def on_system_prompt(
124130
self,
@@ -364,24 +370,9 @@ async def _flush_auto_memory(
364370
snapshots = turn_state["snapshots"]
365371
for marker in submitted:
366372
snapshots.pop(marker, None)
367-
# Best-effort chat title refresh: after a successful auto-memory
368-
# flush the conversation has moved on, so regenerate the session
369-
# title from this recent slice. Spawned as a background task so it
370-
# can never block or break the reply path.
371-
callback = self._title_refresh_callback
372-
if callback is not None:
373-
try:
374-
asyncio.create_task(
375-
callback(
376-
agent,
377-
messages,
378-
session_id=self._agent_session_id(agent),
379-
),
380-
)
381-
except Exception:
382-
logger.exception(
383-
"MemoryMiddleware title refresh scheduling failed",
384-
)
373+
# Signal ChatTitleRefreshMiddleware that a successful auto-memory
374+
# flush just completed.
375+
_auto_memory_flushed_var.set(True)
385376

386377
def _discard_unresolved_pending_markers(
387378
self,
@@ -1021,3 +1012,73 @@ async def on_acting(
10211012
],
10221013
},
10231014
)
1015+
1016+
1017+
class ChatTitleRefreshMiddleware(MiddlewareBase):
1018+
"""Observe conversation lifecycle and refresh chat titles after auto-memory flush.
1019+
1020+
This middleware is the sole consumer of the ``_auto_memory_flushed``
1021+
context variable. When :meth:`on_reply` detects that an auto-memory
1022+
flush just completed, it delegates title generation and persistence to
1023+
the injected :class:`~app.chats.title_refresh_service.ChatTitleRefreshService`.
1024+
1025+
Registered from the application/runtime assembly layer (see
1026+
``runtime/builder.py``) so that both chat services and middleware
1027+
composition are available.
1028+
"""
1029+
1030+
def __init__(
1031+
self,
1032+
*,
1033+
service: Any,
1034+
agent_id: str,
1035+
) -> None:
1036+
self._service = service
1037+
self._agent_id = agent_id
1038+
1039+
async def on_reply(
1040+
self,
1041+
agent: "Agent",
1042+
request: Any,
1043+
next_handler: Any,
1044+
) -> AsyncGenerator[Any, Any]:
1045+
"""Wrap each assistant reply turn.
1046+
1047+
After the inner handler yields (i.e. the model has produced a full
1048+
response), check whether auto-memory was flushed this turn. If so,
1049+
trigger a title refresh via the dedicated service.
1050+
"""
1051+
# Run the inner handler to completion first.
1052+
events: list[Any] = []
1053+
async for event in next_handler():
1054+
events.append(event)
1055+
1056+
# Only proceed if auto-memory was flushed this turn.
1057+
if not _auto_memory_flushed_var.get():
1058+
for ev in events:
1059+
yield ev
1060+
return
1061+
1062+
# Reset the signal for subsequent turns.
1063+
_auto_memory_flushed_var.set(False)
1064+
1065+
# Gather recent messages from the conversation context.
1066+
try:
1067+
from agentscope.message import Msg
1068+
1069+
context = list(agent.state.context)
1070+
recent = [
1071+
m for m in context if isinstance(m, Msg) and getattr(m, "role", "") == "assistant"
1072+
]
1073+
# Take the last few assistant messages as input for title generation.
1074+
recent = recent[-3:] if len(recent) >= 3 else recent
1075+
if recent:
1076+
await self._service.refresh(
1077+
session_id=self._agent_session_id(agent),
1078+
recent_messages=recent,
1079+
)
1080+
except Exception:
1081+
logger.exception("ChatTitleRefreshMiddleware.on_reply failed")
1082+
1083+
for ev in events:
1084+
yield ev
Lines changed: 208 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,208 @@
1+
# -*- coding: utf-8 -*-
2+
"""Chat title refresh service.
3+
4+
Responsible for generating a new title from recent messages and
5+
persisting it via the chat manager's compare-and-set mechanism.
6+
This is a pure chat-layer concern — no knowledge of memory backends.
7+
"""
8+
9+
from __future__ import annotations
10+
11+
import asyncio
12+
import logging
13+
from typing import TYPE_CHECKING, Any
14+
15+
if TYPE_CHECKING:
16+
from ...agents.model_factory import _ModelAndFormatter
17+
from agentscope.message import Msg
18+
19+
logger = logging.getLogger(__name__)
20+
21+
22+
class ChatTitleRefreshService:
23+
"""Generate chat titles and persist them via compare-and-set.
24+
25+
Public API::
26+
27+
async def refresh(
28+
session_id: str,
29+
recent_messages: list[Msg],
30+
) -> None
31+
32+
All failures are logged and swallowed so title refresh never breaks
33+
the request path.
34+
"""
35+
36+
def __init__(
37+
self,
38+
chat_manager: Any,
39+
agent_id: str,
40+
) -> None:
41+
self._chat_manager = chat_manager
42+
self._agent_id = agent_id
43+
44+
async def refresh(
45+
self,
46+
*,
47+
session_id: str,
48+
recent_messages: list[Any],
49+
) -> None:
50+
"""Re-generate a chat title from the recent conversation slice.
51+
52+
Called after each auto-memory flush (see ``ChatTitleRefreshMiddleware``).
53+
The chat is located by runtime session id, the recent messages are fed
54+
to the LLM for a fresh title, and the name is updated compare-and-set
55+
so a user-chosen name is never clobbered.
56+
"""
57+
if not recent_messages:
58+
return
59+
60+
from ...config.config import load_agent_config
61+
from ...exceptions import AppBaseException
62+
63+
try:
64+
cfg = load_agent_config(self._agent_id).running
65+
except (ValueError, AppBaseException) as exc:
66+
logger.info("Auto title refresh skipped: config unavailable (%s)", exc)
67+
return
68+
69+
title_cfg = cfg.auto_title_config
70+
if not title_cfg.enabled or not title_cfg.refresh_on_auto_memory:
71+
logger.info(
72+
"Auto title refresh skipped: refresh_on_auto_memory disabled "
73+
"(enabled=%s refresh=%s)",
74+
title_cfg.enabled,
75+
title_cfg.refresh_on_auto_memory,
76+
)
77+
return
78+
79+
chat = await self._chat_manager.find_chat_by_session_id(session_id)
80+
if chat is None:
81+
logger.info(
82+
"Auto title refresh skipped: no chat for session %s",
83+
session_id,
84+
)
85+
return
86+
87+
chat_id = chat.id
88+
89+
transcript = _messages_to_text(recent_messages)
90+
if not transcript:
91+
await self._record(chat_id, ok=False, reason="empty transcript")
92+
return
93+
94+
try:
95+
from ...agents.model_factory import create_model_and_formatter
96+
from ...utils.model_response import consume_model_response
97+
from ..title_generator import REFRESH_TITLE_PROMPT, _clean_title
98+
from agentscope.message import Msg, TextBlock
99+
100+
try:
101+
model, _ = create_model_and_formatter(
102+
agent_id=self._agent_id,
103+
)
104+
except (ValueError, AppBaseException) as exc:
105+
logger.info(
106+
"Auto title refresh skipped: no model available for chat %s (%s)",
107+
chat_id,
108+
exc,
109+
)
110+
await self._record(chat_id, ok=False, reason=f"no model: {exc}")
111+
return
112+
113+
messages = [
114+
Msg(
115+
name="system",
116+
role="system",
117+
content=[TextBlock(type="text", text=REFRESH_TITLE_PROMPT)],
118+
),
119+
Msg(
120+
name="user",
121+
role="user",
122+
content=[TextBlock(type="text", text=transcript)],
123+
),
124+
]
125+
126+
raw_title = await asyncio.wait_for(
127+
consume_model_response(model, messages),
128+
timeout=title_cfg.timeout_seconds,
129+
)
130+
except Exception:
131+
logger.exception(
132+
"Auto title refresh LLM failed for chat %s",
133+
chat_id,
134+
)
135+
await self._record(chat_id, ok=False, reason="LLM failed")
136+
return
137+
138+
title = _clean_title(raw_title)
139+
if not title:
140+
logger.info(
141+
"Auto title refresh produced empty output for chat %s",
142+
chat_id,
143+
)
144+
await self._record(chat_id, ok=False, reason="empty LLM output")
145+
return
146+
147+
# Compare-and-set: expected name is the last title we set (or the
148+
# current name for chats created before this feature shipped). If the
149+
# user renamed the chat manually, the name no longer matches and the
150+
# update is skipped.
151+
expected_name = chat.meta.get("auto_title_last") or chat.name
152+
updated = await self._chat_manager.set_auto_title(
153+
chat.id,
154+
title,
155+
expected_name=expected_name,
156+
)
157+
if updated is None:
158+
logger.info(
159+
"Auto title refresh skipped: chat %s renamed manually",
160+
chat.id,
161+
)
162+
await self._record(chat_id, ok=False, reason="renamed manually")
163+
return
164+
logger.info(
165+
"Auto-refreshed chat %s title to %r (session %s)",
166+
chat.id,
167+
title,
168+
session_id,
169+
)
170+
await self._record(chat_id, ok=True, reason="ok", title=title)
171+
172+
async def _record(
173+
self,
174+
chat_id: str,
175+
*,
176+
ok: bool,
177+
reason: str = "",
178+
title: str = "",
179+
) -> None:
180+
try:
181+
await self._chat_manager.record_auto_title_refresh(
182+
chat_id,
183+
ok=ok,
184+
reason=reason,
185+
title=title,
186+
)
187+
except Exception:
188+
logger.exception(
189+
"Auto title refresh state record failed for chat %s",
190+
chat_id,
191+
)
192+
193+
194+
def _messages_to_text(messages: list[Any]) -> str:
195+
"""Convert a list of message objects into a plain-text transcript."""
196+
parts: list[str] = []
197+
for msg in messages:
198+
role = getattr(msg, "role", "")
199+
content = getattr(msg, "content", "")
200+
if isinstance(content, list):
201+
# AgentScope 2.0 Msg — extract text blocks
202+
for block in content:
203+
text = getattr(block, "text", "") or ""
204+
if text:
205+
parts.append(f"[{role}] {text}")
206+
elif content:
207+
parts.append(f"[{role}] {content}")
208+
return "\n".join(parts)

0 commit comments

Comments
 (0)