# -*- 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)