deerflow-code/offline-backend-20260512/backend/app/gateway/roundtable_job_executor.py
2026-09-07 18:24:55 +08:00

724 lines
36 KiB
Python
Raw Permalink 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.

"""圆桌后台作业执行器(Phase 3a)。
把 Phase 2 的编排引擎(``run_orchestration``)作为**后台 asyncio 任务**跑起来,
survive 发起请求的生命周期。每个任务:
1. ``set_current_user``:用作业的 user_id 设进程内用户上下文,让引擎里的进度落库
(``JobProgressWriter`` → ``RoundtableJobRepository.update_progress(user_id=AUTO)``)
正确归属(与 scheduled task 后台执行同一套做法);
2. 用 ``gateway_factory`` 构造一个 ``RoundtableGateway``(Phase 3a 默认模拟网关,
Phase 3b 换成进程内真网关),跑完整编排(循环 → 共识 → 报告)。
任务句柄登记在 ``_tasks`` 里,供 ``cancel`` 用;完成后自动摘除。
"""
from __future__ import annotations
import asyncio
import logging
import os
import uuid as _uuid
from collections.abc import Callable
from contextlib import suppress
from dataclasses import dataclass, field
from datetime import UTC, datetime, timedelta
from typing import Any
from deerflow.agents.roundtable_orchestrator import (
JobProgressWriter,
SeatRef,
run_orchestration,
)
from deerflow.persistence.roundtable_diagnostics import DiagnosticsRecorder, get_default_store
from deerflow.persistence.roundtable_drafts.sql import DraftConcurrentWriteError
from deerflow.runtime.user_context import reset_current_user, set_current_user
logger = logging.getLogger(__name__)
# 本 worker 进程的唯一标识(Phase 3 租约持有者)。pid 保证同机多 worker 进程互不相同,
# 追加随机后缀避免 pid 复用导致的极端撞车。跨进程「谁在跑这个作业」由它写进 lease_owner。
WORKER_ID = f"w-{os.getpid()}-{_uuid.uuid4().hex[:8]}"
# 租约 TTL 与心跳间隔(Phase 3)。心跳每 _HEARTBEAT_INTERVAL 秒续租一次,把 lease_until
# 推到 now + _LEASE_TTL_SECONDS。只要 worker 活着,租约永远不过期、别的 worker 抢不走;
# worker 崩溃后心跳停止,最长 _LEASE_TTL_SECONDS 后租约过期,dispatcher 在其它 worker 上
# 重新领取(attempt+1)—— 崩溃可接管的一致性来自 DB 租约而非进程内存。
# 约束:_LEASE_TTL_SECONDS 必须明显大于 _HEARTBEAT_INTERVAL(留足续租余量)。
_LEASE_TTL_SECONDS = 60
_HEARTBEAT_INTERVAL = 20.0
# 跨进程取消轮询间隔(秒)。cancel 端点把 ``status='cancel_requested'`` 写进 DB(权威意图),
# 实际跑作业的 worker 上的看门狗协程定期读 DB 检测到后,收口 ``cancelled`` 终态并对本进程的
# 编排 task 调 ``cancel()``。多 worker 下 cancel 请求可能打到非执行 worker,靠此机制最长延迟
# 一个轮询周期即可让幽灵作业停下。
_CANCEL_POLL_INTERVAL = 5.0
# 草稿乐观锁重试次数。每次尝试都带 expected_version(冲突则重读重试);重试全部耗尽后
# **绝不退化为无版本检查的 last-write-wins 盲写**,而是给作业置 draft_persist_pending
# 标记,由 outbox 补偿器(Phase 3)投影结果到草稿——一致性来自 DB 标记而非运气。
_DRAFT_WRITE_RETRIES = 3
# 「后台挂起」的 Step3 自动产出(方案总结报告 md + 结果绘制 HTML)总开关。
# 历史:2026-06-05 因后台 start_run 路径报告不稳定临时禁用;2026-06-12 定位到根因是
# ``_enforce_agent_access`` 对进程内虚拟 request 的 ``request.state`` 裸访问(AttributeError
# 杀死编排,见 tests/test_roundtable_inprocess_request.py),修复后恢复开启。
# 现 Step2 收口后自动跑:① roundtable-summary 写「方案总结报告.md」② roundtable-report
# 画「方案总览.html」,两段故障隔离(engine._run_step3),产物按手动同款形状写回草稿。
ROUNDTABLE_BACKGROUND_REPORT_ENABLED = True
class _JobUser:
"""最小 CurrentUser:满足 user_context 的 ``.id`` 结构协议。"""
def __init__(self, user_id: str) -> None:
self.id = user_id
# 头像配色:与前端 getAvatarTypeFor 对齐(内置圆桌席位固定色,其余按序轮换)。
_BUILTIN_AVATAR = {
"roundtable-intelligence": "purple",
"roundtable-environment": "emerald",
"roundtable-solution-design": "human",
"roundtable-risk-review": "amber",
"roundtable-execution-plan": "pink",
"roundtable-summary": "blue",
}
_FALLBACK_AVATAR = ["blue", "purple", "emerald", "amber", "pink"]
def _avatar_for(agent_id: str, index: int) -> str:
return _BUILTIN_AVATAR.get(agent_id) or _FALLBACK_AVATAR[index % len(_FALLBACK_AVATAR)]
# 工具步骤卡的回退文案。前端 getStepDisplay 对 tool 步骤是从 key+args 推图标/文案的,
# 存的 label 基本不展示——这里给个合理回退即可(与前端默认文案对齐)。
_STEP_LABEL = {
"agent_orchestration": "派活给子智能体",
"write_file": "写入文件",
"str_replace": "写入文件",
"read_file": "读取文件",
"ls": "列出文件夹",
"bash": "执行命令",
"web_search": "搜索相关信息",
"ask_clarification": "需要你的协助",
"present_files": "展示文件",
}
def _steps_from_tool_calls(tool_calls: list[dict], dialogue_id: object, ts: str) -> list[dict]:
"""把一轮 tool_calls 映射成前端 ``Step2Dialogue.steps``(工具步骤卡)。
- ``key`` = ``tool_calling:<toolName>``——前端 ``getStepDisplay`` 据此推图标 + 文案;
- ``id`` = ``<dialogueId>-<i>``——对齐前端 ``${bubbleId}-${stepIndex}``;
- ``args`` 原样带上(write_file 的 path/content、派活的 agent_name/task 等)。
"""
steps: list[dict] = []
for i, tc in enumerate(tool_calls or []):
name = tc.get("name") or ""
if not name:
continue
steps.append(
{
"id": f"{dialogue_id}-{i}",
"key": f"tool_calling:{name}",
"label": _STEP_LABEL.get(name, f"使用 “{name}” 工具"),
"ts": ts,
"args": tc.get("args") or {},
}
)
return steps
def _seat_stub_messages(dialogues: list[dict]) -> dict[str, list[dict]]:
"""从各对话的 write_file 调用重建 ``seatStubMessages``(threadId → [stub ai message])。
对齐前端 ``upsertSeatToolCall`` 的 stub 形状与 id 约定(消息 ``roundtable-stub-<tid>``、
工具调用 ``tc-<tid>-<i>``),使加载草稿后点开产物的虚拟 URL 能正确解析到内容。
同一 thread 跨轮写多次按 path 去重(后写覆盖),与前端一致。
"""
by_tid: dict[str, dict[str, dict]] = {}
for d in dialogues:
tid = d.get("threadId")
if not tid:
continue
for tc in d.get("toolCalls") or []:
if tc.get("name") != "write_file":
continue
path = str((tc.get("args") or {}).get("path") or "")
by_tid.setdefault(tid, {})[path] = tc # 后写覆盖
out: dict[str, list[dict]] = {}
for tid, by_path in by_tid.items():
tool_calls = [
{"id": f"tc-{tid}-{i}", "name": "write_file", "args": tc.get("args") or {}}
for i, tc in enumerate(by_path.values())
]
out[tid] = [{"id": f"roundtable-stub-{tid}", "type": "ai", "content": "", "tool_calls": tool_calls}]
return out
def _to_step2_dialogue(d: dict, idx_by_agent: dict[str, int]) -> dict:
"""把作业的中性对话记录映射成前端 Step2Dialogue 形状(含 avatarType / steps /
agentThreadId / time),使写进草稿 step2.runs 后,前端 hydrateStep2Run 能像手动跑
一样直接渲染——步骤卡、时间戳、点开产物全部可用。"""
dialogue_id = d.get("id")
time_str = d.get("time") or ""
base: dict = {
"id": dialogue_id,
"content": d.get("content") or "",
"time": time_str,
"streaming": False,
"steps": _steps_from_tool_calls(d.get("toolCalls") or [], dialogue_id, time_str),
}
thread_id = d.get("threadId")
if thread_id:
base["agentThreadId"] = thread_id
if d.get("role") == "leader":
return {**base, "sender": "总控协调", "avatarType": "blue", "confidence": "派活说明"}
aid = d.get("agentId") or ""
return {
**base,
"sender": d.get("name") or aid,
"avatarType": _avatar_for(aid, idx_by_agent.get(aid, 0)),
"confidence": "子智能体交付",
}
@dataclass
class StartParams:
job_id: str
user_id: str
agents: list[SeatRef]
seed_message: str
model: str | None = None
mode: str = "recommend"
intent_text: str = ""
thread_ids: dict[str, str] | None = None
coordinator_name: str = "roundtable-coordinator"
# 关联草稿:作业完成后把 step3 报告写回该草稿,使「查看报告」可用。
draft_id: str | None = None
# 外部任务 id(taskId 深链会商)。非空 = task 作业:draft_id 指向独立的
# roundtable_task_drafts 表,研讨/报告写回 task-draft 存储(不分权)。
task_id: str | None = None
# 是否为「续跑」(resume 澄清)。首启=False → 真网关建全新线程(不复用前端线程,
# 避免与前端本地研讨的活跃 run 撞 409);续跑=True → 复用作业自己的线程。
is_resume: bool = False
# chain 模式来源链条(写进草稿 run 用,展示)。
chain: dict[str, Any] | None = None
# dag 模式分层编排计划(前端 OrchestrationPlan dict):驱动引擎 _run_dag + 写进草稿 run
# 供前端流程图渲染。None = 非 dag(引擎按 mode 回退 chain/recommend)。
orchestration_plan: dict[str, Any] | None = None
# 是否自动产出 Step3(方案总结报告 + 结果绘制)。默认沿用全局开关(现已恢复 True,
# 见 ROUNDTABLE_BACKGROUND_REPORT_ENABLED)。测试可显式覆盖。
enable_report: bool = ROUNDTABLE_BACKGROUND_REPORT_ENABLED
# 席位执行模式(快速/思考/专业问答/多智能体)派生参数,沿用 Step2 选择。仅 special 席位用:
# 覆盖 roundtable_run_policy("seat") 的 thinking/reasoning(None → policy 默认);
# subagent 仅 ultra 档为真。透传给进程内网关(InProcessRoundtableGateway)的席位 run。
seat_thinking_enabled: bool | None = None
seat_reasoning_effort: str | None = None
seat_subagent_enabled: bool = False
# 每席位「技能识别强化」开关映射 {agent_id: bool}(沿用 Step2 所选业务链条里每个席位的
# 配置)。透传给进程内网关,run_seat 按 agent_id 查表,与各席位 agent 自身 config.yaml 的
# 同名开关取 OR。默认空 dict。
seat_skill_directives: dict[str, bool] = field(default_factory=dict)
# 每席位「推理深度覆盖」映射 {agent_id: SeatMode}(值 flash/thinking/pro/ultra,沿用 Step2 所选
# 业务链条里每个席位单独配的推理深度)。透传给进程内网关,run_seat 按 agent_id 查表:命中 → 用该
# 档派生的 thinking/reasoning/subagent 覆盖作业级 seatMode;未命中 → 跟随作业级。默认空 dict。
seat_modes: dict[str, str] = field(default_factory=dict)
# 是否在研讨前先跑一轮**全席位并行『取数』**(业务链条逐条 opt-in):各席位并行只调技能取数、
# 不研讨,产出作为共享资料注入后续研讨轮。透传给 run_orchestration 的 gather_first。默认 False。
gather_first: bool = False
# 深链接业务码(rwfx→6BF 等)。透传给 run_orchestration → build_summary_prompt:写实业务额外注入
# 「业务链完整性核对」,让后台总结报告核对产出相对业务链完不完整。普通会商为空 → 不注入。
business_code: str | None = None
# Phase 3 租约持有者(= dispatcher 领取时写入的 worker_id)。执行器据此:① 所有进度落库
# 带 lease_owner(update_progress 租约门,非持租者的过期写被拒);② 心跳续租;③ 看门狗
# 检测到 cancel_requested 时以持租者身份收口 cancelled。None = 旧路径(不走租约门)。
lease_owner: str | None = None
# ── Phase 3:input_snapshot 序列化 / 反序列化 ─────────────────────────────────
#
# 路由 start/resume 时把整份启动入参序列化成 JSON 存进 ``input_snapshot``;dispatcher
# 在**任意 worker** 领取 queued/租约过期作业后,据此无损重建 StartParams —— 这是「跨进程
# 可接管 / 崩溃可恢复」的前提(不再依赖发起请求的那个 worker 的进程内存)。
#
# 快照 schema 版本 ``_SNAPSHOT_VERSION``:将来字段变更时升版并兼容读旧版。
_SNAPSHOT_VERSION = 1
def build_start_snapshot(
*,
user_id: str,
agents: list[SeatRef],
seed_message: str,
model: str | None = None,
mode: str = "recommend",
intent_text: str = "",
thread_ids: dict[str, str] | None = None,
coordinator_name: str = "roundtable-coordinator",
draft_id: str | None = None,
task_id: str | None = None,
is_resume: bool = False,
chain: dict[str, Any] | None = None,
orchestration_plan: dict[str, Any] | None = None,
enable_report: bool = ROUNDTABLE_BACKGROUND_REPORT_ENABLED,
seat_thinking_enabled: bool | None = None,
seat_reasoning_effort: str | None = None,
seat_subagent_enabled: bool = False,
seat_skill_directives: dict[str, bool] | None = None,
seat_modes: dict[str, str] | None = None,
gather_first: bool = False,
business_code: str | None = None,
) -> dict[str, Any]:
"""把一次启动入参序列化为可 JSON 落库的快照 dict(snake_case,对齐 StartParams)。"""
return {
"v": _SNAPSHOT_VERSION,
"user_id": user_id,
"agents": [{"agent_id": a.agent_id, "name": a.name} for a in agents],
"seed_message": seed_message,
"model": model,
"mode": mode,
"intent_text": intent_text,
"thread_ids": thread_ids,
"coordinator_name": coordinator_name,
"draft_id": draft_id,
"task_id": task_id,
"is_resume": is_resume,
"chain": chain,
"orchestration_plan": orchestration_plan,
"enable_report": enable_report,
"seat_thinking_enabled": seat_thinking_enabled,
"seat_reasoning_effort": seat_reasoning_effort,
"seat_subagent_enabled": seat_subagent_enabled,
"seat_skill_directives": dict(seat_skill_directives or {}),
"seat_modes": dict(seat_modes or {}),
"gather_first": bool(gather_first),
"business_code": business_code,
}
def params_from_snapshot(
snapshot: dict[str, Any] | None, *, job_id: str, lease_owner: str | None
) -> StartParams | None:
"""从 ``input_snapshot`` 无损重建 StartParams。快照缺失 / 损坏返回 None(调用方跳过)。"""
if not snapshot or not isinstance(snapshot, dict):
return None
try:
agents = [
SeatRef(str(a.get("agent_id") or ""), str(a.get("name") or a.get("agent_id") or ""))
for a in (snapshot.get("agents") or [])
if isinstance(a, dict) and a.get("agent_id")
]
if not agents:
return None
return StartParams(
job_id=job_id,
user_id=str(snapshot.get("user_id") or ""),
agents=agents,
seed_message=str(snapshot.get("seed_message") or ""),
model=snapshot.get("model"),
mode=str(snapshot.get("mode") or "recommend"),
intent_text=str(snapshot.get("intent_text") or ""),
thread_ids=snapshot.get("thread_ids"),
coordinator_name=str(snapshot.get("coordinator_name") or "roundtable-coordinator"),
draft_id=snapshot.get("draft_id"),
task_id=snapshot.get("task_id"),
is_resume=bool(snapshot.get("is_resume")),
chain=snapshot.get("chain"),
orchestration_plan=snapshot.get("orchestration_plan"),
enable_report=bool(snapshot.get("enable_report", ROUNDTABLE_BACKGROUND_REPORT_ENABLED)),
seat_thinking_enabled=snapshot.get("seat_thinking_enabled"),
seat_reasoning_effort=snapshot.get("seat_reasoning_effort"),
seat_subagent_enabled=bool(snapshot.get("seat_subagent_enabled")),
seat_skill_directives=dict(snapshot.get("seat_skill_directives") or {}),
seat_modes=dict(snapshot.get("seat_modes") or {}),
gather_first=bool(snapshot.get("gather_first")),
business_code=snapshot.get("business_code"),
lease_owner=lease_owner,
)
except Exception: # noqa: BLE001 — 快照损坏不应让 dispatcher 崩溃
logger.exception("failed to rebuild StartParams from snapshot for job %s", job_id)
return None
# 网关工厂:给定一次启动参数,返回一个 RoundtableGateway 实例。
GatewayFactory = Callable[[StartParams], object]
class RoundtableJobExecutor:
def __init__(self, store, *, gateway_factory: GatewayFactory, draft_store=None, task_draft_store=None) -> None:
self._store = store
self._gateway_factory = gateway_factory
# 草稿仓储(可选):作业完成后把 step3 写回草稿,供前端「查看报告」。
self._draft_store = draft_store
# task 草稿仓储(独立、不分权):task 作业(params.task_id 非空)写回这里而非个人草稿。
self._task_draft_store = task_draft_store
self._tasks: dict[str, asyncio.Task] = {}
def _draft_store_for(self, params: StartParams):
"""task 作业 → 独立 task-draft 存储;普通作业 → 个人草稿存储。"""
return self._task_draft_store if params.task_id else self._draft_store
def start_job(self, params: StartParams) -> bool:
"""启动一个后台编排任务。同一 job 已在跑则忽略(幂等)。返回是否新启动。"""
if params.job_id in self._tasks and not self._tasks[params.job_id].done():
return False
task = asyncio.create_task(self._run(params))
self._tasks[params.job_id] = task
task.add_done_callback(lambda t, jid=params.job_id: self._tasks.pop(jid, None))
return True
async def _run(self, params: StartParams) -> None:
token = set_current_user(_JobUser(params.user_id))
diag = DiagnosticsRecorder(
get_default_store(),
scope="background",
job_id=params.job_id,
draft_id=params.draft_id,
task_id=params.task_id,
user_id=params.user_id,
)
try:
await diag.record(
stage="job", level="info", event="job_started",
message=f"后台挂起作业启动(mode={params.mode}{',续跑' if params.is_resume else ''})",
detail={"mode": params.mode, "is_resume": params.is_resume, "agents": [a.agent_id for a in params.agents]},
)
gateway = self._gateway_factory(params)
thread_ids = params.thread_ids or {}
coordinator_name = params.coordinator_name
# 真网关需要先进程内建好 leader/seat 线程;模拟网关无 init_threads。
# 首启(非 resume)始终建**全新线程** —— 绝不复用前端本地研讨的线程
# (那些线程可能有活跃 run,复用会撞 409)。resume 时复用作业自己的线程。
init = getattr(gateway, "init_threads", None)
if init is not None and not params.is_resume:
thread_ids, coordinator_name = await init()
try:
await self._store.update_progress(
params.job_id,
lease_owner=params.lease_owner,
thread_ids=thread_ids,
coordinator_name=coordinator_name,
)
except Exception: # noqa: BLE001
logger.warning("persist thread_ids for job %s failed", params.job_id, exc_info=True)
# 每轮模型返回后,把对话作为一个带 jobId 的 run 写进草稿 step2.runs[]
# —— 和手动跑同一存取路径(draft_store.update_draft),加载草稿即可显示。
async def _on_dialogue(dialogues: list[dict], consensus: int) -> None:
await self._write_run_to_draft(
params, dialogues=dialogues, status="running", consensus=consensus, thread_ids=thread_ids
)
progress = JobProgressWriter(
self._store,
job_id=params.job_id,
agents=params.agents,
on_dialogue=_on_dialogue if (self._draft_store_for(params) and params.draft_id) else None,
recorder=diag,
# Phase 3:进度落库带租约持有者 —— update_progress 租约门只放行持租者的写,
# 非持租者(被接管/已取消)的过期进度写被 DB 条件 UPDATE 拒绝。
lease_owner=params.lease_owner,
)
# Phase 3 心跳续租:活着就持续把 lease_until 推后,保证本 worker 的租约不过期、
# 别的 worker 抢不走;本协程随编排 task 生死(finally 里 cancel)。崩溃则心跳停,
# 租约最长 _LEASE_TTL_SECONDS 后过期 → dispatcher 在其它 worker 接管。
main_task = asyncio.current_task()
heartbeat = (
asyncio.create_task(self._heartbeat(params.job_id, params.lease_owner, main_task))
if params.lease_owner
else None
)
# 跨进程取消看门狗:run_orchestration 是一个不透明长 await(多轮 LLM),
# 无法在其内部插检查点。改为起一个并行的看门狗协程,定期读 DB 的 job status;
# cancel 端点把 status=cancel_requested 写进 DB(权威意图),看门狗检测到后以**持租者**
# 身份收口 cancelled 终态(清唯一活跃键 + 租约)并对本编排 task 调 cancel() →
# CancelledError 在 run_orchestration 的下一个 await 点抛出并上抛。
# 这让 cancel 在多 worker 下跨进程生效——cancel 打到非执行 worker 也无妨。
watcher = asyncio.create_task(
self._watch_cancellation(params.job_id, main_task, params.lease_owner)
)
try:
result = await run_orchestration(
gateway=gateway,
progress=progress,
agents=params.agents,
thread_ids=thread_ids,
coordinator_name=coordinator_name,
model=params.model,
seed_message=params.seed_message,
mode=params.mode,
intent_text=params.intent_text,
enable_report=params.enable_report,
plan=params.orchestration_plan,
gather_first=params.gather_first,
business_code=params.business_code,
)
finally:
watcher.cancel()
with suppress(asyncio.CancelledError, Exception):
await watcher
if heartbeat is not None:
heartbeat.cancel()
with suppress(asyncio.CancelledError, Exception):
await heartbeat
# 终态:把最终对话 + step3 + jobStatus 落进同一个 run(含 furthest_step=3)。
await self._write_run_to_draft(
params,
dialogues=result.dialogues,
status=result.status,
consensus=100 if result.status == "done" else 95,
step3=result.step3,
thread_ids=thread_ids,
)
# step3 也挂 job(API / 弹窗用)。带租约门:仅持租者可写。
if result.status == "done" and result.step3:
try:
await self._store.update_progress(
params.job_id, lease_owner=params.lease_owner, step3=result.step3
)
except Exception: # noqa: BLE001
logger.warning("persist step3 to job %s failed", params.job_id, exc_info=True)
except asyncio.CancelledError:
# 干净停止即可,**不在被取消的任务里写 DB**:取消态下二次 await 会再抛
# CancelledError、并可能打断进行中的 sqlite 连接。``cancelled`` 终态由
# 取消看门狗(或租约过期后的 dispatcher 回收)权威写入。
logger.info("roundtable job %s cancelled", params.job_id)
raise
except Exception as exc: # noqa: BLE001 — 顶层兜底:run_orchestration 内部已尽量标 error
logger.exception("roundtable job %s crashed", params.job_id)
await diag.record(
stage="job", level="error", event="job_crashed",
message=f"后台挂起作业崩溃:{exc}",
detail={"type": type(exc).__name__, "error": str(exc)},
)
finally:
reset_current_user(token)
async def _heartbeat(
self, job_id: str, lease_owner: str, main_task: asyncio.Task
) -> None:
"""定期续租:把 ``lease_until`` 推到 now + TTL,声明「本 worker 仍在跑、仍持租」。
只要主任务活着且续租成功,别的 worker 的 dispatcher 看到租约未过期就不会抢;主任务
结束(或本协程被 cancel)后不再续租,租约自然到期。续租失败(租约已被接管 / 状态已变)
说明本 worker 已不是持租者,必须立即停止本地编排,避免过期 worker 继续消耗模型资源;
其任何迟到进度写也会被租约门拒绝。
"""
try:
while not main_task.done():
await asyncio.sleep(_HEARTBEAT_INTERVAL)
if main_task.done():
return
try:
ok = await self._store.renew_lease(
job_id,
lease_owner=lease_owner,
lease_until=datetime.now(UTC) + timedelta(seconds=_LEASE_TTL_SECONDS),
)
except Exception: # noqa: BLE001 — DB 抖动不中断心跳,下轮重试
logger.warning("roundtable job %s renew_lease failed", job_id, exc_info=True)
continue
if not ok:
logger.info(
"roundtable job %s: renew_lease rejected (lease lost or status changed)",
job_id,
)
if not main_task.done():
main_task.cancel()
return
except asyncio.CancelledError:
# 编排结束 → _run 的 finally cancel 了心跳 → 安静退出。
return
async def _watch_cancellation(
self, job_id: str, main_task: asyncio.Task, lease_owner: str | None
) -> None:
"""定期轮询 DB 检测跨进程取消标志(``cancel_requested``),并以持租者身份收口。
多 worker 下 cancel 请求可能打到**非执行** worker,该 worker 的 ``cancel()`` 在
``_tasks`` 里找不到此 job(返回 False)。但 cancel 端点始终把 ``status='cancel_requested'``
写进 DB(权威意图)。本看门狗跑在**实际执行作业的 worker** 上,检测到后:
1. 以**持租者**身份 ``finalize_cancel``——条件 UPDATE 把 cancel_requested 收口为终态
``cancelled``,同时清唯一活跃键 + 租约(一致性来自 DB,绝不双收口);
2. 对本地编排 task 调 ``cancel()``,让 ``run_orchestration`` 在下一个 await 点抛
``CancelledError`` 干净退出。
DB 读失败时静默重试(不中断看门狗)。兼容旧数据:若 DB 里已是终态 ``cancelled``
(旧版 cancel 端点直写),直接取消本地 task。
"""
try:
while not main_task.done():
await asyncio.sleep(_CANCEL_POLL_INTERVAL)
if main_task.done():
return
try:
# user_id=None → 不按 user 过滤,纯读 status(看门狗不关心归属)。
row = await self._store.get(job_id, user_id=None)
except Exception: # noqa: BLE001 — DB 读失败不中断看门狗,下轮重试
continue
status = (row.get("status") or "") if row is not None else ""
if status == "cancel_requested":
logger.info(
"roundtable job %s: detected cancel_requested in DB, finalizing "
"cancelled and cancelling orchestration task",
job_id,
)
# 以持租者身份收口;非持租者(理论上不该发生)则退化为无租约条件收口。
try:
await self._store.finalize_cancel(job_id, lease_owner=lease_owner)
except Exception: # noqa: BLE001 — 收口失败不阻断本地取消
logger.warning("roundtable job %s finalize_cancel failed", job_id, exc_info=True)
if not main_task.done():
main_task.cancel()
return
if status == "cancelled":
# 兼容:旧路径已直接把 cancelled 写进 DB —— 本地 task 照样停下。
logger.info(
"roundtable job %s: detected cancelled in DB, cancelling orchestration task",
job_id,
)
if not main_task.done():
main_task.cancel()
return
except asyncio.CancelledError:
# 编排正常结束 → _run 的 finally cancel 了看门狗 → 安静退出。
return
async def _write_run_to_draft(
self,
params: StartParams,
*,
dialogues: list[dict],
status: str,
consensus: int,
step3: dict | None = None,
thread_ids: dict[str, str] | None = None,
) -> None:
"""把本次后台作业作为一个带 jobId 的 run 写进草稿 step2.runs[](同手动跑存取)。
读现有 step2 → 按 run_id=``job-<jobId>`` upsert(保留其它手动 run)→ 设为 activeRun
→ update_draft。前端加载草稿即可像手动跑一样 hydrate 出对话。
带乐观锁重试:**每次**读-改-写都带上 ``expected_version``,若并发写导致版本冲突
(``DraftConcurrentWriteError``)则重新读取最新草稿、重新合并 run 再试。重试全部
耗尽后**绝不**退化为无版本检查的 last-write-wins 盲写,而是给作业置
``draft_persist_pending`` 标记(由 Phase 3 outbox 补偿器投影结果到草稿)——
消除后台作业与前台 PUT 并发时的 step2.runs 丢失更新,且一致性来自 DB 标记而非运气。
"""
draft_store = self._draft_store_for(params)
if draft_store is None or not params.draft_id:
return
# 这些只依赖 params/dialogues,重试间不变——提前算好。
idx_by_agent = {a.agent_id: i for i, a in enumerate(params.agents)}
s2_dialogues = [_to_step2_dialogue(d, idx_by_agent) for d in dialogues]
last_leader = next(
(d.get("content") for d in reversed(dialogues) if d.get("role") == "leader"), ""
)
run_id = f"job-{params.job_id}"
for attempt in range(_DRAFT_WRITE_RETRIES):
# 读最新草稿(冲突重试时重新读,拿到别的小伙伴刚写的 runs)。
try:
draft = await draft_store.get_draft(params.draft_id, params.user_id)
except Exception: # noqa: BLE001
draft = None
expected_version = (draft or {}).get("version")
step2 = dict((draft or {}).get("step2") or {})
runs = list(step2.get("runs") or [])
idx = next((i for i, r in enumerate(runs) if r.get("id") == run_id), -1)
created = runs[idx].get("createdAt") if idx >= 0 else datetime.now(UTC).isoformat()
prev_step3 = runs[idx].get("step3") if idx >= 0 else None
# 本次有效 step3:新报告优先,否则保留上次(running 中间态 step3=None 不抹掉已存的)。
effective_step3 = step3 if step3 is not None else prev_step3
run = {
"id": run_id,
"title": "后台研讨",
"createdAt": created,
"source": params.mode,
"selectedAgents": [{"agent_id": a.agent_id, "name": a.name} for a in params.agents],
"threadIds": (thread_ids or None),
"coordinatorName": params.coordinator_name,
"step2RoundtableDialogues": s2_dialogues,
"lastLeaderContent": last_leader or "",
"hasConsensus": status == "done",
"consensusPercentage": int(consensus),
"budgetLimit": 0,
"seatStubMessages": _seat_stub_messages(dialogues),
"orchestrationMode": params.mode,
"chain": params.chain,
"orchestrationPlan": params.orchestration_plan,
"step3": effective_step3,
"jobId": params.job_id,
"jobStatus": status,
}
runs = [run if i == idx else r for i, r in enumerate(runs)] if idx >= 0 else [*runs, run]
step2["runs"] = runs
step2["activeRunId"] = run_id
has_report = effective_step3 is not None
update_kwargs: dict[str, Any] = {
"step2": step2,
"furthest_step": 3 if (status == "done" and has_report) else 2,
}
if has_report:
update_kwargs["step3"] = effective_step3
# 每次尝试都带 expected_version(乐观锁);绝不降级为无版本检查的盲写。
try:
await draft_store.update_draft(
params.draft_id, params.user_id,
expected_version=expected_version,
**update_kwargs,
)
return # 写入成功
except DraftConcurrentWriteError:
if attempt < _DRAFT_WRITE_RETRIES - 1:
logger.debug(
"roundtable draft %s version conflict (attempt %d/%d), retrying",
params.draft_id, attempt + 1, _DRAFT_WRITE_RETRIES,
)
continue
# 重试全部耗尽:不盲写,置 pending 标记交由 outbox 补偿器投影。
logger.warning(
"roundtable draft %s: version conflicts exhausted after %d attempts, "
"marking job %s draft_persist_pending",
params.draft_id, _DRAFT_WRITE_RETRIES, params.job_id,
)
with suppress(Exception):
await self._store.mark_draft_persist_pending(params.job_id)
return
except Exception as exc: # noqa: BLE001 — 写草稿失败不影响作业本身
logger.warning("write run to draft %s failed", params.draft_id, exc_info=True)
from deerflow.persistence.roundtable_diagnostics import record_diagnostic
await record_diagnostic(
scope="background", stage="job", level="error", event="draft_write_failed",
message=f"研讨结果写回草稿失败(status={status}):{exc}",
detail={"type": type(exc).__name__, "error": str(exc)},
job_id=params.job_id, draft_id=params.draft_id, task_id=params.task_id, user_id=params.user_id,
)
return
def cancel(self, job_id: str) -> bool:
"""请求取消一个在跑的作业。返回是否命中在跑任务。"""
task = self._tasks.get(job_id)
if task is None or task.done():
return False
task.cancel()
return True
def is_running(self, job_id: str) -> bool:
task = self._tasks.get(job_id)
return task is not None and not task.done()
async def await_job(self, job_id: str) -> None:
"""等待某作业任务结束(测试用;生产不调用)。"""
task = self._tasks.get(job_id)
if task is not None:
try:
await task
except (Exception, asyncio.CancelledError): # noqa: BLE001 — 取消是 BaseException
pass