feat(agent): 增加简易 auto 推理强度模式

This commit is contained in:
caoqianming 2026-09-04 09:34:37 +08:00
parent 5b052c464b
commit 8d60e000ce
10 changed files with 532 additions and 12 deletions

View File

@ -136,7 +136,7 @@ Eval 与生产 core 解耦,通过现有 `/v1` API 创建专用任务、监听
默认 `deepseek_v4.flash`;复杂 bug / 终稿升 pro + reasoning_effort=max;fallback 手动切 Claude。成本量级:修 bug flash ~$0.01 / 完整申报书 flash ~$0.30(pro-max ~$1.5,Opus ~$10+)。99% 任务 flash 够用。 默认 `deepseek_v4.flash`;复杂 bug / 终稿升 pro + reasoning_effort=max;fallback 手动切 Claude。成本量级:修 bug flash ~$0.01 / 完整申报书 flash ~$0.30(pro-max ~$1.5,Opus ~$10+)。99% 任务 flash 够用。
模型思考参数由 profile 统一表达:`thinking_enabled` 只表示开关,`thinking_transport` 只表示已验证的传输协议,`reasoning_effort` 只表示开启后的推理强度,`thinking_clear` 表示 provider 是否清除历史思考,`reasoning_replay` 表示状态生命周期(`none` / `tool_turn` / `conversation` / `provider_managed``core/llm_params.py` 是请求参数构造唯一入口,`core/context.py` 是历史消息清洗唯一入口。原始 assistant 响应完整落库provider-bound 副本只向相同生产模型回放未改写的 reasoningDeepSeek V4 仅保留当前用户轮次的工具链状态GLM-5.3 Flash 在同模型会话内保留完整状态,未验证网关明确用 `none`,未来签名/加密 block 走 `provider_managed`。模型切换、上下文折叠和普通压缩都在同一入口应用隔离;上下文统计使用裁剪后的请求视图,原生图片 token 不反向污染 chars/token 校准。 模型思考参数由 profile 统一表达:`thinking_enabled` 只表示开关,`thinking_transport` 只表示已验证的传输协议,`reasoning_effort` 只表示开启后的推理强度,`thinking_clear` 表示 provider 是否清除历史思考,`reasoning_replay` 表示状态生命周期(`none` / `tool_turn` / `conversation` / `provider_managed``core/llm_params.py` 是请求参数构造唯一入口,`core/context.py` 是历史消息清洗唯一入口。档案的 `default_reasoning_effort=auto` 是平台编排值,不进入 provider当前用户轮次首次调用和失败工具步后用 `high`,成功工具步后用 `low`固定档保持原样auto 不选择 `max`。纯 reasoning 从首个推理片段起超过 90 秒且尚无正文/工具调用时关闭当前流并发送 `reasoning_reset`,不持久化半截 assistant随后以 `low` 和仅本次 provider 请求可见的简短约束重试一次,第二次仍超时则明确停止,用户取消始终优先。原始 assistant 响应完整落库provider-bound 副本只向相同生产模型回放未改写的 reasoningDeepSeek V4 仅保留当前用户轮次的工具链状态GLM-5.3 Flash 在同模型会话内保留完整状态,未验证网关明确用 `none`,未来签名/加密 block 走 `provider_managed`。模型切换、上下文折叠和普通压缩都在同一入口应用隔离;上下文统计使用裁剪后的请求视图,原生图片 token 不反向污染 chars/token 校准。
--- ---

View File

@ -2,7 +2,7 @@
> 配合 `DESIGN.md`。本文件只记 phase 状态、决策偏差、文件量、下一步。每条 1-2 句:做了啥 + 关键判断;细节查 `git log` / `git diff` / `DESIGN §7.9` > 配合 `DESIGN.md`。本文件只记 phase 状态、决策偏差、文件量、下一步。每条 1-2 句:做了啥 + 关键判断;细节查 `git log` / `git diff` / `DESIGN §7.9`
最后更新:2026-09-03Unreleased工具健康统一事件口径 最后更新:2026-09-04Unreleased简易 auto 推理强度模式
--- ---
@ -20,6 +20,8 @@
--- ---
## 已完成关键能力 ## 已完成关键能力
- **09-04 / Unreleased / 简易 auto 推理强度模式**:模型档案可保留平台值 `auto`DeepSeek V4 Flash 首次调用和失败工具步后使用 high、成功工具步后使用 lowPro 的 medium 与其他固定档保持原样;纯 reasoning 超过 90 秒会清理直播推理并以 low 和请求内临时约束重试一次,第二次仍超时明确停止,取消优先且半截 assistant 不入库。正常 chat usage 留存配置值、实际档位、决策原因和保护重试标记,熔断复用 `agent_guard`;无 schema、migration、API、前端或版本变化。
- **09-03 / 0.71.0 / 工具链可靠性收敛**tool arguments salvage 新增“至少两个完整且完全一致副本 + 尾部截断副本”的保守恢复,语义不一致仍拒绝。修复 artifacts 部分唯一索引谓词被参数化后 PostgreSQL 无法匹配 `ON CONFLICT` 的问题,并让输出中任意行首 `[GATE FAIL]` 均按质量门记录,避免 SVG/PPT 质检占用真实故障大数。无 schema/migration/API 变化,生产失败历史未改写。 - **09-03 / 0.71.0 / 工具链可靠性收敛**tool arguments salvage 新增“至少两个完整且完全一致副本 + 尾部截断副本”的保守恢复,语义不一致仍拒绝。修复 artifacts 部分唯一索引谓词被参数化后 PostgreSQL 无法匹配 `ON CONFLICT` 的问题,并让输出中任意行首 `[GATE FAIL]` 均按质量门记录,避免 SVG/PPT 质检占用真实故障大数。无 schema/migration/API 变化,生产失败历史未改写。
- **09-03 / 0.71.0 / 工具健康统一写入 usage_events**:普通工具失败在 tool result 落消息后结构化写 `tool_failure`,并新增 `run_stopped`、`agent_guard`、`context_fold_failure`、`quality_gate` 四类任务异常事件;工具健康聚合彻底移除 messages JSONB 回扫,只用一次 usage_events 查询切换前的普通失败历史不迁移。Admin 大数只计真实 failure代理主动停止/重复保护/上下文异常改为独立控制区,质量门继续单列。新增 0040 migration为健康 kind 时间窗和普通失败 message 去重补部分索引。 - **09-03 / 0.71.0 / 工具健康统一写入 usage_events**:普通工具失败在 tool result 落消息后结构化写 `tool_failure`,并新增 `run_stopped`、`agent_guard`、`context_fold_failure`、`quality_gate` 四类任务异常事件;工具健康聚合彻底移除 messages JSONB 回扫,只用一次 usage_events 查询切换前的普通失败历史不迁移。Admin 大数只计真实 failure代理主动停止/重复保护/上下文异常改为独立控制区,质量门继续单列。新增 0040 migration为健康 kind 时间窗和普通失败 message 去重补部分索引。

View File

@ -15,7 +15,7 @@ variants:
thinking_enabled: true thinking_enabled: true
thinking_transport: extra_body thinking_transport: extra_body
reasoning_effort_levels: [low, high, max] reasoning_effort_levels: [low, high, max]
default_reasoning_effort: high default_reasoning_effort: auto
reasoning_replay: tool_turn # 只在当前用户轮次的工具链内原样回传 reasoning reasoning_replay: tool_turn # 只在当前用户轮次的工具链内原样回传 reasoning
code_quality: good code_quality: good
enable_run_python: true enable_run_python: true

View File

@ -15,6 +15,7 @@ REASONING_REPLAY_POLICIES = {
"conversation", "conversation",
"provider_managed", "provider_managed",
} }
REASONING_EFFORT_AUTO = "auto"
def model_profile_of(caps: object) -> str: def model_profile_of(caps: object) -> str:
@ -134,12 +135,20 @@ class ModelCapabilities:
) )
if ( if (
caps.default_reasoning_effort caps.default_reasoning_effort
and caps.default_reasoning_effort != REASONING_EFFORT_AUTO
and caps.default_reasoning_effort not in caps.reasoning_effort_levels and caps.default_reasoning_effort not in caps.reasoning_effort_levels
): ):
raise ValueError( raise ValueError(
f"档案 {path} 的 default_reasoning_effort=" f"档案 {path} 的 default_reasoning_effort="
f"{caps.default_reasoning_effort!r} 不在 reasoning_effort_levels 中" f"{caps.default_reasoning_effort!r} 不在 reasoning_effort_levels 中"
) )
if (
caps.default_reasoning_effort == REASONING_EFFORT_AUTO
and not {"low", "high"}.issubset(caps.reasoning_effort_levels)
):
raise ValueError(
f"档案 {path} 使用 auto 时 reasoning_effort_levels 必须包含 low/high"
)
return caps return caps
@property @property

