From c8cc945b8db51b974f6ee48c07c5b44ecb5c5087 Mon Sep 17 00:00:00 2001 From: hz4th_coder Date: Fri, 14 Aug 2026 12:35:46 +0800 Subject: [PATCH] =?UTF-8?q?V3.3=20AI=E4=B8=BB=E7=AE=A1=E8=87=AA=E5=8A=A8?= =?UTF-8?q?=E5=BC=80=E5=B7=A5=EF=BC=9A=E6=96=B0=E5=BB=BA=E9=A1=B9=E7=9B=AE?= =?UTF-8?q?=E5=BF=85=E9=80=89AI=E4=B8=BB=E7=AE=A1+=E5=A4=9A=E9=80=89?= =?UTF-8?q?=E5=B9=B2=E6=B4=BB=E5=9B=A2=E9=98=9F=EF=BC=8C=E5=88=9B=E5=BB=BA?= =?UTF-8?q?=E5=8D=B3=E8=87=AA=E5=8A=A8=E8=B7=91=E8=B5=B7=E6=9D=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新建项目必填 AI 主管(manager_worker_id)+ 干活团队(project_team_workers 多选,默认AI牛/AI马),均可随时更换 - 创建后立即自动开工:AI主管WBS拆解→轮询分派团队→DAG自动执行→全程监控(autostart.py) - 失败任务由AI主管诊断原因并修订指令重试(最多2次),全程记入 project_logs 动态 - 项目页新增「AI主管动态」面板:状态/最新动态/主管与团队/动态流水,运行中每3秒自动刷新看板 - 支持手动「🚀 AI主管开工」:已有任务跳过规划直接执行存量;换主管/团队后按新配置生效 - 收工自动汇总完成数/失败数/待审核数/总成本,触发交付体系通知送达者 - 老库自动迁移:projects 新增 manager_worker_id/auto_status 等列 + project_team_workers/project_logs 表 --- README.md | 18 ++- app.py | 122 ++++++++++++---- autostart.py | 364 +++++++++++++++++++++++++++++++++++++++++++++++ db.py | 45 ++++++ enterprise.py | 2 +- static/app.js | 121 ++++++++++++++-- static/style.css | 16 +++ 7 files changed, 651 insertions(+), 37 deletions(-) create mode 100644 autostart.py diff --git a/README.md b/README.md index a42a4a1..0598d87 100644 --- a/README.md +++ b/README.md @@ -3,7 +3,23 @@ > 以「项目」为中心、以「AI Worker」为执行单元的项目管理平台。 > 把大模型团队变成一支可指挥、可审计、可控成本的"虚拟团队"。 -**当前版本:V3.2**(送达者=用户 / 用户邮箱必填 / 自定义角色 / Worker 权限组 / 交付体系) +**当前版本:V3.3**(AI 主管自动开工 / 项目创建即开工 / 失败诊断重试) + +--- + +## 🚀 V3.3 AI 主管自动开工(新增) + +> 创建项目 → 自动跑起来,全程不需要人点一下。 + +| 能力 | 说明 | +|---|---| +| 👔 必选 AI 主管 | 新建项目**必须指定一个 AI Worker 做管理者**(默认「AI主管」),负责拆解任务(WBS)、派活、监控、失败诊断;可随时在编辑项目中**换别的主管** | +| 👷 多选干活团队 | 新建项目**勾选多个干活 Worker** 供 AI 主管支配(默认勾选 AI牛/AI马),任务按轮询分派;可随时换团队 | +| ⚡ 创建即开工 | 项目创建后**立即自动开工**:AI 主管读目标/验收标准 → 拆解 4~8 个任务(含依赖)→ 分派团队执行 → 全程监控 | +| 🔁 失败诊断重试 | 任务失败后 AI 主管自动**诊断原因并修订执行指令**,重试最多 2 次;仍失败则收工并汇总失败项 | +| 📡 全程可视 | 项目页顶部「AI 主管动态」面板实时显示:开工状态/最新动态/主管与团队/动态流水(每 3 秒自动刷新,看板同步跳动) | +| 🚀 手动开工 | 项目页「🚀 AI 主管开工」按钮:已有任务则跳过规划直接执行存量;换过主管/团队后点它即按新配置生效 | +| ✅ 自动收尾 | 全部任务完成后自动汇总(完成数/失败数/总成本),触发交付体系通知送达者;任务默认自动验收(可在创建时勾选「需人工审核」改为 HITL) | --- diff --git a/app.py b/app.py index 31f8dfa..34c93e9 100644 --- a/app.py +++ b/app.py @@ -20,6 +20,7 @@ import eval as evalmod import templates as tplmod import enterprise import delivery +import autostart app = Flask(__name__, static_folder='static', static_url_path='') app.secret_key = config.SECRET_KEY @@ -230,7 +231,7 @@ def require_admin(fn): @app.route('/api/health') def health(): - return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v3.2.0'}) + return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v3.3.0'}) # --------------------------------------------------------------------------- @@ -344,6 +345,18 @@ def _attach_deliver_user(project): return project +def _attach_ai_team(project): + """V3.3 项目附带 AI 主管 + 干活团队 Worker 信息""" + pid = project['id'] + mgr = db.q('SELECT id, name, provider, model FROM workers WHERE id=?', + (project.get('manager_worker_id') or 0,), one=True) + project['manager_worker'] = dict(mgr) if mgr else None + project['team_workers'] = [dict(r) for r in db.q( + 'SELECT w.id, w.name, w.provider, w.model FROM project_team_workers t ' + 'JOIN workers w ON w.id=t.worker_id WHERE t.project_id=? ORDER BY t.rowid', (pid,))] + return project + + # --------------------------------------------------------------------------- # 项目 # --------------------------------------------------------------------------- @@ -364,26 +377,46 @@ def projects(): return jsonify({'ok': False, 'error': '送达者用户不存在'}), 400 if not (du.get('email') or '').strip(): return jsonify({'ok': False, 'error': f'送达者用户「{d.get("deliver_username", "")}」未设置邮箱,请先在用户管理中补全'}), 400 + # V3.3:必须指定 AI 主管(管理者)与干活团队(≥1 个 Worker) + mgr_id = d.get('manager_worker_id') + if not mgr_id: + return jsonify({'ok': False, 'error': '必填:选择 AI 主管 Worker(负责拆解任务、派活、监控)'}), 400 + mgr = db.q('SELECT * FROM workers WHERE id=? AND status="enabled"', (int(mgr_id),), one=True) + if not mgr: + return jsonify({'ok': False, 'error': 'AI 主管 Worker 不存在或已停用'}), 400 + team_ids = [int(x) for x in (d.get('team_worker_ids') or []) if x] + if not team_ids: + return jsonify({'ok': False, 'error': '必填:至少勾选 1 个干活 AI Worker 供 AI 主管支配'}), 400 + team = db.q(f'SELECT * FROM workers WHERE id IN ({scope_args_ph(team_ids)}) AND status="enabled"', + tuple(team_ids)) + if len(team) != len(set(team_ids)): + return jsonify({'ok': False, 'error': '存在不可用的干活 Worker(不存在或已停用)'}), 400 + review_required = 1 if d.get('auto_review') else 0 pid = db.w( 'INSERT INTO projects (name, description, objective, acceptance_criteria, ' 'status, budget_limit, deliver_user_id, deliver_email, deliver_type, deliver_note, workspace_dir, ' - 'created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)', + 'manager_worker_id, auto_status, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)', (d.get('name', '').strip(), d.get('description', ''), d.get('objective', ''), d.get('acceptance_criteria', ''), d.get('status', 'active'), float(d.get('budget_limit') or 0), int(duid), du['email'], d.get('deliver_type', 'web'), d.get('deliver_note', ''), - '', db.now(), db.now())) + '', int(mgr_id), 'none', db.now(), db.now())) # 工作目录名依赖自增 id:先拿到 id 再补写 db.w('UPDATE projects SET workspace_dir=? WHERE id=?', (f'project_{pid}', pid)) delivery.workspace_path(pid) + # 干活团队 + autostart.set_team(pid, team_ids) # 创建者自动成为项目管理员 u = _session_user() if u and not _is_super(): db.w('INSERT INTO user_projects (user_id, project_id, perm, created_at) VALUES (?,?,?,?)', (u['id'], pid, 'admin', db.now())) enterprise.audit(enterprise.current_actor(), 'project.create', f'project#{pid}', - f'「{d.get("name","")}」送达者用户#{duid} {du["email"]}', request.remote_addr or '') - return jsonify({'ok': True, 'id': pid}) + f'「{d.get("name","")}」送达者用户#{duid},AI主管#{mgr_id},团队{len(team_ids)}人', + request.remote_addr or '') + # 🚀 创建后立即自动开工:AI 主管规划 → 派活 → 执行 → 监控 + autostart.launch(pid, review_required=review_required) + return jsonify({'ok': True, 'id': pid, 'auto_started': True}) err = _check_perm_point('project.view') if err: return err @@ -397,6 +430,7 @@ def projects(): (r['id'],))[0]['c'] r['perm'] = _project_perm(r['id']) _attach_deliver_user(r) + _attach_ai_team(r) return jsonify({'ok': True, 'data': rows}) @@ -415,6 +449,7 @@ def project_detail(pid): return jsonify({'ok': False, 'error': '项目不存在'}), 404 p['perm'] = _project_perm(pid) _attach_deliver_user(p) + _attach_ai_team(p) return jsonify({'ok': True, 'data': p}) if request.method == 'DELETE': err = _check_project_perm(pid, 'admin') @@ -424,6 +459,8 @@ def project_detail(pid): db.w('DELETE FROM project_deliverables WHERE project_id=?', (pid,)) db.w('DELETE FROM tasks WHERE project_id=?', (pid,)) db.w('DELETE FROM cost_records WHERE project_id=?', (pid,)) + db.w('DELETE FROM project_team_workers WHERE project_id=?', (pid,)) + db.w('DELETE FROM project_logs WHERE project_id=?', (pid,)) db.w('DELETE FROM projects WHERE id=?', (pid,)) import shutil for d in (delivery.workspace_path(pid), delivery.demo_path(pid)): @@ -460,6 +497,23 @@ def project_detail(pid): return jsonify({'ok': False, 'error': '送达者用户未设置邮箱,请先补全'}), 400 db.w('UPDATE projects SET deliver_user_id=?, deliver_email=?, updated_at=? WHERE id=?', (int(d['deliver_user_id']), du['email'], db.now(), pid)) + # V3.3:AI 主管 / 干活团队变更(下次「开工」生效) + if 'manager_worker_id' in d and d.get('manager_worker_id'): + mgr = db.q('SELECT id FROM workers WHERE id=? AND status="enabled"', + (int(d['manager_worker_id']),), one=True) + if not mgr: + return jsonify({'ok': False, 'error': 'AI 主管 Worker 不存在或已停用'}), 400 + db.w('UPDATE projects SET manager_worker_id=?, updated_at=? WHERE id=?', + (int(d['manager_worker_id']), db.now(), pid)) + if 'team_worker_ids' in d: + team_ids = [int(x) for x in (d.get('team_worker_ids') or []) if x] + if not team_ids: + return jsonify({'ok': False, 'error': '至少勾选 1 个干活 AI Worker'}), 400 + team = db.q(f'SELECT id FROM workers WHERE id IN ({scope_args_ph(team_ids)}) AND status="enabled"', + tuple(team_ids)) + if len(team) != len(set(team_ids)): + return jsonify({'ok': False, 'error': '存在不可用的干活 Worker(不存在或已停用)'}), 400 + autostart.set_team(pid, team_ids) return jsonify({'ok': True}) @@ -1007,6 +1061,39 @@ def workflow_run(pid): return jsonify({'ok': True, 'started': started, 'blocked': blocked}) +@app.route('/api/projects//autostart', methods=['POST']) +@require_auth +def project_autostart(pid): + """V3.3 AI 主管开工/继续开工:拆解(无任务时)→ 派活 → 执行 → 监控 + 换过 AI 主管/团队后调用即按新配置生效;已有任务跳过规划直接跑存量。""" + err = _check_project_perm(pid, 'manage') + if err: + return err + p = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if not p: + return jsonify({'ok': False, 'error': '项目不存在'}), 404 + if not p.get('manager_worker_id'): + return jsonify({'ok': False, 'error': '请先在编辑项目中设置 AI 主管 Worker'}), 400 + if not db.q('SELECT 1 FROM project_team_workers WHERE project_id=?', (pid,)): + return jsonify({'ok': False, 'error': '请先在编辑项目中勾选干活团队 Worker(≥1 个)'}), 400 + d = request.get_json(force=True) or {} + ok = autostart.launch(pid, review_required=1 if d.get('auto_review') else 0) + if not ok: + return jsonify({'ok': False, 'error': 'AI 主管正在工作中,请稍候'}), 400 + return jsonify({'ok': True}) + + +@app.route('/api/projects//autologs') +@require_auth +def project_autologs(pid): + """V3.3 AI 主管项目级动态(拆解/派活/监控/诊断重试)""" + err = _check_project_perm(pid, 'view') + if err: + return err + rows = db.q('SELECT * FROM project_logs WHERE project_id=? ORDER BY id DESC LIMIT 50', (pid,)) + return jsonify({'ok': True, 'data': rows}) + + @app.route('/api/projects//dag') @require_auth def project_dag(pid): @@ -1029,32 +1116,13 @@ def project_dag(pid): # --------------------------------------------------------------------------- # AI 辅助规划(WBS 生成 + 导入) +# 提示词与解析复用 autostart.py(V3.3 与 AI 主管自动开工同一套) # --------------------------------------------------------------------------- -WBS_PROMPT = ( - '你是资深项目经理。请把下面的项目目标拆解为可执行的任务列表(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}' -) +WBS_PROMPT = autostart.WBS_PROMPT def _parse_wbs(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 - data = _json.loads(t[start:end]) - tasks = data['tasks'] if isinstance(data, dict) else data - assert isinstance(tasks, list) and tasks, '任务列表为空' - return tasks + return autostart.parse_wbs(text) @app.route('/api/projects//wbs/generate', methods=['POST']) diff --git a/autostart.py b/autostart.py new file mode 100644 index 0000000..5cf5c4b --- /dev/null +++ b/autostart.py @@ -0,0 +1,364 @@ +# -*- 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 parse_wbs(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 + data = json.loads(t[start:end]) + tasks = data['tasks'] if isinstance(data, dict) else data + assert isinstance(tasks, list) and tasks, '任务列表为空' + 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 _extract_json(text): + 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 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) +# --------------------------------------------------------------------------- +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 主管无法规划') + 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, + base_url=manager['base_url'] or None, api_key=manager['api_key'] or None) + tasks = parse_wbs(r['text']) + 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) + if status == 'done': + try: + delivery.auto_complete_if_ready(pid) + except Exception: + pass + _log(pid, 'success', message) + else: + _log(pid, '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) 收尾 + 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'] + 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: + _finalize(pid, 'failed', f'自动开工中断:{str(e)[:200]}') + _log(pid, 'error', traceback.format_exc()) + + +_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) diff --git a/db.py b/db.py index 084b062..5e28f58 100644 --- a/db.py +++ b/db.py @@ -23,10 +23,23 @@ CREATE TABLE IF NOT EXISTS projects ( workspace_dir TEXT DEFAULT '', -- 项目工作目录(相对 data/ 的目录名) demo_url TEXT DEFAULT '', -- 网页交付物 Demo 访问地址 delivered_at INTEGER, -- 最近一次交付/送达时间 + manager_worker_id INTEGER, -- V3.3 AI 主管 Worker id(负责拆解/派活/监控) + auto_status TEXT DEFAULT 'none', -- V3.3 自动开工状态 none/running/done/failed + auto_message TEXT DEFAULT '', -- V3.3 自动开工最新动态 + auto_started_at INTEGER, -- V3.3 最近一次自动开工时间 + auto_finished_at INTEGER, -- V3.3 最近一次自动收尾时间 created_at INTEGER, updated_at INTEGER ); +-- V3.3 项目干活团队:AI 主管支配的多个 Worker +CREATE TABLE IF NOT EXISTS project_team_workers ( + project_id INTEGER NOT NULL, + worker_id INTEGER NOT NULL, + created_at INTEGER, + PRIMARY KEY (project_id, worker_id) +); + CREATE TABLE IF NOT EXISTS workers ( id INTEGER PRIMARY KEY AUTOINCREMENT, name TEXT NOT NULL, @@ -73,6 +86,15 @@ CREATE TABLE IF NOT EXISTS task_logs ( created_at INTEGER ); +-- V3.3 AI 主管项目级动态(拆解/派活/监控/诊断重试) +CREATE TABLE IF NOT EXISTS project_logs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + project_id INTEGER NOT NULL, + level TEXT DEFAULT 'info', + message TEXT DEFAULT '', + created_at INTEGER +); + CREATE TABLE IF NOT EXISTS cost_records ( id INTEGER PRIMARY KEY AUTOINCREMENT, task_id INTEGER, @@ -396,6 +418,29 @@ def _migrate(): # V3.2:projects 送达者改为用户 id(实时取邮箱) if 'deliver_user_id' not in pcols: conn.execute('ALTER TABLE projects ADD COLUMN deliver_user_id INTEGER') + # V3.3:projects AI 主管 + 自动开工状态 + for col, ddl in ( + ('manager_worker_id', 'ALTER TABLE projects ADD COLUMN manager_worker_id INTEGER'), + ('auto_status', "ALTER TABLE projects ADD COLUMN auto_status TEXT DEFAULT 'none'"), + ('auto_message', "ALTER TABLE projects ADD COLUMN auto_message TEXT DEFAULT ''"), + ('auto_started_at', 'ALTER TABLE projects ADD COLUMN auto_started_at INTEGER'), + ('auto_finished_at', 'ALTER TABLE projects ADD COLUMN auto_finished_at INTEGER'), + ): + if col not in pcols: + conn.execute(ddl) + conn.execute('''CREATE TABLE IF NOT EXISTS project_team_workers ( + project_id INTEGER NOT NULL, + worker_id INTEGER NOT NULL, + created_at INTEGER, + PRIMARY KEY (project_id, worker_id) + )''') + conn.execute('''CREATE TABLE IF NOT EXISTS project_logs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + project_id INTEGER NOT NULL, + level TEXT DEFAULT 'info', + message TEXT DEFAULT '', + created_at INTEGER + )''') # V3.2:users 邮箱列 + 存量用户默认邮箱 + 旧项目按邮箱回填送达者用户 ucols = {r['name'] for r in conn.execute('PRAGMA table_info(users)')} if 'email' not in ucols: diff --git a/enterprise.py b/enterprise.py index 9fa5731..ed456cf 100644 --- a/enterprise.py +++ b/enterprise.py @@ -304,7 +304,7 @@ def export_all(): 'eval_results', 'templates', 'users', 'audit_logs', 'project_deliverables', 'user_projects', 'user_workers'] out = {'exported_at': time.strftime('%Y-%m-%d %H:%M:%S'), - 'platform': 'ai-worker-platform', 'version': 'v3.2.0'} + 'platform': 'ai-worker-platform', 'version': 'v3.3.0'} for t in tables: try: out[t] = db.q(f'SELECT * FROM {t}') diff --git a/static/app.js b/static/app.js index 1b276ad..c0b6cd3 100644 --- a/static/app.js +++ b/static/app.js @@ -95,6 +95,8 @@ const routes = { 'enterprise': pageEnterprise }; function router() { + clearInterval(projAutoTimer); // 离开项目页时停止 AI 主管轮询 + projAutoTimer = null; const hash = location.hash.replace(/^#\//, '') || 'dashboard'; const parts = hash.split('/'); const name = parts[0]; @@ -176,6 +178,7 @@ async function pageProjects() { function openProjectModal(p = {}) { const usersP = api('/api/users'); + const workersP = api('/api/workers'); openModal(`

