Files
ai-worker-platform/app.py
T
hz4th_coder 9c45a143b9 DAG画布:文字清晰 + 无边画布 + 删除区回收站
1. 文字看不清修复:.dag-node hover stroke 继承到文字导致白色小字被描边盖糊 → label/sub 加 stroke:none,字号加大(13/11px)加粗
2. 无边画布:bg rect 扩大至±4000且与画布同色 + svg overflow:visible,大幅平移缩放不再露出边缘/裁剪节点
3. 删除区与回收站:
- 画布左下角🗑️删除区,节点拖入即软删除(deleted=1)
- 点击删除区打开回收站:恢复/彻底删除
- 后端 tasks 加 deleted/deleted_at 列,全查询过滤,trash/restore/hard-delete API,看板删除改软删
- 修复画布重建后事件失效:box监听每次重绘重绑、window监听只绑一次且拖拽状态提升为模块级(dagDrag)
2026-08-12 23:44:20 +08:00

1398 lines
60 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 -*-
"""
AI Worker 项目管理平台 - MVP
人派活 → AI 干活 → 人验收 最小闭环
"""
import os
import secrets
import functools
import json as _json
from flask import Flask, request, jsonify, session, send_from_directory, redirect
import config
import db
import engine
import llm_gateway
import rag
import notify
import agents
import eval as evalmod
import templates as tplmod
import enterprise
app = Flask(__name__, static_folder='static', static_url_path='')
app.secret_key = config.SECRET_KEY
db.init_db()
db.recover_stale_runs()
evalmod.create_builtin_datasets()
tplmod.create_builtin_templates()
# ---------------------------------------------------------------------------
# 鉴权(简单口令登录)
# ---------------------------------------------------------------------------
def auth_enabled():
return bool(config.AUTH_PASSWORD)
@app.route('/api/login', methods=['POST'])
def login():
data = request.get_json(force=True)
if not auth_enabled():
session['username'] = 'admin'
session['authed'] = True
return jsonify({'ok': True})
username = (data.get('username') or 'admin').strip()
password = data.get('password') or ''
# 1) 本地用户(含默认 admin)
u = enterprise.verify_local(username, password)
if u:
session['username'] = u['username']
session['role'] = u['role']
session['authed'] = True
enterprise.audit(u['username'], 'auth.login', '本地登录', ip=request.remote_addr or '')
return jsonify({'ok': True, 'user': {'username': u['username'], 'role': u['role'],
'display_name': u['display_name']}})
# 2) LDAP SSO
if enterprise.sso_config().get('ldap_enabled') == '1':
u, err = enterprise.ldap_login(username, password)
if u:
session['username'] = u['username']
session['role'] = u['role']
session['authed'] = True
enterprise.audit(u['username'], 'auth.login', 'LDAP 登录', ip=request.remote_addr or '')
return jsonify({'ok': True, 'user': {'username': u['username'], 'role': u['role'],
'display_name': u['display_name']}})
# 3) 兼容 V1:口令直登(映射为 admin)
if password == config.AUTH_PASSWORD:
session['username'] = 'admin'
session['role'] = 'admin'
session['authed'] = True
enterprise.audit('admin', 'auth.login', '口令登录(V1 兼容)', ip=request.remote_addr or '')
return jsonify({'ok': True, 'user': {'username': 'admin', 'role': 'admin', 'display_name': '管理员'}})
enterprise.audit(username, 'auth.login_failed', '登录失败', ip=request.remote_addr or '')
return jsonify({'ok': False, 'error': '用户名或口令错误'}), 401
@app.route('/api/logout', methods=['POST'])
def logout():
enterprise.audit(session.get('username', '?'), 'auth.logout', '退出登录', ip=request.remote_addr or '')
session.clear()
return jsonify({'ok': True})
@app.route('/api/me')
def me():
authed = not auth_enabled() or session.get('authed')
u = None
if authed:
uname = session.get('username', 'admin')
row = enterprise.get_user(uname)
u = {'username': uname, 'role': session.get('role') or (row['role'] if row else 'admin'),
'display_name': row['display_name'] if row else uname}
return jsonify({'ok': True, 'authed': authed, 'user': u})
def _auth_ok():
"""会话或 API Token 任一通过即可"""
if not auth_enabled():
return True
if session.get('authed'):
return True
hdr = request.headers.get('Authorization', '')
if hdr.startswith('Bearer '):
tok = hdr[7:].strip()
r = db.q('SELECT id FROM api_tokens WHERE token=?', (tok,), one=True)
if r:
db.w('UPDATE api_tokens SET last_used_at=? WHERE id=?', (db.now(), r['id']))
return True
return False
def require_auth(fn):
@functools.wraps(fn)
def wrapper(*args, **kwargs):
if not _auth_ok():
return jsonify({'ok': False, 'error': '未登录或 Token 无效'}), 401
try:
import flask
flask.g.auth_actor = enterprise.current_actor()
except Exception:
pass
return fn(*args, **kwargs)
return wrapper
def require_admin(fn):
"""RBAC:仅管理员(含 token 调用者视为管理员)可访问"""
@functools.wraps(fn)
def wrapper(*args, **kwargs):
if not _auth_ok():
return jsonify({'ok': False, 'error': '未登录或 Token 无效'}), 401
try:
import flask
flask.g.auth_actor = enterprise.current_actor()
except Exception:
pass
actor = enterprise.current_actor()
if actor.startswith('token:'):
return fn(*args, **kwargs)
username = session.get('username')
if username and enterprise.is_admin(username):
return fn(*args, **kwargs)
return jsonify({'ok': False, 'error': '需要管理员权限'}), 403
return wrapper
@app.route('/api/health')
def health():
return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v2.0.0'})
# ---------------------------------------------------------------------------
# 项目
# ---------------------------------------------------------------------------
@app.route('/api/projects', methods=['GET', 'POST'])
@require_auth
def projects():
if request.method == 'POST':
d = request.get_json(force=True)
pid = db.w(
'INSERT INTO projects (name, description, objective, acceptance_criteria, '
'status, budget_limit, 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), db.now(), db.now()))
return jsonify({'ok': True, 'id': pid})
rows = db.q('SELECT * FROM projects ORDER BY id DESC')
for r in rows:
r['task_count'] = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0', (r['id'],))[0]['c']
r['done_count'] = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND status="done" AND deleted=0',
(r['id'],))[0]['c']
return jsonify({'ok': True, 'data': rows})
@app.route('/api/projects/<int:pid>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def project_detail(pid):
if request.method == 'GET':
p = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True)
if not p:
return jsonify({'ok': False, 'error': '项目不存在'}), 404
return jsonify({'ok': True, 'data': p})
if request.method == 'DELETE':
db.w('DELETE FROM tasks WHERE project_id=?', (pid,))
db.w('DELETE FROM cost_records WHERE project_id=?', (pid,))
db.w('DELETE FROM projects WHERE id=?', (pid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['name', 'description', 'objective', 'acceptance_criteria', 'status', 'budget_limit']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if sets:
args.append(db.now())
db.w(f'UPDATE projects SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, pid))
return jsonify({'ok': True})
# ---------------------------------------------------------------------------
# Worker
# ---------------------------------------------------------------------------
@app.route('/api/workers', methods=['GET', 'POST'])
@require_auth
def workers():
if request.method == 'POST':
d = request.get_json(force=True)
wid = db.w(
'INSERT INTO workers (name, description, provider, model, base_url, api_key, '
'system_prompt, temperature, max_tokens, task_cost_limit, monthly_cost_limit, '
'status, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)',
(d.get('name', '').strip(), d.get('description', ''), d.get('provider', ''),
d.get('model', ''), d.get('base_url', ''), d.get('api_key', ''),
d.get('system_prompt', ''), float(d.get('temperature', 0.7)),
int(d.get('max_tokens', 2000)), float(d.get('task_cost_limit') or 0),
float(d.get('monthly_cost_limit') or 0), d.get('status', 'enabled'),
db.now(), db.now()))
return jsonify({'ok': True, 'id': wid})
rows = db.q('SELECT * FROM workers ORDER BY id DESC')
for r in rows:
r['month_cost'] = round(db.monthly_worker_cost(r['id']), 6)
return jsonify({'ok': True, 'data': rows})
@app.route('/api/workers/<int:wid>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def worker_detail(wid):
if request.method == 'GET':
w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True)
return jsonify({'ok': True, 'data': w}) if w else (jsonify({'ok': False, 'error': '不存在'}), 404)
if request.method == 'DELETE':
db.w('UPDATE tasks SET worker_id=NULL WHERE worker_id=?', (wid,))
db.w('DELETE FROM workers WHERE id=?', (wid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['name', 'description', 'provider', 'model', 'base_url', 'api_key',
'system_prompt', 'temperature', 'max_tokens', 'task_cost_limit',
'monthly_cost_limit', 'status']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if sets:
args.append(db.now())
db.w(f'UPDATE workers SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, wid))
return jsonify({'ok': True})
@app.route('/api/workers/<int:wid>/test', methods=['POST'])
@require_auth
def worker_test(wid):
w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True)
if not w:
return jsonify({'ok': False, 'error': '不存在'}), 404
try:
r = llm_gateway.test_connection(w['provider'], w['model'],
base_url=w['base_url'] or None,
api_key=w['api_key'] or None)
return jsonify({'ok': True, 'data': r})
except Exception as e:
return jsonify({'ok': False, 'error': str(e)})
@app.route('/api/workers/<int:wid>/vision_test', methods=['POST'])
@require_auth
def worker_vision_test(wid):
"""视觉智能体能力测试:传入图片(URL 或 base64)与问题,验证多模态理解"""
w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True)
if not w:
return jsonify({'ok': False, 'error': '不存在'}), 404
d = request.get_json(force=True)
image_url = (d.get('image_url') or '').strip()
image_b64 = (d.get('image_base64') or '').strip()
question = (d.get('question') or '请描述这张图片的内容').strip()
if not image_url and not image_b64:
return jsonify({'ok': False, 'error': '请提供图片 URL 或 base64 数据'}), 400
try:
if image_b64:
r = llm_gateway.chat_vision(
w['provider'], w['model'], question, image_url=f'data:image/png;base64,{image_b64}',
temperature=0.3, max_tokens=w['max_tokens'] or 2000,
base_url=w['base_url'] or None, api_key=w['api_key'] or None)
else:
r = llm_gateway.chat_vision(
w['provider'], w['model'], question, image_url=image_url,
temperature=0.3, max_tokens=w['max_tokens'] or 2000,
base_url=w['base_url'] or None, api_key=w['api_key'] or None)
return jsonify({'ok': True, 'data': r})
except Exception as e:
return jsonify({'ok': False, 'error': str(e)})
@app.route('/api/providers')
@require_auth
def providers():
data = []
for k, v in config.PROVIDERS.items():
data.append({
'id': k, 'name': v['name'], 'base_url': v['base_url'],
'has_key': bool(v['api_key']),
'models': [m for m in config.MODEL_PRICING.keys() if m.startswith(
{'doubao': 'doubao', 'deepseek': 'deepseek', 'openai': 'gpt',
'qwen': 'qwen', 'vllm': '', 'autodl': 'qwen3'}.get(k, '__none__'))],
})
return jsonify({'ok': True, 'data': data})
# ---------------------------------------------------------------------------
# 任务
# ---------------------------------------------------------------------------
@app.route('/api/tasks', methods=['GET', 'POST'])
@require_auth
def tasks():
if request.method == 'POST':
d = request.get_json(force=True)
tid = db.w(
'INSERT INTO tasks (project_id, worker_id, title, description, priority, '
'review_required, deadline, depends_on, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?)',
(d.get('project_id'), d.get('worker_id'), d.get('title', '').strip(),
d.get('description', ''), d.get('priority', 'medium'),
1 if d.get('review_required', True) else 0, d.get('deadline', ''),
_json.dumps(d.get('depends_on') or []), db.now(), db.now()))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'info', f'任务创建:{d.get("title","")}', db.now()))
return jsonify({'ok': True, 'id': tid})
pid = request.args.get('project_id')
if pid:
rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id DESC', (int(pid),))
else:
rows = db.q('SELECT * FROM tasks WHERE deleted=0 ORDER BY id DESC LIMIT 200')
return jsonify({'ok': True, 'data': [db.serialize_task(t) for t in rows]})
@app.route('/api/tasks/<int:tid>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def task_detail(tid):
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
if request.method == 'GET':
t = db.serialize_task(t)
t['logs'] = db.q('SELECT * FROM task_logs WHERE task_id=? ORDER BY id', (tid,))
t['costs'] = db.q('SELECT * FROM cost_records WHERE task_id=? ORDER BY id', (tid,))
if t['worker_id']:
t['worker'] = db.q('SELECT id,name,provider,model FROM workers WHERE id=?',
(t['worker_id'],), one=True)
return jsonify({'ok': True, 'data': t})
if request.method == 'DELETE':
# 软删除(进回收站);?hard=1 物理删除
if request.args.get('hard') == '1':
db.w('DELETE FROM task_logs WHERE task_id=?', (tid,))
db.w('DELETE FROM cost_records WHERE task_id=?', (tid,))
db.w('DELETE FROM tasks WHERE id=?', (tid,))
return jsonify({'ok': True})
db.w('UPDATE tasks SET deleted=1, deleted_at=? WHERE id=?', (db.now(), tid))
return jsonify({'ok': True, 'trashed': True})
d = request.get_json(force=True)
# 只允许在非运行中修改基础字段
if t['status'] == 'running':
return jsonify({'ok': False, 'error': '任务执行中,禁止修改'}), 400
fields = ['title', 'description', 'worker_id', 'priority', 'review_required',
'deadline', 'status']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if 'depends_on' in d:
sets.append('depends_on=?')
args.append(_json.dumps(d['depends_on'] or []))
if sets:
args.append(db.now())
db.w(f'UPDATE tasks SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, tid))
return jsonify({'ok': True})
@app.route('/api/tasks/trash')
@require_auth
def tasks_trash():
"""回收站:已软删除的任务列表"""
pid = request.args.get('project_id')
if pid:
rows = db.q('SELECT * FROM tasks WHERE deleted=1 AND project_id=? ORDER BY deleted_at DESC',
(int(pid),))
else:
rows = db.q('SELECT * FROM tasks WHERE deleted=1 ORDER BY deleted_at DESC LIMIT 200')
return jsonify({'ok': True, 'data': [db.serialize_task(t) for t in rows]})
@app.route('/api/tasks/<int:tid>/trash', methods=['POST'])
@require_auth
def task_trash(tid):
"""软删除:移入回收站"""
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
if t['status'] == 'running':
return jsonify({'ok': False, 'error': '任务执行中,禁止删除'}), 400
db.w('UPDATE tasks SET deleted=1, deleted_at=? WHERE id=?', (db.now(), tid))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'warn', '🗑️ 任务已移入回收站', db.now()))
enterprise.audit(enterprise.current_actor(), 'task.trash', f'task#{tid}',
f'「{t["title"]}」移入回收站', request.remote_addr or '')
return jsonify({'ok': True})
@app.route('/api/tasks/<int:tid>/restore', methods=['POST'])
@require_auth
def task_restore(tid):
"""从回收站恢复任务"""
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
db.w('UPDATE tasks SET deleted=0, deleted_at=NULL WHERE id=?', (tid,))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'info', '♻️ 任务已从回收站恢复', db.now()))
enterprise.audit(enterprise.current_actor(), 'task.restore', f'task#{tid}',
f'「{t["title"]}」恢复', request.remote_addr or '')
return jsonify({'ok': True})
@app.route('/api/tasks/<int:tid>/run', methods=['POST'])
@require_auth
def task_run(tid):
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
if t['deleted']:
return jsonify({'ok': False, 'error': '任务已在回收站,请先恢复'}), 400
if t['status'] == 'running':
return jsonify({'ok': False, 'error': '任务已在执行中'}), 400
ok, blockers = engine.check_dependencies(t)
if not ok:
return jsonify({'ok': False, 'error': '前置任务未完成:' + '、'.join(blockers)}), 400
if engine.runner.submit(tid):
enterprise.audit(enterprise.current_actor(), 'task.run', f'task#{tid}',
f'「{t["title"]}」', request.remote_addr or '')
return jsonify({'ok': True})
return jsonify({'ok': False, 'error': '任务已在执行中'}), 400
@app.route('/api/tasks/<int:tid>/cancel', methods=['POST'])
@require_auth
def task_cancel(tid):
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
if t['status'] != 'running':
return jsonify({'ok': False, 'error': '仅执行中的任务可取消'}), 400
# MVP:标记取消(线程无法强杀,完成后会回到 review;这里直接置 cancelled 并忽略结果)
db.w('UPDATE tasks SET status="cancelled", updated_at=?, finished_at=? WHERE id=?',
(db.now(), db.now(), tid))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'warn', '任务已由人工取消(引擎线程将在后台自然结束)', db.now()))
return jsonify({'ok': True})
@app.route('/api/tasks/<int:tid>/review', methods=['POST'])
@require_auth
def task_review(tid):
t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
if not t:
return jsonify({'ok': False, 'error': '任务不存在'}), 404
if t['status'] != 'review':
return jsonify({'ok': False, 'error': '仅待审核状态可审核'}), 400
d = request.get_json(force=True)
action = d.get('action')
reason = (d.get('reason') or '').strip()
if action == 'approve':
db.w('UPDATE tasks SET status="done", updated_at=? WHERE id=?', (db.now(), tid))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'success', '✅ 人工审核通过,任务完成', db.now()))
enterprise.audit(enterprise.current_actor(), 'task.approve', f'task#{tid}',
f'「{t["title"]}」审核通过', request.remote_addr or '')
# DAG:审核放行等同完成,触发下游就绪任务
fresh = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True)
for t in engine._trigger_downstream(fresh):
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(t['id'], 'info', f'🔗 前置任务「{fresh["title"]}」已验收完成,自动触发执行', db.now()))
return jsonify({'ok': True})
if action == 'reject':
if not reason:
return jsonify({'ok': False, 'error': '打回必须填写原因'}), 400
db.w('UPDATE tasks SET status="todo", rejection_count=rejection_count+1, '
'output_text="", updated_at=? WHERE id=?', (db.now(), tid))
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'warn', f'⛔ 人工打回:{reason}(等待返工重跑)', db.now()))
enterprise.audit(enterprise.current_actor(), 'task.reject', f'task#{tid}',
f'「{t["title"]}」打回:{reason}', request.remote_addr or '')
return jsonify({'ok': True})
return jsonify({'ok': False, 'error': '未知操作'}), 400
# ---------------------------------------------------------------------------
# 报表 / 统计
# ---------------------------------------------------------------------------
@app.route('/api/reports/cost')
@require_auth
def report_cost():
group = request.args.get('group', 'project')
if group == 'worker':
rows = db.q(
'SELECT worker_id, provider, model, COUNT(*) runs, SUM(total_tokens) tokens, '
'SUM(cost) cost FROM cost_records GROUP BY worker_id ORDER BY cost DESC')
for r in rows:
w = db.q('SELECT name FROM workers WHERE id=?', (r['worker_id'],), one=True)
r['worker_name'] = w['name'] if w else f'#{r["worker_id"]}'
elif group == 'model':
rows = db.q(
'SELECT provider, model, COUNT(*) runs, SUM(total_tokens) tokens, '
'SUM(cost) cost FROM cost_records GROUP BY model ORDER BY cost DESC')
else:
rows = db.q(
'SELECT project_id, COUNT(*) runs, SUM(total_tokens) tokens, '
'SUM(cost) cost FROM cost_records GROUP BY project_id ORDER BY cost DESC')
for r in rows:
p = db.q('SELECT name FROM projects WHERE id=?', (r['project_id'],), one=True)
r['project_name'] = p['name'] if p else f'#{r["project_id"]}'
return jsonify({'ok': True, 'data': rows})
@app.route('/api/stats')
@require_auth
def stats():
out = {'projects': 0, 'tasks': 0, 'workers': 0, 'total_cost': 0, 'total_tokens': 0,
'by_status': {}, 'recent': [], 'daily_cost': []}
out['projects'] = db.q('SELECT COUNT(*) c FROM projects')[0]['c']
out['workers'] = db.q('SELECT COUNT(*) c FROM workers')[0]['c']
out['tasks'] = db.q('SELECT COUNT(*) c FROM tasks WHERE deleted=0')[0]['c']
for r in db.q('SELECT status, COUNT(*) c FROM tasks WHERE deleted=0 GROUP BY status'):
out['by_status'][r['status']] = r['c']
c = db.q('SELECT COALESCE(SUM(cost),0) cost, COALESCE(SUM(total_tokens),0) tokens '
'FROM cost_records')[0]
out['total_cost'], out['total_tokens'] = round(c['cost'], 4), c['tokens']
# 一次通过率:done 且 rejection_count=0
done = out['by_status'].get('done', 0)
clean = db.q('SELECT COUNT(*) c FROM tasks WHERE status="done" AND rejection_count=0 AND deleted=0')[0]['c']
out['one_pass_rate'] = round(clean / done * 100, 1) if done else None
out['rework_count'] = db.q('SELECT COALESCE(SUM(rejection_count),0) c FROM tasks WHERE deleted=0')[0]['c']
out['recent'] = db.q(
'SELECT t.id, t.title, t.status, t.project_id, p.name AS project_name, t.updated_at '
'FROM tasks t LEFT JOIN projects p ON p.id=t.project_id '
'WHERE t.deleted=0 ORDER BY t.updated_at DESC LIMIT 10')
# 近 7 天成本
import datetime
rows = db.q('SELECT created_at, cost FROM cost_records ORDER BY created_at')
day_map = {}
for r in rows:
d = datetime.datetime.fromtimestamp(r['created_at']).strftime('%m-%d')
day_map[d] = round(day_map.get(d, 0) + r['cost'], 4)
for i in range(6, -1, -1):
d = (datetime.datetime.now() - datetime.timedelta(days=i)).strftime('%m-%d')
out['daily_cost'].append({'day': d, 'cost': day_map.get(d, 0)})
return jsonify({'ok': True, 'data': out})
@app.route('/api/logs')
@require_auth
def logs():
limit = min(int(request.args.get('limit', 100)), 500)
rows = db.q(
'SELECT l.*, t.title AS task_title FROM task_logs l '
'LEFT JOIN tasks t ON t.id=l.task_id ORDER BY l.id DESC LIMIT ?', (limit,))
return jsonify({'ok': True, 'data': rows})
@app.route('/api/projects/<int:pid>/workflow/run', methods=['POST'])
@require_auth
def workflow_run(pid):
"""执行整个工作流:跑所有就绪(无未完成前置)任务"""
rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 AND status IN ("todo","failed")', (pid,))
started, blocked = [], []
for t in rows:
ok, blockers = engine.check_dependencies(t)
if ok:
if engine.runner.submit(t['id']):
started.append({'id': t['id'], 'title': t['title']})
else:
blocked.append({'id': t['id'], 'title': t['title'], 'by': blockers})
return jsonify({'ok': True, 'started': started, 'blocked': blocked})
@app.route('/api/projects/<int:pid>/dag')
@require_auth
def project_dag(pid):
"""DAG 图数据:节点 + 边"""
rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id', (pid,))
nodes, edges, id_map = [], [], {}
for t in rows:
node = db.serialize_task(t)
nodes.append(node)
id_map[t['id']] = node
for t in nodes:
for dep_id in t['depends_on']:
if dep_id in id_map:
edges.append({'from': dep_id, 'to': t['id']})
return jsonify({'ok': True, 'nodes': nodes, 'edges': edges})
# ---------------------------------------------------------------------------
# AI 辅助规划(WBS 生成 + 导入)
# ---------------------------------------------------------------------------
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}'
)
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
@app.route('/api/projects/<int:pid>/wbs/generate', methods=['POST'])
@require_auth
def wbs_generate(pid):
d = request.get_json(force=True) or {}
goal = d.get('goal') or ''
proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True)
if not proj:
return jsonify({'ok': False, 'error': '项目不存在'}), 404
if not goal:
goal = proj.get('objective') or proj.get('name') or ''
if not goal:
return jsonify({'ok': False, 'error': '请提供项目目标'}), 400
worker = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id', one=True)
if not worker:
return jsonify({'ok': False, 'error': '请先注册至少一个 Worker 用于规划'}), 400
try:
r = llm_gateway.chat(worker['provider'], worker['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)
tasks = _parse_wbs(r['text'])
return jsonify({'ok': True, 'data': tasks, 'usage': r['total_tokens'], 'cost': r['cost']})
except Exception as e:
return jsonify({'ok': False, 'error': f'WBS 生成失败:{e}'}), 500
@app.route('/api/projects/<int:pid>/wbs/import', methods=['POST'])
@require_auth
def wbs_import(pid):
d = request.get_json(force=True)
tasks = d.get('tasks') or []
worker_id = d.get('worker_id')
if not tasks:
return jsonify({'ok': False, 'error': '任务列表为空'}), 400
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)]
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, t.get('title', f'任务{i+1}'), t.get('description', ''),
t.get('priority', 'medium'), 1, '', _json.dumps(dep_ids), db.now(), db.now()))
created.append(tid)
db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)',
(tid, 'info', f'AI 规划导入:{t.get("title", "")}', db.now()))
return jsonify({'ok': True, 'created': len(created), 'ids': created})
# ---------------------------------------------------------------------------
# RAG 知识库
# ---------------------------------------------------------------------------
@app.route('/api/projects/<int:pid>/documents', methods=['GET', 'POST'])
@require_auth
def documents(pid):
if request.method == 'POST':
d = request.get_json(force=True)
doc_id = db.w(
'INSERT INTO documents (project_id, name, content, source, created_at, updated_at) '
'VALUES (?,?,?,?,?,?)',
(pid, d.get('name', '未命名文档').strip(), d.get('content', ''),
d.get('source', 'manual'), db.now(), db.now()))
n = rag.rebuild_document(doc_id)
return jsonify({'ok': True, 'id': doc_id, 'chunks': n})
rows = db.q('SELECT * FROM documents WHERE project_id=? ORDER BY id DESC', (pid,))
for r in rows:
r['chunks'] = db.q('SELECT COUNT(*) c FROM doc_chunks WHERE document_id=?', (r['id'],))[0]['c']
return jsonify({'ok': True, 'data': rows})
@app.route('/api/documents/<int:doc_id>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def document_detail(doc_id):
doc = db.q('SELECT * FROM documents WHERE id=?', (doc_id,), one=True)
if not doc:
return jsonify({'ok': False, 'error': '文档不存在'}), 404
if request.method == 'GET':
return jsonify({'ok': True, 'data': doc})
if request.method == 'DELETE':
db.w('DELETE FROM doc_chunks WHERE document_id=?', (doc_id,))
db.w('DELETE FROM documents WHERE id=?', (doc_id,))
return jsonify({'ok': True})
d = request.get_json(force=True)
if 'content' in d:
db.w('UPDATE documents SET content=?, updated_at=? WHERE id=?', (d['content'], db.now(), doc_id))
rag.rebuild_document(doc_id)
if 'name' in d:
db.w('UPDATE documents SET name=? WHERE id=?', (d['name'].strip(), doc_id))
return jsonify({'ok': True})
@app.route('/api/projects/<int:pid>/search')
@require_auth
def kb_search(pid):
q = request.args.get('q', '')
if not q:
return jsonify({'ok': True, 'data': []})
hits, hit = rag.search_project(pid, q, top_k=5)
return jsonify({'ok': True, 'data': hits if hit else [], 'hit': hit})
# ---------------------------------------------------------------------------
# 告警中心
# ---------------------------------------------------------------------------
@app.route('/api/alerts')
@require_auth
def alerts():
limit = min(int(request.args.get('limit', 100)), 500)
rows = db.q('SELECT * FROM alerts ORDER BY id DESC LIMIT ?', (limit,))
return jsonify({'ok': True, 'data': rows})
@app.route('/api/alerts/unread_count')
@require_auth
def alerts_unread():
c = db.q('SELECT COUNT(*) c FROM alerts WHERE read=0')[0]['c']
return jsonify({'ok': True, 'count': c})
@app.route('/api/alerts/<int:aid>/read', methods=['POST'])
@require_auth
def alert_read(aid):
db.w('UPDATE alerts SET read=1 WHERE id=?', (aid,))
return jsonify({'ok': True})
@app.route('/api/alerts/read_all', methods=['POST'])
@require_auth
def alerts_read_all():
db.w('UPDATE alerts SET read=1 WHERE read=0')
return jsonify({'ok': True})
# ---------------------------------------------------------------------------
# 开放 API Token
# ---------------------------------------------------------------------------
@app.route('/api/tokens', methods=['GET', 'POST'])
@require_auth
def api_tokens():
if request.method == 'POST':
d = request.get_json(force=True)
tok = secrets.token_hex(24)
db.w('INSERT INTO api_tokens (name, token, created_at) VALUES (?,?,?)',
(d.get('name', '未命名').strip(), tok, db.now()))
return jsonify({'ok': True, 'token': tok})
rows = db.q('SELECT id, name, created_at, last_used_at FROM api_tokens ORDER BY id DESC')
return jsonify({'ok': True, 'data': rows})
@app.route('/api/tokens/<int:tid>', methods=['DELETE'])
@require_auth
def api_token_delete(tid):
db.w('DELETE FROM api_tokens WHERE id=?', (tid,))
return jsonify({'ok': True})
# ---------------------------------------------------------------------------
# 通知渠道(飞书/企微/邮件)
# ---------------------------------------------------------------------------
@app.route('/api/channels', methods=['GET', 'POST'])
@require_auth
def channels():
if request.method == 'POST':
d = request.get_json(force=True)
cid = db.w(
'INSERT INTO notify_channels (name, type, webhook, email, events, enabled, created_at) '
'VALUES (?,?,?,?,?,?,?)',
(d.get('name', '').strip(), d.get('type', 'feishu'), d.get('webhook', ''),
d.get('email', ''), _json.dumps(d.get('events') or []),
1 if d.get('enabled', True) else 0, db.now()))
return jsonify({'ok': True, 'id': cid})
rows = db.q('SELECT * FROM notify_channels ORDER BY id DESC')
for r in rows:
try:
r['events'] = _json.loads(r['events'] or '[]')
except Exception:
r['events'] = []
return jsonify({'ok': True, 'data': rows})
@app.route('/api/channels/<int:cid>', methods=['PUT', 'DELETE'])
@require_auth
def channel_detail(cid):
if request.method == 'DELETE':
db.w('DELETE FROM notify_channels WHERE id=?', (cid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['name', 'type', 'webhook', 'email', 'enabled']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if 'events' in d:
sets.append('events=?')
args.append(_json.dumps(d['events'] or []))
if sets:
db.w(f'UPDATE notify_channels SET {", ".join(sets)} WHERE id=?', (*args, cid))
return jsonify({'ok': True})
@app.route('/api/channels/<int:cid>/test', methods=['POST'])
@require_auth
def channel_test(cid):
ch = db.q('SELECT * FROM notify_channels WHERE id=?', (cid,), one=True)
if not ch:
return jsonify({'ok': False, 'error': '渠道不存在'}), 404
ok, msg = notify.test_channel(ch)
return jsonify({'ok': ok, 'msg': msg})
@app.route('/api/events')
@require_auth
def events():
return jsonify({'ok': True, 'data': [{'id': k, 'name': v} for k, v in notify.EVENTS.items()]})
# ---------------------------------------------------------------------------
# 设置
# ---------------------------------------------------------------------------
@app.route('/api/settings', methods=['GET', 'PUT'])
@require_auth
def settings():
if request.method == 'PUT':
d = request.get_json(force=True)
for k, v in d.items():
db.set_setting(k, v)
if 'budget_alert_ratio' in d:
config.BUDGET_ALERT_RATIO = float(d['budget_alert_ratio'])
return jsonify({'ok': True})
return jsonify({'ok': True, 'data': {
'budget_alert_ratio': config.BUDGET_ALERT_RATIO,
'auth_enabled': auth_enabled(),
'email_configured': bool(config.EMAIL.get('host')),
}})
# ---------------------------------------------------------------------------
# V2 · 多 Agent 协作
# ---------------------------------------------------------------------------
AGENT_MODES = [
{'id': 'supervisor', 'name': '主管模式', 'icon': '👔',
'desc': '主管拆解任务 → 多个 Worker 并行执行 → 汇总合成最终交付物',
'params': [{'key': 'workers', 'label': '执行 Agent(多选)', 'type': 'multi_worker'},
{'key': 'context', 'label': '背景上下文', 'type': 'textarea'}]},
{'id': 'review', 'name': '评审模式', 'icon': '🔍',
'desc': 'Worker 产出 → Reviewer 打分评审 → 未达标自动返工(最多 N 轮)→ 终稿',
'params': [{'key': 'workers', 'label': '执行 Agent(第1个)', 'type': 'multi_worker'},
{'key': 'rounds', 'label': '最大返工轮次', 'type': 'number', 'default': 2},
{'key': 'pass_score', 'label': '通过分数线(0-100', 'type': 'number', 'default': 70},
{'key': 'context', 'label': '背景上下文', 'type': 'textarea'}]},
{'id': 'debate', 'name': '辩论模式', 'icon': '⚖️',
'desc': '多辩手各自立论 → 互相质询反驳 → 裁判综合裁决输出共识结论',
'params': [{'key': 'workers', 'label': '辩手(最后1个为裁判)', 'type': 'multi_worker'},
{'key': 'rounds', 'label': '质询轮次', 'type': 'number', 'default': 1},
{'key': 'stances', 'label': '立场(逗号分隔,如:正方,反方)', 'type': 'text'},
{'key': 'context', 'label': '背景上下文', 'type': 'textarea'}]},
]
@app.route('/api/agents/modes')
@require_auth
def agent_modes():
return jsonify({'ok': True, 'data': AGENT_MODES})
@app.route('/api/agents/runs', methods=['GET', 'POST'])
@require_auth
def agent_runs():
if request.method == 'POST':
d = request.get_json(force=True)
mode = d.get('mode')
if mode not in ('supervisor', 'review', 'debate'):
return jsonify({'ok': False, 'error': '未知协作模式'}), 400
topic = (d.get('topic') or '').strip()
if not topic:
return jsonify({'ok': False, 'error': '请填写任务主题'}), 400
params = dict(d.get('params') or {})
# 从 params 里拆出 workers 列表与 context
worker_ids = params.pop('workers', None) or d.get('worker_ids') or []
context = params.pop('context', '') or ''
if not worker_ids:
rows = db.q('SELECT id FROM workers WHERE status="enabled" ORDER BY id LIMIT 3')
worker_ids = [r['id'] for r in rows]
if not worker_ids:
return jsonify({'ok': False, 'error': '请先注册至少一个 Worker'}), 400
rid = db.w(
'INSERT INTO agent_runs (mode, title, topic, context, worker_ids, params, status, '
'created_at) VALUES (?,?,?,?,?,?,?,?)',
(mode, d.get('title') or f'{mode} · {topic[:40]}', topic, context,
_json.dumps(worker_ids), _json.dumps(params), 'running', db.now()))
agents.runner.submit(rid)
enterprise.audit(enterprise.current_actor(), 'agent.run', f'run#{rid}',
f'mode={mode} workers={worker_ids}', request.remote_addr or '')
return jsonify({'ok': True, 'id': rid})
rows = db.q('SELECT * FROM agent_runs ORDER BY id DESC LIMIT 100')
return jsonify({'ok': True, 'data': rows})
@app.route('/api/agents/runs/<int:rid>', methods=['GET', 'DELETE'])
@require_auth
def agent_run_detail(rid):
run = db.q('SELECT * FROM agent_runs WHERE id=?', (rid,), one=True)
if not run:
return jsonify({'ok': False, 'error': '不存在'}), 404
if request.method == 'DELETE':
db.w('DELETE FROM agent_steps WHERE run_id=?', (rid,))
db.w('DELETE FROM agent_runs WHERE id=?', (rid,))
return jsonify({'ok': True})
run['steps'] = db.q('SELECT * FROM agent_steps WHERE run_id=? ORDER BY seq', (rid,))
try:
run['worker_ids'] = _json.loads(run['worker_ids'] or '[]')
run['params'] = _json.loads(run['params'] or '{}')
except Exception:
pass
wmap = {w['id']: w['name'] for w in db.q('SELECT id,name FROM workers')}
for s in run['steps']:
s['worker_name'] = wmap.get(s['worker_id'], '—')
return jsonify({'ok': True, 'data': run})
# ---------------------------------------------------------------------------
# V2 · 自动评估体系
# ---------------------------------------------------------------------------
@app.route('/api/eval/datasets', methods=['GET', 'POST'])
@require_auth
def eval_datasets():
if request.method == 'POST':
d = request.get_json(force=True)
did = db.w(
'INSERT INTO eval_datasets (name, description, rubric, tags, is_builtin, created_at, updated_at) '
'VALUES (?,?,?,?,0,?,?)',
(d.get('name', '').strip(), d.get('description', ''), d.get('rubric', ''),
_json.dumps(d.get('tags') or []), db.now(), db.now()))
for c in (d.get('cases') or []):
db.w('INSERT INTO eval_cases (dataset_id, input, expected, tags, created_at) VALUES (?,?,?,?,?)',
(did, c.get('input', ''), c.get('expected', ''), _json.dumps(d.get('tags') or []), db.now()))
return jsonify({'ok': True, 'id': did})
rows = db.q('SELECT * FROM eval_datasets ORDER BY id DESC')
for r in rows:
r['case_count'] = db.q('SELECT COUNT(*) c FROM eval_cases WHERE dataset_id=?', (r['id'],))[0]['c']
r['best_score'] = db.q('SELECT MAX(score) s FROM eval_runs WHERE dataset_id=? AND status="done"',
(r['id'],))[0]['s']
r['tags'] = _json.loads(r.get('tags') or '[]')
return jsonify({'ok': True, 'data': rows})
@app.route('/api/eval/datasets/<int:did>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def eval_dataset_detail(did):
ds = db.q('SELECT * FROM eval_datasets WHERE id=?', (did,), one=True)
if not ds:
return jsonify({'ok': False, 'error': '数据集不存在'}), 404
if request.method == 'GET':
ds['cases'] = db.q('SELECT * FROM eval_cases WHERE dataset_id=? ORDER BY id', (did,))
ds['tags'] = _json.loads(ds.get('tags') or '[]')
return jsonify({'ok': True, 'data': ds})
if request.method == 'DELETE':
db.w('DELETE FROM eval_results WHERE case_id IN (SELECT id FROM eval_cases WHERE dataset_id=?)', (did,))
db.w('DELETE FROM eval_cases WHERE dataset_id=?', (did,))
db.w('DELETE FROM eval_runs WHERE dataset_id=?', (did,))
db.w('DELETE FROM eval_datasets WHERE id=?', (did,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['name', 'description', 'rubric']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if 'tags' in d:
sets.append('tags=?')
args.append(_json.dumps(d['tags']))
if sets:
args.append(db.now())
db.w(f'UPDATE eval_datasets SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, did))
return jsonify({'ok': True})
@app.route('/api/eval/datasets/<int:did>/cases', methods=['POST'])
@require_auth
def eval_case_create(did):
d = request.get_json(force=True)
cid = db.w('INSERT INTO eval_cases (dataset_id, input, expected, tags, created_at) VALUES (?,?,?,?,?)',
(did, d.get('input', ''), d.get('expected', ''), _json.dumps(d.get('tags') or []), db.now()))
return jsonify({'ok': True, 'id': cid})
@app.route('/api/eval/cases/<int:cid>', methods=['PUT', 'DELETE'])
@require_auth
def eval_case_detail(cid):
case = db.q('SELECT * FROM eval_cases WHERE id=?', (cid,), one=True)
if not case:
return jsonify({'ok': False, 'error': '用例不存在'}), 404
if request.method == 'DELETE':
db.w('DELETE FROM eval_results WHERE case_id=?', (cid,))
db.w('DELETE FROM eval_cases WHERE id=?', (cid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['input', 'expected', 'tags']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(_json.dumps(d[f]) if f == 'tags' else d[f])
if sets:
db.w(f'UPDATE eval_cases SET {", ".join(sets)} WHERE id=?', (*args, cid))
return jsonify({'ok': True})
@app.route('/api/eval/datasets/<int:did>/sink', methods=['POST'])
@require_auth
def eval_sink(did):
"""把已验收任务产出沉淀为该数据集用例"""
n = evalmod.sink_done_tasks()
return jsonify({'ok': True, 'sunk': n})
@app.route('/api/eval/runs', methods=['GET', 'POST'])
@require_auth
def eval_runs():
if request.method == 'POST':
d = request.get_json(force=True)
dataset_id = d.get('dataset_id')
worker_id = d.get('worker_id')
if not dataset_id or not worker_id:
return jsonify({'ok': False, 'error': '请选择数据集与 Worker'}), 400
rid = db.w(
'INSERT INTO eval_runs (dataset_id, worker_id, status, cases_total, created_at) '
'VALUES (?,?,?,?,?)',
(dataset_id, worker_id, 'running',
db.q('SELECT COUNT(*) c FROM eval_cases WHERE dataset_id=?', (dataset_id,))[0]['c'], db.now()))
evalmod.runner.submit(rid)
enterprise.audit(enterprise.current_actor(), 'eval.run', f'run#{rid}',
f'dataset={dataset_id} worker={worker_id}', request.remote_addr or '')
return jsonify({'ok': True, 'id': rid})
rows = db.q('SELECT * FROM eval_runs ORDER BY id DESC LIMIT 100')
wmap = {w['id']: w['name'] for w in db.q('SELECT id,name FROM workers')}
dmap = {d['id']: d['name'] for d in db.q('SELECT id,name FROM eval_datasets')}
for r in rows:
r['worker_name'] = wmap.get(r['worker_id'], f'#{r["worker_id"]}')
r['dataset_name'] = dmap.get(r['dataset_id'], f'#{r["dataset_id"]}')
return jsonify({'ok': True, 'data': rows})
@app.route('/api/eval/runs/<int:rid>', methods=['GET', 'DELETE'])
@require_auth
def eval_run_detail(rid):
run = db.q('SELECT * FROM eval_runs WHERE id=?', (rid,), one=True)
if not run:
return jsonify({'ok': False, 'error': '不存在'}), 404
if request.method == 'DELETE':
db.w('DELETE FROM eval_results WHERE run_id=?', (rid,))
db.w('DELETE FROM eval_runs WHERE id=?', (rid,))
return jsonify({'ok': True})
run['results'] = db.q('SELECT * FROM eval_results WHERE run_id=? ORDER BY id', (rid,))
cid_map = {c['id']: c for c in db.q('SELECT * FROM eval_cases WHERE dataset_id=?', (run['dataset_id'],))}
for r in run['results']:
c = cid_map.get(r['case_id'])
r['input'] = c['input'][:200] if c else ''
r['expected'] = (c['expected'] or '')[:200] if c else ''
return jsonify({'ok': True, 'data': run})
@app.route('/api/eval/leaderboard')
@require_auth
def eval_leaderboard():
"""Worker 评估排行榜:每个 worker 在完成 run 上的平均分"""
rows = db.q(
'SELECT worker_id, COUNT(*) runs, ROUND(AVG(score),2) avg_score, '
'ROUND(SUM(cost),4) cost, SUM(total_tokens) tokens '
'FROM eval_runs WHERE status="done" GROUP BY worker_id ORDER BY avg_score DESC')
wmap = {w['id']: w for w in db.q('SELECT id,name,provider,model FROM workers')}
for r in rows:
w = wmap.get(r['worker_id'])
r['worker_name'] = w['name'] if w else f'#{r["worker_id"]}'
r['model'] = f"{w['provider']}/{w['model']}" if w else ''
return jsonify({'ok': True, 'data': rows})
# ---------------------------------------------------------------------------
# V2 · 模板市场
# ---------------------------------------------------------------------------
@app.route('/api/templates', methods=['GET', 'POST'])
@require_auth
def templates_api():
if request.method == 'POST':
d = request.get_json(force=True)
tid = db.w(
'INSERT INTO templates (type, name, description, content, tags, author, is_builtin, '
'usage_count, created_at, updated_at) VALUES (?,?,?,?,?,?,0,0,?,?)',
(d.get('type'), d.get('name', '').strip(), d.get('description', ''),
_json.dumps(d.get('content') or {}), _json.dumps(d.get('tags') or []),
enterprise.current_actor(), db.now(), db.now()))
return jsonify({'ok': True, 'id': tid})
ttype = request.args.get('type', '')
if ttype:
rows = db.q('SELECT * FROM templates WHERE type=? ORDER BY usage_count DESC, id DESC', (ttype,))
else:
rows = db.q('SELECT * FROM templates ORDER BY type, usage_count DESC, id DESC')
for r in rows:
try:
r['content'] = _json.loads(r['content'] or '{}')
r['tags'] = _json.loads(r.get('tags') or '[]')
except Exception:
pass
return jsonify({'ok': True, 'data': rows})
@app.route('/api/templates/<int:tid>', methods=['GET', 'PUT', 'DELETE'])
@require_auth
def template_detail(tid):
t = tplmod.get_template(tid)
if not t:
return jsonify({'ok': False, 'error': '模板不存在'}), 404
if request.method == 'GET':
try:
t['content'] = _json.loads(t['content'] or '{}')
t['tags'] = _json.loads(t.get('tags') or '[]')
except Exception:
pass
return jsonify({'ok': True, 'data': t})
if request.method == 'DELETE':
if t['is_builtin']:
return jsonify({'ok': False, 'error': '内置模板不可删除'}), 400
db.w('DELETE FROM templates WHERE id=?', (tid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
fields = ['name', 'description', 'type']
sets, args = [], []
for f in fields:
if f in d:
sets.append(f'{f}=?')
args.append(d[f])
if 'content' in d:
sets.append('content=?')
args.append(_json.dumps(d['content']))
if 'tags' in d:
sets.append('tags=?')
args.append(_json.dumps(d['tags']))
if sets:
args.append(db.now())
db.w(f'UPDATE templates SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, tid))
return jsonify({'ok': True})
@app.route('/api/templates/<int:tid>/apply', methods=['POST'])
@require_auth
def template_apply(tid):
t = tplmod.get_template(tid)
if not t:
return jsonify({'ok': False, 'error': '模板不存在'}), 404
d = request.get_json(force=True)
variables = d.get('variables') or {}
worker_id = d.get('worker_id')
try:
if t['type'] == 'task':
pid = d.get('project_id')
if not pid:
return jsonify({'ok': False, 'error': '请选择目标项目'}), 400
new_id = tplmod.apply_task_template(t, variables, pid, worker_id)
return jsonify({'ok': True, 'kind': 'task', 'id': new_id})
if t['type'] == 'project':
pid, ids = tplmod.apply_project_template(t, variables, worker_id)
return jsonify({'ok': True, 'kind': 'project', 'id': pid, 'task_ids': ids})
if t['type'] == 'team':
ids = tplmod.apply_team_template(t, variables, provider=d.get('provider', 'deepseek'),
model=d.get('model', 'deepseek-v4-flash'))
return jsonify({'ok': True, 'kind': 'team', 'ids': ids})
return jsonify({'ok': False, 'error': '未知模板类型'}), 400
except Exception as e:
return jsonify({'ok': False, 'error': f'应用模板失败: {e}'}), 500
@app.route('/api/templates/builtin/restore', methods=['POST'])
@require_admin
def templates_restore():
n = tplmod.create_builtin_templates()
return jsonify({'ok': True, 'note': '内置模板已就绪' if n is None else '已存在,跳过'})
# ---------------------------------------------------------------------------
# V2 · 企业版:用户 / SSO / 审计 / 合规
# ---------------------------------------------------------------------------
@app.route('/api/enterprise/users', methods=['GET', 'POST'])
@require_admin
def ent_users():
if request.method == 'POST':
d = request.get_json(force=True)
username = (d.get('username') or '').strip()
if not username:
return jsonify({'ok': False, 'error': '用户名不能为空'}), 400
if d.get('source') in ('oidc', 'ldap'):
if enterprise.get_user(username):
return jsonify({'ok': False, 'error': '用户已存在'}), 400
db.w('INSERT INTO users (username, password_hash, display_name, role, source, status, created_at) '
'VALUES (?,?,?,?,?,?,?)',
(username, '', d.get('display_name', username), d.get('role', 'member'),
d.get('source'), 'active', db.now()))
else:
if not d.get('password'):
return jsonify({'ok': False, 'error': '本地用户必须设置密码'}), 400
uid, err = enterprise.create_local_user(username, d['password'],
d.get('display_name', ''), d.get('role', 'member'))
if err:
return jsonify({'ok': False, 'error': err}), 400
return jsonify({'ok': True})
rows = db.q('SELECT id, username, display_name, role, source, status, last_login_at, created_at '
'FROM users ORDER BY id')
return jsonify({'ok': True, 'data': rows})
@app.route('/api/enterprise/users/<int:uid>', methods=['PUT', 'DELETE'])
@require_admin
def ent_user_detail(uid):
u = db.q('SELECT * FROM users WHERE id=?', (uid,), one=True)
if not u:
return jsonify({'ok': False, 'error': '用户不存在'}), 404
if request.method == 'DELETE':
if u['username'] == session.get('username'):
return jsonify({'ok': False, 'error': '不能删除自己'}), 400
if u['username'] == 'admin' and u['source'] == 'local':
return jsonify({'ok': False, 'error': '内置管理员不可删除'}), 400
db.w('DELETE FROM users WHERE id=?', (uid,))
return jsonify({'ok': True})
d = request.get_json(force=True)
if 'role' in d and d['role'] in ('admin', 'member', 'auditor'):
db.w('UPDATE users SET role=? WHERE id=?', (d['role'], uid))
if 'status' in d and d['status'] in ('active', 'disabled'):
db.w('UPDATE users SET status=? WHERE id=?', (d['status'], uid))
if 'display_name' in d:
db.w('UPDATE users SET display_name=? WHERE id=?', (d['display_name'], uid))
if d.get('password'):
db.w('UPDATE users SET password_hash=? WHERE id=?', (db.hash_password(d['password']), uid))
return jsonify({'ok': True})
@app.route('/api/enterprise/audit')
@require_auth
def ent_audit():
limit = min(int(request.args.get('limit', 200)), 1000)
rows = db.q('SELECT * FROM audit_logs ORDER BY id DESC LIMIT ?', (limit,))
return jsonify({'ok': True, 'data': rows})
@app.route('/api/enterprise/audit/export')
@require_admin
def ent_audit_export():
from flask import Response
csv_text = enterprise.export_audit_csv()
enterprise.audit(enterprise.current_actor(), 'compliance.audit_export', '导出审计日志',
ip=request.remote_addr or '')
return Response(csv_text, mimetype='text/csv',
headers={'Content-Disposition': 'attachment; filename=audit_logs.csv'})
@app.route('/api/enterprise/sso', methods=['GET', 'PUT'])
@require_admin
def ent_sso():
if request.method == 'PUT':
d = request.get_json(force=True)
enterprise.save_sso_config(d)
return jsonify({'ok': True})
cfg = enterprise.sso_config()
# 密钥打码返回
for k in ('oidc_client_secret', 'ldap_bind_password'):
if cfg.get(k):
cfg[k] = '********'
return jsonify({'ok': True, 'data': cfg})
@app.route('/api/enterprise/sso/oidc/login', methods=['POST'])
def ent_oidc_login():
"""发起 OIDC 授权(未登录也允许,SSO 登录页用)"""
try:
state = secrets.token_hex(16)
session['oidc_state'] = state
url = enterprise.oidc_authorize_url(state)
return jsonify({'ok': True, 'redirect': url})
except Exception as e:
return jsonify({'ok': False, 'error': str(e)}), 500
@app.route('/api/enterprise/sso/oidc/callback')
def ent_oidc_callback():
code = request.args.get('code')
state = request.args.get('state')
if not code or state != session.get('oidc_state'):
return 'SSO 回调校验失败(state 不匹配或缺少 code', 400
try:
userinfo = enterprise.oidc_exchange(code)
u, err = enterprise.sso_login(userinfo)
if err:
return f'SSO 登录失败: {err}', 403
session['username'] = u['username']
session['role'] = u['role']
session['authed'] = True
enterprise.audit(u['username'], 'auth.sso_login', 'OIDC SSO 登录', ip=request.remote_addr or '')
return redirect('/#/settings')
except Exception as e:
return f'SSO 登录异常: {str(e)[:200]}', 500
@app.route('/api/enterprise/sso/ldap/test', methods=['POST'])
@require_admin
def ent_ldap_test():
d = request.get_json(force=True)
username = d.get('username', '')
password = d.get('password', '')
if not username or not password:
return jsonify({'ok': False, 'error': '请填写测试账号密码'}), 400
info, err = enterprise.ldap_authenticate(username, password)
if err:
return jsonify({'ok': False, 'error': err}), 400
return jsonify({'ok': True, 'data': info})
@app.route('/api/enterprise/export')
@require_admin
def ent_export():
"""合规:全量数据导出(JSON"""
from flask import Response
data = enterprise.export_all()
enterprise.audit(enterprise.current_actor(), 'compliance.export', '全量数据导出',
ip=request.remote_addr or '')
return Response(_json.dumps(data, ensure_ascii=False, indent=2), mimetype='application/json',
headers={'Content-Disposition': 'attachment; filename=aiworker_export.json'})
@app.route('/api/enterprise/compliance', methods=['GET', 'POST'])
@require_admin
def ent_compliance():
if request.method == 'POST':
d = request.get_json(force=True)
if 'retention_days' in d:
db.set_setting('compliance_retention_days', int(d['retention_days']))
if 'mask_pii' in d:
db.set_setting('compliance_mask_pii', '1' if d['mask_pii'] else '0')
if d.get('apply_retention'):
enterprise.apply_retention()
return jsonify({'ok': True})
return jsonify({'ok': True, 'data': {
'retention_days': int(db.get_setting('compliance_retention_days', '0') or 0),
'mask_pii': db.get_setting('compliance_mask_pii', '0') == '1',
'audit_count': db.q('SELECT COUNT(*) c FROM audit_logs')[0]['c'],
'agent_runs_count': db.q('SELECT COUNT(*) c FROM agent_runs')[0]['c'],
'eval_runs_count': db.q('SELECT COUNT(*) c FROM eval_runs')[0]['c'],
'consent_token': enterprise.generate_consent_token(),
}})
# ---------------------------------------------------------------------------
# 前端
# ---------------------------------------------------------------------------
@app.route('/')
def index():
return send_from_directory(app.static_folder, 'index.html')
if __name__ == '__main__':
print(f'AI Worker 项目管理平台启动: http://0.0.0.0:{config.PORT}')
app.run(host=config.HOST, port=config.PORT, debug=False, threaded=True)