mirror of
https://github.com/agentscope-ai/ReMe.git
synced 2026-10-05 02:41:43 +00:00
refactor(plugins): move BaseAgenticAnswerStep out of the core benchmark steps package into the lme/beam plugins (#576)
* refactor(plugins): move BaseAgenticAnswerStep into lme/beam plugins Move the shared benchmark answer base class from the core reme.steps.benchmark package into each plugin's own src tree (reme_lme.base_agentic_answer / reme_beam.base_agentic_answer) with absolute imports, drop the core benchmark steps package, and update plugin READMEs accordingly. * refactor(plugins): simplify benchmark answer steps * docs(plugins): clarify benchmark answer step migration --------- Co-authored-by: jinli.yl <jinli.yl@alibaba-inc.com>
This commit is contained in:
parent
30227e509f
commit
67936d5a43
12 changed files with 213 additions and 169 deletions
|
|
@ -47,7 +47,7 @@ and concise documentation together.
|
|||
- `reme/components/`: agent wrappers, model adapters, stores, catalogs, graphs, indexes, clients, tokenizers, and
|
||||
outbound proxies.
|
||||
- `reme/steps/`: registered job steps grouped by common, file I/O, index, evolve, cookbook, and transfer
|
||||
concerns, plus shared benchmark base classes under `benchmark/`.
|
||||
concerns.
|
||||
- `reme/utils/`: shared utilities, including service discovery, logging, web-static resolution, session I/O, token
|
||||
accounting, and wikilink handling.
|
||||
- `tests/unit/`: primary fast, isolated validation suite.
|
||||
|
|
|
|||
|
|
@ -77,7 +77,6 @@ reme/
|
|||
steps/
|
||||
base_step.py # BaseStep, Ref, dispatch_steps
|
||||
common/ # version, help, health_check, status, chat
|
||||
benchmark/ # LongMemEval / BEAM evaluation steps
|
||||
file_io/ # read/write/edit/delete/move/frontmatter/daily
|
||||
index/ # watch/init/update/search/traverse
|
||||
evolve/ # auto_memory, auto_resource, auto_dream, proactive
|
||||
|
|
|
|||
|
|
@ -72,7 +72,6 @@ reme/
|
|||
steps/
|
||||
base_step.py # BaseStep、Ref、dispatch_steps
|
||||
common/ # version、help、health_check、status、chat
|
||||
benchmark/ # LongMemEval / BEAM 评测步骤
|
||||
file_io/ # read/write/edit/delete/move/frontmatter/daily
|
||||
index/ # watch/init/update/search/traverse
|
||||
evolve/ # auto_memory、auto_resource、auto_dream、proactive
|
||||
|
|
|
|||
|
|
@ -31,10 +31,12 @@ The existing `auto_memory`, `agentic_answer`, `answer_judge`, `bench` and `judge
|
|||
names and model environment variables are unchanged. Explicit application/CLI overrides
|
||||
still take precedence. Installing this plugin does not start an evaluation.
|
||||
|
||||
The shared answer base class lives in `reme.steps.benchmark.base_agentic_answer`.
|
||||
The old core-owned `reme.steps.benchmark.beam` Python import path is removed.
|
||||
Custom Python callers should import memory, search and answer Steps from `reme_beam`, and install
|
||||
`beam-judge` before importing the judge Step from `judge_beam`. After uninstalling,
|
||||
Applications and CLI services must omit the plugin until it is installed again.
|
||||
`BeamAgenticAnswerStep` now implements the answer behavior directly. The former
|
||||
`reme.steps.benchmark.BaseAgenticAnswerStep` import is gone; custom subclasses can
|
||||
extend `reme_beam.BeamAgenticAnswerStep` for BEAM behavior, or implement their own
|
||||
Step using `reme.steps.base_step.BaseStep`.
|
||||
Uninstallation never removes datasets, workspaces or results.
|
||||
Restart an existing service after changing plugins.
|
||||
|
|
|
|||
|
|
@ -28,8 +28,11 @@ editable 安装会注册 `beam` entry point,并让源码修改立即生效。r
|
|||
均保持关闭。原有 `auto_memory`、`agentic_answer`、`answer_judge`、`bench`、`judge` 名称及模型环境变量
|
||||
保持不变,显式应用参数和 CLI 覆盖仍优先。安装或启用插件不会自动开始评测。
|
||||
|
||||
共享回答基类位于 `reme.steps.benchmark.base_agentic_answer`。
|
||||
原 `reme.steps.benchmark.beam` Python 导入路径已移除。自定义 Python 调用应从 `reme_beam`
|
||||
导入记忆、搜索和回答 Step;安装 `beam-judge` 后再从 `judge_beam` 导入评判 Step。
|
||||
自定义 Python 调用应从 `reme_beam` 导入记忆、搜索和回答 Step;安装 `beam-judge` 后再从
|
||||
`judge_beam` 导入评判 Step。
|
||||
`BeamAgenticAnswerStep` 现在直接实现回答逻辑。原来的
|
||||
`reme.steps.benchmark.BaseAgenticAnswerStep` 导入路径已移除;自定义子类可继承
|
||||
`reme_beam.BeamAgenticAnswerStep` 以复用 BEAM 回答行为,或基于
|
||||
`reme.steps.base_step.BaseStep` 自行实现 Step。
|
||||
卸载插件后,Application 和 CLI 服务必须移除插件选择,直到再次安装。
|
||||
卸载不会删除数据集、工作区或结果。修改插件后需重启已有服务。
|
||||
|
|
|
|||
|
|
@ -1,14 +1,101 @@
|
|||
"""BEAM agentic answer step – ReAct agent that answers questions using the search tool."""
|
||||
"""BEAM agentic-answer step."""
|
||||
|
||||
from reme.steps.benchmark import BaseAgenticAnswerStep
|
||||
import os
|
||||
|
||||
from reme.enumeration import ChunkEnum
|
||||
from reme.steps.base_step import BaseStep
|
||||
from reme.steps.index._dedup import _ToolContextDedupMixin
|
||||
from reme.utils.counter import global_counter_inc
|
||||
|
||||
|
||||
class BeamAgenticAnswerStep(BaseAgenticAnswerStep):
|
||||
"""Answer a BEAM probing question via ReAct agent with access to the search tool.
|
||||
class BeamAgenticAnswerStep(BaseStep):
|
||||
"""Answer a BEAM probing question with the ReAct agent and workspace tools."""
|
||||
|
||||
The agent uses the ``agent_wrapper`` component in ReAct mode, calling the
|
||||
``search`` job tool to retrieve relevant memory chunks before generating
|
||||
a final answer.
|
||||
"""
|
||||
# Reasoning-round budget. The AgentScope wrapper converts it to the
|
||||
# backend's iteration-counting semantics; other backends ignore it.
|
||||
MAX_ITERATION = 10
|
||||
TOOL_CONTEXT_PREFIX: str = "beam_agentic_answer"
|
||||
JOB_TOOLS: list[str] = ["search", "add_draft", "read_all_draft", "read"]
|
||||
INJECTED_JOB_KWARGS: dict = {"read_step_format_session": True}
|
||||
|
||||
TOOL_CONTEXT_PREFIX = "beam_agentic_answer"
|
||||
def _injected_job_kwargs(self, query: str) -> dict:
|
||||
"""Enable query-aware session compression when requested."""
|
||||
injected = dict(self.INJECTED_JOB_KWARGS)
|
||||
if self.context is not None and self.context.get("compress_session"):
|
||||
injected["_search"] = {
|
||||
"_compress": {"session": "true"},
|
||||
"queries": [query],
|
||||
"type": "query-aware",
|
||||
}
|
||||
return injected
|
||||
|
||||
async def execute(self):
|
||||
"""Answer the current query and store the result in the response."""
|
||||
assert self.context is not None
|
||||
query: str = self.context.get("query", "")
|
||||
query_time: str | None = self.context.get("query_time")
|
||||
|
||||
if not query:
|
||||
self.context.response.success = False
|
||||
self.context.response.answer = "Skipped: empty query"
|
||||
return self.context.response
|
||||
|
||||
# Build system prompt with optional temporal context
|
||||
sys_prompt = self.get_prompt("system_prompt")
|
||||
if query_time:
|
||||
sys_prompt += "\n" + self.prompt_format("temporal_hint", query_time=query_time)
|
||||
|
||||
if self.app_context is not None:
|
||||
tool_context_id = (
|
||||
f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_"
|
||||
f"{global_counter_inc(self.app_context.metadata, [self.TOOL_CONTEXT_PREFIX])}"
|
||||
)
|
||||
else:
|
||||
tool_context_id = f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_local"
|
||||
wrapper_kwargs = {
|
||||
"system_prompt": sys_prompt,
|
||||
"job_tools": list(self.JOB_TOOLS),
|
||||
"react_config": {"max_iters": self.MAX_ITERATION},
|
||||
"tool_context_id": tool_context_id,
|
||||
}
|
||||
if injected_job_kwargs := self._injected_job_kwargs(f"{query}(query time: {query_time})"):
|
||||
wrapper_kwargs["injected_job_kwargs"] = injected_job_kwargs
|
||||
|
||||
if self.context.stream:
|
||||
text = await self._stream_reply(query, **wrapper_kwargs)
|
||||
else:
|
||||
result = await self.agent_wrapper.reply(query, **wrapper_kwargs)
|
||||
text = (result.get("result") or "").strip()
|
||||
|
||||
self.logger.debug(f"[{self.name}] response: {text!r}")
|
||||
|
||||
self.context.response.success = True
|
||||
self.context.response.answer = text
|
||||
self.context.response.metadata.update(
|
||||
{
|
||||
"query": query,
|
||||
"query_time": query_time,
|
||||
"sys_prompt": sys_prompt,
|
||||
"response": text,
|
||||
},
|
||||
)
|
||||
|
||||
if self.app_context is not None:
|
||||
self.app_context.metadata.get(_ToolContextDedupMixin.TOOL_CONTEXTS_KEY, {}).pop(tool_context_id, None)
|
||||
return self.context.response
|
||||
|
||||
async def _stream_reply(self, query: str, **wrapper_kwargs) -> str:
|
||||
"""Stream unified chunks to the context stream queue."""
|
||||
assert self.context is not None
|
||||
text_parts: list[str] = []
|
||||
|
||||
async for chunk in self.agent_wrapper.reply_stream(query, **wrapper_kwargs):
|
||||
await self.context.add_stream_string(chunk.chunk, chunk.chunk_type)
|
||||
|
||||
if chunk.chunk_type == ChunkEnum.CONTENT and isinstance(chunk.chunk, str):
|
||||
text_parts.append(chunk.chunk)
|
||||
|
||||
if chunk.session_id:
|
||||
self.context.response.metadata["session_id"] = chunk.session_id
|
||||
|
||||
return "".join(text_parts).strip()
|
||||
|
|
|
|||
|
|
@ -31,10 +31,12 @@ The existing `auto_memory`, `agentic_answer`, `answer_judge`, `bench` and `judge
|
|||
names and model environment variables are unchanged. Explicit application/CLI overrides
|
||||
still take precedence. Installing this plugin does not start an evaluation.
|
||||
|
||||
The shared answer base class lives in `reme.steps.benchmark.base_agentic_answer`.
|
||||
The old core-owned `reme.steps.benchmark.lme` Python import path is removed.
|
||||
Custom Python callers should import memory, search and answer Steps from `reme_lme`, and install
|
||||
`lme-judge` before importing the judge Step from `judge_lme`. After uninstalling,
|
||||
Applications and CLI services must omit the plugin until it is installed again.
|
||||
`LmeAgenticAnswerStep` now implements the answer behavior directly. The former
|
||||
`reme.steps.benchmark.BaseAgenticAnswerStep` import is gone; custom subclasses can
|
||||
extend `reme_lme.LmeAgenticAnswerStep` for LongMemEval behavior, or implement their
|
||||
own Step using `reme.steps.base_step.BaseStep`.
|
||||
Uninstallation never removes datasets, workspaces or results.
|
||||
Restart an existing service after changing plugins.
|
||||
|
|
|
|||
|
|
@ -28,8 +28,11 @@ editable 安装会注册 `lme` entry point,并让源码修改立即生效。ru
|
|||
均保持关闭。原有 `auto_memory`、`agentic_answer`、`answer_judge`、`bench`、`judge` 名称及模型环境变量
|
||||
保持不变,显式应用参数和 CLI 覆盖仍优先。安装或启用插件不会自动开始评测。
|
||||
|
||||
共享回答基类位于 `reme.steps.benchmark.base_agentic_answer`。
|
||||
原 `reme.steps.benchmark.lme` Python 导入路径已移除。自定义 Python 调用应从 `reme_lme`
|
||||
导入记忆、搜索和回答 Step;安装 `lme-judge` 后再从 `judge_lme` 导入评判 Step。
|
||||
自定义 Python 调用应从 `reme_lme` 导入记忆、搜索和回答 Step;安装 `lme-judge` 后再从
|
||||
`judge_lme` 导入评判 Step。
|
||||
`LmeAgenticAnswerStep` 现在直接实现回答逻辑。原来的
|
||||
`reme.steps.benchmark.BaseAgenticAnswerStep` 导入路径已移除;自定义子类可继承
|
||||
`reme_lme.LmeAgenticAnswerStep` 以复用 LongMemEval 回答行为,或基于
|
||||
`reme.steps.base_step.BaseStep` 自行实现 Step。
|
||||
卸载插件后,Application 和 CLI 服务必须移除插件选择,直到再次安装。
|
||||
卸载不会删除数据集、工作区或结果。修改插件后需重启已有服务。
|
||||
|
|
|
|||
|
|
@ -1,18 +1,101 @@
|
|||
"""LongMemEval agentic answer step – ReAct agent that answers questions using the search tool."""
|
||||
"""LongMemEval agentic-answer step."""
|
||||
|
||||
from reme.steps.benchmark import BaseAgenticAnswerStep
|
||||
import os
|
||||
|
||||
from reme.enumeration import ChunkEnum
|
||||
from reme.steps.base_step import BaseStep
|
||||
from reme.steps.index._dedup import _ToolContextDedupMixin
|
||||
from reme.utils.counter import global_counter_inc
|
||||
|
||||
|
||||
class LmeAgenticAnswerStep(BaseAgenticAnswerStep):
|
||||
"""Answer a LongMemEval query via ReAct agent with access to the search tool.
|
||||
class LmeAgenticAnswerStep(BaseStep):
|
||||
"""Answer a LongMemEval query with the ReAct agent and workspace tools."""
|
||||
|
||||
The agent uses the ``agent_wrapper`` component in ReAct mode, calling the
|
||||
``search`` job tool to retrieve relevant memory chunks before generating
|
||||
a final answer.
|
||||
# Reasoning-round budget. The AgentScope wrapper converts it to the
|
||||
# backend's iteration-counting semantics; other backends ignore it.
|
||||
MAX_ITERATION = 10
|
||||
TOOL_CONTEXT_PREFIX: str = "lme_agentic_answer"
|
||||
JOB_TOOLS: list[str] = ["search", "add_draft", "read_all_draft", "read"]
|
||||
INJECTED_JOB_KWARGS: dict = {"read_step_format_session": True}
|
||||
|
||||
Session-transcript compression in the plugin's search Step is controlled by the
|
||||
``compress_session`` flag in the runtime context (set by the benchmark
|
||||
runner from ``evaluation.compress_session``); it is off by default.
|
||||
"""
|
||||
def _injected_job_kwargs(self, query: str) -> dict:
|
||||
"""Enable query-aware session compression when requested."""
|
||||
injected = dict(self.INJECTED_JOB_KWARGS)
|
||||
if self.context is not None and self.context.get("compress_session"):
|
||||
injected["_search"] = {
|
||||
"_compress": {"session": "true"},
|
||||
"queries": [query],
|
||||
"type": "query-aware",
|
||||
}
|
||||
return injected
|
||||
|
||||
TOOL_CONTEXT_PREFIX = "lme_agentic_answer"
|
||||
async def execute(self):
|
||||
"""Answer the current query and store the result in the response."""
|
||||
assert self.context is not None
|
||||
query: str = self.context.get("query", "")
|
||||
query_time: str | None = self.context.get("query_time")
|
||||
|
||||
if not query:
|
||||
self.context.response.success = False
|
||||
self.context.response.answer = "Skipped: empty query"
|
||||
return self.context.response
|
||||
|
||||
# Build system prompt with optional temporal context
|
||||
sys_prompt = self.get_prompt("system_prompt")
|
||||
if query_time:
|
||||
sys_prompt += "\n" + self.prompt_format("temporal_hint", query_time=query_time)
|
||||
|
||||
if self.app_context is not None:
|
||||
tool_context_id = (
|
||||
f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_"
|
||||
f"{global_counter_inc(self.app_context.metadata, [self.TOOL_CONTEXT_PREFIX])}"
|
||||
)
|
||||
else:
|
||||
tool_context_id = f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_local"
|
||||
wrapper_kwargs = {
|
||||
"system_prompt": sys_prompt,
|
||||
"job_tools": list(self.JOB_TOOLS),
|
||||
"react_config": {"max_iters": self.MAX_ITERATION},
|
||||
"tool_context_id": tool_context_id,
|
||||
}
|
||||
if injected_job_kwargs := self._injected_job_kwargs(f"{query}(query time: {query_time})"):
|
||||
wrapper_kwargs["injected_job_kwargs"] = injected_job_kwargs
|
||||
|
||||
if self.context.stream:
|
||||
text = await self._stream_reply(query, **wrapper_kwargs)
|
||||
else:
|
||||
result = await self.agent_wrapper.reply(query, **wrapper_kwargs)
|
||||
text = (result.get("result") or "").strip()
|
||||
|
||||
self.logger.debug(f"[{self.name}] response: {text!r}")
|
||||
|
||||
self.context.response.success = True
|
||||
self.context.response.answer = text
|
||||
self.context.response.metadata.update(
|
||||
{
|
||||
"query": query,
|
||||
"query_time": query_time,
|
||||
"sys_prompt": sys_prompt,
|
||||
"response": text,
|
||||
},
|
||||
)
|
||||
|
||||
if self.app_context is not None:
|
||||
self.app_context.metadata.get(_ToolContextDedupMixin.TOOL_CONTEXTS_KEY, {}).pop(tool_context_id, None)
|
||||
return self.context.response
|
||||
|
||||
async def _stream_reply(self, query: str, **wrapper_kwargs) -> str:
|
||||
"""Stream unified chunks to the context stream queue."""
|
||||
assert self.context is not None
|
||||
text_parts: list[str] = []
|
||||
|
||||
async for chunk in self.agent_wrapper.reply_stream(query, **wrapper_kwargs):
|
||||
await self.context.add_stream_string(chunk.chunk, chunk.chunk_type)
|
||||
|
||||
if chunk.chunk_type == ChunkEnum.CONTENT and isinstance(chunk.chunk, str):
|
||||
text_parts.append(chunk.chunk)
|
||||
|
||||
if chunk.session_id:
|
||||
self.context.response.metadata["session_id"] = chunk.session_id
|
||||
|
||||
return "".join(text_parts).strip()
|
||||
|
|
|
|||
|
|
@ -1,11 +1,10 @@
|
|||
"""steps"""
|
||||
|
||||
from . import benchmark, common, evolve, file_io, index, transfer
|
||||
from . import common, evolve, file_io, index, transfer
|
||||
from .base_step import BaseStep
|
||||
|
||||
__all__ = [
|
||||
"BaseStep",
|
||||
"benchmark",
|
||||
"common",
|
||||
"evolve",
|
||||
"file_io",
|
||||
|
|
|
|||
|
|
@ -1,5 +0,0 @@
|
|||
"""Shared benchmark steps; concrete implementations live in plugins."""
|
||||
|
||||
from .base_agentic_answer import BaseAgenticAnswerStep
|
||||
|
||||
__all__ = ["BaseAgenticAnswerStep"]
|
||||
|
|
@ -1,128 +0,0 @@
|
|||
"""Shared base class for benchmark agentic-answer steps."""
|
||||
|
||||
import os
|
||||
|
||||
from ..base_step import BaseStep
|
||||
from ..index._dedup import _ToolContextDedupMixin
|
||||
from ...enumeration import ChunkEnum
|
||||
from ...utils.counter import global_counter_inc
|
||||
|
||||
|
||||
class BaseAgenticAnswerStep(BaseStep):
|
||||
"""ReAct-agent answer implementation shared by benchmark plugins.
|
||||
|
||||
Subclasses only need to set:
|
||||
TOOL_CONTEXT_PREFIX (str): prefix used to build the unique tool_context_id.
|
||||
JOB_TOOLS (list[str]): tools exposed to the ReAct agent; override to customize.
|
||||
INJECTED_JOB_KWARGS (dict): server-owned kwargs injected into every job
|
||||
tool call via ``injected_job_kwargs``; override the attribute or the
|
||||
``_injected_job_kwargs`` hook to customize.
|
||||
|
||||
Concrete subclasses are registered by the plugin manifest.
|
||||
|
||||
Inputs (from RuntimeContext):
|
||||
query (str, required): The question to answer.
|
||||
query_time (str, optional): ISO timestamp representing the query time,
|
||||
used to ground the agent's temporal context.
|
||||
|
||||
Output (written to context.response.answer):
|
||||
The agent's final answer text.
|
||||
"""
|
||||
|
||||
# Reasoning-round budget. The AgentScope wrapper converts it to the
|
||||
# backend's iteration-counting semantics; other backends ignore it.
|
||||
MAX_ITERATION = 10
|
||||
TOOL_CONTEXT_PREFIX: str = "content_agentic_answer"
|
||||
JOB_TOOLS: list[str] = ["search", "add_draft", "read_all_draft", "read"]
|
||||
INJECTED_JOB_KWARGS: dict = {"read_step_format_session": True}
|
||||
|
||||
def _injected_job_kwargs(self, query: str) -> dict: # pylint: disable=unused-argument
|
||||
"""Server-owned kwargs injected into every job tool call.
|
||||
|
||||
Overridable hook: subclasses can extend the static
|
||||
``INJECTED_JOB_KWARGS`` with per-request values derived from ``query``.
|
||||
|
||||
When the runtime context carries a truthy ``compress_session`` flag,
|
||||
session-transcript compression is enabled in the benchmark plugin's search Step by
|
||||
injecting a ``_search._compress.session`` marker plus the current
|
||||
``query`` as the query-aware relevance filter. Default (falsy) leaves
|
||||
session chunks uncompressed.
|
||||
"""
|
||||
injected = dict(self.INJECTED_JOB_KWARGS)
|
||||
if self.context is not None and self.context.get("compress_session"):
|
||||
injected["_search"] = {
|
||||
"_compress": {"session": "true"},
|
||||
"queries": [query],
|
||||
"type": "query-aware",
|
||||
}
|
||||
return injected
|
||||
|
||||
async def execute(self):
|
||||
assert self.context is not None
|
||||
query: str = self.context.get("query", "")
|
||||
query_time: str | None = self.context.get("query_time")
|
||||
|
||||
if not query:
|
||||
self.context.response.success = False
|
||||
self.context.response.answer = "Skipped: empty query"
|
||||
return self.context.response
|
||||
|
||||
# Build system prompt with optional temporal context
|
||||
sys_prompt = self.get_prompt("system_prompt")
|
||||
if query_time:
|
||||
sys_prompt += "\n" + self.prompt_format("temporal_hint", query_time=query_time)
|
||||
|
||||
if self.app_context is not None:
|
||||
tool_context_id = (
|
||||
f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_"
|
||||
f"{global_counter_inc(self.app_context.metadata, [self.TOOL_CONTEXT_PREFIX])}"
|
||||
)
|
||||
else:
|
||||
tool_context_id = f"{self.TOOL_CONTEXT_PREFIX}_{os.getpid()}_local"
|
||||
wrapper_kwargs = {
|
||||
"system_prompt": sys_prompt,
|
||||
"job_tools": list(self.JOB_TOOLS),
|
||||
"react_config": {"max_iters": self.MAX_ITERATION},
|
||||
"tool_context_id": tool_context_id,
|
||||
}
|
||||
if injected_job_kwargs := self._injected_job_kwargs(f"{query}(query time: {query_time})"):
|
||||
wrapper_kwargs["injected_job_kwargs"] = injected_job_kwargs
|
||||
|
||||
if self.context.stream:
|
||||
text = await self._stream_reply(query, **wrapper_kwargs)
|
||||
else:
|
||||
result = await self.agent_wrapper.reply(query, **wrapper_kwargs)
|
||||
text = (result.get("result") or "").strip()
|
||||
|
||||
self.logger.debug(f"[{self.name}] response: {text!r}")
|
||||
|
||||
self.context.response.success = True
|
||||
self.context.response.answer = text
|
||||
self.context.response.metadata.update(
|
||||
{
|
||||
"query": query,
|
||||
"query_time": query_time,
|
||||
"sys_prompt": sys_prompt,
|
||||
"response": text,
|
||||
},
|
||||
)
|
||||
|
||||
if self.app_context is not None:
|
||||
self.app_context.metadata.get(_ToolContextDedupMixin.TOOL_CONTEXTS_KEY, {}).pop(tool_context_id, None)
|
||||
return self.context.response
|
||||
|
||||
async def _stream_reply(self, query: str, **wrapper_kwargs) -> str:
|
||||
"""Stream unified chunks to the context stream queue."""
|
||||
assert self.context is not None
|
||||
text_parts: list[str] = []
|
||||
|
||||
async for chunk in self.agent_wrapper.reply_stream(query, **wrapper_kwargs):
|
||||
await self.context.add_stream_string(chunk.chunk, chunk.chunk_type)
|
||||
|
||||
if chunk.chunk_type == ChunkEnum.CONTENT and isinstance(chunk.chunk, str):
|
||||
text_parts.append(chunk.chunk)
|
||||
|
||||
if chunk.session_id:
|
||||
self.context.response.metadata["session_id"] = chunk.session_id
|
||||
|
||||
return "".join(text_parts).strip()
|
||||
Loading…
Add table
Reference in a new issue