zcbot/web/runs.py

251 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""BG run worker + 渠道入站对话核心(从 app.py 析出,2026-07-23 拆分)。
`run_agent_bg` 是「一次 run」的工作线程本体(build_agent → agent.run → 写终态),
web 路由与渠道回调都经它起 run;`run_channel_conversation` 是微信/企微共用的
入站对话编排(§8.7)。依赖 web 层 broker/sink 桥事件,故住 web/ 不下沉 core。
"""
from __future__ import annotations
import asyncio
from uuid import UUID
from sqlalchemy import select, update
from core.paths import from_db_path
from core.storage import session_scope
from core.storage.models import Task
from core.storage.telemetry import record_run_error
from core.toolfail import alert_provider_critical
from .broker import broker
from .sinks import WebEventSink
def run_agent_bg(
task_id: UUID, user_id: UUID, user_message: str,
image_variant: str = "", video_variant: str = "",
scheduled: bool = False,
*, user_message_persisted: bool = False,
) -> None:
"""工作线程:`build_agent(resume=True)` → 装 WebEventSink + cancel_check → `agent.run` → 写 tasks.run_status。
sink 通过 broker.emit 桥事件回 asyncio loop;agent.run 是 sync,所以在 to_thread 跑。
user_id 必须从 JWT 那侧透传过来 —— 决定 memory_block 读哪个 per-user 子树。
cancel_check 桥 broker.is_cancelled,loop 在 stream chunk 间 + 工具调用之间 poll;
cancel 延迟 ~ 单 chunk 间隔(100ms 级);seedance 轮询间也读这个 cancel_check 用于
用户停止按钮(必须在 build_agent 阶段就传进去,因为 SeedanceTool ctor 持有它,
不能像以前那样 build_agent 返回后再赋 agent.cancel_check)。
`ok` 收尾回 `idle`;`cancelled`(用户停止)与 `error` 一样落持久终态 —— 前端据此
补「已停止」/ 错误卡(扛过收尾重渲),下次起新 run(post_message 写 running)覆盖清掉。
image_variant / video_variant:本 run 用哪个 image/video variant 装 tool(空 → yaml 第一个)。
随消息 POST 传进来,不入 DB —— UI 下拉的选择就跟在这一条消息上生效。
user_message_persisted=True 供 Web / 渠道 / 定时生命周期入口:用户消息已和
tasks.running 同事务提交worker 从 Session 恢复后直接处理最后一条 user
避免重复落库。CLI 等非 Web 入口保持旧行为。
"""
from core.agent_builder import build_agent, sync_task_tokens
cancel_check = lambda tid=task_id: broker.is_cancelled(tid)
try:
broker.emit(task_id, {"type": "run_start"})
agent, session, sid, task_state, task_dir = build_agent(
session_id=str(task_id), resume=True, user_id=user_id,
image_variant=image_variant,
video_variant=video_variant,
cancel_check=cancel_check,
scheduled_run=scheduled,
)
agent.sink = WebEventSink(broker, task_id)
result = (
agent.run_persisted_turn()
if user_message_persisted
else agent.run(user_message)
)
sync_task_tokens(task_state)
# 收尾终态:agent.run 在任一取消路径都 return "[cancelled]"(loop.py)——
# 用户停止 → 落持久 cancelled(前端 renderPersistedRunTerminal 据此补「已停止」卡,
# 扛过收尾 loadMessages 整屏重建);正常完成 → 回 idle。两者都清 run_error。
# cancelled 与 error 同为持久终态,下次起新 run(post_message 写 running)自然覆盖。
final_status = "cancelled" if result == "[cancelled]" else "idle"
with session_scope() as s:
s.execute(
update(Task).where(Task.task_id == task_id).values(
run_status=final_status, run_error=None,
)
)
# 当前 working_dir 不能在 agent 仍执行工具时改名executor/system prompt 都握
# 旧 cwd。rename_working_dir 工具只登记意图;正常完成并退出 running 后,
# 在 done 事件前统一走 DB-aware 服务。取消时不执行,避免用户停止后仍发生变更。
rename_to = agent.deferred_actions.rename_working_dir_to
if final_status == "idle" and rename_to:
from core.working_dirs import WorkingDirRenameError, rename_working_dir
try:
renamed = rename_working_dir(
user_id=user_id,
old_path=task_dir,
new_path=task_dir.parent / rename_to,
)
broker.emit(task_id, {
"type": "warn",
"level": "info",
"msg": f"工作目录已重命名为 {renamed.new_path.name}",
})
except WorkingDirRenameError as e:
broker.emit(task_id, {
"type": "warn",
"msg": f"工作目录改名未执行:{e}",
})
except Exception as e:
# 延迟附加动作失败不应把已经成功的对话 run 改标 error保留目录原状
# 给当前订阅者明确提示,下轮可重试。
broker.emit(task_id, {
"type": "warn",
"msg": (
"工作目录改名未执行:"
f"{type(e).__name__}: {e}"
),
})
except Exception as e:
err = f"{type(e).__name__}: {e}"
broker.emit(task_id, {"type": "error", "msg": err})
mp = ""
try:
with session_scope() as s:
mp = s.execute(
select(Task.model_profile).where(Task.task_id == task_id)
).scalar_one_or_none() or ""
s.execute(
update(Task).where(Task.task_id == task_id).values(
run_status="error", run_error=err,
)
)
except Exception:
pass # 已 emit error 给前端,DB 写失败不放大噪声
# 留痕 + 告警(0.58.21,反 2a1bc25d 教训:Zai 余额不足连挂 3 次续跑无人知):
# run_error 列只留最后一次,usage_events(kind=run_error)才是聚合面板/巡检
# 邮件的完整数据源;余额/认证类 provider 级错误另走即时邮件(6h 签名冷却)。
# 两路都静默失败 —— 留痕/告警绝不能在错误路径上再抛。
try:
record_run_error(
task_id=task_id, user_id=user_id, model_profile=mp, error=err,
)
except Exception:
pass
alert_provider_critical(err, task_id=task_id, model_profile=mp)
finally:
broker.clear_cancel(task_id)
broker.close(task_id)
async def run_channel_conversation(app, uid, text, attachments, *, channel):
"""渠道无关的入站对话核心(§8.7):解析/建该用户该渠道常驻 task → 落盘附件 → 抢 run 锁
→ run_agent_bg → 取 assistant 回复文本。两渠道各一张会话 task,互不串扰。
channel:'wechat'(个人微信 ClawBot,绑定快照取 chat_task_id)| 'wecom'(企业微信,
wecom 绑定行取 chat_task_id)。attachments:已下载解密的入站附件(可空,wecom 暂只收文本)。
返回回复文本(供 ClawBot 回流 / wecom 主动推回)。
"""
from core.wechat import service as _wx
from core.wechat.ilink import attachment_basename
from core.wechat.inbound import extract_last_assistant_text
from .run_lifecycle import (
RunScheduleError,
RunTaskBusy,
RunTaskNotFound,
claim_run_with_message,
schedule_claimed_run,
)
# 解析/建该渠道常驻 chat task(不存在自动建)—— 与 push 记录(send_to_user)共用
# ensure_channel_chat_task,避免两条建 task 路径漂移。wechat 无 binding → 返回 None。
tid = await asyncio.to_thread(_wx.ensure_channel_chat_task, uid, channel)
if tid is None:
return ""
# 手动「新话题」命令:硬重置上下文窗口(base=总数),不跑 agent,直接回执。之前的
# 对话全留 DB(网页端可翻),只是不再喂模型。纯文本命令,有附件则不当命令处理。
if not attachments and text.strip() in _wx.NEW_TOPIC_COMMANDS:
await asyncio.to_thread(_wx.reset_channel_context, tid, hard=True)
return "已开启新话题,之前的对话已归档(网页端仍可查看完整历史)。"
# 快捷指令展开(渠道无关,见 core/shortcuts.py):整条精确命中触发词 → 文本换成完整指令
# 再照常跑;不进上下文、不问模型。放「新话题」命令之后、附件/gap 之前:展开后的文本仍会
# 被下面的附件行追加,故打「简报」+ 附图也成立。
from core.agent_builder import resolve_workspace as _resolve_ws
from core import shortcuts as _shortcuts
_ws = await asyncio.to_thread(_resolve_ws, None)
text, _hit = await asyncio.to_thread(_shortcuts.expand, _ws, uid, text)
if _hit:
print(f"[shortcut] {str(uid)[:8]} '{_hit}' expanded")
# 自动分段:距上次消息超过 gap 阈值 → 软重置(base=最后一条 user 消息 idx,保留上一轮
# 原文做续聊锚点)。在入站消息落库前判断,故 last_at 取的是上一轮的时间。push 不走这。
from core.agent_builder import load_config as _load_config
gap_hours = float(
(_load_config().get("channel") or {}).get(
"session_gap_hours", _wx.SESSION_GAP_HOURS_DEFAULT
)
)
await asyncio.to_thread(_wx.maybe_gap_reset, tid, gap_hours)
# 落盘入站附件到 <wd>/inbound/,拼 [用户上传的...] 行进 text(复用 web 端粘贴图约定)
if attachments:
from datetime import datetime
with session_scope() as s:
wd_db = s.execute(
select(Task.working_dir).where(Task.task_id == tid)
).scalar_one()
inbound_dir = from_db_path(wd_db) / "inbound"
inbound_dir.mkdir(parents=True, exist_ok=True)
ts = datetime.now().strftime("%Y%m%d-%H%M%S")
lines: list[str] = []
for i, att in enumerate(attachments):
if not att.data:
continue
base = attachment_basename(att)
name = f"{ts}-{i}-{base}"
(inbound_dir / name).write_bytes(att.data)
rel = f"inbound/{name}"
tag = "[用户上传的参考图]" if att.kind == "image" else "[用户上传的文件]"
lines.append(f"{tag} {rel}")
if lines:
extra = "\n".join(lines)
text = f"{text}\n\n{extra}" if text.strip() else extra
# 渠道消息也与 running 原子提交,消除进程在 worker 启动前退出时的丢输入窗口。
try:
claim_run_with_message(tid, uid, text)
except RunTaskNotFound:
return "[出错] 对话 task 不存在"
except RunTaskBusy:
return "上一条还在处理中,请稍候再发。"
try:
runner = schedule_claimed_run(app, tid, uid, text)
except RunScheduleError as e:
return f"[出错] {e}"
await runner
with session_scope() as s:
st = s.execute(
select(Task.run_status, Task.run_error).where(Task.task_id == tid)
).first()
if st is not None and st.run_status == "error":
return f"[出错] {st.run_error}"
reply = await asyncio.to_thread(extract_last_assistant_text, tid)
return reply or "(本轮无文本回复)"
async def transcribe_wecom_voice(media_id: str) -> str:
"""企业微信入站语音 → 文本:media/get 下 AMR → `core.audio.transcribe_voice`
(ffmpeg 解 16k PCM → 讯飞 IAT 整段转写,60s 截断在里面)。
返回转写文本(空串=没听出内容)。失败抛异常,message 可直接给用户看
(XfyunASRNotConfigured / AudioNotConfigured 自带配置指引)。
"""
from core import audio
from core.wechat import wecom
data, _fname = await asyncio.to_thread(wecom.download_media, media_id)
return await audio.transcribe_voice(data)