- 项目级验收:全部任务完成后进入「待验收」,邮件通知负责人,验收通过才算完成;公开验收链接 /review/<token> 免登录一键通过/打回(打回必填原因) - 打回自动返工:自动生成含负责人意见的返工任务,AI主管立即重做并再次提交验收;平台内也可验收 - 真实网页交付:任务HTML产出自动落盘工作目录(剥离围栏/前置叙述),Demo展示真实页面;返工产出覆盖入口页;老项目已回填 - 事件记录面板:合并AI主管动态+任务日志+交付记录,可展开/收起(记忆状态),任务事件可点击定位 - 邮件通道修复:补 Date/Message-ID 头(amavisd 拒收 invalid header section 根因),打包附件/通知恢复送达;notify.py 改标准 MIME - 规划重试:WBS 拆解失败自动重试3次,仍失败邮件通知负责人 - 项目状态新增 review(待验收),全端展示
401 lines
18 KiB
Python
401 lines
18 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""
|
||
V3.3 AI 主管自动开工引擎
|
||
========================
|
||
新建项目后立即自动运转:
|
||
1. AI 主管(项目 manager_worker_id)读取项目目标/验收标准,拆解 WBS 任务
|
||
2. 把任务分派给团队 Worker(project_team_workers,默认 AI牛/AI马)轮询分配
|
||
3. 自动执行整个工作流(DAG 就绪任务逐个跑)
|
||
4. 监控:失败任务由 AI 主管诊断并修订指令重试(最多 2 次)
|
||
5. 全部完成后收尾(交付物自动部署/打包并通知送达者)
|
||
|
||
对外接口:
|
||
- WBS_PROMPT / parse_wbs() —— WBS 生成提示词与解析(app.py 复用)
|
||
- launch(pid) —— 后台线程启动自动开工
|
||
- rerun(pid) —— 手动重新开工(已有任务则跳过规划直接跑)
|
||
"""
|
||
import json
|
||
import threading
|
||
import time
|
||
import traceback
|
||
|
||
import db
|
||
import llm_gateway
|
||
import engine
|
||
import delivery
|
||
import notify
|
||
|
||
WBS_PROMPT = (
|
||
'你是资深项目经理(AI 主管)。请把下面的项目目标拆解为可执行的任务列表(WBS),'
|
||
'要求:\n1. 输出严格 JSON,格式 {{"tasks": [{{"title": "任务标题", '
|
||
'"description": "给AI Worker的执行指令(含要求与输出格式)", "depends_on": [0,2]}}]}}\n'
|
||
'2. depends_on 是前置任务的数组下标(无依赖填 []),下标从 0 开始\n'
|
||
'3. 4~8 个任务,逻辑清晰,可并行任务并行,不要输出 JSON 以外的任何内容\n\n'
|
||
'项目目标:{goal}\n'
|
||
'项目验收标准:{accept}'
|
||
)
|
||
|
||
REVISE_PROMPT = (
|
||
'你是 AI 主管。下面这个任务执行失败了,请诊断原因并输出修订后的执行指令。\n'
|
||
'输出严格 JSON:{{"instruction": "修订后的完整执行指令(含要求与输出格式)", '
|
||
'"reason": "一句话失败原因诊断"}}\n'
|
||
'只输出 JSON。\n\n'
|
||
'任务标题:{title}\n'
|
||
'任务指令:{instruction}\n'
|
||
'失败原因:{error}\n'
|
||
'修订原则:把容易出错的步骤写得更具体(明确格式、步骤、约束),必要时拆小范围。'
|
||
)
|
||
|
||
|
||
def _extract_json(text):
|
||
"""用 raw_decode 稳健提取第一个完整 JSON 值(自动忽略尾部杂质/多余 JSON)"""
|
||
t = text.strip()
|
||
if t.startswith('```'):
|
||
t = t.strip('`')
|
||
if t.startswith('json'):
|
||
t = t[4:]
|
||
t = t.strip()
|
||
candidates = [i for i in (t.find('{'), t.find('[')) if i >= 0]
|
||
last_err = None
|
||
for i in sorted(candidates):
|
||
try:
|
||
obj, _ = json.JSONDecoder().raw_decode(t[i:])
|
||
return obj
|
||
except Exception as e:
|
||
last_err = e
|
||
raise last_err or ValueError('未找到 JSON 内容')
|
||
|
||
|
||
def parse_wbs(text):
|
||
"""从 LLM 输出中稳健提取 JSON 任务列表(兼容任务键名/多余内容)"""
|
||
try:
|
||
data = _extract_json(text)
|
||
except Exception:
|
||
raise
|
||
if isinstance(data, dict):
|
||
for key in ('tasks', 'subtasks', 'task_list', 'items', 'task'):
|
||
if key in data:
|
||
data = data[key]
|
||
break
|
||
tasks = data if isinstance(data, list) else []
|
||
if not tasks:
|
||
raise ValueError('任务列表为空或格式不正确')
|
||
return tasks
|
||
|
||
|
||
def _log(pid, level, message):
|
||
try:
|
||
db.w('INSERT INTO project_logs (project_id, level, message, created_at) '
|
||
'VALUES (?,?,?,?)', (pid, level, message, db.now()))
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def _set_auto(pid, status, message, started=False, finished=False):
|
||
sets, args = ['auto_status=?', 'auto_message=?'], [status, (message or '')[:500]]
|
||
if started:
|
||
sets.append('auto_started_at=?')
|
||
args.append(db.now())
|
||
if finished:
|
||
sets.append('auto_finished_at=?')
|
||
args.append(db.now())
|
||
sets.append('updated_at=?')
|
||
args.append(db.now())
|
||
db.w(f'UPDATE projects SET {", ".join(sets)} WHERE id=?', (*args, pid))
|
||
|
||
|
||
def _chat_json(worker, messages, max_tokens=3000, temperature=0.3):
|
||
"""调用 LLM 并解析 JSON;解析失败自动加大 max_tokens 重试"""
|
||
last_err = None
|
||
for attempt in range(3):
|
||
mt = max_tokens * (attempt + 1)
|
||
r = llm_gateway.chat(
|
||
worker['provider'], worker['model'], messages,
|
||
temperature=temperature, max_tokens=mt,
|
||
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
|
||
try:
|
||
return _extract_json(r['text']), r
|
||
except Exception as e:
|
||
last_err = e
|
||
raise last_err or ValueError('JSON 解析失败')
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 团队读取
|
||
# ---------------------------------------------------------------------------
|
||
def get_manager(pid):
|
||
"""项目 AI 主管 Worker"""
|
||
p = db.q('SELECT manager_worker_id FROM projects WHERE id=?', (pid,), one=True)
|
||
if not p or not p.get('manager_worker_id'):
|
||
return None
|
||
return db.q('SELECT * FROM workers WHERE id=? AND status="enabled"',
|
||
(p['manager_worker_id'],), one=True)
|
||
|
||
|
||
def get_team(pid):
|
||
"""项目干活团队 Worker 列表(按加入顺序)"""
|
||
rows = db.q(
|
||
'SELECT w.* FROM project_team_workers t JOIN workers w ON w.id=t.worker_id '
|
||
'WHERE t.project_id=? AND w.status="enabled" ORDER BY t.rowid', (pid,))
|
||
return rows
|
||
|
||
|
||
def set_team(pid, worker_ids):
|
||
"""整组替换团队(先清后插)"""
|
||
db.w('DELETE FROM project_team_workers WHERE project_id=?', (pid,))
|
||
for wid in worker_ids or []:
|
||
db.w('INSERT OR IGNORE INTO project_team_workers (project_id, worker_id, created_at) '
|
||
'VALUES (?,?,?)', (pid, int(wid), db.now()))
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 规划:AI 主管拆解 WBS → 建任务(轮询分派团队 Worker)
|
||
# V3.4:解析失败自动重试 3 次(每次重新调用 LLM),仍失败抛出异常由上层通知负责人
|
||
# ---------------------------------------------------------------------------
|
||
def plan_tasks(pid, manager, team, review_required=0):
|
||
"""AI 主管规划并创建任务。返回 (created_ids, plan_message)"""
|
||
proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True)
|
||
goal = (proj.get('objective') or proj.get('name') or '').strip()
|
||
if not goal:
|
||
raise ValueError('项目缺少目标(Objective),AI 主管无法规划')
|
||
last_err = None
|
||
for attempt in range(1, 4):
|
||
try:
|
||
r = llm_gateway.chat(
|
||
manager['provider'], manager['model'],
|
||
[{'role': 'system', 'content': '你只输出 JSON,不输出任何解释文字。'},
|
||
{'role': 'user', 'content': WBS_PROMPT.format(
|
||
goal=goal, accept=proj.get('acceptance_criteria') or '—')}],
|
||
temperature=0.3, max_tokens=3000 * attempt,
|
||
base_url=manager['base_url'] or None, api_key=manager['api_key'] or None)
|
||
tasks = parse_wbs(r['text'])
|
||
break
|
||
except Exception as e:
|
||
last_err = e
|
||
_log(pid, 'warn', f'任务拆解第 {attempt}/3 次失败({str(e)[:100]}),重新规划中…')
|
||
else:
|
||
raise ValueError(f'AI 主管连续 3 次拆解任务失败({str(last_err)[:120]}),已中断并通知负责人')
|
||
created = []
|
||
for i, t in enumerate(tasks):
|
||
dep_idx = t.get('depends_on') or []
|
||
dep_ids = [created[idx] for idx in dep_idx
|
||
if isinstance(idx, int) and 0 <= idx < len(created)]
|
||
worker = team[i % len(team)] if team else None
|
||
tid = db.w(
|
||
'INSERT INTO tasks (project_id, worker_id, title, description, priority, '
|
||
'review_required, deadline, depends_on, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?)',
|
||
(pid, worker['id'] if worker else None,
|
||
t.get('title', f'任务{i+1}'), t.get('description', ''),
|
||
t.get('priority', 'medium'), 1 if review_required else 0, '',
|
||
json.dumps(dep_ids), db.now(), db.now()))
|
||
created.append(tid)
|
||
who = f'→ {worker["name"]}' if worker else '(自动路由)'
|
||
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
|
||
(tid, 'info', f'AI 主管规划分派{who}:{t.get("title","")}', db.now()))
|
||
return created, f'AI 主管拆解 {len(tasks)} 个任务,分派给 {len(team)} 个团队 Worker'
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 监控:等待执行完 + AI 主管诊断重试失败任务
|
||
# ---------------------------------------------------------------------------
|
||
def _retry_count(task_id):
|
||
rows = db.q("SELECT COUNT(*) c FROM task_logs WHERE task_id=? AND message LIKE 'AI主管重试%'",
|
||
(task_id,))
|
||
return rows[0]['c'] if rows else 0
|
||
|
||
|
||
def _revise_and_retry(pid, task, manager, team):
|
||
"""AI 主管诊断失败任务:修订指令 → 重置 todo → 重新执行"""
|
||
n = _retry_count(task['id'])
|
||
if n >= 2 or not manager:
|
||
_log(pid, 'warn', f'任务「{task["title"]}」重试已达上限,保持失败状态')
|
||
return False
|
||
try:
|
||
data, r = _chat_json(
|
||
manager,
|
||
[{'role': 'system', 'content': '你只输出 JSON。你是严谨的 AI 主管。'},
|
||
{'role': 'user', 'content': REVISE_PROMPT.format(
|
||
title=task['title'], instruction=task['description'] or task['title'],
|
||
error=(task.get('error') or '')[:500])}],
|
||
max_tokens=2000, temperature=0.3)
|
||
if not isinstance(data, dict) or not data.get('instruction'):
|
||
raise ValueError('修订指令为空')
|
||
new_instr = data['instruction']
|
||
reason = data.get('reason') or '未知原因'
|
||
except Exception as e:
|
||
_log(pid, 'warn', f'任务「{task["title"]}」AI 主管诊断失败:{e},直接重试一次')
|
||
new_instr, reason = task['description'] or task['title'], 'AI主管诊断失败,原样重试'
|
||
db.w('UPDATE tasks SET description=?, status="todo", error="", updated_at=? WHERE id=?',
|
||
(new_instr, db.now(), task['id']))
|
||
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
|
||
(task['id'], 'warn', f'AI主管重试第{n+1}次(诊断:{reason}):修订指令后重新执行', db.now()))
|
||
_log(pid, 'info', f'AI 主管诊断「{task["title"]}」失败:{reason},修订后重试第 {n+1} 次')
|
||
engine.runner.submit(task['id'])
|
||
return True
|
||
|
||
|
||
def _running_tasks(pid):
|
||
return db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 AND status="running"', (pid,))
|
||
|
||
|
||
def _pending_tasks(pid):
|
||
return db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 AND status IN ("todo","failed")', (pid,))
|
||
|
||
|
||
def _blocked_dep(task):
|
||
"""todo 任务的前置依赖是否有未完成/失败的(导致永远无法就绪)"""
|
||
try:
|
||
deps = json.loads(task.get('depends_on') or '[]')
|
||
except Exception:
|
||
deps = []
|
||
for dep_id in deps:
|
||
dep = db.q('SELECT status FROM tasks WHERE id=?', (dep_id,), one=True)
|
||
if not dep or dep['status'] != 'done':
|
||
return True
|
||
return False
|
||
|
||
|
||
def monitor(pid, manager, team, timeout=1800):
|
||
"""轮询监控:任务全跑完或超时为止;失败任务交给 AI 主管诊断重试(最多 2 次)"""
|
||
deadline = time.time() + timeout
|
||
while time.time() < deadline:
|
||
# 失败任务:逐个请 AI 主管诊断重试(重试后转 todo 重新入队)
|
||
for t in [x for x in _pending_tasks(pid) if x['status'] == 'failed']:
|
||
try:
|
||
_revise_and_retry(pid, t, manager, team)
|
||
except Exception:
|
||
_log(pid, 'warn', f'任务「{t["title"]}」重试触发异常')
|
||
running = _running_tasks(pid)
|
||
pending = _pending_tasks(pid)
|
||
todo = [x for x in pending if x['status'] == 'todo']
|
||
retryable = [x for x in pending if x['status'] == 'failed' and _retry_count(x['id']) < 2]
|
||
# 出口:无运行中任务,且没有可重试的失败任务;
|
||
# 剩余 todo 均为被失败/未完成依赖卡死的任务 → 视为收尾(后续可手动「继续开工」)
|
||
if not running and not retryable:
|
||
if not todo or all(_blocked_dep(t) for t in todo):
|
||
return 'done'
|
||
time.sleep(5)
|
||
return 'timeout'
|
||
|
||
|
||
def _finalize(pid, status, message):
|
||
_set_auto(pid, status, message, finished=True)
|
||
_log(pid, 'success' if status == 'done' else 'warn', message)
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 主流程
|
||
# ---------------------------------------------------------------------------
|
||
def autostart_project(pid, review_required=0):
|
||
"""自动开工主流程(在后台线程中执行)"""
|
||
proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True)
|
||
if not proj:
|
||
return
|
||
_set_auto(pid, 'running', 'AI 主管已接单,正在规划…', started=True)
|
||
db.w('UPDATE projects SET status="active", updated_at=? WHERE id=?', (db.now(), pid))
|
||
try:
|
||
manager = get_manager(pid)
|
||
team = get_team(pid)
|
||
if not manager:
|
||
raise ValueError('未设置 AI 主管 Worker(项目需指定 manager_worker_id)')
|
||
if not team:
|
||
raise ValueError('未设置干活团队 Worker(至少勾选 1 个)')
|
||
_log(pid, 'info', f'AI 主管「{manager["name"]}」接管项目,团队:'
|
||
+ '、'.join(w['name'] for w in team))
|
||
|
||
# 1) 规划(已有任务则跳过,直接执行存量)
|
||
existing = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0', (pid,))[0]['c']
|
||
if existing == 0:
|
||
_set_auto(pid, 'running', 'AI 主管正在拆解任务(WBS 规划)…')
|
||
created, msg = plan_tasks(pid, manager, team, review_required=review_required)
|
||
_set_auto(pid, 'running', msg)
|
||
else:
|
||
msg = f'项目已有 {existing} 个任务,跳过规划直接开工'
|
||
_log(pid, 'info', msg)
|
||
_set_auto(pid, 'running', msg)
|
||
|
||
# 2) 执行整个工作流
|
||
rows = _pending_tasks(pid)
|
||
started = 0
|
||
for t in rows:
|
||
ok, blockers = engine.check_dependencies(t)
|
||
if ok and engine.runner.submit(t['id']):
|
||
started += 1
|
||
if started == 0 and not _running_tasks(pid):
|
||
raise ValueError('没有可执行的任务(请检查任务依赖或 Worker 状态)')
|
||
_set_auto(pid, 'running', f'已派活 {started} 个任务,AI 主管全程监控中…')
|
||
|
||
# 3) 监控 + 失败诊断重试
|
||
result = monitor(pid, manager, team)
|
||
if result == 'timeout':
|
||
raise ValueError('监控超时,仍有任务未完成(可在任务就绪后点「继续开工」)')
|
||
|
||
# 4) 收尾:全部任务完成 → 提交负责人验收(AI 无权宣布项目完成)
|
||
by = {r['status']: r['c'] for r in db.q(
|
||
'SELECT status, COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0 GROUP BY status', (pid,))}
|
||
done = by.get('done', 0)
|
||
failed = by.get('failed', 0)
|
||
review = by.get('review', 0)
|
||
total = sum(by.values())
|
||
costs = db.q('SELECT COALESCE(SUM(cost),0) c FROM cost_records WHERE project_id=?', (pid,))[0]['c']
|
||
if failed == 0 and review == 0 and total > 0:
|
||
# 转待验收:delivery.auto_complete_if_ready → enter_review(发邮件给负责人)
|
||
try:
|
||
r = delivery.auto_complete_if_ready(pid)
|
||
if r and r.get('ok'):
|
||
msg = f'AI 主管收工:{done}/{total} 个任务完成,已提交负责人验收,总成本 ¥{costs:.4f}'
|
||
else:
|
||
msg = f'AI 主管收工:{done}/{total} 个任务完成,但验收通知异常:{((r or {}).get("msg") or "")[:80]}'
|
||
except Exception as e:
|
||
msg = f'AI 主管收工:{done}/{total} 个任务完成,验收提交异常:{str(e)[:80]}'
|
||
_finalize(pid, 'done', msg)
|
||
else:
|
||
msg = f'AI 主管收工:{done}/{total} 个任务完成' + (f',{failed} 个失败' if failed else '') + \
|
||
(f',{review} 个待人工审核' if review else '') + f',总成本 ¥{costs:.4f}'
|
||
_finalize(pid, 'done' if failed == 0 else 'failed', msg)
|
||
try:
|
||
notify.notify('task_done', f'项目自动开工完成:{proj["name"]}', msg, save_alert=False)
|
||
except Exception:
|
||
pass
|
||
except Exception as e:
|
||
msg = f'自动开工中断:{str(e)[:200]}'
|
||
_finalize(pid, 'failed', msg)
|
||
_log(pid, 'error', traceback.format_exc())
|
||
# V3.4:连续失败/中断 → 邮件通知项目负责人(送达者)
|
||
try:
|
||
proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True)
|
||
if proj:
|
||
email = delivery.deliver_email_of(proj)
|
||
if email:
|
||
delivery.notify_deliverer(
|
||
proj, f'⚠️ 项目自动开工中断:{proj["name"]}',
|
||
f'AI 主管自动开工中断:{msg}\n\n'
|
||
f'请到平台查看项目详情(任务/日志),或稍后在项目页点击「🚀 AI 主管开工」重试。')
|
||
else:
|
||
_log(pid, 'warn', '项目未配置送达者邮箱,中断通知无法发送')
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
_threads = {}
|
||
|
||
|
||
def launch(pid, review_required=0):
|
||
"""后台线程启动自动开工(幂等:已有运行中线程则忽略)"""
|
||
t = _threads.get(pid)
|
||
if t and t.is_alive():
|
||
return False
|
||
proj = db.q('SELECT auto_status FROM projects WHERE id=?', (pid,), one=True)
|
||
if proj and proj.get('auto_status') == 'running':
|
||
# 数据库标记运行中但线程已死(服务重启遗留)→ 强制接管
|
||
db.w('UPDATE projects SET auto_status="none" WHERE id=?', (pid,))
|
||
t = threading.Thread(target=autostart_project, args=(pid, review_required), daemon=True)
|
||
_threads[pid] = t
|
||
t.start()
|
||
return True
|
||
|
||
|
||
def rerun(pid):
|
||
"""手动重新开工(换过主管/团队后生效;已有任务跳过规划直接执行)"""
|
||
return launch(pid)
|