251 lines
12 KiB
Python
251 lines
12 KiB
Python
"""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 .common import INSTANCE
|
||
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 POST:用户消息已和 tasks.running 同事务提交,
|
||
worker 从 Session 恢复后直接处理最后一条 user,避免重复落库。其余入口保持旧行为。
|
||
"""
|
||
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
|
||
|
||
# 解析/建该渠道常驻 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
|
||
|
||
# 抢 run 锁:正忙 → 提示稍候(同用户串行;ClawBot loop 本就串行,wecom 回调靠此挡并发)
|
||
with session_scope() as s:
|
||
row = s.execute(
|
||
select(Task.run_status).where(Task.task_id == tid).with_for_update()
|
||
).first()
|
||
if row is None:
|
||
return "[出错] 对话 task 不存在"
|
||
if row.run_status in ("running", "cancelling"):
|
||
return "上一条还在处理中,请稍候再发。"
|
||
s.execute(update(Task).where(Task.task_id == tid).values(
|
||
run_status="running", run_error=None, run_owner=INSTANCE or None))
|
||
|
||
broker.start(tid)
|
||
runner = asyncio.create_task(asyncio.to_thread(
|
||
run_agent_bg, tid, uid, text, "", "", False,
|
||
))
|
||
app.state.inflight[runner] = tid
|
||
runner.add_done_callback(lambda t: app.state.inflight.pop(t, None))
|
||
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)
|