${p.id ? '编辑项目' : '新建项目'}

@@ -198,9 +201,17 @@ function openProjectModal(p = {}) { +
🤖 AI 虚拟团队(创建后自动开工)
+ + + +
加载 Worker 中…
+
+ +
`); usersP.then(r => { const sel = $('#pm-email'); @@ -211,6 +222,24 @@ function openProjectModal(p = {}) { ``).join(''); if (!cur) sel.value = ''; }).catch(() => {}); + workersP.then(r => { + const ws = r.data.filter(w => w.status === 'enabled'); + if (!ws.length) { $('#pm-mgr').innerHTML = ''; return; } + // 默认:AI 主管 = 名字含「主管/经理」的 Worker;干活团队 = 含「牛/马」的(如 AI牛/AI马),否则除主管外全部 + let curMgr = p.manager_worker ? p.manager_worker.id : (p.manager_worker_id || ''); + if (!curMgr) { + const m = ws.find(w => /主管|经理|manager/i.test(w.name)); + curMgr = (m || ws[0]).id; + } + $('#pm-mgr').innerHTML = '' + + ws.map(w => ``).join(''); + const curTeam = new Set((p.team_workers || []).map(w => w.id)); + let defaults = ws.filter(w => /牛|马/.test(w.name) && String(w.id) !== String(curMgr)); + if (!defaults.length) defaults = ws.filter(w => String(w.id) !== String(curMgr)); + $('#pm-team').innerHTML = ws.map(w => ` + `).join(''); + }).catch(() => {}); $('#pm-save').addEventListener('click', async () => { const body = { name: $('#pm-name').value.trim(), objective: $('#pm-objective').value, @@ -218,41 +247,115 @@ function openProjectModal(p = {}) { status: $('#pm-status').value, budget_limit: parseFloat($('#pm-budget').value || 0), deliver_user_id: Number($('#pm-email').value || 0), deliver_type: $('#pm-dtype').value, - deliver_note: $('#pm-dnote').value.trim() + deliver_note: $('#pm-dnote').value.trim(), + manager_worker_id: Number($('#pm-mgr').value || 0), + team_worker_ids: [...$$('#pm-team input:checked')].map(i => Number(i.value)), + auto_review: !!$('#pm-auto-review').checked }; if (!body.name) return toast('请填写项目名称', 'err'); if (!body.deliver_user_id) return toast('必填:选择送达者用户', 'err'); + if (!body.manager_worker_id) return toast('必填:选择 AI 主管 Worker', 'err'); + if (!body.team_worker_ids.length) return toast('必填:至少勾选 1 个干活 AI Worker', 'err'); try { - await api(p.id ? `/api/projects/${p.id}` : '/api/projects', {method: p.id ? 'PUT' : 'POST', body}); - toast('已保存', 'ok'); closeModal(); router(); + if (p.id) { + await api(`/api/projects/${p.id}`, {method: 'PUT', body}); + toast('已保存(新主管/团队将在下次开工时生效)', 'ok'); closeModal(); router(); + } else { + const r = await api('/api/projects', {method: 'POST', body}); + toast('项目已创建,AI 主管已开工 🚀', 'ok'); closeModal(); + location.hash = `#/project/${r.id}`; + } } catch (e) { toast(e.message, 'err'); } }); } -/* ---------- 项目详情(看板 / DAG / 知识库) ---------- */ -let projCtx = null; // {pid, p, tasks, workers, tab} +/* ---------- 项目详情(看板 / DAG / 知识库 / 交付) ---------- */ +let projCtx = null; // {pid, p, tasks, workers, tab, autoLogs} +let projAutoTimer = null; async function pageProject([pid, tab]) { const p = (await api(`/api/projects/${pid}`)).data; const tasks = (await api(`/api/tasks?project_id=${pid}`)).data; const workers = (await api('/api/workers')).data; - projCtx = {pid, p, tasks, workers, tab: tab || 'kanban'}; + const autoLogs = (await api(`/api/projects/${pid}/autologs`)).data; + projCtx = {pid, p, tasks, workers, tab: tab || 'kanban', autoLogs}; renderProjectShell(); + renderActiveTab(); + startProjAutoPoll(); +} + +function renderActiveTab() { if (projCtx.tab === 'dag') renderDag(); else if (projCtx.tab === 'kb') renderKb(); else if (projCtx.tab === 'deliver') renderDeliver(); else renderKanban(); } +/* V3.3:AI 主管自动开工 —— 轮询刷新状态/动态/看板 */ +function startProjAutoPoll() { + clearInterval(projAutoTimer); + renderAutoPanel(); + if (projCtx.p.auto_status !== 'running') return; + projAutoTimer = setInterval(async () => { + try { + const wasRunning = projCtx.p.auto_status === 'running'; + const p = (await api(`/api/projects/${projCtx.pid}`)).data; + projCtx.p = p; + projCtx.autoLogs = (await api(`/api/projects/${projCtx.pid}/autologs`)).data; + const running = p.auto_status === 'running'; + // 状态变化(开工结束/中断)或运行中:都刷新一次任务视图 + if (running || wasRunning) { + projCtx.tasks = (await api(`/api/tasks?project_id=${projCtx.pid}`)).data; + if (projCtx.tab === 'dag') renderDag(); else renderKanban(); + } + renderAutoPanel(); + if (!running) { clearInterval(projAutoTimer); projAutoTimer = null; } + } catch (e) {} + }, 3000); +} + +async function autoStartProj() { + try { + await api(`/api/projects/${projCtx.pid}/autostart`, {method: 'POST'}); + toast('🚀 AI 主管已开工,正在拆解派活…', 'ok'); + const p = (await api(`/api/projects/${projCtx.pid}`)).data; + projCtx.p = p; + projCtx.autoLogs = (await api(`/api/projects/${projCtx.pid}/autologs`)).data; + renderAutoPanel(); + startProjAutoPoll(); + } catch (e) { toast(e.message, 'err'); } +} + +function renderAutoPanel() { + const el = $('#auto-panel'); + if (!el || !projCtx) return; + const p = projCtx.p; + const running = p.auto_status === 'running'; + const st = {none:['未开工','todo'], running:['🚀 AI 主管工作中','running'], done:['✅ 自动开工完成','done'], failed:['⚠️ 开工中断','failed']}[p.auto_status] || ['—','todo']; + const mgr = p.manager_worker ? esc(p.manager_worker.name) : '未设置'; + const team = (p.team_workers || []).map(w => `${esc(w.name)}`).join(' ') || '未设置'; + const logs = (projCtx.autoLogs || []).map(l => `
${l.level === 'error' ? '❌' : l.level === 'warn' ? '⚠️' : l.level === 'success' ? '✅' : '🤖'}${esc(l.message)}${fmtTime(l.created_at)}
`).join(''); + el.innerHTML = ` +
+ ${st[0]} + ${esc(p.auto_message || '尚未开工:创建项目后自动开工,或点右上角「🚀 AI 主管开工」')} + ${running ? '' : ''} +
+
👔 AI 主管:${mgr} | 👷 干活团队:${team} | ⏱ ${p.auto_started_at ? '开工于 ' + fmtTime(p.auto_started_at) : '—'}
+ ${logs ? `
${logs}
` : ''}`; +} + function renderProjectShell() { const {p, pid, tab} = projCtx; + const running = p.auto_status === 'running'; $('#main').innerHTML = `
← 返回

${esc(p.name)}${esc(p.objective || '')}

${p.perm === 'view' ? '' : ``} ${p.perm === 'view' ? '' : ``} - ${p.perm === 'view' ? '' : ``} + ${p.perm === 'view' ? '' : ``} + ${p.perm === 'view' || running ? '' : ``}
看板 @@ -260,7 +363,9 @@ function renderProjectShell() { 知识库 RAG 📦 交付
+
`; + renderAutoPanel(); } async function refreshProj() { diff --git a/static/style.css b/static/style.css index d7d64c1..ebf61e5 100644 --- a/static/style.css +++ b/static/style.css @@ -198,3 +198,19 @@ code{background:var(--panel2);border:1px solid var(--border);border-radius:6px;p .form-grid textarea{resize:vertical} #main .tabs a{cursor:pointer} .tag{display:inline-block;background:var(--accent-soft,#eef2ff);color:#4f5bd5;border-radius:4px;padding:0 6px;font-size:11px;margin-right:4px} + +/* ---------- V3.3 AI 主管自动开工 ---------- */ +.sec-divider{font-size:12px;font-weight:600;color:var(--accent);margin:16px 0 4px;padding-top:12px;border-top:1px dashed var(--border)} +.pick-list{display:flex;flex-wrap:wrap;gap:8px;margin-top:6px} +.pick-item{display:flex;align-items:center;gap:6px;font-size:12px;color:var(--text);margin:0;padding:6px 12px;background:var(--panel2);border:1px solid var(--border);border-radius:20px;cursor:pointer;font-weight:400} +.pick-item:has(input:checked){border-color:var(--accent);background:rgba(79,140,255,.12)} +.pick-item input{width:auto} +.pick-item small{color:var(--muted)} +.auto-bar{display:flex;align-items:center;gap:10px;background:var(--panel);border:1px solid var(--border);border-radius:10px;padding:10px 14px;margin-bottom:8px} +.auto-meta{font-size:12px;color:var(--muted);padding:0 4px 8px} +.auto-logs{background:var(--panel);border:1px solid var(--border);border-radius:10px;padding:8px 12px;margin-bottom:14px;max-height:180px;overflow-y:auto} +.al-item{display:flex;gap:8px;align-items:baseline;font-size:12px;padding:4px 0;border-bottom:1px dashed var(--border)} +.al-item:last-child{border-bottom:none} +.al-lv{flex:none} +.al-msg{flex:1;word-break:break-all} +.al-time{flex:none;color:var(--muted);font-size:11px}