- 全部 AI Worker 切换到 deepseek/deepseek-v4-flash(api.deepseek.com,实测757ms/次,成本降40倍) - 团队模板默认 provider/model 同步改为 deepseek/deepseek-v4-flash - 新增视觉智能体「视觉分析师」:autodl/qwen3.6-plus 多模态 - llm_gateway: chat_vision 多模态调用(URL/base64) + 空content自动重试 + 聚合后端路由容错 - engine: 任务描述支持  图片注入(视觉任务直接派活) - API: POST /api/workers/<id>/vision_test;前端 Worker 页新增🖼️视觉测试按钮 - 实测: 截图结构化分析 ✅ 雪羊图片识别 ✅ 辩论模式DeepSeek 55秒完成
1347 lines
57 KiB
Python
1347 lines
57 KiB
Python
# -*- 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=?', (r['id'],))[0]['c']
|
||
r['done_count'] = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND status="done"',
|
||
(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=? ORDER BY id DESC', (int(pid),))
|
||
else:
|
||
rows = db.q('SELECT * FROM tasks 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':
|
||
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})
|
||
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/<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['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')[0]['c']
|
||
for r in db.q('SELECT status, COUNT(*) c FROM tasks 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')[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')[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 '
|
||
'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 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=? 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)
|