Files
hz4th_coder f28d73f1bb fix: 多Agent协作 JSON 截断导致主管规划解析失败
- 根因: deepseek-v4-flash 规划输出超 max_tokens=2000 被截断,JSON 不完整解析失败
- agents.py: 新增 _chat_json() 统一处理 JSON 类调用,解析失败自动加大 max_tokens 重试(×1/×2/×3)
- 主管规划/辩论质询/辩论裁决/评审 全部改用 _chat_json 并提高初始上限(3000/1200/2500/2000)
- 实测: 主管(计算器网页版4子任务) 评审(92分) 辩论(指定阵容)
2026-08-12 17:08:16 +08:00

442 lines
20 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.
# -*- coding: utf-8 -*-
"""
V2 多 Agent 协作引擎
三种模式:
- supervisor 主管模式:主管拆解任务 → 委派多个 Worker 并行执行 → 汇总合成最终产出
- review 评审模式:Worker 产出 → Reviewer 评审打分 → 未达标返工(最多 N 轮)→ 终稿
- debate 辩论模式:多位辩手各自立论 → 互相质询(可多轮)→ 裁判综合裁决
"""
import json
import threading
import traceback
import db
import llm_gateway
import notify
def _log(run_id, role, worker_id, stage, content, usage=None):
tokens = usage['total_tokens'] if usage else 0
cost = usage['cost'] if usage else 0.0
db.w(
'INSERT INTO agent_steps (run_id, role, worker_id, seq, stage, content, tokens, cost, created_at) '
'VALUES (?,?,?,?,?,?,?,?,?)',
(run_id, role, worker_id,
db.q('SELECT COALESCE(MAX(seq),0)+1 s FROM agent_steps WHERE run_id=?', (run_id,))[0]['s'],
stage, content, tokens, cost, db.now()))
def _update_run(run_id, **fields):
if not fields:
return
fields['finished_at'] = db.now()
sets = ', '.join(f'{k}=?' for k in fields)
db.w(f'UPDATE agent_runs SET {sets} WHERE id=?', (*fields.values(), run_id))
def _chat_worker(worker, messages, temperature=None, max_tokens=None):
"""调用某个 Worker 的模型,返回 (text, usage)"""
r = llm_gateway.chat(
worker['provider'], worker['model'], messages,
temperature=temperature if temperature is not None else worker['temperature'],
max_tokens=max_tokens or worker['max_tokens'] or 2000,
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
return r['text'], r
def _workers_from_ids(ids):
"""按 id 列表取 Worker;自动路由兜底(id 为空或不存在时)"""
out = []
for wid in ids or []:
w = db.q('SELECT * FROM workers WHERE id=? AND status="enabled"', (wid,), one=True)
if w:
out.append(w)
if not out:
rows = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id LIMIT 3')
out = rows
return out
def _extract_json(text):
"""从 LLM 输出中稳健提取 JSON"""
t = text.strip()
if t.startswith('```'):
t = t.strip('`')
if t.startswith('json'):
t = t[4:]
t = t.strip()
start = min([i for i in (t.find('{'), t.find('[')) if i >= 0] or [0])
end = max(t.rfind('}'), t.rfind(']')) + 1
if end <= start:
raise ValueError('未找到 JSON 内容')
return json.loads(t[start:end])
def _extract_score(text):
"""从评审 JSON 中提取分数,失败则启发式解析"""
try:
data = _extract_json(text)
if isinstance(data, dict):
for k in ('score', '总分', '评分'):
if k in data:
return float(data[k]), data.get('judgment') or data.get('意见') or text
except Exception:
pass
# 启发式:找 0-100 数字
import re
m = re.search(r'(?:score|总分|评分)[^\d]*(\d{1,3})', text, re.I)
if m:
return float(m.group(1)), text
m = re.search(r'(\d{1,3})\s*/\s*100', text)
if m:
return float(m.group(1)), text
return 60.0, text
def _chat_json(worker, messages, max_tokens=2000, temperature=None):
"""调用 LLM 并解析 JSON;解析失败自动加大 max_tokens 重试(防止长文截断)。
返回 (data, raw_text, usage);全部失败抛最后异常。"""
last_err = None
for attempt in range(3):
mt = max_tokens * (attempt + 1) # 2000 → 4000 → 6000
text, usage = _chat_worker(worker, messages, temperature=temperature, max_tokens=mt)
try:
data = _extract_json(text)
if isinstance(data, dict) and not data:
raise ValueError('空 JSON')
return data, text, usage
except Exception as e:
last_err = e
raise last_err or ValueError('JSON 解析失败')
# ---------------------------------------------------------------------------
# 主管模式
# ---------------------------------------------------------------------------
SUPERVISOR_PLAN_PROMPT = (
'你是资深项目经理(主管 Agent)。请把下面的任务拆解为 2~5 个子任务,分配给团队协作完成。\n'
'输出严格 JSON{{"subtasks": [{{"title": "子任务标题", "instruction": "给执行 Agent 的完整指令(含要求与输出格式)"}}]}}\n'
'只输出 JSON,不要任何解释。\n\n'
'任务主题:{topic}\n'
'背景上下文:{context}'
)
SUPERVISOR_SYNTH_PROMPT = (
'你是主管 Agent。团队已完成以下子任务,请汇总合成一份完整、连贯、高质量的最终交付物。\n'
'要求:结构清晰、覆盖所有子任务要点、去除重复、补足衔接,直接输出最终成果(不要解释过程)。\n\n'
'原始任务:{topic}\n'
'子任务成果:\n{parts}'
)
def run_supervisor(run_id, run, workers):
_log(run_id, 'supervisor', None, 'plan',
f'主管开始规划:{run["topic"][:200]}')
# 1) 主管拆解(JSON 解析失败自动加大 max_tokens 重试)
try:
data, _, u1 = _chat_json(
workers[0],
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': SUPERVISOR_PLAN_PROMPT.format(
topic=run['topic'], context=run['context'] or '无')}],
max_tokens=3000, temperature=0.3)
subtasks = data.get('subtasks') or data.get('tasks') or []
if isinstance(data, list):
subtasks = data
except Exception as e:
_update_run(run_id, status='failed', error=f'主管规划解析失败: {e}')
_log(run_id, 'supervisor', None, 'plan', f'❌ 规划解析失败: {e}')
return
_log(run_id, 'supervisor', workers[0]['id'], 'plan',
f'规划完成,拆解为 {len(subtasks)} 个子任务:' + ''.join(s.get('title', '?')[:30] for s in subtasks),
u1)
# 2) 委派并行执行(轮询分配 Worker)
parts = []
total_tokens, total_cost = u1['total_tokens'], u1['cost']
for i, st in enumerate(subtasks):
w = workers[i % len(workers)]
title = st.get('title', f'子任务{i+1}')
instr = st.get('instruction') or st.get('description') or title
_log(run_id, 'worker', w['id'], 'delegate', f'委派子任务「{title}」→ {w["name"]}')
try:
text, u = _chat_worker(
w, [{'role': 'system', 'content': w['system_prompt'] or '你是高效可靠的执行 Agent。'},
{'role': 'user', 'content': instr}],
max_tokens=2500)
parts.append(f'【子任务{i+1}: {title}\n{text}')
total_tokens += u['total_tokens']; total_cost += u['cost']
_log(run_id, 'worker', w['id'], 'produce',
f'子任务「{title}」完成({u["total_tokens"]} tokens):\n{text[:600]}', u)
except Exception as e:
parts.append(f'【子任务{i+1}: {title}\n(执行失败: {str(e)[:200]}')
_log(run_id, 'worker', w['id'], 'produce', f'❌ 子任务「{title}」失败: {e}')
# 3) 主管合成
_log(run_id, 'supervisor', None, 'synthesize', '汇总合成最终产出…')
try:
final_text, u2 = _chat_worker(
workers[0],
[{'role': 'system', 'content': '你是资深主管 Agent,负责最终交付物合成。'},
{'role': 'user', 'content': SUPERVISOR_SYNTH_PROMPT.format(
topic=run['topic'], parts='\n\n'.join(parts))}],
temperature=0.4, max_tokens=3000)
total_tokens += u2['total_tokens']; total_cost += u2['cost']
_log(run_id, 'supervisor', workers[0]['id'], 'synthesize', f'✅ 最终产出({u2["total_tokens"]} tokens):\n{final_text[:800]}', u2)
except Exception as e:
final_text = '\n\n'.join(parts)
_log(run_id, 'supervisor', None, 'synthesize', f'⚠️ 合成失败,退回拼接结果: {e}')
_update_run(run_id, status='done', result=final_text,
summary=f'拆解 {len(subtasks)} 个子任务,由 {len(workers)} 个 Agent 协作完成',
total_tokens=total_tokens, cost=round(total_cost, 6))
_finish_notify(run, '主管协作', final_text)
# ---------------------------------------------------------------------------
# 评审模式
# ---------------------------------------------------------------------------
REVIEWER_PROMPT = (
'你是严格的评审专家(Reviewer)。请对下面的执行成果进行评审,输出严格 JSON:\n'
'{{"score": 0到100的整数, "judgment": "总体评价", "issues": ["问题1", "问题2"], "suggestions": "具体修改建议"}}\n'
'评分标准:{rubric}\n'
'只输出 JSON。\n\n'
'原始任务:{topic}\n'
'执行成果:\n{output}'
)
REVISE_PROMPT = (
'你是执行 Agent。请根据评审意见修改你的成果,输出修改后的最终版本(直接输出成果内容,不要解释过程)。\n\n'
'原始任务:{topic}\n'
'上一版成果:\n{output}\n\n'
'评审意见:\n{review}'
)
def run_review(run_id, run, workers):
main_w = workers[0]
reviewer_w = workers[1] if len(workers) > 1 else workers[0]
params = json.loads(run.get('params') or '{}')
max_rounds = int(params.get('rounds', 2))
pass_score = float(params.get('pass_score', 70))
rubric = params.get('rubric') or '内容准确性 40%,完整性 30%,清晰度与格式 30%。'
_log(run_id, 'worker', main_w['id'], 'produce', f'执行 Agent 开始产出:{run["topic"][:200]}')
output, u1 = _chat_worker(
main_w, [{'role': 'system', 'content': main_w['system_prompt'] or '你是高质量执行 Agent。'},
{'role': 'user', 'content': run['topic'] + (('\n\n上下文:' + run['context']) if run['context'] else '')}],
max_tokens=2500)
total_tokens, total_cost = u1['total_tokens'], u1['cost']
_log(run_id, 'worker', main_w['id'], 'produce', f'初稿完成({u1["total_tokens"]} tokens):\n{output[:600]}', u1)
final_output, last_review, final_score = output, '', 0.0
for round_i in range(max_rounds + 1):
# 评审
review_usage = None
try:
review_data, review_text, u2 = _chat_json(
reviewer_w,
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': REVIEWER_PROMPT.format(
topic=run['topic'], output=final_output, rubric=rubric)}],
max_tokens=2000, temperature=0.2)
review_usage = u2
total_tokens += u2['total_tokens']; total_cost += u2['cost']
if isinstance(review_data, dict):
score = float(review_data.get('score', review_data.get('总分', 60)))
judgment = review_data.get('judgment') or review_data.get('意见') or review_text
else:
score, judgment = _extract_score(review_text)
except Exception as e:
score, judgment = 0, f'评审调用失败: {e}'
final_score = score
last_review = judgment
_log(run_id, 'reviewer', reviewer_w['id'], 'critique',
f'第 {round_i+1} 轮评审:得分 {score}/100\n{judgment[:500]}', review_usage)
if score >= pass_score or round_i >= max_rounds:
_log(run_id, 'reviewer', reviewer_w['id'], 'critique',
f'{"✅ 评审通过" if score >= pass_score else "⚠️ 已达最大轮次,采纳当前稿"}(得分 {score}')
break
# 返工
try:
revised, u3 = _chat_worker(
main_w,
[{'role': 'system', 'content': main_w['system_prompt'] or '你是高质量执行 Agent。'},
{'role': 'user', 'content': REVISE_PROMPT.format(
topic=run['topic'], output=final_output, review=judgment)}],
max_tokens=2500)
total_tokens += u3['total_tokens']; total_cost += u3['cost']
final_output = revised
_log(run_id, 'worker', main_w['id'], 'revise',
f'第 {round_i+1} 轮返工完成:\n{revised[:600]}', u3)
except Exception as e:
_log(run_id, 'worker', main_w['id'], 'revise', f'❌ 返工失败: {e}')
break
_update_run(run_id, status='done', result=final_output,
summary=f'评审得分 {final_score}/100(阈值 {pass_score}),共评审 {min(max_rounds+1, 3)} 轮',
total_tokens=total_tokens, cost=round(total_cost, 6))
_finish_notify(run, '评审协作', final_output)
# ---------------------------------------------------------------------------
# 辩论模式
# ---------------------------------------------------------------------------
DEBATE_OPEN_PROMPT = (
'你是辩手(立场:{stance})。针对议题给出你的立场陈述与核心论点。\n'
'要求:论点鲜明、论据充分、逻辑严谨,控制在 400 字以内。\n\n'
'议题:{topic}\n'
'背景:{context}'
)
DEBATE_REBUT_PROMPT = (
'你是辩手(立场:{stance})。请阅读其他辩手的观点,进行质询与反驳,同时回应对方对你观点的质疑。\n'
'输出严格 JSON{{"rebuttal": "你的反驳与回应", "refine": "是否坚持原立场 true/false"}}\n'
'只输出 JSON。\n\n'
'议题:{topic}\n'
'你的上一轮立场陈述:\n{my_view}\n\n'
'其他辩手观点:\n{others}'
)
JUDGE_PROMPT = (
'你是首席裁判(Judge)。请综合以下辩论内容,给出最终裁决。\n'
'输出严格 JSON{{"winner": "胜出方或共识结论", "consensus": "最终结论(可直接作为交付物)", "summary": "辩论过程摘要"}}\n'
'只输出 JSON。\n\n'
'议题:{topic}\n'
'辩论记录:\n{transcript}'
)
def run_debate(run_id, run, workers):
params = json.loads(run.get('params') or '{}')
rounds = int(params.get('rounds', 1))
stances = params.get('stances') or []
judge_w = workers[-1]
debaters = workers[:-1] if len(workers) > 1 else workers
# 立场分配
if len(stances) < len(debaters):
default_stances = ['正方(支持)', '反方(反对)', '中立(分析利弊)', '补充视角']
stances = stances + default_stances[len(stances):]
views = {}
total_tokens, total_cost = 0, 0.0
# 第一轮:立论
for i, w in enumerate(debaters):
stance = stances[i % len(stances)]
try:
text, u = _chat_worker(
w, [{'role': 'system', 'content': w['system_prompt'] or '你是犀利严谨的辩手。'},
{'role': 'user', 'content': DEBATE_OPEN_PROMPT.format(
stance=stance, topic=run['topic'], context=run['context'] or '无')}],
temperature=0.7, max_tokens=800)
views[w['id']] = {'stance': stance, 'view': text}
total_tokens += u['total_tokens']; total_cost += u['cost']
_log(run_id, 'debater', w['id'], 'produce', f'【{stance}{w["name"]} 立论:\n{text[:500]}', u)
except Exception as e:
views[w['id']] = {'stance': stance, 'view': f'(发言失败: {e}'}
_log(run_id, 'debater', w['id'], 'produce', f'❌ {w["name"]} 立论失败: {e}')
# 质询轮
for rnd in range(rounds):
for i, w in enumerate(debaters):
mine = views[w['id']]
others = '\n\n'.join(
f"【{v['stance']}】(#{k}){v['view'][:500]}"
for k, v in views.items() if k != w['id'])
try:
data, _, u = _chat_json(
w,
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': DEBATE_REBUT_PROMPT.format(
stance=mine['stance'], topic=run['topic'],
my_view=mine['view'], others=others)}],
max_tokens=1200, temperature=0.7)
rebut = data.get('rebuttal') if isinstance(data, dict) else text
mine['view'] = mine['view'] + '\n\n【质询回应】' + str(rebut)
total_tokens += u['total_tokens']; total_cost += u['cost']
_log(run_id, 'debater', w['id'], 'produce',
f'第 {rnd+1} 轮质询回应({w["name"]}):\n{str(rebut)[:500]}', u)
except Exception as e:
_log(run_id, 'debater', w['id'], 'produce', f'❌ {w["name"]} 质询失败: {e}')
# 裁判裁决
transcript = '\n\n'.join(
f"【{v['stance']}{v['view']}" for v in views.values())
_log(run_id, 'judge', judge_w['id'], 'verdict', '裁判综合裁决中…')
try:
data, verdict, u4 = _chat_json(
judge_w,
[{'role': 'system', 'content': '你只输出 JSON。你是公正严明的首席裁判。'},
{'role': 'user', 'content': JUDGE_PROMPT.format(topic=run['topic'], transcript=transcript)}],
max_tokens=2500, temperature=0.3)
if isinstance(data, dict):
consensus = data.get('consensus') or data.get('结论') or verdict
summary = data.get('summary') or data.get('摘要') or ''
winner = data.get('winner') or ''
else:
consensus, summary, winner = verdict, '', ''
total_tokens += u4['total_tokens']; total_cost += u4['cost']
_log(run_id, 'judge', judge_w['id'], 'verdict',
f'🏆 裁决:{winner}\n共识结论:{str(consensus)[:800]}', u4)
except Exception as e:
consensus, summary, winner = transcript, '裁判裁决失败: ' + str(e), ''
_log(run_id, 'judge', judge_w['id'], 'verdict', f'❌ 裁决失败: {e}')
_update_run(run_id, status='done', result=str(consensus),
summary=f'辩论参与者 {len(debaters)} 名,质询 {rounds} 轮' + (f';胜出:{winner}' if winner else ''),
total_tokens=total_tokens, cost=round(total_cost, 6))
_finish_notify(run, '辩论协作', str(consensus))
def _finish_notify(run, mode_label, result_preview):
try:
notify.notify('task_done', f'多 Agent 协作完成:{run["title"]}',
f'模式:{mode_label}\n主题:{run["topic"][:100]}\n'
f'产出预览:{str(result_preview)[:200]}',
save_alert=False)
except Exception:
pass
def execute_run(run_id):
"""后台线程:执行一次多 Agent 协作"""
run = db.q('SELECT * FROM agent_runs WHERE id=?', (run_id,), one=True)
if not run:
return
try:
workers = _workers_from_ids(json.loads(run.get('worker_ids') or '[]'))
if not workers:
_update_run(run_id, status='failed', error='没有可用的 Worker')
_log(run_id, 'system', None, 'plan', '❌ 无可用 Worker')
return
mode = run['mode']
if mode == 'supervisor':
run_supervisor(run_id, run, workers)
elif mode == 'review':
run_review(run_id, run, workers)
elif mode == 'debate':
run_debate(run_id, run, workers)
else:
_update_run(run_id, status='failed', error=f'未知模式: {mode}')
except Exception as e:
_update_run(run_id, status='failed', error=f'引擎异常: {str(e)[:300]}')
_log(run_id, 'system', None, 'plan', '❌ ' + traceback.format_exc())
class AgentRunner:
def __init__(self):
self._threads = {}
def submit(self, run_id):
if run_id in self._threads and self._threads[run_id].is_alive():
return False
t = threading.Thread(target=execute_run, args=(run_id,), daemon=True)
self._threads[run_id] = t
t.start()
return True
runner = AgentRunner()