View File

@ -19,6 +19,8 @@ def build_thinking_kwargs(
``extra_body`` 对应当前 DeepSeekGLM 与方舟 ChatCompletions 的共同协议 ``extra_body`` 对应当前 DeepSeekGLM 与方舟 ChatCompletions 的共同协议
effort 仅在开启且档案提供非空值时发送 effort 仅在开启且档案提供非空值时发送
""" """
if reasoning_effort == "auto":
raise ValueError("reasoning_effort=auto 必须在调用编排层解析后才能发送")
if transport == "none": if transport == "none":
return {} return {}
if transport != "extra_body": if transport != "extra_body":

View File

@ -23,7 +23,11 @@ import litellm
from . import pptx_guard from . import pptx_guard
from .artifacts import MAX_ARTIFACTS_PER_MESSAGE from .artifacts import MAX_ARTIFACTS_PER_MESSAGE
from .attachments import materialize_native_images from .attachments import materialize_native_images
from .capabilities import ModelCapabilities, model_profile_of from .capabilities import (
REASONING_EFFORT_AUTO,
ModelCapabilities,
model_profile_of,
)
from .context import ( from .context import (
CHARS_PER_TOKEN, CHARS_PER_TOKEN,
COMPACT_CONTEXT_RATIO, COMPACT_CONTEXT_RATIO,
@ -55,12 +59,49 @@ from .storage import (
record_tool_failure, record_tool_failure,
) )
from .task_actions import DeferredTaskActions from .task_actions import DeferredTaskActions
from .tool_failure import structured_failure
# 产物机检只挂能落盘的执行类工具(fs 写工具不适合造 pptx,机检无意义) # 产物机检只挂能落盘的执行类工具(fs 写工具不适合造 pptx,机检无意义)
_PPTX_GUARD_TOOLS = ("shell", "run_python") _PPTX_GUARD_TOOLS = ("shell", "run_python")
_CANCELLED_TOOL_PLACEHOLDER = "[cancelled by user]" _CANCELLED_TOOL_PLACEHOLDER = "[cancelled by user]"
_REASONING_RETRY_INSTRUCTION = (
"本次重试请压缩内部推理,尽快给出正文或发起必要的工具调用。"
)
class ReasoningPhaseTimeout(RuntimeError):
"""当前流只输出 reasoning超过保护时限后已被关闭。"""
class ReasoningGuardExhausted(RuntimeError):
"""纯 reasoning 保护重试仍超时,本轮应明确停止。"""
def resolve_reasoning_effort(
configured: str,
*,
first_call: bool,
previous_tools_succeeded: Optional[bool],
) -> Tuple[Optional[str], str]:
"""把档案配置解析为本次 provider 档位;``auto`` 永不直接出站。"""
if configured != REASONING_EFFORT_AUTO:
return configured or None, "configured"
if first_call:
return "high", "first_call"
if previous_tools_succeeded is False:
return "high", "previous_tools_failed"
return "low", "previous_tools_succeeded"
def tool_result_succeeded_for_reasoning(
name: str, arguments: Any, result: str, *, productive: bool
) -> bool:
"""auto 使用的工具成功口径:净产出且没有错误或质量门失败。"""
if not productive or "[产物机检 ERROR]" in result:
return False
return structured_failure(name, result, arguments=arguments) is None
# 错误签名归一:同一类工具报错在不同参数/路径/数字下抹平,让「反复撞同一堵墙」 # 错误签名归一:同一类工具报错在不同参数/路径/数字下抹平,让「反复撞同一堵墙」
@ -265,6 +306,11 @@ class AgentLoop:
# Structured deliverables accumulated across tool steps in the current user turn. # Structured deliverables accumulated across tool steps in the current user turn.
# They are persisted on the final assistant message, not mixed into provider payloads. # They are persisted on the final assistant message, not mixed into provider payloads.
self._pending_artifact_refs: list[dict] = [] self._pending_artifact_refs: list[dict] = []
# auto reasoning 只在当前用户轮次内决策;档案值本身不直接发给 provider。
self._llm_call_count = 0
self._previous_tool_step_succeeded: Optional[bool] = None
self._active_reasoning_effort: Optional[str] = None
self._reasoning_usage: dict[str, Any] = {}
def _emit(self, event: dict) -> None: def _emit(self, event: dict) -> None:
if self.sink is not None: if self.sink is not None:
@ -296,6 +342,10 @@ class AgentLoop:
def _run(self, user_message: Optional[str]) -> str: def _run(self, user_message: Optional[str]) -> str:
self._pending_artifact_refs = [] self._pending_artifact_refs = []
self._llm_call_count = 0
self._previous_tool_step_succeeded = None
self._active_reasoning_effort = None
self._reasoning_usage = {}
self._maybe_fold_context() self._maybe_fold_context()
if user_message is not None: if user_message is not None:
self.session.append({"role": "user", "content": user_message}) self.session.append({"role": "user", "content": user_message})
@ -306,7 +356,22 @@ class AgentLoop:
return "[cancelled]" return "[cancelled]"
start = time.monotonic() start = time.monotonic()
try:
response, cancelled_mid_stream = self._stream_llm() response, cancelled_mid_stream = self._stream_llm()
except ReasoningGuardExhausted:
# 用户取消优先于自动保护的终态;两者竞态时保持停止按钮语义。
if self._is_cancelled():
self._emit({"type": "cancelled"})
return "[cancelled]"
self._emit({
"type": "warn",
"msg": (
"模型连续两次仅输出推理且超过 90 秒,已停止本轮以免继续空耗。"
"回复「继续」可重新尝试。"
),
})
self._emit({"type": "done"})
return "[stopped: reasoning timeout]"
elapsed = time.monotonic() - start elapsed = time.monotonic() - start
if cancelled_mid_stream: if cancelled_mid_stream:
@ -347,9 +412,13 @@ class AgentLoop:
cache_hit_cny_per_mtoken=self.caps.cache_hit_cny_per_mtoken, cache_hit_cny_per_mtoken=self.caps.cache_hit_cny_per_mtoken,
pricing=self.caps.pricing, pricing=self.caps.pricing,
extra_units={ extra_units={
k: v for k, v in usage_details.items() **{
k: v
for k, v in usage_details.items()
if k not in ("tokens_in", "tokens_out") and v if k not in ("tokens_in", "tokens_out") and v
}, },
**self._reasoning_usage,
},
response=response, response=response,
) )
except Exception as e: except Exception as e:
@ -381,6 +450,7 @@ class AgentLoop:
return content return content
step_productive = False step_productive = False
step_reasoning_success = True
for i, tc in enumerate(tool_calls): for i, tc in enumerate(tool_calls):
if self._is_cancelled(): if self._is_cancelled():
self._fill_cancelled_tool_results(tool_calls[i:]) self._fill_cancelled_tool_results(tool_calls[i:])
@ -389,6 +459,15 @@ class AgentLoop:
result, productive, artifacts = self._execute_tool_call(tc) result, productive, artifacts = self._execute_tool_call(tc)
self._remember_artifacts(artifacts) self._remember_artifacts(artifacts)
step_productive = step_productive or productive step_productive = step_productive or productive
succeeded_for_reasoning = tool_result_succeeded_for_reasoning(
tc.function.name,
tc.function.arguments,
result,
productive=productive,
)
step_reasoning_success = (
step_reasoning_success and succeeded_for_reasoning
)
message_id = self.session.append( message_id = self.session.append(
{ {
"role": "tool", "role": "tool",
@ -412,6 +491,10 @@ class AgentLoop:
# 可观测性留痕失败不能打断主对话。 # 可观测性留痕失败不能打断主对话。
pass pass
# 下一次 auto 决策只看刚完成的整步:所有工具均有净产出且未报错/未过
# 质量门才视为成功;任一错误、门失败或整步无净产出都回到 high。
self._previous_tool_step_succeeded = step_reasoning_success
# ask_user:本步调用了人工选择工具 → 提前结束本轮,等用户点选项 / 文字讨论, # ask_user:本步调用了人工选择工具 → 提前结束本轮,等用户点选项 / 文字讨论,
# 不回灌 LLM。选项已随该 tool_call 的 arguments 流给前端渲染成选项卡;tool 结果 # 不回灌 LLM。选项已随该 tool_call 的 arguments 流给前端渲染成选项卡;tool 结果
# 只是占位,下轮用户回复(点选项 = 发选项 label 文本)后模型自然接上。 # 只是占位,下轮用户回复(点选项 = 发选项 label 文本)后模型自然接上。
@ -471,6 +554,7 @@ class AgentLoop:
# 一整段 2k+ token 生成),首败即降级非流式(provider 服务端拼 tool_calls,绕开 # 一整段 2k+ token 生成),首败即降级非流式(provider 服务端拼 tool_calls,绕开
# 流式 delta 错位),历史数据里非流式兜底从未再畸形。 # 流式 delta 错位),历史数据里非流式兜底从未再畸形。
_MAX_MALFORMED_ATTEMPTS = 3 _MAX_MALFORMED_ATTEMPTS = 3
_REASONING_PHASE_TIMEOUT_S = 90.0
# DeepSeek 的长 write/edit arguments 在流式 delta 中偶发错位。function.name 首包 # DeepSeek 的长 write/edit arguments 在流式 delta 中偶发错位。function.name 首包
# 到达时立即关流并非流式重发;正文和其他工具仍走流式。只限定已实证的模型族, # 到达时立即关流并非流式重发;正文和其他工具仍走流式。只限定已实证的模型族,
@ -556,6 +640,24 @@ class AgentLoop:
返回 (response, cancelled_mid_stream);语义见 robust_stream docstring 返回 (response, cancelled_mid_stream);语义见 robust_stream docstring
""" """
configured = str(getattr(self.caps, "default_reasoning_effort", "") or "")
call_count = int(getattr(self, "_llm_call_count", 0))
effort, reason = resolve_reasoning_effort(
configured,
first_call=call_count == 0,
previous_tools_succeeded=getattr(
self, "_previous_tool_step_succeeded", None
),
)
self._llm_call_count = call_count + 1
self._active_reasoning_effort = effort
self._reasoning_usage = {
"reasoning_config": configured,
"reasoning_effort": effort or "",
"reasoning_reason": reason,
"reasoning_guard_retry": False,
}
# 上下文压力门槛按当前模型 reliable_context 折算:体量未到阈值前不压缩(缓存全暖 + 不丢信息)。 # 上下文压力门槛按当前模型 reliable_context 折算:体量未到阈值前不压缩(缓存全暖 + 不丢信息)。
# 换算比值走校准态(实报 usage 优先,回退 2.5)—— 门槛语义是 token 口径,chars 只是载体。 # 换算比值走校准态(实报 usage 优先,回退 2.5)—— 门槛语义是 token 口径,chars 只是载体。
ratio = self._context_ratio() ratio = self._context_ratio()
@ -591,6 +693,47 @@ class AgentLoop:
sc["function"]["name"]: (sc["function"].get("parameters") or {}).get("required") or [] sc["function"]["name"]: (sc["function"].get("parameters") or {}).get("required") or []
for sc in self.executor.schemas() for sc in self.executor.schemas()
} }
try:
return self._run_robust_stream(
llm_messages=llm_messages,
required_by_tool=required_by_tool,
llm_start_event=llm_start_event,
)
except ReasoningPhaseTimeout:
if self._is_cancelled():
return None, True
self._record_reasoning_guard("reasoning_timeout_retry", count=1)
self._emit({
"type": "warn",
"level": "info",
"msg": "推理阶段超过 90 秒,已降为 low 并压缩推理重试一次",
})
self._active_reasoning_effort = "low"
self._reasoning_usage.update({
"reasoning_effort": "low",
"reasoning_reason": "reasoning_guard_retry",
"reasoning_guard_retry": True,
})
retry_messages = self._with_reasoning_retry_instruction(llm_messages)
try:
return self._run_robust_stream(
llm_messages=retry_messages,
required_by_tool=required_by_tool,
llm_start_event=llm_start_event,
)
except ReasoningPhaseTimeout as exc:
if self._is_cancelled():
return None, True
self._record_reasoning_guard("reasoning_timeout_stop", count=2)
raise ReasoningGuardExhausted from exc
def _run_robust_stream(
self,
*,
llm_messages: List[dict],
required_by_tool: Dict[str, List[str]],
llm_start_event: dict,
) -> Tuple[Optional[Any], bool]:
return robust_stream( return robust_stream(
collect_stream=self._collect_stream_once, collect_stream=self._collect_stream_once,
nonstream=self._nonstream_once, nonstream=self._nonstream_once,
@ -605,6 +748,34 @@ class AgentLoop:
max_attempts=self._MAX_MALFORMED_ATTEMPTS, max_attempts=self._MAX_MALFORMED_ATTEMPTS,
) )
@staticmethod
def _with_reasoning_retry_instruction(llm_messages: List[dict]) -> List[dict]:
"""只改本次 provider 请求副本,不进入 Session/DB。"""
retry_messages = list(llm_messages)
insert_at = len(retry_messages)
for i in range(len(retry_messages) - 1, -1, -1):
if retry_messages[i].get("role") == "user":
insert_at = i
break
retry_messages.insert(insert_at, {
"role": "system",
"content": _REASONING_RETRY_INSTRUCTION,
})
return retry_messages
def _record_reasoning_guard(self, guard: str, *, count: int) -> None:
try:
record_agent_guard(
task_id=self.session.task_id,
user_id=self.user_id,
model_profile=model_profile_of(self.caps),
tool="(llm)",
guard=guard,
count=count,
)
except Exception:
pass
def _try_salvage_response(self, response: Any) -> bool: def _try_salvage_response(self, response: Any) -> bool:
"""尝试就地抢救本轮所有畸形 tool_call 的 arguments;成功才改写并返回 True。 """尝试就地抢救本轮所有畸形 tool_call 的 arguments;成功才改写并返回 True。
@ -679,11 +850,32 @@ class AgentLoop:
"""跑一次流式:攒 chunk + content delta 即时 emit,拼回完整 response。 """跑一次流式:攒 chunk + content delta 即时 emit,拼回完整 response。
返回 (response, cancelled_mid_stream)""" 返回 (response, cancelled_mid_stream)"""
chunks: List[Any] = [] chunks: List[Any] = []
reasoning_started_at: Optional[float] = None
reasoning_seen = False
output_seen = False
reasoning_timed_out = False
def _stop_stream() -> bool:
nonlocal reasoning_timed_out
# 取消始终优先,竞态时不把用户主动停止误记成自动保护。
if self._is_cancelled():
return True
if (
reasoning_seen
and not output_seen
and reasoning_started_at is not None
and time.monotonic() - reasoning_started_at
>= self._REASONING_PHASE_TIMEOUT_S
):
reasoning_timed_out = True
return True
return False
stream = self.llm.chat_stream( stream = self.llm.chat_stream(
messages=llm_messages, messages=llm_messages,
tools=self.executor.schemas(), tools=self.executor.schemas(),
reasoning_effort=self.caps.default_reasoning_effort or None, reasoning_effort=getattr(self, "_active_reasoning_effort", None),
cancel_check=self._is_cancelled, cancel_check=_stop_stream,
) )
cancelled = False cancelled = False
may_reroute = self.caps.family == "deepseek_v4" may_reroute = self.caps.family == "deepseek_v4"
@ -695,6 +887,12 @@ class AgentLoop:
break break
chunks.append(chunk) chunks.append(chunk)
tool_names = extract_delta_tool_names(chunk) tool_names = extract_delta_tool_names(chunk)
try:
has_tool_delta = bool(chunk.choices[0].delta.tool_calls)
except (AttributeError, IndexError, TypeError):
has_tool_delta = False
if tool_names or has_tool_delta:
output_seen = True
if ( if (
may_reroute may_reroute
and any(name in self._DEEPSEEK_NONSTREAM_TOOLS for name in tool_names) and any(name in self._DEEPSEEK_NONSTREAM_TOOLS for name in tool_names)
@ -712,11 +910,15 @@ class AgentLoop:
# _execute_tool_call 时机发更直观)。 # _execute_tool_call 时机发更直观)。
delta_text = extract_delta_content(chunk) delta_text = extract_delta_content(chunk)
if delta_text: if delta_text:
output_seen = True
self._emit({"type": "text", "delta": delta_text}) self._emit({"type": "text", "delta": delta_text})
# thinking 模型的推理 delta 也实时流出(reasoning 事件):深度推理可达 # thinking 模型的推理 delta 也实时流出(reasoning 事件):深度推理可达
# 分钟级,不发的话前端全程静止"思考中",用户以为卡死。 # 分钟级,不发的话前端全程静止"思考中",用户以为卡死。
delta_reasoning = extract_delta_reasoning(chunk) delta_reasoning = extract_delta_reasoning(chunk)
if delta_reasoning: if delta_reasoning:
if reasoning_started_at is None:
reasoning_started_at = time.monotonic()
reasoning_seen = True
self._emit({"type": "reasoning", "delta": delta_reasoning}) self._emit({"type": "reasoning", "delta": delta_reasoning})
reasoning_emitted = True reasoning_emitted = True
# interruptible stream 会在无新 chunk 的等待期直接因 cancel 结束迭代; # interruptible stream 会在无新 chunk 的等待期直接因 cancel 结束迭代;
@ -731,6 +933,10 @@ class AgentLoop:
if cancelled: if cancelled:
return None, True return None, True
if reasoning_timed_out:
if reasoning_emitted:
self._emit({"type": "reasoning_reset"})
raise ReasoningPhaseTimeout
# 用 litellm 官方 helper 拼回完整 response(包括 tool_calls 拼接 + usage)。 # 用 litellm 官方 helper 拼回完整 response(包括 tool_calls 拼接 + usage)。
# messages 参数仅用于失败时回填 prompt token 估算,正常路径 stream_options.include_usage # messages 参数仅用于失败时回填 prompt token 估算,正常路径 stream_options.include_usage
@ -756,7 +962,7 @@ class AgentLoop:
box["resp"] = self.llm.chat( box["resp"] = self.llm.chat(
messages=llm_messages, messages=llm_messages,
tools=self.executor.schemas(), tools=self.executor.schemas(),
reasoning_effort=self.caps.default_reasoning_effort or None, reasoning_effort=getattr(self, "_active_reasoning_effort", None),
) )
except BaseException as e: # noqa: BLE001 — 原样转抛回主线程 except BaseException as e: # noqa: BLE001 — 原样转抛回主线程
box["exc"] = e box["exc"] = e

View File

@ -14,7 +14,7 @@ from __future__ import annotations
from dataclasses import dataclass, field from dataclasses import dataclass, field
from typing import Any, List, Optional from typing import Any, List, Optional
from .capabilities import ModelCapabilities from .capabilities import REASONING_EFFORT_AUTO, ModelCapabilities
from .llm import LLM from .llm import LLM
@ -147,6 +147,8 @@ def probe_thinking(llm: LLM, caps: ModelCapabilities) -> ProbeResult:
caps.default_reasoning_effort caps.default_reasoning_effort
or (caps.reasoning_effort_levels[0] if caps.reasoning_effort_levels else None) or (caps.reasoning_effort_levels[0] if caps.reasoning_effort_levels else None)
) )
if effort == REASONING_EFFORT_AUTO:
effort = "high"
try: try:
resp = llm.chat( resp = llm.chat(
messages=[{"role": "user", "content": "Briefly: what is 17 * 23?"}], messages=[{"role": "user", "content": "Briefly: what is 17 * 23?"}],

View File

@ -1,4 +1,5 @@
import os import os
import tempfile
import unittest import unittest
from pathlib import Path from pathlib import Path
from unittest.mock import patch from unittest.mock import patch
@ -101,6 +102,15 @@ class LLMKwargsTests(unittest.TestCase):
[{"role": "user", "content": "hello"}], None, None, "high" [{"role": "user", "content": "hello"}], None, None, "high"
) )
def test_auto_is_rejected_at_provider_boundary(self) -> None:
llm = self._llm(
family="deepseek_v4", thinking_enabled=True, thinking_transport="extra_body"
)
with self.assertRaisesRegex(ValueError, "auto"):
llm._build_kwargs(
[{"role": "user", "content": "hello"}], None, None, "auto"
)
def test_flash_profile_matches_0731_capabilities(self) -> None: def test_flash_profile_matches_0731_capabilities(self) -> None:
caps = ModelCapabilities.load( caps = ModelCapabilities.load(
"deepseek_v4.flash", Path(__file__).resolve().parents[1] / "config" / "models" "deepseek_v4.flash", Path(__file__).resolve().parents[1] / "config" / "models"
@ -108,7 +118,7 @@ class LLMKwargsTests(unittest.TestCase):
self.assertTrue(caps.thinking_enabled) self.assertTrue(caps.thinking_enabled)
self.assertEqual(caps.reasoning_effort_levels, ["low", "high", "max"]) self.assertEqual(caps.reasoning_effort_levels, ["low", "high", "max"])
self.assertEqual(caps.default_reasoning_effort, "high") self.assertEqual(caps.default_reasoning_effort, "auto")
self.assertEqual(caps.max_output, 8192) self.assertEqual(caps.max_output, 8192)
self.assertEqual(caps.output_cny_per_mtoken, 4.752) self.assertEqual(caps.output_cny_per_mtoken, 4.752)
self.assertEqual(caps.cache_hit_cny_per_mtoken, 0.0504) self.assertEqual(caps.cache_hit_cny_per_mtoken, 0.0504)
@ -118,6 +128,25 @@ class LLMKwargsTests(unittest.TestCase):
self.assertEqual(caps.thinking_transport, "extra_body") self.assertEqual(caps.thinking_transport, "extra_body")
self.assertEqual(caps.reasoning_replay, "tool_turn") self.assertEqual(caps.reasoning_replay, "tool_turn")
pro = ModelCapabilities.load(
"deepseek_v4.pro", Path(__file__).resolve().parents[1] / "config" / "models"
)
self.assertEqual(pro.default_reasoning_effort, "medium")
def test_auto_profile_requires_low_and_high_levels(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "test.yaml"
path.write_text(
"family: test\nvariants:\n bad:\n"
" thinking_enabled: true\n"
" thinking_transport: extra_body\n"
" reasoning_effort_levels: [low, max]\n"
" default_reasoning_effort: auto\n",
encoding="utf-8",
)
with self.assertRaisesRegex(ValueError, "low/high"):
ModelCapabilities.load("test.bad", Path(tmp))
def test_other_controllable_profiles_declare_transport(self) -> None: def test_other_controllable_profiles_declare_transport(self) -> None:
models_dir = Path(__file__).resolve().parents[1] / "config" / "models" models_dir = Path(__file__).resolve().parents[1] / "config" / "models"

View File

@ -37,6 +37,7 @@ def _make_loop(stream_results, nonstream_results):
loop.caps = SimpleNamespace(reliable_context=64_000, family="test", variant="t") loop.caps = SimpleNamespace(reliable_context=64_000, family="test", variant="t")
loop.session = SimpleNamespace(messages=[], task_id="test-task") loop.session = SimpleNamespace(messages=[], task_id="test-task")
loop.user_id = "test-user" # 无 DB:_log_empty_response 落库路径静默跳过 loop.user_id = "test-user" # 无 DB:_log_empty_response 落库路径静默跳过
loop.user_root = None
loop.executor = SimpleNamespace(schemas=lambda: []) # required_by_tool 取值用,空即可 loop.executor = SimpleNamespace(schemas=lambda: []) # required_by_tool 取值用,空即可
loop.events = [] loop.events = []
loop._emit = loop.events.append loop._emit = loop.events.append

View File

@ -0,0 +1,269 @@
from __future__ import annotations
import unittest
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
from uuid import uuid4
from core.loop import (
AgentLoop,
ReasoningGuardExhausted,
ReasoningPhaseTimeout,
resolve_reasoning_effort,
tool_result_succeeded_for_reasoning,
)
from core.probe import probe_thinking
def _text_response(text: str = "ok"):
return SimpleNamespace(
choices=[SimpleNamespace(message=SimpleNamespace(
content=text, tool_calls=None,
))],
usage=None,
)
def _reasoning_chunk(text: str = "thinking"):
return SimpleNamespace(choices=[SimpleNamespace(delta=SimpleNamespace(
reasoning_content=text, content=None, tool_calls=None,
))])
def _stream_loop(*, configured: str = "auto", cancelled=False) -> AgentLoop:
loop = object.__new__(AgentLoop)
loop.caps = SimpleNamespace(
default_reasoning_effort=configured,
reliable_context=64_000,
family="deepseek_v4",
variant="flash",
reasoning_replay="tool_turn",
native_image_input=False,
)
loop.session = SimpleNamespace(
messages=[{"role": "user", "content": "hello"}], task_id="task",
)
loop.user_id = "user"
loop.user_root = None
loop.working_dir = Path(".")
loop.executor = SimpleNamespace(schemas=lambda: [])
loop.cancel_check = (lambda: cancelled)
loop.events = []
loop._emit = loop.events.append
loop._ctx_chars_per_token = 2.5
loop._last_sent_chars = 0
loop._last_had_native_images = False
loop._llm_call_count = 0
loop._previous_tool_step_succeeded = None
loop._reasoning_usage = {}
return loop
class ReasoningDecisionTests(unittest.TestCase):
def test_auto_first_call_is_high(self) -> None:
self.assertEqual(
resolve_reasoning_effort(
"auto", first_call=True, previous_tools_succeeded=None
),
("high", "first_call"),
)
def test_auto_success_is_low_and_failure_is_high(self) -> None:
self.assertEqual(
resolve_reasoning_effort(
"auto", first_call=False, previous_tools_succeeded=True
),
("low", "previous_tools_succeeded"),
)
self.assertEqual(
resolve_reasoning_effort(
"auto", first_call=False, previous_tools_succeeded=False
),
("high", "previous_tools_failed"),
)
def test_fixed_effort_is_unchanged(self) -> None:
self.assertEqual(
resolve_reasoning_effort(
"max", first_call=False, previous_tools_succeeded=True
),
("max", "configured"),
)
def test_tool_success_error_and_quality_gate_signals(self) -> None:
self.assertTrue(tool_result_succeeded_for_reasoning(
"read", '{}', "file content", productive=True,
))
self.assertFalse(tool_result_succeeded_for_reasoning(
"read", '{}', "[Error] missing", productive=False,
))
self.assertFalse(tool_result_succeeded_for_reasoning(
"shell", '{"command":"build"}',
"created\n[产物机检 ERROR] 发现整页贴图", productive=True,
))
def test_probe_resolves_auto_before_provider_call(self) -> None:
efforts = []
def chat(**kwargs):
efforts.append(kwargs["reasoning_effort"])
return SimpleNamespace(choices=[SimpleNamespace(message=SimpleNamespace(
reasoning_content="brief reasoning", content="391",
))])
caps = SimpleNamespace(
thinking_enabled=True,
default_reasoning_effort="auto",
reasoning_effort_levels=["low", "high", "max"],
)
result = probe_thinking(SimpleNamespace(chat=chat), caps)
self.assertEqual(efforts, ["high"])
self.assertEqual(result.status, "ok")
class ReasoningGuardTests(unittest.TestCase):
def test_pure_reasoning_timeout_resets_stream(self) -> None:
loop = _stream_loop()
def chat_stream(**kwargs):
yield _reasoning_chunk()
self.assertTrue(kwargs["cancel_check"]())
loop.llm = SimpleNamespace(chat_stream=chat_stream)
loop._active_reasoning_effort = "high"
loop._REASONING_PHASE_TIMEOUT_S = 0
with self.assertRaises(ReasoningPhaseTimeout):
loop._collect_stream_once(loop.session.messages)
self.assertEqual(loop.events, [
{"type": "reasoning", "delta": "thinking"},
{"type": "reasoning_reset"},
])
@patch("core.loop.record_agent_guard")
def test_timeout_retries_once_with_low_and_ephemeral_instruction(self, guard) -> None:
loop = _stream_loop()
calls = []
def run_robust(**kwargs):
calls.append((loop._active_reasoning_effort, kwargs["llm_messages"]))
if len(calls) == 1:
raise ReasoningPhaseTimeout
return _text_response(), False
loop._run_robust_stream = run_robust
response, cancelled = loop._stream_llm()
self.assertFalse(cancelled)
self.assertEqual(response.choices[0].message.content, "ok")
self.assertEqual([call[0] for call in calls], ["high", "low"])
self.assertEqual(len(calls[0][1]), 1)
self.assertEqual(len(calls[1][1]), 2)
self.assertEqual(calls[1][1][0]["role"], "system")
self.assertEqual(loop.session.messages, [{"role": "user", "content": "hello"}])
self.assertEqual(loop._reasoning_usage, {
"reasoning_config": "auto",
"reasoning_effort": "low",
"reasoning_reason": "reasoning_guard_retry",
"reasoning_guard_retry": True,
})
guard.assert_called_once()
@patch("core.loop.record_agent_guard")
def test_second_timeout_stops(self, guard) -> None:
loop = _stream_loop()
loop._run_robust_stream = MagicMock(side_effect=[
ReasoningPhaseTimeout(), ReasoningPhaseTimeout(),
])
with self.assertRaises(ReasoningGuardExhausted):
loop._stream_llm()
self.assertEqual(guard.call_count, 2)
@patch("core.loop.record_agent_guard")
def test_user_cancel_wins_over_guard_retry(self, guard) -> None:
loop = _stream_loop(cancelled=True)
loop._run_robust_stream = MagicMock(side_effect=ReasoningPhaseTimeout())
response, cancelled = loop._stream_llm()
self.assertIsNone(response)
self.assertTrue(cancelled)
guard.assert_not_called()
def test_nonstream_fallback_uses_resolved_effort(self) -> None:
loop = _stream_loop()
efforts = []
def chat(**kwargs):
efforts.append(kwargs["reasoning_effort"])
return _text_response()
loop.llm = SimpleNamespace(chat=chat)
loop._active_reasoning_effort = "low"
response = loop._nonstream_once(loop.session.messages)
self.assertEqual(response.choices[0].message.content, "ok")
self.assertEqual(efforts, ["low"])
class _Session:
def __init__(self):
self.task_id = uuid4()
self.messages = [{"role": "user", "content": "hello"}]
self.appended = []
def append(self, message, **_kwargs):
self.messages.append(message)
self.appended.append(message)
return uuid4()
class ReasoningPersistenceTests(unittest.TestCase):
def test_exhausted_guard_does_not_persist_partial_assistant(self) -> None:
session = _Session()
loop = AgentLoop(
llm=MagicMock(), executor=MagicMock(), session=session,
capabilities=SimpleNamespace(max_iterations=1),
user_id=uuid4(), working_dir=Path("."),
)
loop._maybe_fold_context = MagicMock()
loop._stream_llm = MagicMock(side_effect=ReasoningGuardExhausted())
result = loop.run_persisted_turn()
self.assertEqual(result, "[stopped: reasoning timeout]")
self.assertEqual(session.appended, [])
self.assertEqual(loop.events if hasattr(loop, "events") else [], [])
@patch("core.loop.record_chat_usage")
def test_successful_chat_records_reasoning_metadata(self, record_usage) -> None:
session = _Session()
caps = SimpleNamespace(
max_iterations=1, family="deepseek_v4", variant="flash",
input_cny_per_mtoken=0, output_cny_per_mtoken=0,
cache_hit_cny_per_mtoken=0, pricing={},
)
loop = AgentLoop(
llm=MagicMock(), executor=MagicMock(), session=session,
capabilities=caps, user_id=uuid4(), working_dir=Path("."),
)
loop._maybe_fold_context = MagicMock()
def stream():
loop._reasoning_usage = {
"reasoning_config": "auto",
"reasoning_effort": "high",
"reasoning_reason": "first_call",
"reasoning_guard_retry": False,
}
return _text_response("done"), False
loop._stream_llm = MagicMock(side_effect=stream)
self.assertEqual(loop.run_persisted_turn(), "done")
units = record_usage.call_args.kwargs["extra_units"]
self.assertEqual(units["reasoning_config"], "auto")
self.assertEqual(units["reasoning_effort"], "high")
self.assertEqual(units["reasoning_reason"], "first_call")
self.assertFalse(units["reasoning_guard_retry"])
if __name__ == "__main__":
unittest.main()