feat(runtime): 抽出 RuntimeService 公共任务入口,任务线程拆除 Flask 隐式上下文
Gateway 化阶段一/二实施(公共任务入口),计划与验收记录见 cache_research/gateway/phase12_implementation_plan.md;契约文档 docs/runtime_contract.md 为本地文档按仓库惯例不入库 - 新增 server/runtime/:RuntimeContext 三层模型(TrustedPrincipal/TaskParams/InternalDirectives)+ RuntimeService(create_task/cancel/guidance/queue/get_task_events) - create_chat_task 强制显式 session_data(缺失 ValueError,i18n key tasks.missing_session_data);6 处调用点迁移至 runtime_service.create_task()(tasks/api、api_v1、workflow_runtime_api x2、chat_flow_task_main x2),门闸 token 移交与 notice 互斥豁免语义保留 - _run_chat_task 拆除 test_request_context 桥,任务线程全程 RuntimeIdentity 驱动 - 审批超时经 terminal._approval_timeout_seconds 透传至工具循环 4 个 _wait_* 调用点(默认 3600s 语义不变) - 附带:config/_load_dotenv 容忍沙箱禁读 .env;AGENTS.md 同步结构 测试:改造相关 24/24 全绿(冒烟 6 + runtime_service 10 + model_persistence 4 + identity 路由 4);真实环境验证已由用户人工完成 Co-authored-by: Astrion powered by Kimi-K3 <astrion-agent@users.noreply.github.com> Co-authored-by: Codex powered by ChatGPT-6-Astra <codex@example.com>
This commit is contained in:
parent
7d290d1790
commit
652c160640
@ -23,7 +23,8 @@
|
|||||||
- `server/app.py`: 推荐的 Web 服务入口(封装并转发到 `server/app_legacy.py`)
|
- `server/app.py`: 推荐的 Web 服务入口(封装并转发到 `server/app_legacy.py`)
|
||||||
- `web_server.py`: 兼容入口,已标记 deprecated,但仍可启动
|
- `web_server.py`: 兼容入口,已标记 deprecated,但仍可启动
|
||||||
- **后端核心目录**
|
- **后端核心目录**
|
||||||
- `server/`: Flask 业务主线(chat/task/status 已拆分为子包:`server/chat/`、`server/status/`、`server/tasks/`;REST 任务轮询为主,Socket.IO 主要用于兼容与实时辅助通道)
|
- `server/`: Flask 业务主线(chat/task/status/context 已拆分为子包:`server/chat/`、`server/status/`、`server/tasks/`、`server/context/`(用户资源/身份/广播/个性化,原 `server/context.py` 已拆包并保留兼容 re-export);REST 任务轮询为主,Socket.IO 主要用于兼容与实时辅助通道)
|
||||||
|
- `server/runtime/`: 公共任务入口(2026-09 Gateway 化阶段二新增):`context.py` 定义 RuntimeContext 三层模型(TrustedPrincipal/TaskParams/InternalDirectives),`service.py` 提供 RuntimeService(create_task/cancel/guidance/queue/get_task_events);契约见 `docs/runtime_contract.md`
|
||||||
- `core/`: 终端与工具编排(`main_terminal.py`、`web_terminal.py`、`main_terminal_parts/*`;其中 `main_terminal_parts/context/` 和 `main_terminal_parts/tools_definition` 已拆分为 base + mixin 子包)
|
- `core/`: 终端与工具编排(`main_terminal.py`、`web_terminal.py`、`main_terminal_parts/*`;其中 `main_terminal_parts/context/` 和 `main_terminal_parts/tools_definition` 已拆分为 base + mixin 子包)
|
||||||
- `modules/`: 可复用能力模块(terminal/file/memory/sub_agent/upload_security/user 等;`file_manager`、`persistent_terminal`、`terminal_ops`、`mcp_client_manager` 已拆分为子包)
|
- `modules/`: 可复用能力模块(terminal/file/memory/sub_agent/upload_security/user 等;`file_manager`、`persistent_terminal`、`terminal_ops`、`mcp_client_manager` 已拆分为子包)
|
||||||
- `config/`: 配置拆分(`api.py`, `limits.py`, `terminal.py`, `paths.py` ...),由 `config/__init__.py` 聚合并加载 `.env`
|
- `config/`: 配置拆分(`api.py`, `limits.py`, `terminal.py`, `paths.py` ...),由 `config/__init__.py` 聚合并加载 `.env`
|
||||||
|
|||||||
92
cache_research/gateway/phase12_implementation_plan.md
Normal file
92
cache_research/gateway/phase12_implementation_plan.md
Normal file
@ -0,0 +1,92 @@
|
|||||||
|
# 阶段一/二实施计划记录
|
||||||
|
|
||||||
|
> 日期:2026-09-07
|
||||||
|
> 上游文档:`gateway_work_plan.md`(三阶段路线)、`eval_summary.md`(范围与复杂度评估,含 R1-R6 审阅修订)
|
||||||
|
> 范围:**阶段一(契约文档)+ 阶段二(公共任务入口)**;阶段三定时任务未讨论功能与实现,明确不做。
|
||||||
|
|
||||||
|
## 目标
|
||||||
|
|
||||||
|
按 work plan 阶段一/二的完成标准交付:
|
||||||
|
|
||||||
|
- 阶段一:`docs/runtime_contract.md` 落地——状态责任表、概念对齐、公共入口上下文契约、调用方迁移表、回归用例清单。目标入口的状态修改路径和执行裁决可定位;已有保障成为明确约束。
|
||||||
|
- 阶段二:不创建浏览器会话、不伪造 HTTP 请求,能通过显式身份和资源上下文启动一次受控任务、读取结果并取消;同对话重入仍受门闸保护;现有 Web/CLI 行为兼容。
|
||||||
|
|
||||||
|
## 阶段二设计决策(按 R4 审阅意见)
|
||||||
|
|
||||||
|
### RuntimeContext 三层分离
|
||||||
|
|
||||||
|
```
|
||||||
|
TrustedPrincipal 可信身份与资源范围:username / workspace_id / host_mode /
|
||||||
|
host_workspace_id / is_api_user / role
|
||||||
|
——只能由适配层认证后构造,客户端不可自报
|
||||||
|
TaskParams 本次任务参数:message / images / videos / files / model_key /
|
||||||
|
run_mode / thinking_mode / max_iterations / conversation_id /
|
||||||
|
goal_mode / skill_context_messages / message_source
|
||||||
|
InternalDirectives 内部执行信息:main_task_gate_token / auto_user_message_event /
|
||||||
|
auto_user_message_payload / preceding_user_notices /
|
||||||
|
approval_timeout_seconds(仅透传机制,默认语义不变)
|
||||||
|
——普通客户端不可提交,仅内部调用方(通知链/工作流/派发器)使用
|
||||||
|
```
|
||||||
|
|
||||||
|
默认值解析优先级:**本次显式传参 > 对话元数据绑定 > 会话/用户偏好快照 > 系统默认**。
|
||||||
|
会话配置恢复(模型/模式)仍在 `get_user_resources` + 会话加载链路内完成,RuntimeContext 携带的是「本次覆盖」与「身份快照」,不替代既有恢复逻辑。
|
||||||
|
|
||||||
|
### RuntimeService 最小接口集
|
||||||
|
|
||||||
|
- `create_task(principal, params, directives) -> task_id`(受理:互斥裁决 + 登记 + 起执行)
|
||||||
|
- `cancel_task(username, task_id)`
|
||||||
|
- `enqueue_runtime_guidance / enqueue_runtime_pending_message / remove_runtime_pending_message / promote_runtime_pending_to_guidance`
|
||||||
|
- `get_task / get_task_events(username, task_id, offset)`(内部查询接口,CLI/定时任务不必走 HTTP)
|
||||||
|
- 审批回答复用现有三个 manager,通过明确关联接入,不进 RuntimeService 首版。
|
||||||
|
|
||||||
|
### 两步走
|
||||||
|
|
||||||
|
- **步骤①(兼容期)**:RuntimeContext + RuntimeService 骨架落地,`create_chat_task` 接受显式上下文;6 处调用点迁移;`test_request_context` 桥**保留**作为兜底(显式快照灌回 session 的既有行为不变)。
|
||||||
|
- **步骤②(拆桥)**:`get_user_resources` 增加显式参数变体(principal + 快照参数),`_apply_workspace_personalization_preferences` 参数化,`auth_helpers` 提供显式 record/role 入口;`_run_chat_task` 改走显式路径后删除 test_request_context 桥;`ensure_conversation_loaded` 的 session 回写移出任务路径。
|
||||||
|
|
||||||
|
## 关键风险与对策
|
||||||
|
|
||||||
|
1. `get_user_resources` 参数化(最高风险):host_mode/is_api_user 分支选错 → 静默串工作区。对策:显式变体与现有 web 路径并存,先任务线程单点切换,回归验证后再推广。
|
||||||
|
2. 门闸 token 移交语义:InternalDirectives 原样承载;通知链「预占→移交→认领→失败回滚」不变。
|
||||||
|
3. `task_type="notice"` 互斥豁免保留。
|
||||||
|
4. 异常回退路径(chat_flow_task_main.py:645-675 直接执行 handle_task_with_sender):保留语义,验收覆盖其门闸/事件/取消生命周期。
|
||||||
|
5. 事件前台干扰:阶段二不加 source 字段(阶段三产品决策),现有行为不变。
|
||||||
|
|
||||||
|
## 验收
|
||||||
|
|
||||||
|
- `python3 -m py_compile` 触及文件全过
|
||||||
|
- `python -m pytest test/test_server_refactor_smoke.py -q` 通过
|
||||||
|
- 新增测试:显式上下文(无 HTTP 请求)受理任务 → 拒绝同对话并发 chat 任务 → 取消;不触达真实模型调用
|
||||||
|
- 服务重启与人工验证交用户执行(不擅自重启 8091/8092)
|
||||||
|
|
||||||
|
## 实施结果(2026-09-07 完成,待 commit)
|
||||||
|
|
||||||
|
**状态:阶段一/二代码工作全部完成,改造相关测试 24/24 全绿;工作区改动未提交。**
|
||||||
|
|
||||||
|
### 阶段一交付
|
||||||
|
|
||||||
|
- `docs/runtime_contract.md`:概念对齐(Session/Run/Schedule/Occurrence/Event)、10 项状态责任表、已有保障约束清单、RuntimeService 契约(含 RuntimeContext 三层模型定义)、调用方迁移表、12 项回归用例 T01-T12。
|
||||||
|
|
||||||
|
### 阶段二交付
|
||||||
|
|
||||||
|
- 新建 `server/runtime/` 包:
|
||||||
|
- `context.py`:RuntimeContext 三层模型(TrustedPrincipal / TaskParams / InternalDirectives),`from_terminal()`、`principal_from_session_snapshot()`、`to_session_data()`。
|
||||||
|
- `service.py`:`RuntimeService`(create_task / cancel / guidance / queue / get_task_events)+ 进程级单例 `runtime_service`。
|
||||||
|
- `TaskManager.create_chat_task` 强制显式 session_data(缺失即 `ValueError`,i18n key `tasks.missing_session_data`)。
|
||||||
|
- 6 处调用点全部迁移到 `runtime_service.create_task()`:`server/tasks/api.py`、`server/api_v1.py`、`server/workflow_runtime_api.py`×2、`server/chat_flow_task_main.py`×2(完成通知派发 :613、多智能体 idle 派发 :1416)。门闸 token 移交、`task_type="notice"` 互斥豁免语义原样保留。
|
||||||
|
- `server/context.py`(989 行)拆分为 `server/context/` 子包:identity / broadcast / personalization / usage / upload / conversation / resources / decorators / reaper 共 9 模块 + `__init__.py` 兼容 re-export(外部 import 路径不变)。
|
||||||
|
- `get_user_resources` 参数化:新增 `RuntimeIdentity` 显式身份快照(定义在 `server/context/identity.py`,runtime 包引用之,依赖方向自下而上无循环)。
|
||||||
|
- **test_request_context 桥已拆除**:`server/tasks/models.py::_run_chat_task` 不再建立 Flask 请求上下文,任务线程全程 RuntimeIdentity 驱动。
|
||||||
|
- 审批超时透传管道:terminal 属性 `_approval_timeout_seconds`(默认语义 3600s 不变),`_run_chat_task` setattr → `server/chat_flow_tool_loop.py::_approval_timeout_for()` helper → 4 个 `_wait_*` 调用点。
|
||||||
|
- 附带修复:`config/_load_dotenv` 对禁读 `.env` 的沙箱环境加 try/except(不再 PermissionError 崩溃)。
|
||||||
|
|
||||||
|
### 测试验收
|
||||||
|
|
||||||
|
- 改造相关 24/24 全绿:`test_server_refactor_smoke`(6)+ `test_runtime_service`(10,新增)+ `test_conversation_model_persistence`(4,patch 目标随迁)+ `test_runtime_identity_resources`(4,新增,覆盖 get_user_resources 的 web/host/api 身份路由)。
|
||||||
|
- 存量失败 4 项(conversation_workspace_storage / host_workspace_manager / skills_manager / token_usage_extractor)经甄别与本次改动零相关,未修。
|
||||||
|
|
||||||
|
### 遗留待办
|
||||||
|
|
||||||
|
1. 真实运行环境验证(Web 聊天 / 停止 / 审批 / workflow 激活 / 多智能体派发)需用户重启服务后人工完成——get_user_resources 分支选错会静默串工作区,这是最高风险点。
|
||||||
|
2. 审批条目 task_id 恒 None(静态疑点,待运行时验证)。
|
||||||
|
3. socket 软 stop 不打断审批等待(REST 硬取消可以),留待阶段三或独立决策。
|
||||||
@ -23,7 +23,13 @@ def _load_dotenv():
|
|||||||
else:
|
else:
|
||||||
env_path = Path(__file__).resolve().parents[1] / '.env'
|
env_path = Path(__file__).resolve().parents[1] / '.env'
|
||||||
env_from_file: dict = {}
|
env_from_file: dict = {}
|
||||||
if env_path.exists():
|
try:
|
||||||
|
# 权限受限环境(只读沙箱白名单禁读 .env)下 exists() 会抛 PermissionError,
|
||||||
|
# 与「文件不存在」同等降级处理,而不是崩溃。
|
||||||
|
env_exists = env_path.exists()
|
||||||
|
except Exception:
|
||||||
|
env_exists = False
|
||||||
|
if env_exists:
|
||||||
try:
|
try:
|
||||||
for raw_line in env_path.read_text(encoding="utf-8").splitlines():
|
for raw_line in env_path.read_text(encoding="utf-8").splitlines():
|
||||||
line = raw_line.strip()
|
line = raw_line.strip()
|
||||||
|
|||||||
@ -102,6 +102,10 @@ MESSAGES = {
|
|||||||
"zh-CN": "当前对话已有运行中的任务,请稍后再试。",
|
"zh-CN": "当前对话已有运行中的任务,请稍后再试。",
|
||||||
"en-US": "This conversation already has a running task. Please try again later.",
|
"en-US": "This conversation already has a running task. Please try again later.",
|
||||||
},
|
},
|
||||||
|
"tasks.missing_session_data": {
|
||||||
|
"zh-CN": "缺少任务运行上下文(session_data),请通过公共任务入口提交",
|
||||||
|
"en-US": "Missing runtime context (session_data); please submit via the runtime service entry",
|
||||||
|
},
|
||||||
"tasks.system_not_initialized": {
|
"tasks.system_not_initialized": {
|
||||||
"zh-CN": "系统未初始化",
|
"zh-CN": "系统未初始化",
|
||||||
"en-US": "System not initialized",
|
"en-US": "System not initialized",
|
||||||
|
|||||||
@ -10,6 +10,7 @@ from flask import Blueprint, request, jsonify, send_file, session
|
|||||||
from .api_auth import api_token_required
|
from .api_auth import api_token_required
|
||||||
from .security import rate_limited
|
from .security import rate_limited
|
||||||
from .tasks import task_manager
|
from .tasks import task_manager
|
||||||
|
from server.runtime import RuntimeContext, TaskParams, principal_from_session_snapshot, runtime_service
|
||||||
from .context import get_user_resources, ensure_conversation_loaded, get_upload_guard, apply_conversation_overrides
|
from .context import get_user_resources, ensure_conversation_loaded, get_upload_guard, apply_conversation_overrides
|
||||||
from .utils_common import sanitize_filename_preserve_unicode
|
from .utils_common import sanitize_filename_preserve_unicode
|
||||||
from .utils_common import debug_log
|
from .utils_common import debug_log
|
||||||
@ -316,18 +317,22 @@ def send_message_api(workspace_id: str):
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
return jsonify({"success": False, "error": tr("api_v1.custom_params_error", error=exc)}), 400
|
return jsonify({"success": False, "error": tr("api_v1.custom_params_error", error=exc)}), 400
|
||||||
|
|
||||||
|
# 公共任务入口:principal 携带 api_token_required 注入的身份
|
||||||
|
# (is_api_user=True / role="api"),丢失会静默串资源管理器(契约 §5)。
|
||||||
try:
|
try:
|
||||||
rec = task_manager.create_chat_task(
|
ctx = RuntimeContext(
|
||||||
username=username,
|
principal=principal_from_session_snapshot(session, ws.workspace_id, username=username),
|
||||||
workspace_id=ws.workspace_id,
|
params=TaskParams(
|
||||||
message=message,
|
message=message,
|
||||||
images=images,
|
images=images,
|
||||||
conversation_id=conversation_id,
|
conversation_id=conversation_id,
|
||||||
model_key=model_key,
|
model_key=model_key,
|
||||||
thinking_mode=thinking_mode,
|
thinking_mode=thinking_mode,
|
||||||
run_mode=run_mode,
|
run_mode=run_mode,
|
||||||
max_iterations=max_iterations,
|
max_iterations=max_iterations,
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
rec = runtime_service.create_task(ctx)
|
||||||
except ValueError as exc:
|
except ValueError as exc:
|
||||||
return jsonify({"success": False, "error": str(exc)}), 400
|
return jsonify({"success": False, "error": str(exc)}), 400
|
||||||
except RuntimeError as exc:
|
except RuntimeError as exc:
|
||||||
|
|||||||
@ -43,6 +43,7 @@ from modules.upload_security import UploadSecurityError
|
|||||||
from modules.user_manager import UserWorkspace
|
from modules.user_manager import UserWorkspace
|
||||||
from modules.usage_tracker import QUOTA_DEFAULTS
|
from modules.usage_tracker import QUOTA_DEFAULTS
|
||||||
from modules.sub_agent import TERMINAL_STATUSES
|
from modules.sub_agent import TERMINAL_STATUSES
|
||||||
|
from server.runtime import InternalDirectives, RuntimeContext, TaskParams, runtime_service
|
||||||
from modules.multi_agent.debug_logger import ma_debug
|
from modules.multi_agent.debug_logger import ma_debug
|
||||||
from modules.versioning_manager import ConversationVersioningManager, VersioningError
|
from modules.versioning_manager import ConversationVersioningManager, VersioningError
|
||||||
from modules.shallow_versioning import ShallowVersioningManager
|
from modules.shallow_versioning import ShallowVersioningManager
|
||||||
@ -564,39 +565,21 @@ async def _dispatch_completion_user_notice(
|
|||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
from .tasks import task_manager
|
|
||||||
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
|
||||||
host_mode = bool(getattr(workspace, "username", None) == "host")
|
|
||||||
session_data = {
|
|
||||||
"username": username,
|
|
||||||
"role": getattr(web_terminal, "user_role", "user"),
|
|
||||||
"is_api_user": getattr(web_terminal, "user_role", "") == "api",
|
|
||||||
"host_mode": host_mode,
|
|
||||||
"host_workspace_id": workspace_id if host_mode else None,
|
|
||||||
"workspace_id": workspace_id,
|
|
||||||
"run_mode": getattr(web_terminal, "run_mode", None),
|
|
||||||
"thinking_mode": getattr(web_terminal, "thinking_mode", None),
|
|
||||||
"model_key": getattr(web_terminal, "model_key", None),
|
|
||||||
"message_source": message_source,
|
|
||||||
}
|
|
||||||
# 派发方预占的主任务门闸 token 随任务移交,由任务线程认领(见 process_message_task)
|
|
||||||
if main_task_gate_token:
|
|
||||||
session_data["main_task_gate_token"] = main_task_gate_token
|
|
||||||
ui_defaults = _user_message_ui_defaults(
|
ui_defaults = _user_message_ui_defaults(
|
||||||
message_source,
|
message_source,
|
||||||
auto_user_message_event=True,
|
auto_user_message_event=True,
|
||||||
)
|
)
|
||||||
# 关键:通知类后台任务需要把 user_message 写入任务事件流,
|
# 关键:通知类后台任务需要把 user_message 写入任务事件流,
|
||||||
# 否则前端轮询只会看到 AI/tool 事件,看不到 user_message。
|
# 否则前端轮询只会看到 AI/tool 事件,看不到 user_message。
|
||||||
session_data["auto_user_message_event"] = True
|
auto_payload = {
|
||||||
session_data["auto_user_message_payload"] = {
|
|
||||||
**dict(extra_payload or {}),
|
**dict(extra_payload or {}),
|
||||||
**ui_defaults,
|
**ui_defaults,
|
||||||
"timestamp": datetime.now().isoformat(),
|
"timestamp": datetime.now().isoformat(),
|
||||||
}
|
}
|
||||||
# 轮询客户端通过后续任务事件流回放前置通知(在线客户端已由上面的 socketio 回显覆盖)。
|
# 轮询客户端通过后续任务事件流回放前置通知(在线客户端已由上面的 socketio 回显覆盖)。
|
||||||
|
preceding_list = None
|
||||||
if preceding_notices:
|
if preceding_notices:
|
||||||
session_data["preceding_user_notices"] = [
|
preceding_list = [
|
||||||
{
|
{
|
||||||
"message": str(item.get("message") or ""),
|
"message": str(item.get("message") or ""),
|
||||||
"payload": dict(item.get("payload") or {}),
|
"payload": dict(item.get("payload") or {}),
|
||||||
@ -604,23 +587,35 @@ async def _dispatch_completion_user_notice(
|
|||||||
for item in preceding_notices
|
for item in preceding_notices
|
||||||
if str(item.get("message") or "").strip()
|
if str(item.get("message") or "").strip()
|
||||||
]
|
]
|
||||||
|
# 公共任务入口:内部调用方,从 terminal 构造可信上下文;
|
||||||
|
# 模型/模式覆盖取 terminal 当前值(与偏好快照同值,对齐既有通知链语义)。
|
||||||
|
ctx = RuntimeContext.from_terminal(
|
||||||
|
web_terminal,
|
||||||
|
workspace,
|
||||||
|
username,
|
||||||
|
params=TaskParams(
|
||||||
|
message=user_message,
|
||||||
|
conversation_id=conversation_id,
|
||||||
|
message_source=message_source,
|
||||||
|
model_key=getattr(web_terminal, "model_key", None),
|
||||||
|
run_mode=getattr(web_terminal, "run_mode", None),
|
||||||
|
thinking_mode=getattr(web_terminal, "thinking_mode", None),
|
||||||
|
),
|
||||||
|
directives=InternalDirectives(
|
||||||
|
# 派发方预占的主任务门闸 token 随任务移交,由任务线程认领
|
||||||
|
main_task_gate_token=main_task_gate_token,
|
||||||
|
auto_user_message_event=True,
|
||||||
|
auto_user_message_payload=auto_payload,
|
||||||
|
preceding_user_notices=preceding_list,
|
||||||
|
),
|
||||||
|
)
|
||||||
ma_debug(
|
ma_debug(
|
||||||
"dispatch_completion_create_chat_task",
|
"dispatch_completion_create_chat_task",
|
||||||
conversation_id=conversation_id,
|
conversation_id=conversation_id,
|
||||||
user_message_preview=user_message[:300],
|
user_message_preview=user_message[:300],
|
||||||
preceding_count=len(preceding_notices),
|
preceding_count=len(preceding_notices),
|
||||||
)
|
)
|
||||||
rec = task_manager.create_chat_task(
|
rec = runtime_service.create_task(ctx)
|
||||||
username,
|
|
||||||
workspace_id,
|
|
||||||
user_message,
|
|
||||||
[],
|
|
||||||
conversation_id,
|
|
||||||
model_key=session_data.get("model_key"),
|
|
||||||
thinking_mode=session_data.get("thinking_mode"),
|
|
||||||
run_mode=session_data.get("run_mode"),
|
|
||||||
session_data=session_data,
|
|
||||||
)
|
|
||||||
ma_debug(
|
ma_debug(
|
||||||
"dispatch_completion_task_created",
|
"dispatch_completion_task_created",
|
||||||
conversation_id=conversation_id,
|
conversation_id=conversation_id,
|
||||||
@ -1347,25 +1342,11 @@ async def _dispatch_multi_agent_idle_messages(
|
|||||||
ui_defaults["starts_work"] = False
|
ui_defaults["starts_work"] = False
|
||||||
|
|
||||||
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
||||||
host_mode = bool(getattr(workspace, "username", None) == "host")
|
|
||||||
session_data = {
|
|
||||||
"username": username,
|
|
||||||
"role": getattr(web_terminal, "user_role", "user"),
|
|
||||||
"is_api_user": getattr(web_terminal, "user_role", "") == "api",
|
|
||||||
"host_mode": host_mode,
|
|
||||||
"host_workspace_id": workspace_id if host_mode else None,
|
|
||||||
"workspace_id": workspace_id,
|
|
||||||
"run_mode": getattr(web_terminal, "run_mode", None),
|
|
||||||
"thinking_mode": getattr(web_terminal, "thinking_mode", None),
|
|
||||||
"model_key": getattr(web_terminal, "model_key", None),
|
|
||||||
"message_source": message_source,
|
|
||||||
}
|
|
||||||
ma_auto_type = _auto_message_type_for_multi_agent_subtype(last["subtype"])
|
ma_auto_type = _auto_message_type_for_multi_agent_subtype(last["subtype"])
|
||||||
# 重要:ui_defaults 默认会给出 visibility="compact"/starts_work=False,
|
# 重要:ui_defaults 默认会给出 visibility="compact"/starts_work=False,
|
||||||
# 多智能体主消息是正常聊天消息,必须在 **ui_defaults 之后覆写为 chat,
|
# 多智能体主消息是正常聊天消息,必须在 **ui_defaults 之后覆写为 chat,
|
||||||
# 否则后端历史 metadata.visibility=compact 会让加载的消息走通知渲染。
|
# 否则后端历史 metadata.visibility=compact 会让加载的消息走通知渲染。
|
||||||
session_data["auto_user_message_event"] = True
|
auto_payload = {
|
||||||
session_data["auto_user_message_payload"] = {
|
|
||||||
**ui_defaults,
|
**ui_defaults,
|
||||||
"message_source": message_source,
|
"message_source": message_source,
|
||||||
"sub_agent_notice": True,
|
"sub_agent_notice": True,
|
||||||
@ -1384,7 +1365,7 @@ async def _dispatch_multi_agent_idle_messages(
|
|||||||
workspace_id=workspace_id,
|
workspace_id=workspace_id,
|
||||||
username=username,
|
username=username,
|
||||||
)
|
)
|
||||||
session_data["auto_user_message_payload"].setdefault("metadata", {
|
auto_payload.setdefault("metadata", {
|
||||||
**ui_defaults,
|
**ui_defaults,
|
||||||
"message_source": message_source,
|
"message_source": message_source,
|
||||||
"auto_message_type": ma_auto_type,
|
"auto_message_type": ma_auto_type,
|
||||||
@ -1394,8 +1375,9 @@ async def _dispatch_multi_agent_idle_messages(
|
|||||||
"starts_work": False,
|
"starts_work": False,
|
||||||
"visibility": "chat",
|
"visibility": "chat",
|
||||||
})
|
})
|
||||||
|
preceding_list = None
|
||||||
if len(parsed_messages) > 1:
|
if len(parsed_messages) > 1:
|
||||||
session_data["preceding_user_notices"] = [
|
preceding_list = [
|
||||||
{
|
{
|
||||||
"message": item["text"],
|
"message": item["text"],
|
||||||
"payload": {
|
"payload": {
|
||||||
@ -1413,18 +1395,26 @@ async def _dispatch_multi_agent_idle_messages(
|
|||||||
]
|
]
|
||||||
|
|
||||||
try:
|
try:
|
||||||
rec = task_manager.create_chat_task(
|
ctx = RuntimeContext.from_terminal(
|
||||||
|
web_terminal,
|
||||||
|
workspace,
|
||||||
username,
|
username,
|
||||||
workspace_id,
|
params=TaskParams(
|
||||||
last["text"],
|
message=last["text"],
|
||||||
[],
|
conversation_id=conversation_id,
|
||||||
conversation_id,
|
message_source=message_source,
|
||||||
model_key=session_data.get("model_key"),
|
model_key=getattr(web_terminal, "model_key", None),
|
||||||
thinking_mode=session_data.get("thinking_mode"),
|
run_mode=getattr(web_terminal, "run_mode", None),
|
||||||
run_mode=session_data.get("run_mode"),
|
thinking_mode=getattr(web_terminal, "thinking_mode", None),
|
||||||
session_data=session_data,
|
task_type="notice",
|
||||||
task_type="notice",
|
),
|
||||||
|
directives=InternalDirectives(
|
||||||
|
auto_user_message_event=True,
|
||||||
|
auto_user_message_payload=auto_payload,
|
||||||
|
preceding_user_notices=preceding_list,
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
rec = runtime_service.create_task(ctx)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
ma_debug(
|
ma_debug(
|
||||||
"dispatch_ma_idle_create_task_exception",
|
"dispatch_ma_idle_create_task_exception",
|
||||||
|
|||||||
@ -229,6 +229,22 @@ def _format_rejected_tool_text(reason: str) -> str:
|
|||||||
return tr("tool_loop.tool_call_rejected", reason=clean_reason)
|
return tr("tool_loop.tool_call_rejected", reason=clean_reason)
|
||||||
|
|
||||||
|
|
||||||
|
def _approval_timeout_for(web_terminal) -> Optional[float]:
|
||||||
|
"""任务级审批/提问等待超时(秒),来自受理时的 RuntimeContext 透传。
|
||||||
|
|
||||||
|
未设置/非法值 = None → 调用点回退默认 3600s(既有语义不变)。
|
||||||
|
超时后的语义(拒绝当前动作继续 vs 结束任务)属阶段三产品决策。
|
||||||
|
"""
|
||||||
|
try:
|
||||||
|
value = getattr(web_terminal, "_approval_timeout_seconds", None)
|
||||||
|
if value is None:
|
||||||
|
return None
|
||||||
|
value = float(value)
|
||||||
|
return value if value > 0 else None
|
||||||
|
except Exception:
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
async def _wait_for_tool_approval(*, approval_id: str, username: str, timeout_seconds: float = 3600.0) -> Dict[str, Any]:
|
async def _wait_for_tool_approval(*, approval_id: str, username: str, timeout_seconds: float = 3600.0) -> Dict[str, Any]:
|
||||||
started = time.time()
|
started = time.time()
|
||||||
while True:
|
while True:
|
||||||
@ -432,7 +448,11 @@ async def _handle_submit_plan(*, web_terminal, arguments: Dict[str, Any], sender
|
|||||||
})
|
})
|
||||||
|
|
||||||
# 4. 阻塞等待用户决定
|
# 4. 阻塞等待用户决定
|
||||||
resolved = await _wait_for_plan_approval(approval_id=str(approval.get("approval_id") or ""), username=username)
|
resolved = await _wait_for_plan_approval(
|
||||||
|
approval_id=str(approval.get("approval_id") or ""),
|
||||||
|
username=username,
|
||||||
|
timeout_seconds=_approval_timeout_for(web_terminal) or 3600.0,
|
||||||
|
)
|
||||||
status = str(resolved.get("status") or "")
|
status = str(resolved.get("status") or "")
|
||||||
comment = str(resolved.get("comment") or "").strip()
|
comment = str(resolved.get("comment") or "").strip()
|
||||||
sender('plan_approval_resolved', {
|
sender('plan_approval_resolved', {
|
||||||
@ -574,6 +594,7 @@ async def _execute_tool_calls_impl(*, web_terminal, tool_calls, sender, messages
|
|||||||
wait_answers = await _wait_for_user_questions(
|
wait_answers = await _wait_for_user_questions(
|
||||||
question_ids=[str(q.get("question_id") or "") for q in created_questions],
|
question_ids=[str(q.get("question_id") or "") for q in created_questions],
|
||||||
username=username,
|
username=username,
|
||||||
|
timeout_seconds=_approval_timeout_for(web_terminal) or 3600.0,
|
||||||
)
|
)
|
||||||
for question in created_questions:
|
for question in created_questions:
|
||||||
qid = str(question.get("question_id") or "")
|
qid = str(question.get("question_id") or "")
|
||||||
@ -883,6 +904,7 @@ async def _execute_tool_calls_impl(*, web_terminal, tool_calls, sender, messages
|
|||||||
wait_result = await _wait_for_tool_approval(
|
wait_result = await _wait_for_tool_approval(
|
||||||
approval_id=approval_item.get("approval_id"),
|
approval_id=approval_item.get("approval_id"),
|
||||||
username=username,
|
username=username,
|
||||||
|
timeout_seconds=_approval_timeout_for(web_terminal) or 3600.0,
|
||||||
)
|
)
|
||||||
sender('tool_approval_resolved', {
|
sender('tool_approval_resolved', {
|
||||||
'approval_id': approval_item.get("approval_id"),
|
'approval_id': approval_item.get("approval_id"),
|
||||||
@ -1153,6 +1175,7 @@ async def _execute_tool_calls_impl(*, web_terminal, tool_calls, sender, messages
|
|||||||
wait_result = await _wait_for_tool_approval(
|
wait_result = await _wait_for_tool_approval(
|
||||||
approval_id=approval_item.get("approval_id"),
|
approval_id=approval_item.get("approval_id"),
|
||||||
username=username,
|
username=username,
|
||||||
|
timeout_seconds=_approval_timeout_for(web_terminal) or 3600.0,
|
||||||
)
|
)
|
||||||
sender('tool_approval_resolved', {
|
sender('tool_approval_resolved', {
|
||||||
'approval_id': approval_item.get("approval_id"),
|
'approval_id': approval_item.get("approval_id"),
|
||||||
|
|||||||
19
server/runtime/__init__.py
Normal file
19
server/runtime/__init__.py
Normal file
@ -0,0 +1,19 @@
|
|||||||
|
# server/runtime/__init__.py - 公共任务入口包(契约 docs/runtime_contract.md)
|
||||||
|
from server.runtime.context import (
|
||||||
|
InternalDirectives,
|
||||||
|
RuntimeContext,
|
||||||
|
TaskParams,
|
||||||
|
TrustedPrincipal,
|
||||||
|
principal_from_session_snapshot,
|
||||||
|
)
|
||||||
|
from server.runtime.service import RuntimeService, runtime_service
|
||||||
|
|
||||||
|
__all__ = [
|
||||||
|
"InternalDirectives",
|
||||||
|
"RuntimeContext",
|
||||||
|
"RuntimeService",
|
||||||
|
"TaskParams",
|
||||||
|
"TrustedPrincipal",
|
||||||
|
"principal_from_session_snapshot",
|
||||||
|
"runtime_service",
|
||||||
|
]
|
||||||
197
server/runtime/context.py
Normal file
197
server/runtime/context.py
Normal file
@ -0,0 +1,197 @@
|
|||||||
|
"""运行时上下文模型(RuntimeContext 三层分离)。
|
||||||
|
|
||||||
|
契约见 docs/runtime_contract.md §4.1:
|
||||||
|
|
||||||
|
- ``TrustedPrincipal``:可信身份与资源范围。只能由适配层在完成认证后构造;
|
||||||
|
禁止接受客户端或未来 Schedule payload 自报的 role/is_api_user。
|
||||||
|
- ``TaskParams``:本次任务参数(消息、媒体、模型/模式覆盖等)。
|
||||||
|
- ``InternalDirectives``:内部执行信息(门闸 token、通知回放等),
|
||||||
|
仅内部调用方(通知链/工作流/派发器)使用,普通客户端不可提交。
|
||||||
|
|
||||||
|
默认值解析优先级:本次显式传参(TaskParams)> 对话元数据绑定 >
|
||||||
|
用户偏好快照(TrustedPrincipal.preferred_*)> 系统默认。
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from dataclasses import dataclass, field
|
||||||
|
from typing import Any, Dict, List, Optional
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class TrustedPrincipal:
|
||||||
|
"""可信身份与资源范围(适配层认证后构造)。
|
||||||
|
|
||||||
|
preferred_* 是「用户偏好快照」:资源装配(get_user_resources)新建 terminal 时
|
||||||
|
的默认值来源;与 TaskParams 中的「本次覆盖」字段语义不同,不要混用。
|
||||||
|
"""
|
||||||
|
|
||||||
|
username: str
|
||||||
|
workspace_id: str
|
||||||
|
role: str = "user"
|
||||||
|
is_api_user: bool = False
|
||||||
|
host_mode: bool = False
|
||||||
|
host_workspace_id: Optional[str] = None
|
||||||
|
preferred_model_key: Optional[str] = None
|
||||||
|
preferred_run_mode: Optional[str] = None
|
||||||
|
preferred_thinking_mode: Optional[bool] = None
|
||||||
|
|
||||||
|
def validate(self) -> None:
|
||||||
|
if not str(self.username or "").strip():
|
||||||
|
raise ValueError("runtime_context: username 不能为空")
|
||||||
|
if not str(self.workspace_id or "").strip():
|
||||||
|
raise ValueError("runtime_context: workspace_id 不能为空")
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class TaskParams:
|
||||||
|
"""本次任务参数(优先级最高的覆盖层)。"""
|
||||||
|
|
||||||
|
message: str = ""
|
||||||
|
images: List[Any] = field(default_factory=list)
|
||||||
|
videos: List[Any] = field(default_factory=list)
|
||||||
|
files: List[str] = field(default_factory=list)
|
||||||
|
conversation_id: Optional[str] = None
|
||||||
|
model_key: Optional[str] = None
|
||||||
|
run_mode: Optional[str] = None
|
||||||
|
thinking_mode: Optional[bool] = None
|
||||||
|
max_iterations: Optional[int] = None
|
||||||
|
goal_mode: bool = False
|
||||||
|
skill_context_messages: List[Dict[str, str]] = field(default_factory=list)
|
||||||
|
message_source: Optional[str] = None
|
||||||
|
task_type: str = "chat"
|
||||||
|
# 审批/提问等待超时透传(秒)。None = 保持既有默认语义(3600s)不变;
|
||||||
|
# 超时后的语义(拒绝当前动作继续 vs 结束任务)属阶段三产品决策,
|
||||||
|
# 本阶段仅建立透传机制,适配层暂不向 Web 客户端暴露该字段。
|
||||||
|
approval_timeout_seconds: Optional[int] = None
|
||||||
|
|
||||||
|
def validate(self) -> None:
|
||||||
|
if self.run_mode is not None:
|
||||||
|
normalized = str(self.run_mode).lower()
|
||||||
|
if normalized not in {"fast", "thinking", "deep"}:
|
||||||
|
raise ValueError("runtime_context: run_mode 只支持 fast/thinking/deep")
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class InternalDirectives:
|
||||||
|
"""内部执行信息。仅内部调用方使用;HTTP 适配层禁止从请求体构造本结构。"""
|
||||||
|
|
||||||
|
# 派发方预占的主任务门闸 token,随任务移交由任务线程认领
|
||||||
|
main_task_gate_token: Optional[str] = None
|
||||||
|
# 通知类任务:把 user 消息写入任务事件流,保证轮询客户端/刷新后可见
|
||||||
|
auto_user_message_event: bool = False
|
||||||
|
auto_user_message_payload: Optional[Dict[str, Any]] = None
|
||||||
|
# 本批通知池中除触发消息外的前置通知,随任务事件流回放
|
||||||
|
preceding_user_notices: Optional[List[Dict[str, Any]]] = None
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True)
|
||||||
|
class RuntimeContext:
|
||||||
|
"""一轮 Run 的完整显式上下文:可信身份 + 任务参数 + 内部指令。"""
|
||||||
|
|
||||||
|
principal: TrustedPrincipal
|
||||||
|
params: TaskParams
|
||||||
|
directives: InternalDirectives = field(default_factory=InternalDirectives)
|
||||||
|
|
||||||
|
def validate(self) -> None:
|
||||||
|
self.principal.validate()
|
||||||
|
self.params.validate()
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def from_terminal(
|
||||||
|
cls,
|
||||||
|
terminal: Any,
|
||||||
|
workspace: Any,
|
||||||
|
username: str,
|
||||||
|
params: TaskParams,
|
||||||
|
directives: Optional[InternalDirectives] = None,
|
||||||
|
) -> "RuntimeContext":
|
||||||
|
"""从对话级 terminal/工作区构造(通知链、多智能体派发等内部调用方)。
|
||||||
|
|
||||||
|
与既有 session_data 手工构造逐字段对齐:
|
||||||
|
host_mode 取 workspace.username == "host";偏好快照取 terminal 当前值。
|
||||||
|
"""
|
||||||
|
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
||||||
|
host_mode = bool(getattr(workspace, "username", None) == "host")
|
||||||
|
role = getattr(terminal, "user_role", None) or "user"
|
||||||
|
return cls(
|
||||||
|
principal=TrustedPrincipal(
|
||||||
|
username=username,
|
||||||
|
workspace_id=workspace_id,
|
||||||
|
role=role,
|
||||||
|
is_api_user=(role == "api"),
|
||||||
|
host_mode=host_mode,
|
||||||
|
host_workspace_id=workspace_id if host_mode else None,
|
||||||
|
preferred_model_key=getattr(terminal, "model_key", None),
|
||||||
|
preferred_run_mode=getattr(terminal, "run_mode", None),
|
||||||
|
preferred_thinking_mode=getattr(terminal, "thinking_mode", None),
|
||||||
|
),
|
||||||
|
params=params,
|
||||||
|
directives=directives or InternalDirectives(),
|
||||||
|
)
|
||||||
|
|
||||||
|
def to_session_data(self) -> Dict[str, Any]:
|
||||||
|
"""兼容转换:合并为现有 ``create_chat_task`` 的 session_data 快照 dict。
|
||||||
|
|
||||||
|
快照在受理时固化、随任务线程传递:身份/偏好由 ``_run_chat_task`` 还原为
|
||||||
|
``RuntimeIdentity`` 驱动资源装配(无 Flask 隐式上下文);门闸移交
|
||||||
|
(main_task_gate_token)、事件注入(auto_user_message_*)、terminal 属性
|
||||||
|
设置(message_source/goal_mode/skill_context_messages)语义不变。
|
||||||
|
|
||||||
|
其中 run_mode/thinking_mode/model_key 取「用户偏好快照」层
|
||||||
|
(principal.preferred_*),供资源装配新建 terminal 时恢复默认值;
|
||||||
|
本次覆盖值(params.*)由 RuntimeService.create_task 走显式参数传递,
|
||||||
|
不进入本快照。
|
||||||
|
"""
|
||||||
|
p = self.principal
|
||||||
|
session_data: Dict[str, Any] = {
|
||||||
|
"username": p.username,
|
||||||
|
"role": p.role,
|
||||||
|
"is_api_user": p.is_api_user,
|
||||||
|
"host_mode": p.host_mode,
|
||||||
|
"host_workspace_id": p.host_workspace_id or (p.workspace_id if p.host_mode else None),
|
||||||
|
"workspace_id": p.workspace_id,
|
||||||
|
"run_mode": p.preferred_run_mode,
|
||||||
|
"thinking_mode": p.preferred_thinking_mode,
|
||||||
|
"model_key": p.preferred_model_key,
|
||||||
|
"message_source": self.params.message_source,
|
||||||
|
"goal_mode": bool(self.params.goal_mode),
|
||||||
|
"skill_context_messages": list(self.params.skill_context_messages or []),
|
||||||
|
}
|
||||||
|
if self.params.approval_timeout_seconds is not None:
|
||||||
|
session_data["approval_timeout_seconds"] = int(self.params.approval_timeout_seconds)
|
||||||
|
d = self.directives
|
||||||
|
if d.main_task_gate_token:
|
||||||
|
session_data["main_task_gate_token"] = d.main_task_gate_token
|
||||||
|
if d.auto_user_message_event:
|
||||||
|
session_data["auto_user_message_event"] = True
|
||||||
|
if d.auto_user_message_payload:
|
||||||
|
session_data["auto_user_message_payload"] = dict(d.auto_user_message_payload)
|
||||||
|
if d.preceding_user_notices:
|
||||||
|
session_data["preceding_user_notices"] = list(d.preceding_user_notices)
|
||||||
|
return session_data
|
||||||
|
|
||||||
|
|
||||||
|
def principal_from_session_snapshot(
|
||||||
|
session_snapshot: Dict[str, Any],
|
||||||
|
workspace_id: str,
|
||||||
|
username: Optional[str] = None,
|
||||||
|
) -> TrustedPrincipal:
|
||||||
|
"""适配层工具:从 Flask session(或其 dict 快照)构造可信 principal。
|
||||||
|
|
||||||
|
调用方必须在完成认证后调用;``session_snapshot`` 由适配层显式传入
|
||||||
|
(本模块不 import flask,保持服务层可测试、无隐式上下文依赖)。
|
||||||
|
"""
|
||||||
|
snap = session_snapshot or {}
|
||||||
|
resolved_workspace = str(workspace_id or snap.get("workspace_id") or "default")
|
||||||
|
host_mode = bool(snap.get("host_mode"))
|
||||||
|
return TrustedPrincipal(
|
||||||
|
username=str(username or snap.get("username") or ""),
|
||||||
|
workspace_id=resolved_workspace,
|
||||||
|
role=str(snap.get("role") or "user"),
|
||||||
|
is_api_user=bool(snap.get("is_api_user")),
|
||||||
|
host_mode=host_mode,
|
||||||
|
host_workspace_id=str(snap.get("host_workspace_id") or (resolved_workspace if host_mode else "") or "") or None,
|
||||||
|
preferred_model_key=snap.get("model_key"),
|
||||||
|
preferred_run_mode=snap.get("run_mode"),
|
||||||
|
preferred_thinking_mode=snap.get("thinking_mode"),
|
||||||
|
)
|
||||||
114
server/runtime/service.py
Normal file
114
server/runtime/service.py
Normal file
@ -0,0 +1,114 @@
|
|||||||
|
"""RuntimeService:公共任务受理与控制入口(契约 docs/runtime_contract.md §4.2)。
|
||||||
|
|
||||||
|
定位:Web(HTTP 适配层)、CLI、未来的定时触发器等调用方共用的任务入口。
|
||||||
|
本服务只做「受理裁决 + 显式上下文转发 + 控制委托」,不持有任务状态——
|
||||||
|
任务记录、事件流、门闸、保存保护仍由既有 TaskManager / main_task_gate /
|
||||||
|
conversation_manager 承载(契约 §2 状态责任表不变)。
|
||||||
|
|
||||||
|
兼容期说明:create_task 内部把 RuntimeContext 转换为既有
|
||||||
|
session_data 快照传入 create_chat_task;任务线程已改为显式 RuntimeIdentity
|
||||||
|
驱动资源装配(test_request_context 桥已拆除),session_data 快照仍承载
|
||||||
|
门闸 token、事件回放等内部指令的跨线程传递。
|
||||||
|
"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
from typing import Any, Dict, List, Optional, Tuple
|
||||||
|
|
||||||
|
from modules.i18n import tr
|
||||||
|
from server.runtime.context import RuntimeContext
|
||||||
|
|
||||||
|
|
||||||
|
class RuntimeService:
|
||||||
|
"""公共任务入口。无状态:全部状态委托给 task_manager 单例。"""
|
||||||
|
|
||||||
|
# ---- 受理 ----
|
||||||
|
|
||||||
|
def create_task(self, ctx: RuntimeContext):
|
||||||
|
"""受理一轮 Run:校验显式上下文 → 互斥裁决(create_chat_task 内)→ 登记 → 起执行。
|
||||||
|
|
||||||
|
抛错契约(适配层负责映射 HTTP 状态码):
|
||||||
|
- ValueError:上下文/参数非法(400)
|
||||||
|
- RuntimeError:同对话已有运行中 chat 任务等业务冲突(409)
|
||||||
|
"""
|
||||||
|
ctx.validate()
|
||||||
|
# 延迟导入避免循环:server.tasks 的 blueprint 链不依赖本包
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
params = ctx.params
|
||||||
|
return task_manager.create_chat_task(
|
||||||
|
ctx.principal.username,
|
||||||
|
ctx.principal.workspace_id,
|
||||||
|
params.message,
|
||||||
|
list(params.images or []),
|
||||||
|
params.conversation_id,
|
||||||
|
videos=list(params.videos or []),
|
||||||
|
model_key=params.model_key,
|
||||||
|
thinking_mode=params.thinking_mode,
|
||||||
|
run_mode=params.run_mode,
|
||||||
|
max_iterations=params.max_iterations,
|
||||||
|
session_data=ctx.to_session_data(),
|
||||||
|
message_source=params.message_source,
|
||||||
|
goal_mode=params.goal_mode,
|
||||||
|
skill_context_messages=list(params.skill_context_messages or []),
|
||||||
|
files=list(params.files or []),
|
||||||
|
task_type=params.task_type,
|
||||||
|
)
|
||||||
|
|
||||||
|
# ---- 控制 ----
|
||||||
|
|
||||||
|
def cancel_task(self, username: str, task_id: str) -> bool:
|
||||||
|
"""停止主 Run(不触碰后台工作者;后台任务有独立控制入口)。"""
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.cancel_task(username, task_id)
|
||||||
|
|
||||||
|
def enqueue_runtime_guidance(
|
||||||
|
self, username: str, task_id: str, message: str, source: Optional[str] = None
|
||||||
|
) -> Dict[str, Any]:
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.enqueue_runtime_guidance(username, task_id, message, source=source)
|
||||||
|
|
||||||
|
def enqueue_runtime_pending_message(
|
||||||
|
self, username: str, task_id: str, message: str, files: Optional[List[str]] = None
|
||||||
|
) -> Dict[str, Any]:
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.enqueue_runtime_pending_message(username, task_id, message, files=files)
|
||||||
|
|
||||||
|
def remove_runtime_pending_message(self, username: str, task_id: str, message_id: str) -> Dict[str, Any]:
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.remove_runtime_pending_message(username, task_id, message_id)
|
||||||
|
|
||||||
|
def promote_runtime_pending_to_guidance(self, username: str, task_id: str, message_id: str) -> Dict[str, Any]:
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.promote_runtime_pending_to_guidance(username, task_id, message_id)
|
||||||
|
|
||||||
|
# ---- 观察(内部接口;后台调用方不必为观察任务发 HTTP 请求)----
|
||||||
|
|
||||||
|
def get_task(self, username: str, task_id: str):
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
return task_manager.get_task(username, task_id)
|
||||||
|
|
||||||
|
def get_task_events(
|
||||||
|
self, username: str, task_id: str, offset: int
|
||||||
|
) -> Tuple[Optional[List[Dict[str, Any]]], Optional[int], Optional[str]]:
|
||||||
|
"""按 offset 增量读取任务事件流(idx/offset 协议,与 REST 轮询同一语义)。
|
||||||
|
|
||||||
|
返回 (events, next_offset, error)。error 非空表示任务不存在或无权访问。
|
||||||
|
"""
|
||||||
|
from server.tasks import task_manager
|
||||||
|
|
||||||
|
rec = task_manager.get_task(username, task_id)
|
||||||
|
if not rec:
|
||||||
|
return None, None, tr("tasks.task_not_found")
|
||||||
|
events = task_manager.get_events_since(rec, max(0, int(offset or 0)))
|
||||||
|
next_offset = events[-1]["idx"] + 1 if events else max(0, int(offset or 0))
|
||||||
|
return events, next_offset, None
|
||||||
|
|
||||||
|
|
||||||
|
# 进程级单例(无状态,可安全共享)
|
||||||
|
runtime_service = RuntimeService()
|
||||||
@ -24,6 +24,7 @@ from utils.host_workspace_debug import write_host_workspace_debug
|
|||||||
from config import DATA_DIR, WORKSPACE_SKILLS_DIRNAME
|
from config import DATA_DIR, WORKSPACE_SKILLS_DIRNAME
|
||||||
from modules.goal_state_manager import GoalStateManager, REASON_USER_CANCEL
|
from modules.goal_state_manager import GoalStateManager, REASON_USER_CANCEL
|
||||||
from server.tasks import task_manager
|
from server.tasks import task_manager
|
||||||
|
from server.runtime import RuntimeContext, TaskParams, principal_from_session_snapshot, runtime_service
|
||||||
from server.tasks.skills import _build_skill_context_messages
|
from server.tasks.skills import _build_skill_context_messages
|
||||||
from server.tasks.helpers import _task_public_payload
|
from server.tasks.helpers import _task_public_payload
|
||||||
from server.tasks.media import _normalize_media_payload, _normalize_files_payload
|
from server.tasks.media import _normalize_media_payload, _normalize_files_payload
|
||||||
@ -196,23 +197,29 @@ def create_task_api():
|
|||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
debug_log(f"[TaskAPI] 补建对话失败(继续按无 cid 处理): {exc}")
|
debug_log(f"[TaskAPI] 补建对话失败(继续按无 cid 处理): {exc}")
|
||||||
|
|
||||||
|
# 公共任务入口(契约 docs/runtime_contract.md §4):适配层在完成认证后
|
||||||
|
# 从 Flask session 显式构造可信 principal,服务层不再回退读隐式上下文。
|
||||||
try:
|
try:
|
||||||
rec = task_manager.create_chat_task(
|
ctx = RuntimeContext(
|
||||||
username,
|
principal=principal_from_session_snapshot(session, workspace_id, username=username),
|
||||||
workspace_id,
|
params=TaskParams(
|
||||||
message,
|
message=message,
|
||||||
images,
|
images=images,
|
||||||
conversation_id,
|
videos=videos,
|
||||||
videos=videos,
|
files=files or [],
|
||||||
model_key=model_key,
|
conversation_id=conversation_id,
|
||||||
thinking_mode=thinking_mode,
|
model_key=model_key,
|
||||||
run_mode=run_mode,
|
thinking_mode=thinking_mode,
|
||||||
max_iterations=max_iterations,
|
run_mode=run_mode,
|
||||||
message_source=message_source,
|
max_iterations=max_iterations,
|
||||||
goal_mode=goal_mode,
|
message_source=message_source,
|
||||||
skill_context_messages=skill_context_messages,
|
goal_mode=goal_mode,
|
||||||
files=files,
|
skill_context_messages=skill_context_messages,
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
rec = runtime_service.create_task(ctx)
|
||||||
|
except ValueError as exc:
|
||||||
|
return jsonify({"success": False, "error": str(exc)}), 400
|
||||||
except RuntimeError as exc:
|
except RuntimeError as exc:
|
||||||
return jsonify({"success": False, "error": str(exc)}), 409
|
return jsonify({"success": False, "error": str(exc)}), 409
|
||||||
return jsonify({
|
return jsonify({
|
||||||
|
|||||||
@ -11,11 +11,7 @@ from collections import deque
|
|||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Dict, Any, Optional, List
|
from typing import Dict, Any, Optional, List
|
||||||
|
|
||||||
from flask import Blueprint, request, jsonify
|
from server.context import RuntimeIdentity, get_user_resources, ensure_conversation_loaded
|
||||||
from flask import current_app, session
|
|
||||||
|
|
||||||
from server.auth_helpers import api_login_required, get_current_username
|
|
||||||
from server.context import get_user_resources, ensure_conversation_loaded
|
|
||||||
from server.chat_flow import run_chat_task_sync
|
from server.chat_flow import run_chat_task_sync
|
||||||
from server.main_task_gate import release_main_task_gate
|
from server.main_task_gate import release_main_task_gate
|
||||||
from server.work_timer import finalize_conversation_work_timer
|
from server.work_timer import finalize_conversation_work_timer
|
||||||
@ -172,41 +168,19 @@ class TaskManager:
|
|||||||
raise RuntimeError(tr("tasks.task_already_running"))
|
raise RuntimeError(tr("tasks.task_already_running"))
|
||||||
task_id = str(uuid.uuid4())
|
task_id = str(uuid.uuid4())
|
||||||
record = TaskRecord(task_id, username, workspace_id, message, conversation_id, model_key, thinking_mode, run_mode, max_iterations, task_type=normalized_task_type)
|
record = TaskRecord(task_id, username, workspace_id, message, conversation_id, model_key, thinking_mode, run_mode, max_iterations, task_type=normalized_task_type)
|
||||||
# 记录当前 session 快照,便于后台线程内使用
|
# 运行上下文快照(RuntimeContext.to_session_data 产物,契约 docs/runtime_contract.md §4.1)。
|
||||||
if session_data is not None:
|
# 必须显式传入:禁止在受理层回退读 Flask session(隐式上下文在后台线程中
|
||||||
snapshot = dict(session_data)
|
# 不可靠且会静默丢身份)。调用方统一走 RuntimeService.create_task 构造快照。
|
||||||
snapshot.setdefault("workspace_id", workspace_id)
|
if session_data is None:
|
||||||
if message_source is not None:
|
raise ValueError(tr("tasks.missing_session_data"))
|
||||||
snapshot.setdefault("message_source", str(message_source))
|
snapshot = dict(session_data)
|
||||||
snapshot["goal_mode"] = bool(goal_mode)
|
snapshot.setdefault("workspace_id", workspace_id)
|
||||||
if skill_context_messages:
|
if message_source is not None:
|
||||||
snapshot["skill_context_messages"] = list(skill_context_messages)
|
snapshot.setdefault("message_source", str(message_source))
|
||||||
try:
|
snapshot["goal_mode"] = bool(goal_mode)
|
||||||
snapshot.setdefault("host_mode", session.get("host_mode"))
|
if skill_context_messages:
|
||||||
if snapshot.get("host_mode"):
|
snapshot["skill_context_messages"] = list(skill_context_messages)
|
||||||
snapshot.setdefault("host_workspace_id", session.get("host_workspace_id") or workspace_id)
|
record.session_data = snapshot
|
||||||
except Exception:
|
|
||||||
if snapshot.get("host_mode"):
|
|
||||||
snapshot.setdefault("host_workspace_id", workspace_id)
|
|
||||||
record.session_data = snapshot
|
|
||||||
else:
|
|
||||||
try:
|
|
||||||
record.session_data = {
|
|
||||||
"username": session.get("username"),
|
|
||||||
"role": session.get("role"),
|
|
||||||
"is_api_user": session.get("is_api_user"),
|
|
||||||
"host_mode": session.get("host_mode"),
|
|
||||||
"host_workspace_id": session.get("host_workspace_id") or workspace_id,
|
|
||||||
"workspace_id": workspace_id,
|
|
||||||
"run_mode": session.get("run_mode"),
|
|
||||||
"thinking_mode": session.get("thinking_mode"),
|
|
||||||
"model_key": session.get("model_key"),
|
|
||||||
"message_source": str(message_source) if message_source is not None else None,
|
|
||||||
"goal_mode": bool(goal_mode),
|
|
||||||
"skill_context_messages": list(skill_context_messages or []),
|
|
||||||
}
|
|
||||||
except Exception:
|
|
||||||
record.session_data = {}
|
|
||||||
with self._lock:
|
with self._lock:
|
||||||
self._tasks[task_id] = record
|
self._tasks[task_id] = record
|
||||||
thread = threading.Thread(target=self._run_chat_task, args=(record, images, videos or [], files or []), daemon=True)
|
thread = threading.Thread(target=self._run_chat_task, args=(record, images, videos or [], files or []), daemon=True)
|
||||||
@ -789,25 +763,34 @@ class TaskManager:
|
|||||||
workspace = None
|
workspace = None
|
||||||
stop_hint = False
|
stop_hint = False
|
||||||
try:
|
try:
|
||||||
# 为后台线程构造最小请求上下文,填充 session
|
# 显式运行上下文(契约 docs/runtime_contract.md §4.1):身份与偏好快照在
|
||||||
from server.app import app as flask_app
|
# 受理时由 RuntimeContext.to_session_data 固化,这里还原为 RuntimeIdentity
|
||||||
with flask_app.test_request_context():
|
# 直接驱动资源装配——不再伪造 Flask 请求上下文(原 test_request_context
|
||||||
try:
|
# 桥已拆除,任务线程全程无隐式上下文)。
|
||||||
for k, v in (rec.session_data or {}).items():
|
sd = rec.session_data or {}
|
||||||
if v is not None:
|
identity = RuntimeIdentity(
|
||||||
session[k] = v
|
host_mode=bool(sd.get("host_mode")),
|
||||||
if session.get("host_mode"):
|
host_workspace_id=sd.get("host_workspace_id"),
|
||||||
session["workspace_id"] = workspace_id
|
is_api_user=bool(sd.get("is_api_user")),
|
||||||
session["host_workspace_id"] = session.get("host_workspace_id") or workspace_id
|
role=sd.get("role"),
|
||||||
write_host_workspace_debug(
|
preferred_model_key=sd.get("model_key"),
|
||||||
"tasks.run_chat_task.apply_host_session",
|
preferred_run_mode=sd.get("run_mode"),
|
||||||
task_id=rec.task_id,
|
preferred_thinking_mode=sd.get("thinking_mode"),
|
||||||
workspace_id=workspace_id,
|
)
|
||||||
host_workspace_id=session.get("host_workspace_id"),
|
if identity.host_mode:
|
||||||
)
|
write_host_workspace_debug(
|
||||||
except Exception:
|
"tasks.run_chat_task.apply_host_session",
|
||||||
pass
|
task_id=rec.task_id,
|
||||||
terminal, workspace = get_user_resources(username, workspace_id=workspace_id, conversation_id=rec.conversation_id)
|
workspace_id=workspace_id,
|
||||||
|
host_workspace_id=identity.host_workspace_id or workspace_id,
|
||||||
|
)
|
||||||
|
terminal, workspace = get_user_resources(
|
||||||
|
username,
|
||||||
|
workspace_id=workspace_id,
|
||||||
|
update_session=False,
|
||||||
|
conversation_id=rec.conversation_id,
|
||||||
|
identity=identity,
|
||||||
|
)
|
||||||
if not terminal or not workspace:
|
if not terminal or not workspace:
|
||||||
raise RuntimeError(tr("tasks.system_not_initialized"))
|
raise RuntimeError(tr("tasks.system_not_initialized"))
|
||||||
stop_hint = bool(stop_flags.get(rec.task_id, {}).get("stop"))
|
stop_hint = bool(stop_flags.get(rec.task_id, {}).get("stop"))
|
||||||
@ -850,7 +833,7 @@ class TaskManager:
|
|||||||
# 确保会话加载
|
# 确保会话加载
|
||||||
conversation_id = rec.conversation_id
|
conversation_id = rec.conversation_id
|
||||||
try:
|
try:
|
||||||
conversation_id, _ = ensure_conversation_loaded(terminal, conversation_id, workspace=workspace)
|
conversation_id, _ = ensure_conversation_loaded(terminal, conversation_id, workspace=workspace, update_session=False)
|
||||||
rec.conversation_id = conversation_id
|
rec.conversation_id = conversation_id
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
raise RuntimeError(tr("tasks.conversation_load_failed", error=exc)) from exc
|
raise RuntimeError(tr("tasks.conversation_load_failed", error=exc)) from exc
|
||||||
@ -974,6 +957,13 @@ class TaskManager:
|
|||||||
"_skill_context_messages",
|
"_skill_context_messages",
|
||||||
list((rec.session_data or {}).get("skill_context_messages") or []),
|
list((rec.session_data or {}).get("skill_context_messages") or []),
|
||||||
)
|
)
|
||||||
|
# 审批/提问等待超时透传(契约 §6):None = 保持既有默认语义(3600s),
|
||||||
|
# 超时后的语义属阶段三产品决策,本阶段仅建立透传机制。
|
||||||
|
setattr(
|
||||||
|
terminal,
|
||||||
|
"_approval_timeout_seconds",
|
||||||
|
(rec.session_data or {}).get("approval_timeout_seconds"),
|
||||||
|
)
|
||||||
except Exception:
|
except Exception:
|
||||||
previous_auto_user_event = None
|
previous_auto_user_event = None
|
||||||
previous_message_source = None
|
previous_message_source = None
|
||||||
@ -1009,6 +999,7 @@ class TaskManager:
|
|||||||
setattr(terminal, "_skill_context_messages", previous_skill_context_messages)
|
setattr(terminal, "_skill_context_messages", previous_skill_context_messages)
|
||||||
else:
|
else:
|
||||||
setattr(terminal, "_skill_context_messages", [])
|
setattr(terminal, "_skill_context_messages", [])
|
||||||
|
setattr(terminal, "_approval_timeout_seconds", None)
|
||||||
if terminal and getattr(terminal, "context_manager", None):
|
if terminal and getattr(terminal, "context_manager", None):
|
||||||
terminal.context_manager.set_web_terminal_callback(previous_ctx_callback)
|
terminal.context_manager.set_web_terminal_callback(previous_ctx_callback)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
|||||||
@ -14,6 +14,7 @@ from flask import Blueprint, jsonify, request
|
|||||||
|
|
||||||
from server.auth_helpers import api_login_required
|
from server.auth_helpers import api_login_required
|
||||||
from server.context import make_terminal_callback, with_terminal
|
from server.context import make_terminal_callback, with_terminal
|
||||||
|
from server.runtime import InternalDirectives, RuntimeContext, TaskParams, runtime_service
|
||||||
from modules.i18n import tr
|
from modules.i18n import tr
|
||||||
|
|
||||||
workflow_runtime_bp = Blueprint("workflow_runtime", __name__)
|
workflow_runtime_bp = Blueprint("workflow_runtime", __name__)
|
||||||
@ -162,34 +163,33 @@ def api_activate_workflow(terminal, workspace, username):
|
|||||||
f"{activation_text}"
|
f"{activation_text}"
|
||||||
)
|
)
|
||||||
|
|
||||||
from .tasks import task_manager
|
# 公共任务入口:内部调用方,从 terminal 构造可信上下文;
|
||||||
|
# 门闸 token 与 user 消息事件回放走 InternalDirectives(客户端不可提交)。
|
||||||
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
ctx = RuntimeContext.from_terminal(
|
||||||
session_data = {
|
terminal,
|
||||||
"username": username,
|
workspace,
|
||||||
"message_source": "workflow",
|
username,
|
||||||
# 门闸 token 随任务移交,由任务线程认领(见 process_message_task)
|
params=TaskParams(
|
||||||
"main_task_gate_token": gate_token,
|
message=prompt,
|
||||||
# 让任务事件流携带该 user 消息,保证轮询客户端/刷新后可见
|
conversation_id=conversation_id,
|
||||||
"auto_user_message_event": True,
|
|
||||||
"auto_user_message_payload": {
|
|
||||||
"message_source": "workflow",
|
|
||||||
"workflow_activate": True,
|
|
||||||
"visibility": "chat",
|
|
||||||
"starts_work": True,
|
|
||||||
"timestamp": datetime.now().isoformat(),
|
|
||||||
},
|
|
||||||
}
|
|
||||||
try:
|
|
||||||
rec = task_manager.create_chat_task(
|
|
||||||
username,
|
|
||||||
workspace_id,
|
|
||||||
prompt,
|
|
||||||
[],
|
|
||||||
conversation_id,
|
|
||||||
message_source="workflow",
|
message_source="workflow",
|
||||||
session_data=session_data,
|
),
|
||||||
)
|
directives=InternalDirectives(
|
||||||
|
# 门闸 token 随任务移交,由任务线程认领(见 process_message_task)
|
||||||
|
main_task_gate_token=gate_token,
|
||||||
|
# 让任务事件流携带该 user 消息,保证轮询客户端/刷新后可见
|
||||||
|
auto_user_message_event=True,
|
||||||
|
auto_user_message_payload={
|
||||||
|
"message_source": "workflow",
|
||||||
|
"workflow_activate": True,
|
||||||
|
"visibility": "chat",
|
||||||
|
"starts_work": True,
|
||||||
|
"timestamp": datetime.now().isoformat(),
|
||||||
|
},
|
||||||
|
),
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
rec = runtime_service.create_task(ctx)
|
||||||
except RuntimeError as exc:
|
except RuntimeError as exc:
|
||||||
release_main_task_gate(terminal, gate_token)
|
release_main_task_gate(terminal, gate_token)
|
||||||
return jsonify({"error": str(exc)}), 409
|
return jsonify({"error": str(exc)}), 409
|
||||||
@ -245,31 +245,28 @@ def api_deactivate_workflow(terminal, workspace, username):
|
|||||||
str(n.get("message") or "").strip() for n in notices if str(n.get("message") or "").strip()
|
str(n.get("message") or "").strip() for n in notices if str(n.get("message") or "").strip()
|
||||||
)
|
)
|
||||||
if notice_text:
|
if notice_text:
|
||||||
from .tasks import task_manager
|
ctx = RuntimeContext.from_terminal(
|
||||||
|
terminal,
|
||||||
workspace_id = getattr(workspace, "workspace_id", None) or "default"
|
workspace,
|
||||||
session_data = {
|
|
||||||
"username": username,
|
|
||||||
"message_source": "workflow",
|
|
||||||
"main_task_gate_token": gate_token,
|
|
||||||
"auto_user_message_event": True,
|
|
||||||
"auto_user_message_payload": {
|
|
||||||
"message_source": "workflow",
|
|
||||||
"workflow_notice": True,
|
|
||||||
"visibility": "chat",
|
|
||||||
"starts_work": True,
|
|
||||||
"timestamp": datetime.now().isoformat(),
|
|
||||||
},
|
|
||||||
}
|
|
||||||
task_manager.create_chat_task(
|
|
||||||
username,
|
username,
|
||||||
workspace_id,
|
params=TaskParams(
|
||||||
notice_text,
|
message=notice_text,
|
||||||
[],
|
conversation_id=conversation_id,
|
||||||
conversation_id,
|
message_source="workflow",
|
||||||
message_source="workflow",
|
),
|
||||||
session_data=session_data,
|
directives=InternalDirectives(
|
||||||
|
main_task_gate_token=gate_token,
|
||||||
|
auto_user_message_event=True,
|
||||||
|
auto_user_message_payload={
|
||||||
|
"message_source": "workflow",
|
||||||
|
"workflow_notice": True,
|
||||||
|
"visibility": "chat",
|
||||||
|
"starts_work": True,
|
||||||
|
"timestamp": datetime.now().isoformat(),
|
||||||
|
},
|
||||||
|
),
|
||||||
)
|
)
|
||||||
|
runtime_service.create_task(ctx)
|
||||||
dispatched = True
|
dispatched = True
|
||||||
except Exception: # noqa: BLE001
|
except Exception: # noqa: BLE001
|
||||||
# 任何失败:通知放回池(等轮询器/工具循环消费),不静默丢失
|
# 任何失败:通知放回池(等轮询器/工具循环消费),不静默丢失
|
||||||
|
|||||||
195
test/test_runtime_service.py
Normal file
195
test/test_runtime_service.py
Normal file
@ -0,0 +1,195 @@
|
|||||||
|
"""阶段二公共任务入口(RuntimeService)的显式上下文验收测试。
|
||||||
|
|
||||||
|
契约用例(docs/runtime_contract.md §7):
|
||||||
|
- T01:显式上下文受理(无 HTTP 请求、无 test_request_context)
|
||||||
|
- T02:同对话并发 chat 互斥 / notice 豁免
|
||||||
|
- T04:取消(受理层语义;执行线程以 no-op 替身阻断,不触达模型调用)
|
||||||
|
- 快照完整性:to_session_data 携带资源装配所需的全部身份与偏好字段
|
||||||
|
|
||||||
|
本测试全程不创建 Flask 应用/请求上下文——这本身就是
|
||||||
|
「公共入口不依赖隐式 Web 环境」的直接证明。
|
||||||
|
"""
|
||||||
|
import sys
|
||||||
|
import unittest
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
|
||||||
|
|
||||||
|
from server.runtime import ( # noqa: E402
|
||||||
|
InternalDirectives,
|
||||||
|
RuntimeContext,
|
||||||
|
TaskParams,
|
||||||
|
TrustedPrincipal,
|
||||||
|
principal_from_session_snapshot,
|
||||||
|
runtime_service,
|
||||||
|
)
|
||||||
|
from server.tasks import task_manager # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
def _make_ctx(conversation_id="conv_test_rt", task_type="chat", username="tester"):
|
||||||
|
return RuntimeContext(
|
||||||
|
principal=TrustedPrincipal(
|
||||||
|
username=username,
|
||||||
|
workspace_id="default",
|
||||||
|
role="user",
|
||||||
|
is_api_user=False,
|
||||||
|
host_mode=True,
|
||||||
|
host_workspace_id="default",
|
||||||
|
preferred_model_key="kimi-test",
|
||||||
|
preferred_run_mode="fast",
|
||||||
|
preferred_thinking_mode=False,
|
||||||
|
),
|
||||||
|
params=TaskParams(
|
||||||
|
message="hello",
|
||||||
|
conversation_id=conversation_id,
|
||||||
|
run_mode="fast",
|
||||||
|
task_type=task_type,
|
||||||
|
),
|
||||||
|
directives=InternalDirectives(main_task_gate_token="tok123"),
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
class RuntimeContextModelTest(unittest.TestCase):
|
||||||
|
def test_validate_rejects_empty_username(self):
|
||||||
|
ctx = _make_ctx()
|
||||||
|
object.__setattr__(ctx.principal, "username", "") # frozen dataclass 测试绕过
|
||||||
|
with self.assertRaises(ValueError):
|
||||||
|
ctx.validate()
|
||||||
|
|
||||||
|
def test_to_session_data_carries_identity_and_preferences(self):
|
||||||
|
ctx = _make_ctx()
|
||||||
|
snap = ctx.to_session_data()
|
||||||
|
# 身份与资源范围
|
||||||
|
self.assertEqual(snap["username"], "tester")
|
||||||
|
self.assertEqual(snap["workspace_id"], "default")
|
||||||
|
self.assertEqual(snap["role"], "user")
|
||||||
|
self.assertFalse(snap["is_api_user"])
|
||||||
|
self.assertTrue(snap["host_mode"])
|
||||||
|
self.assertEqual(snap["host_workspace_id"], "default")
|
||||||
|
# 偏好快照层(非本次覆盖)
|
||||||
|
self.assertEqual(snap["model_key"], "kimi-test")
|
||||||
|
self.assertEqual(snap["run_mode"], "fast")
|
||||||
|
self.assertFalse(snap["thinking_mode"])
|
||||||
|
# 内部指令
|
||||||
|
self.assertEqual(snap["main_task_gate_token"], "tok123")
|
||||||
|
# 默认不出现事件回放键
|
||||||
|
self.assertNotIn("auto_user_message_event", snap)
|
||||||
|
|
||||||
|
def test_to_session_data_directives(self):
|
||||||
|
ctx = RuntimeContext(
|
||||||
|
principal=_make_ctx().principal,
|
||||||
|
params=TaskParams(message="m", conversation_id="c1", approval_timeout_seconds=120),
|
||||||
|
directives=InternalDirectives(
|
||||||
|
auto_user_message_event=True,
|
||||||
|
auto_user_message_payload={"visibility": "chat"},
|
||||||
|
preceding_user_notices=[{"message": "n1", "payload": {}}],
|
||||||
|
),
|
||||||
|
)
|
||||||
|
snap = ctx.to_session_data()
|
||||||
|
self.assertTrue(snap["auto_user_message_event"])
|
||||||
|
self.assertEqual(snap["auto_user_message_payload"], {"visibility": "chat"})
|
||||||
|
self.assertEqual(len(snap["preceding_user_notices"]), 1)
|
||||||
|
# 超时透传机制(默认 None 时不出现,语义不变)
|
||||||
|
self.assertEqual(snap["approval_timeout_seconds"], 120)
|
||||||
|
|
||||||
|
def test_principal_from_session_snapshot(self):
|
||||||
|
snap = {
|
||||||
|
"username": "u1",
|
||||||
|
"role": "api",
|
||||||
|
"is_api_user": True,
|
||||||
|
"host_mode": False,
|
||||||
|
"workspace_id": "ws1",
|
||||||
|
"model_key": "m1",
|
||||||
|
}
|
||||||
|
p = principal_from_session_snapshot(snap, "ws1")
|
||||||
|
self.assertEqual(p.username, "u1")
|
||||||
|
self.assertEqual(p.role, "api")
|
||||||
|
self.assertTrue(p.is_api_user)
|
||||||
|
self.assertFalse(p.host_mode)
|
||||||
|
self.assertIsNone(p.host_workspace_id)
|
||||||
|
self.assertEqual(p.preferred_model_key, "m1")
|
||||||
|
|
||||||
|
|
||||||
|
class RuntimeServiceAdmissionTest(unittest.TestCase):
|
||||||
|
"""受理层行为:用 no-op 替身阻断执行线程,不触达模型与工作区装配。"""
|
||||||
|
|
||||||
|
def setUp(self):
|
||||||
|
self._orig_run = task_manager._run_chat_task
|
||||||
|
task_manager._run_chat_task = lambda *a, **k: None # 线程即刻结束
|
||||||
|
self._created = []
|
||||||
|
|
||||||
|
def tearDown(self):
|
||||||
|
task_manager._run_chat_task = self._orig_run
|
||||||
|
for task_id in self._created:
|
||||||
|
with task_manager._lock:
|
||||||
|
task_manager._tasks.pop(task_id, None)
|
||||||
|
|
||||||
|
def _create(self, **kw):
|
||||||
|
rec = runtime_service.create_task(_make_ctx(**kw))
|
||||||
|
self._created.append(rec.task_id)
|
||||||
|
return rec
|
||||||
|
|
||||||
|
def test_t01_create_task_without_http_context(self):
|
||||||
|
rec = self._create()
|
||||||
|
self.assertTrue(rec.task_id)
|
||||||
|
self.assertEqual(rec.username, "tester")
|
||||||
|
self.assertEqual(rec.conversation_id, "conv_test_rt")
|
||||||
|
# 快照经显式上下文进入,未触碰 Flask session
|
||||||
|
self.assertEqual(rec.session_data["username"], "tester")
|
||||||
|
self.assertEqual(rec.session_data["main_task_gate_token"], "tok123")
|
||||||
|
|
||||||
|
def test_t02_same_conversation_chat_mutex_and_notice_exempt(self):
|
||||||
|
# 第一个任务保持 running(线程 no-op 但 status 已被置 running)
|
||||||
|
rec1 = self._create()
|
||||||
|
self.assertEqual(rec1.status, "running")
|
||||||
|
# 同对话第二个 chat 被拒
|
||||||
|
with self.assertRaises(RuntimeError):
|
||||||
|
self._create()
|
||||||
|
# notice 类型豁免互斥
|
||||||
|
rec2 = self._create(task_type="notice")
|
||||||
|
self.assertEqual(rec2.task_type, "notice")
|
||||||
|
# 不同对话不受影响
|
||||||
|
rec3 = self._create(conversation_id="conv_other")
|
||||||
|
self.assertTrue(rec3.task_id)
|
||||||
|
|
||||||
|
def test_create_task_validates_context(self):
|
||||||
|
bad = _make_ctx()
|
||||||
|
object.__setattr__(bad.principal, "workspace_id", "")
|
||||||
|
with self.assertRaises(ValueError):
|
||||||
|
runtime_service.create_task(bad)
|
||||||
|
|
||||||
|
def test_create_chat_task_requires_explicit_session_data(self):
|
||||||
|
# 直调底层入口且不带快照 → 明确拒绝(不再静默读 Flask session)
|
||||||
|
with self.assertRaises(ValueError):
|
||||||
|
task_manager.create_chat_task(
|
||||||
|
"tester", "default", "msg", [], "conv_x",
|
||||||
|
)
|
||||||
|
|
||||||
|
def test_t04_cancel_task(self):
|
||||||
|
rec = self._create()
|
||||||
|
ok = runtime_service.cancel_task("tester", rec.task_id)
|
||||||
|
self.assertTrue(ok)
|
||||||
|
# 他人不可取消
|
||||||
|
rec2 = self._create(conversation_id="conv_cancel2")
|
||||||
|
self.assertFalse(runtime_service.cancel_task("someone_else", rec2.task_id))
|
||||||
|
|
||||||
|
def test_get_task_events_offset_protocol(self):
|
||||||
|
rec = self._create()
|
||||||
|
task_manager._append_event(rec, "system_message", {"text": "a"})
|
||||||
|
task_manager._append_event(rec, "text_chunk", {"text": "b"})
|
||||||
|
events, next_offset, err = runtime_service.get_task_events("tester", rec.task_id, 0)
|
||||||
|
self.assertIsNone(err)
|
||||||
|
self.assertEqual([e["idx"] for e in events], [0, 1])
|
||||||
|
self.assertEqual(next_offset, 2)
|
||||||
|
# offset 续读
|
||||||
|
events2, next2, _ = runtime_service.get_task_events("tester", rec.task_id, 1)
|
||||||
|
self.assertEqual([e["idx"] for e in events2], [1])
|
||||||
|
self.assertEqual(next2, 2)
|
||||||
|
# 无权/不存在
|
||||||
|
events3, _, err3 = runtime_service.get_task_events("someone_else", rec.task_id, 0)
|
||||||
|
self.assertIsNone(events3)
|
||||||
|
self.assertTrue(err3)
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
Loading…
Reference in New Issue
Block a user