Files
ai-worker-platform/autostart.py
T
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

365 lines
16 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""
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)