From 7fb7ba23a4fdfd54c9e1858085edca5fa196fdb7 Mon Sep 17 00:00:00 2001 From: YUZHEthefool <2804776511@qq.com> Date: Thu, 10 Sep 2026 23:34:44 +0800 Subject: [PATCH 1/3] feat(btw): extract work execution and task state Reuse the existing Agent executor with bounded execution concurrency, terminal retention, cancellation propagation and generic failure state. Related: #125 AI-Generated: true Generated-At: 2026-09-10T15:34:43Z --- astrbot/core/agent/btw/__init__.py | 1 + astrbot/core/agent/btw/i18n.py | 50 +++++ astrbot/core/agent/btw/types.py | 65 +++++++ astrbot/core/agent/btw/work_loop.py | 173 ++++++++++++++++++ astrbot/core/agent/btw/work_sessions.py | 106 +++++++++++ astrbot/core/agent/conversation_loop.py | 27 +++ astrbot/core/config/default.py | 24 ++- .../en-US/features/config-metadata.json | 16 ++ .../zh-CN/features/config-metadata.json | 16 ++ docs/en/dev/astrbot-config.md | 2 + docs/zh/dev/astrbot-config.md | 2 + tests/unit/test_btw_work_loop.py | 146 +++++++++++++++ tests/unit/test_conversation_loop.py | 25 +++ 13 files changed, 652 insertions(+), 1 deletion(-) create mode 100644 astrbot/core/agent/btw/__init__.py create mode 100644 astrbot/core/agent/btw/i18n.py create mode 100644 astrbot/core/agent/btw/types.py create mode 100644 astrbot/core/agent/btw/work_loop.py create mode 100644 astrbot/core/agent/btw/work_sessions.py create mode 100644 tests/unit/test_btw_work_loop.py diff --git a/astrbot/core/agent/btw/__init__.py b/astrbot/core/agent/btw/__init__.py new file mode 100644 index 0000000000..71280bdf32 --- /dev/null +++ b/astrbot/core/agent/btw/__init__.py @@ -0,0 +1 @@ +# BTW runtime primitives; keep package imports inert. diff --git a/astrbot/core/agent/btw/i18n.py b/astrbot/core/agent/btw/i18n.py new file mode 100644 index 0000000000..33d5239da1 --- /dev/null +++ b/astrbot/core/agent/btw/i18n.py @@ -0,0 +1,50 @@ +"""Locale-aware user-facing strings for the BTW work loop. + +The work loop runs in core without a plugin context, so it resolves the +locale from the event extra/session the same way ``PluginContext._locale`` +does, then looks the string up in these bundles. Missing locales fall back +to ``zh-CN``. +""" + +LOCALES: dict[str, dict[str, str]] = { + "zh-CN": { + "btw.work.started": "🔧 工作任务已开始处理。", + "btw.work.status.pending": "工作任务正在排队。", + "btw.work.status.running": "工作任务正在执行。", + "btw.work.status.completed": "工作任务已完成。", + "btw.work.status.failed": "工作任务执行失败。", + "btw.work.status.cancelled": "工作任务已取消。", + }, + "en-US": { + "btw.work.started": "🔧 Work task started.", + "btw.work.status.pending": "The work task is queued.", + "btw.work.status.running": "The work task is running.", + "btw.work.status.completed": "The work task is completed.", + "btw.work.status.failed": "The work task failed.", + "btw.work.status.cancelled": "The work task was cancelled.", + }, +} + +_FALLBACK_LOCALE = "zh-CN" + + +def resolve_event_locale(event) -> str: + """Return the locale for an event (extra first, then the stored session).""" + getter = getattr(event, "get_extra", None) + if callable(getter): + try: + extra = getter("locale") + except Exception: # noqa: BLE001 + extra = None + if extra: + return str(extra) + return _FALLBACK_LOCALE + + +def text(locale: str, key: str) -> str: + """Return one BTW string for a locale, falling back to zh-CN then key.""" + bundle = LOCALES.get(locale) or LOCALES[_FALLBACK_LOCALE] + value = bundle.get(key) + if value is None: + value = LOCALES[_FALLBACK_LOCALE].get(key, key) + return value diff --git a/astrbot/core/agent/btw/types.py b/astrbot/core/agent/btw/types.py new file mode 100644 index 0000000000..ca600b8d87 --- /dev/null +++ b/astrbot/core/agent/btw/types.py @@ -0,0 +1,65 @@ +"""Types shared by the BTW conversation and work loops.""" + +from collections.abc import Mapping +from dataclasses import dataclass, field +from datetime import UTC, datetime +from enum import StrEnum +from uuid import uuid4 + + +def is_work_loop_enabled(config: object) -> bool: + """Return whether the profile explicitly enables BTW and work.""" + if not isinstance(config, Mapping): + return False + btw = config.get("btw", {}) + if not isinstance(btw, Mapping) or not btw.get("enabled", False): + return False + work = btw.get("work_loop", {}) + return isinstance(work, Mapping) and bool(work.get("enabled", False)) + + +class TaskType(StrEnum): + """The execution loop selected for a user request.""" + + CONVERSATION = "conversation" + WORK = "work" + + +class WorkSessionStatus(StrEnum): + """Lifecycle states for one work-loop request.""" + + PENDING = "pending" + RUNNING = "running" + COMPLETED = "completed" + FAILED = "failed" + CANCELLED = "cancelled" + + +@dataclass(slots=True) +class WorkSession: + """Runtime state shared by the conversation and work loops.""" + + origin: str + request: str + task_type: TaskType = TaskType.WORK + id: str = field(default_factory=lambda: uuid4().hex) + status: WorkSessionStatus = WorkSessionStatus.PENDING + created_at: datetime = field(default_factory=lambda: datetime.now(UTC)) + updated_at: datetime = field(default_factory=lambda: datetime.now(UTC)) + error: str | None = None + + def update_status( + self, + status: WorkSessionStatus, + *, + error: str | None = None, + ) -> None: + """Record a status transition. + + Args: + status: The new work-session status. + error: A safe diagnostic for failed work, when available. + """ + self.status = status + self.error = error + self.updated_at = datetime.now(UTC) diff --git a/astrbot/core/agent/btw/work_loop.py b/astrbot/core/agent/btw/work_loop.py new file mode 100644 index 0000000000..c2d565056e --- /dev/null +++ b/astrbot/core/agent/btw/work_loop.py @@ -0,0 +1,173 @@ +"""The BTW work-loop prototype backed by the existing Agent tool loop.""" + +import asyncio +from collections.abc import AsyncGenerator, Awaitable, Callable +from typing import Protocol + +from astrbot.core.message.message_event_result import MessageEventResult +from astrbot.core.platform.astr_message_event import AstrMessageEvent +from astrbot.core.utils.task_utils import create_tracked_task + +from . import i18n as work_i18n +from .types import WorkSessionStatus +from .work_sessions import WorkSessionManager + + +class AgentRequestExecutor(Protocol): + """The existing Agent request path required by the work loop.""" + + def process(self, event: AstrMessageEvent) -> AsyncGenerator[None]: + """Yield pipeline progress markers for one event. + + Protocol stub; concrete implementations are the pipeline's Agent + request sub-stage. The body raises so the statement is effectful + (CodeQL py/ineffectual-statement); the unreachable ``yield`` keeps + the declared ``AsyncGenerator`` return type type-checkable. + """ + raise NotImplementedError + yield # noqa: B901 -- unreachable marker for the type checker + + +ResultDispatcher = Callable[[AstrMessageEvent], Awaitable[None]] +EventFinalizer = Callable[[AstrMessageEvent], Awaitable[None]] + + +class WorkLoop: + """Run classified work with the current Agent and tool infrastructure.""" + + def __init__( + self, + executor: AgentRequestExecutor, + sessions: WorkSessionManager, + *, + max_concurrent: int = 2, + ) -> None: + self.executor = executor + self.sessions = sessions + self._semaphore = asyncio.Semaphore(max(1, max_concurrent)) + self._background_tasks: set[asyncio.Task] | None = None + self._result_dispatcher: ResultDispatcher | None = None + self._event_finalizer: EventFinalizer | None = None + + def configure_detached_execution( + self, + *, + background_tasks: set[asyncio.Task], + result_dispatcher: ResultDispatcher, + event_finalizer: EventFinalizer, + ) -> None: + """Attach runtime-owned background execution services. + + Args: + background_tasks: Runtime task registry cancelled during shutdown. + result_dispatcher: Delivers a generated work result through the + configured result-decorate and response stages. + event_finalizer: Releases the event after detached work finishes. + """ + self._background_tasks = background_tasks + self._result_dispatcher = result_dispatcher + self._event_finalizer = event_finalizer + + async def process(self, event: AstrMessageEvent) -> AsyncGenerator[None]: + """Execute one work-loop request inline. + + Args: + event: The classified message event. + + Yields: + Pipeline progress markers emitted by the existing Agent executor. + """ + session = await self.sessions.create( + event.unified_msg_origin, event.message_str + ) + self._prepare_event(event, session.id) + async for progress in self._execute(event, session.id): + yield progress + + async def submit(self, event: AstrMessageEvent) -> AsyncGenerator[None]: + """Acknowledge work, then run it without retaining the request pipeline. + + Falls back to inline execution when no runtime task registry is + attached, which keeps the primitive usable in isolated tests. + """ + if ( + self._background_tasks is None + or self._result_dispatcher is None + or self._event_finalizer is None + ): + async for progress in self.process(event): + yield progress + return + + session = await self.sessions.create( + event.unified_msg_origin, event.message_str + ) + self._prepare_event(event, session.id) + event.set_result( + MessageEventResult().message( + work_i18n.text( + work_i18n.resolve_event_locale(event), "btw.work.started" + ) + ) + ) + yield + + # The first yield returns only after the normal response stages deliver + # the acknowledgement. Marking it here prevents the scheduler from + # releasing event-owned temporary files before the worker needs them. + event.set_extra("btw_detached_work", True) + create_tracked_task( + self._background_tasks, + self._run_detached(event, session.id), + name=f"btw_work:{session.id}", + ) + + @staticmethod + def _prepare_event(event: AstrMessageEvent, session_id: str) -> None: + """Mark an event so Agent assembly uses the work-loop policy.""" + event.set_extra("btw_work_session_id", session_id) + event.set_extra("btw_loop", "work") + event.set_extra("btw_agent_lock_key", f"{event.unified_msg_origin}:work") + + async def _execute( + self, + event: AstrMessageEvent, + session_id: str, + ) -> AsyncGenerator[None]: + """Run one already-created work session and update its lifecycle.""" + try: + async with self._semaphore: + await self.sessions.update_status( + session_id, + WorkSessionStatus.RUNNING, + ) + async for progress in self.executor.process(event): + yield progress + except asyncio.CancelledError: + await self.sessions.update_status( + session_id, + WorkSessionStatus.CANCELLED, + ) + raise + except Exception: + await self.sessions.update_status( + session_id, + WorkSessionStatus.FAILED, + error="Work task failed.", + ) + raise + else: + await self.sessions.update_status( + session_id, + WorkSessionStatus.COMPLETED, + ) + + async def _run_detached(self, event: AstrMessageEvent, session_id: str) -> None: + """Run work in the runtime task registry and deliver each result.""" + assert self._result_dispatcher is not None + assert self._event_finalizer is not None + try: + async for _ in self._execute(event, session_id): + await self._result_dispatcher(event) + finally: + await self._event_finalizer(event) diff --git a/astrbot/core/agent/btw/work_sessions.py b/astrbot/core/agent/btw/work_sessions.py new file mode 100644 index 0000000000..7c919f9efc --- /dev/null +++ b/astrbot/core/agent/btw/work_sessions.py @@ -0,0 +1,106 @@ +"""In-memory runtime ownership for BTW work sessions.""" + +import asyncio +from datetime import UTC, datetime, timedelta + +from .types import WorkSession, WorkSessionStatus + + +class WorkSessionManager: + """Own active and recent work-loop session state for one pipeline. + + ``get_for_origin`` returns the *latest* session per origin: a new task + replaces the previous session's status-query target, while older + sessions stay addressable by id until they expire. Work-loop + concurrency is bounded by the work loop's semaphore, not here. + """ + + def __init__(self, *, max_age_seconds: int = 3600) -> None: + self._by_origin: dict[str, WorkSession] = {} + self._by_id: dict[str, WorkSession] = {} + self._lock = asyncio.Lock() + self.set_max_age_seconds(max_age_seconds) + + def set_max_age_seconds(self, value: int) -> None: + """Set the retention period for terminal sessions. + + Args: + value: Number of seconds to keep a completed, failed, or cancelled + session before a later manager operation removes it. + """ + self._max_age_seconds = max(1, value) if type(value) is int else 3600 + + async def create(self, origin: str, request: str) -> WorkSession: + """Create and register a work session. + + Args: + origin: The unified message origin that owns the work. + request: The user request being processed. + + Returns: + The newly created work session. + """ + session = WorkSession(origin=origin, request=request) + async with self._lock: + self._cleanup_expired_locked() + self._by_origin[origin] = session + self._by_id[session.id] = session + return session + + async def get_for_origin(self, origin: str) -> WorkSession | None: + """Return the most recent work session for an origin.""" + async with self._lock: + self._cleanup_expired_locked() + return self._by_origin.get(origin) + + async def get_by_id(self, session_id: str) -> WorkSession | None: + """Return a recent work session by its identifier. + + Args: + session_id: The generated work-session identifier. + + Returns: + The matching session, or ``None`` after it has expired. + """ + async with self._lock: + self._cleanup_expired_locked() + return self._by_id.get(session_id) + + async def update_status( + self, + session_id: str, + status: WorkSessionStatus, + *, + error: str | None = None, + ) -> WorkSession | None: + """Transition one known work session. + + Args: + session_id: The work-session identifier. + status: The new lifecycle status. + error: A safe failure message, when applicable. + + Returns: + The updated session, or ``None`` when it has expired. + """ + async with self._lock: + self._cleanup_expired_locked() + session = self._by_id.get(session_id) + if session is not None: + session.update_status(status, error=error) + return session + + def _cleanup_expired_locked(self) -> None: + """Remove old terminal sessions while the manager lock is held.""" + cutoff = datetime.now(UTC) - timedelta(seconds=self._max_age_seconds) + expired_ids = [ + session_id + for session_id, session in self._by_id.items() + if session.status + not in {WorkSessionStatus.PENDING, WorkSessionStatus.RUNNING} + and session.updated_at < cutoff + ] + for session_id in expired_ids: + session = self._by_id.pop(session_id) + if self._by_origin.get(session.origin) is session: + self._by_origin.pop(session.origin, None) diff --git a/astrbot/core/agent/conversation_loop.py b/astrbot/core/agent/conversation_loop.py index 237befeca5..6a1914ff84 100644 --- a/astrbot/core/agent/conversation_loop.py +++ b/astrbot/core/agent/conversation_loop.py @@ -3,6 +3,9 @@ from collections.abc import AsyncGenerator from typing import TYPE_CHECKING +from astrbot.core.agent.btw.types import is_work_loop_enabled +from astrbot.core.agent.btw.work_loop import WorkLoop +from astrbot.core.agent.btw.work_sessions import WorkSessionManager from astrbot.core.platform.astr_message_event import AstrMessageEvent if TYPE_CHECKING: @@ -18,6 +21,8 @@ class ConversationLoop: def __init__(self, agent_request: AgentRequestSubStage) -> None: self.agent_request = agent_request self._btw_enabled = False + self.work_sessions = WorkSessionManager() + self.work_loop: WorkLoop | None = None async def initialize(self, ctx: PipelineContext) -> None: """Initialize the shared Agent executor for this profile.""" @@ -25,9 +30,31 @@ async def initialize(self, ctx: PipelineContext) -> None: btw = self.astrbot_config.get("btw", {}) self._btw_enabled = isinstance(btw, dict) and bool(btw.get("enabled", False)) await self.agent_request.initialize(ctx) + btw = btw if isinstance(btw, dict) else {} + work = btw.get("work_loop", {}) + work = work if isinstance(work, dict) else {} + retention = btw.get("work_session", {}) + retention = retention if isinstance(retention, dict) else {} + self.work_sessions.set_max_age_seconds(retention.get("max_age_seconds", 3600)) + concurrency = work.get("max_concurrent", 2) + self.work_loop = WorkLoop( + self.agent_request, + self.work_sessions, + max_concurrent=concurrency if type(concurrency) is int else 2, + ) async def process(self, event: AstrMessageEvent) -> AsyncGenerator[None]: """Process one admitted conversation using the current Agent path.""" + if ( + self._btw_enabled + and event.get_extra("btw_force_work") + and is_work_loop_enabled(self.astrbot_config) + ): + if self.work_loop is None: + raise RuntimeError("ConversationLoop is not initialized") + async for response in self.work_loop.submit(event): + yield response + return if self._btw_enabled: event.set_extra("btw_loop", "conversation") async for response in self.agent_request.process(event): diff --git a/astrbot/core/config/default.py b/astrbot/core/config/default.py index 74bbbd48fb..76e7848718 100644 --- a/astrbot/core/config/default.py +++ b/astrbot/core/config/default.py @@ -188,7 +188,11 @@ ), "agents": [], }, - "btw": {"enabled": False}, + "btw": { + "enabled": False, + "work_loop": {"enabled": False, "max_concurrent": 2}, + "work_session": {"max_age_seconds": 3600}, + }, "provider_stt_settings": { "enable": False, "provider_id": "", @@ -4695,6 +4699,24 @@ "type": "bool", "hint": "实验功能,默认关闭。开启后,普通 AI 请求通过对话循环进入现有 Agent。", }, + "btw.work_loop.enabled": { + "description": "启用工作循环", + "type": "bool", + "hint": "默认关闭;允许显式工作请求使用工作执行器。", + "condition": {"btw.enabled": True}, + }, + "btw.work_loop.max_concurrent": { + "description": "工作任务执行并发", + "type": "int", + "hint": "同时执行的工作任务数,默认 2。此值不是等待队列的长度限制。", + "condition": {"btw.work_loop.enabled": True}, + }, + "btw.work_session.max_age_seconds": { + "description": "终态工作会话保留秒数", + "type": "int", + "hint": "已完成、失败或取消的工作会话保留时间,默认 3600 秒。", + "condition": {"btw.enabled": True}, + }, }, } diff --git a/dashboard/src/i18n/locales/en-US/features/config-metadata.json b/dashboard/src/i18n/locales/en-US/features/config-metadata.json index bf01745949..d1ee999b84 100644 --- a/dashboard/src/i18n/locales/en-US/features/config-metadata.json +++ b/dashboard/src/i18n/locales/en-US/features/config-metadata.json @@ -1181,6 +1181,22 @@ "enabled": { "description": "Enable BTW dual loops", "hint": "Experimental and disabled by default. Ordinary admitted AI requests use the conversation entry over the existing Agent." + }, + "work_loop": { + "enabled": { + "description": "Enable work loop", + "hint": "Disabled by default. Allow explicit work execution." + }, + "max_concurrent": { + "description": "Concurrent work execution", + "hint": "Active execution limit, default 2. This is not a waiting-queue length limit." + } + }, + "work_session": { + "max_age_seconds": { + "description": "Terminal work retention (seconds)", + "hint": "Keep completed, failed or cancelled records for 3600 seconds by default." + } } } } diff --git a/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json b/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json index 364895e686..9ce3d70f91 100644 --- a/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json +++ b/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json @@ -1175,6 +1175,22 @@ "enabled": { "description": "启用 BTW 双循环", "hint": "实验功能,默认关闭。普通且已通过准入的 AI 请求经对话入口使用现有 Agent。" + }, + "work_loop": { + "enabled": { + "description": "??????", + "hint": "????????????????" + }, + "max_concurrent": { + "description": "????????", + "hint": "???????????? 2????????????" + } + }, + "work_session": { + "max_age_seconds": { + "description": "??????????", + "hint": "??????????????? 3600 ??" + } } } } diff --git a/docs/en/dev/astrbot-config.md b/docs/en/dev/astrbot-config.md index 497cae6ca4..a6567540e3 100644 --- a/docs/en/dev/astrbot-config.md +++ b/docs/en/dev/astrbot-config.md @@ -205,6 +205,8 @@ Alkaid [Long-term Memory](../use/long-term-memory) currently has no enable/disab Automatic classifier candidates are evaluated separately. Enabling this entry does not select an automatic routing strategy. +The work executor additionally requires `btw.work_loop.enabled`, also `false` by default. It reuses the Agent executor and records pending, running, completed, failed, and cancelled task states. `btw.work_loop.max_concurrent` limits active execution (default `2`); it does not impose a waiting-queue length limit. `btw.work_session.max_age_seconds` retains terminal states for `3600` seconds by default; active tasks do not expire, and expired terminal records are removed during the next session operation. Runtime-owned background services perform task execution and cleanup when attached by the scheduler. + ## WebUI and authentication Important `dashboard` defaults: diff --git a/docs/zh/dev/astrbot-config.md b/docs/zh/dev/astrbot-config.md index 835dc20251..e3fb243319 100644 --- a/docs/zh/dev/astrbot-config.md +++ b/docs/zh/dev/astrbot-config.md @@ -207,6 +207,8 @@ Alkaid [长期记忆](../use/long-term-memory) 当前没有对应的启停配置 自动分类器候选将分别评估。开启此入口不会选定自动路由方案。 +工作执行器还需要开启 `btw.work_loop.enabled`,默认同样为 `false`。它复用 Agent 执行器并记录排队、运行、完成、失败、取消状态。`btw.work_loop.max_concurrent` 限制正在执行的任务数,默认 `2`,不限制等待队列长度。`btw.work_session.max_age_seconds` 默认保留终态记录 `3600` 秒;活动任务不会过期,终态过期记录在下次会话操作时清除。调度器接入后台服务后,由运行时拥有工作任务的执行和清理。 + ## WebUI 与认证 `dashboard` 的关键默认值: diff --git a/tests/unit/test_btw_work_loop.py b/tests/unit/test_btw_work_loop.py new file mode 100644 index 0000000000..89d4b0952f --- /dev/null +++ b/tests/unit/test_btw_work_loop.py @@ -0,0 +1,146 @@ +import asyncio +from datetime import UTC, datetime, timedelta +from unittest.mock import AsyncMock + +import pytest + +from astrbot.core.agent.btw.types import WorkSessionStatus, is_work_loop_enabled +from astrbot.core.agent.btw.work_loop import WorkLoop +from astrbot.core.agent.btw.work_sessions import WorkSessionManager + + +class FakeEvent: + def __init__(self, message: str) -> None: + self.unified_msg_origin = "umo-1" + self.message_str = message + self.extras = {} + self.result = None + + def set_extra(self, key, value) -> None: + self.extras[key] = value + + def get_extra(self, key): + return self.extras.get(key) + + def set_result(self, value) -> None: + self.result = value + + +class FailingExecutor: + async def process(self, event): + del event + raise RuntimeError("provider token leaked") + yield + + +class BlockingExecutor: + def __init__(self) -> None: + self.started = asyncio.Event() + self.release = asyncio.Event() + + async def process(self, event): + del event + self.started.set() + await self.release.wait() + yield "done" + + +@pytest.mark.asyncio +async def test_work_loop_marks_failures_without_exposing_executor_error(): + sessions = WorkSessionManager() + work_loop = WorkLoop(FailingExecutor(), sessions) + event = FakeEvent("执行命令") + + with pytest.raises(RuntimeError, match="provider token leaked"): + _ = [item async for item in work_loop.process(event)] + + session = await sessions.get_for_origin(event.unified_msg_origin) + assert session is not None + assert session.status is WorkSessionStatus.FAILED + assert "provider token leaked" not in (session.error or "") + + +@pytest.mark.asyncio +async def test_work_session_manager_expires_terminal_sessions(): + sessions = WorkSessionManager(max_age_seconds=60) + session = await sessions.create("umo-1", "执行命令") + await sessions.update_status(session.id, WorkSessionStatus.COMPLETED) + session.updated_at = datetime.now(UTC) - timedelta(seconds=61) + + assert await sessions.get_by_id(session.id) is None + assert await sessions.get_for_origin(session.origin) is None + + +@pytest.mark.asyncio +async def test_work_loop_acknowledges_then_runs_in_background(): + executor = BlockingExecutor() + sessions = WorkSessionManager() + work_loop = WorkLoop(executor, sessions) + background_tasks: set[asyncio.Task] = set() + result_dispatcher = AsyncMock() + event_finalizer = AsyncMock() + work_loop.configure_detached_execution( + background_tasks=background_tasks, + result_dispatcher=result_dispatcher, + event_finalizer=event_finalizer, + ) + event = FakeEvent("执行命令") + + output = [item async for item in work_loop.submit(event)] + + assert output == [None] + assert event.result.get_plain_text() == "🔧 工作任务已开始处理。" + assert len(background_tasks) == 1 + [task] = background_tasks + await asyncio.wait_for(executor.started.wait(), timeout=1) + session = await sessions.get_for_origin(event.unified_msg_origin) + assert session is not None + assert session.status is WorkSessionStatus.RUNNING + + executor.release.set() + await asyncio.wait_for(task, timeout=5) + + assert session.status is WorkSessionStatus.COMPLETED + result_dispatcher.assert_awaited_once_with(event) + event_finalizer.assert_awaited_once_with(event) + + +@pytest.mark.parametrize( + "config, enabled", + [ + ({"btw": {"enabled": True, "work_loop": {"enabled": True}}}, True), + ({"btw": {"enabled": False, "work_loop": {"enabled": True}}}, False), + ({"btw": {"enabled": True, "work_loop": None}}, False), + ({"btw": None}, False), + (None, False), + ], +) +def test_work_admission_requires_both_valid_switches(config, enabled): + assert is_work_loop_enabled(config) is enabled + + +@pytest.mark.asyncio +async def test_cancelled_work_retains_cancelled_state_and_propagates(): + executor = BlockingExecutor() + sessions = WorkSessionManager() + loop = WorkLoop(executor, sessions) + event = FakeEvent("cancel this work") + + async def run(): + return [item async for item in loop.process(event)] + + task = asyncio.create_task(run()) + await asyncio.wait_for(executor.started.wait(), timeout=1) + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + session = await sessions.get_for_origin(event.unified_msg_origin) + assert session.status is WorkSessionStatus.CANCELLED + + +@pytest.mark.asyncio +async def test_retention_does_not_expire_active_work(): + sessions = WorkSessionManager(max_age_seconds=1) + session = await sessions.create("origin", "still running") + session.updated_at = datetime.now(UTC) - timedelta(seconds=120) + assert await sessions.get_by_id(session.id) is session diff --git a/tests/unit/test_conversation_loop.py b/tests/unit/test_conversation_loop.py index fb8e318e9b..a7321dcb13 100644 --- a/tests/unit/test_conversation_loop.py +++ b/tests/unit/test_conversation_loop.py @@ -24,6 +24,9 @@ def __init__(self) -> None: def set_extra(self, key, value) -> None: self.extras[key] = value + def get_extra(self, key): + return self.extras.get(key) + @pytest.mark.asyncio @pytest.mark.parametrize("btw", [{"enabled": True}, {"enabled": False}, {}, None]) @@ -40,3 +43,25 @@ async def test_conversation_entry_preserves_agent_execution_and_disabled_metadat assert event.extras == ( {"btw_loop": "conversation"} if btw and btw["enabled"] else {} ) + + +@pytest.mark.asyncio +async def test_explicit_work_uses_the_work_executor_without_a_classifier(): + executor = FakeAgentRequest() + loop = ConversationLoop(executor) + await loop.initialize( + SimpleNamespace( + astrbot_config={ + "btw": {"enabled": True, "work_loop": {"enabled": True}}, + } + ) + ) + event = FakeEvent() + event.message_str = "inspect the workspace" + event.unified_msg_origin = "origin" + event.set_extra("btw_force_work", True) + + assert [item async for item in loop.process(event)] == ["first", "second"] + assert event.get_extra("btw_loop") == "work" + session = await loop.work_sessions.get_for_origin("origin") + assert session.status.value == "completed" From 0851fb122af9433a8f7514631978dde8cd8b2819 Mon Sep 17 00:00:00 2001 From: YUZHEthefool <2804776511@qq.com> Date: Fri, 11 Sep 2026 00:37:36 +0800 Subject: [PATCH 2/3] fix(btw): restore Chinese work runtime labels Restore the six work enablement, concurrency, and retention translations from their runtime metadata descriptions and hints. AI-Generated: true Generated-At: 2026-09-10T16:37:36Z --- .../i18n/locales/zh-CN/features/config-metadata.json | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json b/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json index 9ce3d70f91..4a42cbf145 100644 --- a/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json +++ b/dashboard/src/i18n/locales/zh-CN/features/config-metadata.json @@ -1178,18 +1178,18 @@ }, "work_loop": { "enabled": { - "description": "??????", - "hint": "????????????????" + "description": "启用工作循环", + "hint": "默认关闭;允许显式工作请求使用工作执行器。" }, "max_concurrent": { - "description": "????????", - "hint": "???????????? 2????????????" + "description": "工作任务执行并发", + "hint": "同时执行的工作任务数,默认 2。此值不是等待队列的长度限制。" } }, "work_session": { "max_age_seconds": { - "description": "??????????", - "hint": "??????????????? 3600 ??" + "description": "终态工作会话保留秒数", + "hint": "已完成、失败或取消的工作会话保留时间,默认 3600 秒。" } } } From e2f20cc598258595b382cb7ba0f5a141692556d1 Mon Sep 17 00:00:00 2001 From: BegoniaHe Date: Thu, 10 Sep 2026 21:00:29 +0200 Subject: [PATCH 3/3] fix(btw): redact detached work failures Prevent raw executor exceptions from reaching tracked-task logging and cover the redaction boundary. AI-Generated: true Generated-At: 2026-09-10T19:02:24Z --- astrbot/core/agent/btw/work_loop.py | 8 +++++++ tests/unit/test_btw_work_loop.py | 33 +++++++++++++++++++++++++++++ 2 files changed, 41 insertions(+) diff --git a/astrbot/core/agent/btw/work_loop.py b/astrbot/core/agent/btw/work_loop.py index c2d565056e..7cc353f5ca 100644 --- a/astrbot/core/agent/btw/work_loop.py +++ b/astrbot/core/agent/btw/work_loop.py @@ -4,8 +4,10 @@ from collections.abc import AsyncGenerator, Awaitable, Callable from typing import Protocol +from astrbot import logger from astrbot.core.message.message_event_result import MessageEventResult from astrbot.core.platform.astr_message_event import AstrMessageEvent +from astrbot.core.utils.error_redaction import safe_error from astrbot.core.utils.task_utils import create_tracked_task from . import i18n as work_i18n @@ -169,5 +171,11 @@ async def _run_detached(self, event: AstrMessageEvent, session_id: str) -> None: try: async for _ in self._execute(event, session_id): await self._result_dispatcher(event) + except asyncio.CancelledError: + raise + except Exception as exc: + # The task registry logs unhandled exceptions with their traceback. + # Consume executor failures here so provider details never reach it. + logger.error("BTW work task failed: %s", safe_error("", exc)) finally: await self._event_finalizer(event) diff --git a/tests/unit/test_btw_work_loop.py b/tests/unit/test_btw_work_loop.py index 89d4b0952f..6a5b179846 100644 --- a/tests/unit/test_btw_work_loop.py +++ b/tests/unit/test_btw_work_loop.py @@ -33,6 +33,13 @@ async def process(self, event): yield +class SecretFailingExecutor: + async def process(self, event): + del event + raise RuntimeError("provider failed: api_key=btw-work-loop-secret") + yield + + class BlockingExecutor: def __init__(self) -> None: self.started = asyncio.Event() @@ -105,6 +112,32 @@ async def test_work_loop_acknowledges_then_runs_in_background(): event_finalizer.assert_awaited_once_with(event) +@pytest.mark.asyncio +async def test_detached_work_redacts_executor_failure_from_logs(caplog): + sessions = WorkSessionManager() + work_loop = WorkLoop(SecretFailingExecutor(), sessions) + background_tasks: set[asyncio.Task] = set() + event_finalizer = AsyncMock() + work_loop.configure_detached_execution( + background_tasks=background_tasks, + result_dispatcher=AsyncMock(), + event_finalizer=event_finalizer, + ) + event = FakeEvent("execute work") + + with caplog.at_level("ERROR", logger="astrbot"): + _ = [item async for item in work_loop.submit(event)] + [task] = background_tasks + await asyncio.wait_for(task, timeout=5) + + session = await sessions.get_for_origin(event.unified_msg_origin) + assert session is not None + assert session.status is WorkSessionStatus.FAILED + assert session.error == "Work task failed." + assert "btw-work-loop-secret" not in caplog.text + event_finalizer.assert_awaited_once_with(event) + + @pytest.mark.parametrize( "config, enabled", [