Compare commits

...
1 Commits
Author SHA1 Message Date
hz4th_coder c8cc945b8d V3.3 AI主管自动开工:新建项目必选AI主管+多选干活团队,创建即自动跑起来
- 新建项目必填 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 表
2026-08-14 12:35:46 +08:00
7 changed files with 651 additions and 37 deletions
+17 -1
View File
@@ -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) |
---
+95 -27
View File
@@ -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/<int:pid>/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/<int:pid>/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/<int:pid>/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/<int:pid>/wbs/generate', methods=['POST'])
+364
View File
@@ -0,0 +1,364 @@
# -*- coding: utf-8 -*-
"""
V3.3 AI 主管自动开工引擎
========================
新建项目后立即自动运转:
1. AI 主管(项目 manager_worker_id)读取项目目标/验收标准,拆解 WBS 任务
2. 把任务分派给团队 Workerproject_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)
+45
View File
@@ -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.2projects 送达者改为用户 id(实时取邮箱)
if 'deliver_user_id' not in pcols:
conn.execute('ALTER TABLE projects ADD COLUMN deliver_user_id INTEGER')
# V3.3projects 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:
+1 -1
View File
@@ -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}')
+113 -8
View File
@@ -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(`
<h3>${p.id ? '编辑项目' : '新建项目'}<button class="close" data-close>×</button></h3>
<label>项目名称 *</label><input id="pm-name" value="${esc(p.name || '')}" placeholder="例如:营销内容批量生产">
@@ -198,9 +201,17 @@ function openProjectModal(p = {}) {
</select></div>
</div>
<label>交付说明</label><input id="pm-dnote" value="${esc(p.deliver_note || '')}" placeholder="给送达者的交付说明(可选)">
<div class="sec-divider">🤖 AI 虚拟团队(创建后自动开工)</div>
<label>👔 AI 主管(管理者)* <small style="color:var(--muted)">负责拆解任务、派活、监控,默认「AI主管」</small></label>
<select id="pm-mgr"><option value="">加载 Worker 中…</option></select>
<label>👷 干活团队(AI 主管支配的 Worker,可多选)* <small style="color:var(--muted)">默认 AI牛 / AI马</small></label>
<div id="pm-team" class="pick-list"><div class="empty" style="padding:8px">加载 Worker 中…</div></div>
<div class="row">
<label style="display:flex;align-items:center;gap:6px;font-weight:normal"><input type="checkbox" id="pm-auto-review"> 任务完成后需人工审核(默认不勾选 = AI 主管全自动跑到完工)</label>
</div>
<div class="modal-actions">
<button class="btn" data-close>取消</button>
<button class="btn primary" id="pm-save">保存</button>
<button class="btn primary" id="pm-save">${p.id ? '保存' : '创建并开工 🚀'}</button>
</div>`);
usersP.then(r => {
const sel = $('#pm-email');
@@ -211,6 +222,24 @@ function openProjectModal(p = {}) {
`<option value="${u.id}" ${String(u.id) === String(cur) ? 'selected' : ''}>${esc(u.display_name || u.username)}${esc(u.username)} · ${esc(u.email || '无邮箱')}</option>`).join('');
if (!cur) sel.value = '';
}).catch(() => {});
workersP.then(r => {
const ws = r.data.filter(w => w.status === 'enabled');
if (!ws.length) { $('#pm-mgr').innerHTML = '<option value="">请先在 Worker 管理中注册 AI Worker</option>'; 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 = '<option value="">— 请选择 AI 主管 —</option>' +
ws.map(w => `<option value="${w.id}" ${String(w.id) === String(curMgr) ? 'selected' : ''}>${esc(w.name)}${esc(w.provider)}/${esc(w.model)}</option>`).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 => `
<label class="pick-item"><input type="checkbox" value="${w.id}" ${curTeam.has(w.id) || (!p.id && defaults.includes(w)) ? 'checked' : ''}>
<span>${esc(w.name)}</span><small>${esc(w.provider)}/${esc(w.model)}</small></label>`).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) : '<span style="color:var(--warn)">未设置</span>';
const team = (p.team_workers || []).map(w => `<span class="pill" style="color:var(--accent)">${esc(w.name)}</span>`).join(' ') || '<span style="color:var(--warn)">未设置</span>';
const logs = (projCtx.autoLogs || []).map(l => `<div class="al-item"><span class="al-lv">${l.level === 'error' ? '❌' : l.level === 'warn' ? '⚠️' : l.level === 'success' ? '✅' : '🤖'}</span><span class="al-msg">${esc(l.message)}</span><span class="al-time">${fmtTime(l.created_at)}</span></div>`).join('');
el.innerHTML = `
<div class="auto-bar">
<span class="badge st-${st[1]}">${st[0]}</span>
<span style="flex:1;color:var(--muted)">${esc(p.auto_message || '尚未开工:创建项目后自动开工,或点右上角「🚀 AI 主管开工」')}</span>
${running ? '<span class="spin"></span>' : ''}
</div>
<div class="auto-meta">👔 AI 主管:${mgr} 👷 干活团队:${team} ${p.auto_started_at ? '开工于 ' + fmtTime(p.auto_started_at) : '—'}</div>
${logs ? `<div class="auto-logs">${logs}</div>` : ''}`;
}
function renderProjectShell() {
const {p, pid, tab} = projCtx;
const running = p.auto_status === 'running';
$('#main').innerHTML = `
<div class="toolbar">
<a class="btn" href="#/projects">← 返回</a>
<h1 class="page-title" style="margin:0;flex:1">${esc(p.name)}<small>${esc(p.objective || '')}</small></h1>
${p.perm === 'view' ? '' : `<button class="btn" onclick="openProjectModal(${JSON.stringify(p).replace(/"/g, '&quot;')})">编辑项目</button>`}
${p.perm === 'view' ? '' : `<button class="btn" onclick="openWbsModal()">🧠 AI 规划(WBS)</button>`}
${p.perm === 'view' ? '' : `<button class="btn primary" onclick="openTaskModal()"> 新建任务</button>`}
${p.perm === 'view' ? '' : `<button class="btn" onclick="openTaskModal()"> 新建任务</button>`}
${p.perm === 'view' || running ? '' : `<button class="btn primary" onclick="autoStartProj()">🚀 AI 主管开工</button>`}
</div>
<div class="tabs">
<a href="#/project/${pid}/kanban" class="${tab === 'kanban' ? 'active' : ''}">看板</a>
@@ -260,7 +363,9 @@ function renderProjectShell() {
<a href="#/project/${pid}/kb" class="${tab === 'kb' ? 'active' : ''}">知识库 RAG</a>
<a href="#/project/${pid}/deliver" class="${tab === 'deliver' ? 'active' : ''}">📦 交付</a>
</div>
<div id="auto-panel"></div>
<div id="tab-body"></div>`;
renderAutoPanel();
}
async function refreshProj() {
+16
View File
@@ -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}