diff --git a/README.md b/README.md index cc499a6..2a76265 100644 --- a/README.md +++ b/README.md @@ -1,22 +1,33 @@ -# AI Worker 项目管理平台(MVP) +# AI Worker 项目管理平台 > 以「项目」为中心、以「AI Worker」为执行单元的项目管理平台。 > 把大模型团队变成一支可指挥、可审计、可控成本的"虚拟团队"。 -**MVP 目标**:跑通「人派活 → AI 干活 → 人验收」最小闭环。 +**当前版本:V1**(MVP 闭环 + 多 Agent 编排 / RAG / 预算告警 / 开放 API / 飞书企微集成) -## 功能清单(MVP) +--- +## 功能总览 + +### MVP(已交付 v1.0) | 模块 | 能力 | |---|---| -| 📁 项目管理 | 项目 CRUD、目标/验收标准/预算上限、状态流转(规划/进行/完成/归档) | -| 🧑💻 Worker 管理 | 注册虚拟员工(供应商 + 模型 + 角色提示词 + 温度/Token + 成本上限)、连通性测试、启停用 | -| 📋 任务编排 | 任务 CRUD、看板(拖拽改状态)、优先级、截止时间、**自动路由**(按价格选最便宜可用 Worker) | -| 🚀 执行引擎 | 后台线程执行、统一模型网关(多供应商 OpenAI 兼容)、失败重试、超时、成本上限预检、断点状态 | -| ✅ 人工审核 HITL | 产出自动进入「待审核」,人工通过/打回(必填原因)→ 返工重跑,AI 不可绕过 | -| 💰 成本记账 | Token/成本按 项目×Worker×模型 多维核算、预算硬约束、近 7 天趋势 | -| 📊 质量指标 | 一次通过率、返工次数、失败数 | -| 📜 全链路日志 | 任务级执行日志 + 全局日志流,每个结论可溯源 | +| 📁 项目管理 | 项目 CRUD、目标/验收标准/预算上限、状态流转 | +| 🧑💻 Worker 管理 | 虚拟员工档案(供应商+模型+角色提示词+温度/Token+成本上限)、连通性测试、自动路由 | +| 📋 看板 | 6 列看板拖拽、优先级、截止时间、返工计数 | +| 🚀 执行引擎 | 统一模型网关(doubao/deepseek/openai/qwen/vLLM)、失败重试、超时、成本预检 | +| ✅ HITL 审核 | 产出进「待审核」,通过/打回(必填原因)→ 返工,AI 不可绕过 | +| 💰 成本记账 | Token/成本按 项目×Worker×模型 核算、预算硬约束、一次通过率 | + +### V1(本次交付 v2.0) +| 模块 | 能力 | +|---|---| +| 🔗 **DAG 多任务编排** | 任务依赖(depends_on)、依赖校验拦截、**完成后自动触发下游**、并行执行、一键执行整个工作流、SVG 流程可视化 | +| 🧠 **AI 辅助规划 WBS** | 输入目标 → LLM 自动生成 4~8 个任务的任务分解(含依赖关系)→ 可编辑预览 → 一键导入项目 | +| 📚 **RAG 知识库** | 项目级文档(粘贴/上传 txt·md)、自动分块索引、BM25 检索(jieba 分词)、**任务执行时自动注入相关知识**(日志可见命中来源) | +| 🚨 **预算控制与告警** | 项目预算使用率阈值告警(默认 80%)、成本上限拦截、**告警中心**(未读角标/标记已读)、去重防刷屏 | +| 🔌 **开放 API** | Bearer Token 认证、Token 管理、外部系统可建任务/派活/查结果/审核(curl 示例内置) | +| 📨 **通知集成** | 飞书/企微群机器人 Webhook + 邮件 SMTP,订阅事件:待审核/完成/失败/预算告警/Worker 异常/规划完成,渠道连通性测试 | ## 快速开始 @@ -26,45 +37,49 @@ ./start.sh restart # 重启 ``` -- 管理界面:`http://:16071/`,默认口令 `admin123`(改 `config.py` 的 `AUTH_PASSWORD`,留空则免登录) -- 首次启动自动灌入演示数据(1 项目 / 1 Worker / 3 任务),可用 `python3 seed.py` 重新初始化 +- 管理界面:`http://:16071/`,默认口令 `admin123`(改 `config.py` 的 `AUTH_PASSWORD`) +- 首次启动自动灌入演示数据;`python3 seed.py` 重新初始化 +- V1 全链路验证脚本:`python3 dag_verify.py <终点任务ID>`(自动审核放行直到链路完成) ## 技术栈 -- 后端:Python + Flask + SQLite(WAL),无外部依赖,纯标准库 + flask/requests -- 前端:原生 HTML/JS/CSS 单页应用(无 CDN,离线可用) -- 模型网关:OpenAI 兼容协议,内置 doubao / deepseek / openai / qwen / vLLM 供应商,可扩展 +- 后端:Python + Flask + SQLite(WAL),依赖仅 flask/requests/jieba +- 前端:原生 HTML/JS/CSS SPA(无 CDN,离线可用) +- RAG:本地 BM25(jieba 分词),零依赖;升级路径 = 商用 embedding + pgvector/Qdrant +- 编排:自研轻量 DAG 引擎(线程级,断点状态落库);升级路径 = Temporal 持久执行 + LangGraph 多 Agent 推理 ## 目录结构 ``` ai-worker-platform/ -├── app.py # Flask 应用 + REST API -├── config.py # 供应商/定价/鉴权/端口配置 -├── db.py # SQLite 数据层(projects/workers/tasks/logs/cost_records) -├── engine.py # 任务执行引擎(后台线程 + 自动路由 + 预算预检) -├── llm_gateway.py # 统一模型网关(调用/重试/计量/计价) +├── app.py # Flask 应用 + REST API(含 V1 路由) +├── config.py # 供应商/定价/鉴权/告警阈值/邮件 SMTP +├── db.py # SQLite 数据层 + 自动迁移 +├── engine.py # 执行引擎:DAG 依赖/自动触发/RAG 注入/告警/通知 +├── llm_gateway.py # 统一模型网关 +├── rag.py # 知识库:分块 + BM25 检索 +├── notify.py # 飞书/企微/邮件通知 +├── dag_verify.py # DAG 全链路验证脚本 ├── seed.py # 演示数据 ├── start.sh # 启停脚本 ├── static/ # 前端 SPA +├── docs/ # 界面截图 └── data/ # SQLite 库;logs/ 运行日志 ``` -## API 摘要 +## 开放 API 摘要 -- `POST /api/login` 登录;`GET /api/me` 会话检查 -- `GET|POST /api/projects`、`GET|PUT|DELETE /api/projects/` -- `GET|POST /api/workers`、`GET|PUT|DELETE /api/workers/`、`POST /api/workers//test` -- `GET|POST /api/tasks`、`GET|PUT|DELETE /api/tasks/` -- `POST /api/tasks//run` 派活、`/cancel` 取消、`/review` 审核(approve/reject+reason) -- `GET /api/reports/cost?group=project|worker|model`、`GET /api/stats`、`GET /api/logs` +所有接口支持 `Authorization: Bearer `(管理界面「开放 API」页生成): -## 定价配置 +- `POST /api/projects` 建项目 | `GET /api/projects` +- `POST /api/tasks` 建任务(body 含 project_id/title/description/depends_on…) +- `POST /api/tasks//run` 派活 | `GET /api/tasks/` 查结果 +- `POST /api/tasks//review` 审核 `{"action":"approve"}` 或 `{"action":"reject","reason":"…"}` +- `POST /api/projects//workflow/run` 执行整个工作流 +- `POST /api/projects//wbs/generate` AI 生成任务分解 | `/wbs/import` 导入 +- `GET /api/reports/cost?group=project|worker|model` 成本报表 +- `GET /api/alerts` 告警 | `GET /api/stats` 统计 -`config.py` 中 `MODEL_PRICING` 按「元 / 1M tokens」登记各模型单价,未登记的走默认价。 -成本 = 输入 tokens/1e6 × 输入单价 + 输出 tokens/1e6 × 输出单价。 +## 路线图 -## 路线图(超出 MVP) - -- V1:DAG 多任务编排(Temporal 持久执行)、RAG 知识库、开放 API、飞书/企微集成 -- V2:多 Agent 协作(主管/评审/辩论模式)、自动评估体系、模板市场、企业版 +- **V2**:多 Agent 协作模式(主管/评审/辩论)、Temporal 持久执行、自动评估体系(eval)、模板市场、企业版(私有化/SSO/审计) diff --git a/app.py b/app.py index 87e97e0..68fbe47 100644 --- a/app.py +++ b/app.py @@ -4,13 +4,17 @@ AI Worker 项目管理平台 - MVP 人派活 → AI 干活 → 人验收 最小闭环 """ import os +import secrets import functools +import json as _json from flask import Flask, request, jsonify, session, send_from_directory import config import db import engine import llm_gateway +import rag +import notify app = Flask(__name__, static_folder='static', static_url_path='') app.secret_key = config.SECRET_KEY @@ -46,11 +50,27 @@ def me(): return jsonify({'ok': True, 'authed': not auth_enabled() or session.get('authed')}) +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 auth_enabled() and not session.get('authed'): - return jsonify({'ok': False, 'error': '未登录'}), 401 + if not _auth_ok(): + return jsonify({'ok': False, 'error': '未登录或 Token 无效'}), 401 return fn(*args, **kwargs) return wrapper @@ -199,11 +219,11 @@ def tasks(): d = request.get_json(force=True) tid = db.w( 'INSERT INTO tasks (project_id, worker_id, title, description, priority, ' - 'review_required, deadline, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?)', + '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', ''), - db.now(), db.now())) + _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}) @@ -245,6 +265,9 @@ def task_detail(tid): 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)) @@ -259,6 +282,9 @@ def task_run(tid): 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): return jsonify({'ok': True}) return jsonify({'ok': False, 'error': '任务已在执行中'}), 400 @@ -295,6 +321,11 @@ def task_review(tid): 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())) + # 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: @@ -330,7 +361,7 @@ def report_cost(): '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) + 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}) @@ -383,6 +414,307 @@ def logs(): return jsonify({'ok': True, 'data': rows}) +@app.route('/api/projects//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//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//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//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//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/', 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//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//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/', 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/', 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//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')), + }}) + + # --------------------------------------------------------------------------- # 前端 # --------------------------------------------------------------------------- diff --git a/config.py b/config.py index d375934..c5f2a54 100644 --- a/config.py +++ b/config.py @@ -77,5 +77,18 @@ DEFAULT_PRICE = {'input': 2.0, 'output': 8.0} MAX_RETRY = 1 # 失败重试次数(429/5xx/网络错误) TASK_TIMEOUT = 600 # 单任务超时(秒) +# 预算告警阈值(项目预算使用率 >= 该值触发告警) +BUDGET_ALERT_RATIO = 0.8 + +# 邮件通知配置(可选,留空 host 则禁用) +EMAIL = { + 'host': os.environ.get('SMTP_HOST', 'mail.tphai.com'), + 'port': int(os.environ.get('SMTP_PORT', 587)), + 'starttls': os.environ.get('SMTP_STARTTLS', '0') == '1', + 'user': os.environ.get('SMTP_USER', 'hz4th_coder@tphai.com'), + 'password': os.environ.get('SMTP_PASSWORD', 'hz4th_coder@!'), + 'from_name': 'AI Worker 平台', +} + # 自动路由:按模型单价升序挑选可用 Worker AUTO_ROUTE_POOL = 'enabled' # enabled | all diff --git a/dag_verify.py b/dag_verify.py new file mode 100644 index 0000000..58f51c6 --- /dev/null +++ b/dag_verify.py @@ -0,0 +1,58 @@ +#!/usr/bin/env python3 +"""DAG 全链路自动验证:审核放行 → 观察自动触发 → 直至终点任务 done""" +import json, sys, time, urllib.request + +BASE = 'http://127.0.0.1:16071' +COOKIE = '/tmp/aw_cookies.txt' + +def req(path, method='GET', body=None): + r = urllib.request.Request(BASE + path, method=method) + cookie = '' + for line in open(COOKIE): + parts = line.strip().split('\t') + if len(parts) >= 7 and parts[5] and parts[6]: + cookie += f'{parts[5]}={parts[6]}; ' + r.add_header('Cookie', cookie.strip()) + if body is not None: + r.add_header('Content-Type', 'application/json') + data = json.dumps(body).encode() + else: + data = None + with urllib.request.urlopen(r, data=data) as resp: + return json.loads(resp.read()) + +def task_status(tid): + d = req(f'/api/tasks/{tid}')['data'] + return d['status'], d.get('costs', []) + +# 终点任务 +END = int(sys.argv[1]) if len(sys.argv) > 1 else 9 +chain = [4, 5, 6, 7, 8, 9] +approved = set() +deadline = time.time() + 900 + +print(f'开始 DAG 全链路验证,终点任务 #{END}') +while time.time() < deadline: + all_done = True + for tid in chain: + st, costs = task_status(tid) + mark = f'¥{sum(c["cost"] for c in costs):.4f}' if costs else '-' + print(f' #{tid}: {st:9s} {mark}', flush=True) + if st == 'review' and tid not in approved: + req(f'/api/tasks/{tid}/review', 'POST', {'action': 'approve'}) + approved.add(tid) + print(f' ✅ 审核通过 #{tid}', flush=True) + if st != 'done': + all_done = False + if all_done: + print(f'🎉 全链路完成!共审核 {len(approved)} 个任务') + break + time.sleep(20) +else: + print('⏰ 超时未完成') + +total = 0 +for tid in chain: + _, costs = task_status(tid) + total += sum(c['cost'] for c in costs) +print(f'总成本: ¥{total:.4f}') diff --git a/db.py b/db.py index 095dfee..64b76df 100644 --- a/db.py +++ b/db.py @@ -80,13 +80,81 @@ CREATE TABLE IF NOT EXISTS cost_records ( created_at INTEGER ); +CREATE TABLE IF NOT EXISTS documents ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + project_id INTEGER NOT NULL, + name TEXT NOT NULL, + content TEXT DEFAULT '', + source TEXT DEFAULT 'manual', -- manual/file/url + chunk_size INTEGER DEFAULT 0, + created_at INTEGER, + updated_at INTEGER +); + +CREATE TABLE IF NOT EXISTS doc_chunks ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + document_id INTEGER NOT NULL, + idx INTEGER DEFAULT 0, + content TEXT DEFAULT '', + tokens INTEGER DEFAULT 0 +); + +CREATE TABLE IF NOT EXISTS alerts ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + type TEXT DEFAULT 'system', -- budget/task_failed/worker/limit/notify/plan + level TEXT DEFAULT 'info', -- info/warn/critical + title TEXT DEFAULT '', + detail TEXT DEFAULT '', + read INTEGER DEFAULT 0, + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS api_tokens ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + token TEXT NOT NULL UNIQUE, + created_at INTEGER, + last_used_at INTEGER +); + +CREATE TABLE IF NOT EXISTS notify_channels ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + type TEXT NOT NULL, -- feishu/wecom/email + webhook TEXT DEFAULT '', + email TEXT DEFAULT '', + events TEXT DEFAULT '[]', -- JSON: task_review/task_done/task_failed/budget_alert/worker_alert + enabled INTEGER DEFAULT 1, + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS settings ( + key TEXT PRIMARY KEY, + value TEXT DEFAULT '' +); + CREATE INDEX IF NOT EXISTS idx_tasks_project ON tasks(project_id); CREATE INDEX IF NOT EXISTS idx_tasks_status ON tasks(status); CREATE INDEX IF NOT EXISTS idx_logs_task ON task_logs(task_id); CREATE INDEX IF NOT EXISTS idx_cost_task ON cost_records(task_id); CREATE INDEX IF NOT EXISTS idx_cost_project ON cost_records(project_id); +CREATE INDEX IF NOT EXISTS idx_docs_project ON documents(project_id); +CREATE INDEX IF NOT EXISTS idx_chunks_doc ON doc_chunks(document_id); +CREATE INDEX IF NOT EXISTS idx_alerts_read ON alerts(read); """ +# --------------------------------------------------------------------------- +# 迁移:给旧表补列(V1) +# --------------------------------------------------------------------------- +def _migrate(): + conn = get_conn() + cols = {r['name'] for r in conn.execute('PRAGMA table_info(tasks)')} + if 'depends_on' not in cols: + conn.execute("ALTER TABLE tasks ADD COLUMN depends_on TEXT DEFAULT '[]'") + conn.execute('CREATE INDEX IF NOT EXISTS idx_tasks_depends ON tasks(depends_on)') + conn.commit() + conn.close() + def get_conn(): conn = sqlite3.connect(DB_PATH, timeout=30) @@ -101,6 +169,7 @@ def init_db(): conn.executescript(SCHEMA) conn.commit() conn.close() + _migrate() def q(sql, args=(), one=False): @@ -154,4 +223,23 @@ def task_worker_cost(worker_id): def serialize_task(t): t = dict(t) t['review_required'] = bool(t['review_required']) + try: + t['depends_on'] = json.loads(t.get('depends_on') or '[]') + except Exception: + t['depends_on'] = [] return t + + +def get_setting(key, default=''): + r = q('SELECT value FROM settings WHERE key=?', (key,), one=True) + return r['value'] if r else default + + +def set_setting(key, value): + conn = get_conn() + try: + conn.execute('INSERT INTO settings (key, value) VALUES (?,?) ' + 'ON CONFLICT(key) DO UPDATE SET value=excluded.value', (key, str(value))) + conn.commit() + finally: + conn.close() diff --git a/docs/aw_dag.png b/docs/aw_dag.png new file mode 100644 index 0000000..3179c04 Binary files /dev/null and b/docs/aw_dag.png differ diff --git a/docs/aw_kb.png b/docs/aw_kb.png new file mode 100644 index 0000000..10bd525 Binary files /dev/null and b/docs/aw_kb.png differ diff --git a/engine.py b/engine.py index 82580a2..ce603ae 100644 --- a/engine.py +++ b/engine.py @@ -1,13 +1,19 @@ # -*- coding: utf-8 -*- """ -任务执行引擎:后台线程执行单次 LLM 调用,记录日志与成本, -完成后进入「待审核」或直接「已完成」。 +任务执行引擎 V1: +- DAG 依赖校验 + 下游自动触发(串行/并行) +- RAG 知识库上下文注入 +- 预算告警(项目预算使用率阈值) +- 事件通知(待审核/完成/失败)与告警记录 """ +import json import threading import traceback import db import llm_gateway import config +import rag +import notify def _log(task_id, level, message): @@ -36,6 +42,25 @@ def _cost_record(task, worker, usage): usage['total_tokens'], usage['cost'], db.now())) +def _deps(task): + try: + return json.loads(task.get('depends_on') or '[]') + except Exception: + return [] + + +def check_dependencies(task): + """DAG 依赖检查:返回 (ok, blockers)""" + blockers = [] + for dep_id in _deps(task): + dep = db.q('SELECT id, title, status FROM tasks WHERE id=?', (dep_id,), one=True) + if not dep: + blockers.append(f'# {dep_id}(已删除)') + elif dep['status'] != 'done': + blockers.append(f'「{dep["title"]}」(#{dep_id}) {dep["status"]}') + return (not blockers), blockers + + def pick_worker_auto(task): """自动路由:按模型输入单价升序挑选 enabled Worker""" rows = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id') @@ -63,17 +88,71 @@ def check_worker_limits(worker): return True, '' +def _project_cost(project_id): + rows = db.q('SELECT COALESCE(SUM(cost),0) AS t FROM cost_records WHERE project_id=?', + (project_id,)) + return rows[0]['t'] if rows else 0.0 + + def _check_project_budget(task): proj = db.q('SELECT * FROM projects WHERE id=?', (task['project_id'],), one=True) if proj and proj['budget_limit'] and proj['budget_limit'] > 0: - rows = db.q('SELECT COALESCE(SUM(cost),0) AS t FROM cost_records WHERE project_id=?', - (task['project_id'],)) - used = rows[0]['t'] if rows else 0 + used = _project_cost(task['project_id']) if used >= proj['budget_limit']: return False, f'项目预算已用完({used:.2f}/{proj["budget_limit"]:.2f} 元)' return True, '' +def _budget_alert(project_id): + """预算使用率告警(每次任务完成后检查,避免重复刷屏)""" + proj = db.q('SELECT * FROM projects WHERE id=?', (project_id,), one=True) + if not proj or not proj['budget_limit'] or proj['budget_limit'] <= 0: + return + used = _project_cost(project_id) + ratio = used / proj['budget_limit'] + if ratio >= config.BUDGET_ALERT_RATIO: + # 同项目 1 小时内只告警一次,避免刷屏 + dup = db.q('SELECT COUNT(*) c FROM alerts WHERE type="budget" AND detail LIKE ? ' + 'AND created_at > ?', (f'项目「{proj["name"]}」%', db.now() - 3600)) + if dup[0]['c'] == 0: + notify.notify('budget_alert', + f'预算告警:项目「{proj["name"]}」已使用 {ratio*100:.0f}%', + f'已花费 ¥{used:.2f} / 预算 ¥{proj["budget_limit"]:.2f},' + f'超过阈值 {config.BUDGET_ALERT_RATIO*100:.0f}%,请关注成本控制。', + save_alert=True, level='warn', atype='budget') + + +def _build_messages(task, worker): + """构造提示词:任务指令 + RAG 知识库上下文""" + messages = [] + if worker['system_prompt']: + messages.append({'role': 'system', 'content': worker['system_prompt']}) + user_text = task['description'] or task['title'] + ctx, refs = rag.build_context(task['project_id'], user_text) + if ctx: + user_text = f'{ctx}\n\n----\n\n任务指令:{user_text}' + _log(task['id'], 'info', f'📚 RAG 知识库命中 {len(refs)} 个片段:' + ';'.join(refs[:5])) + messages.append({'role': 'user', 'content': user_text}) + return messages + + +def _trigger_downstream(task): + """DAG:任务完成后自动触发所有就绪的下游任务""" + rows = db.q('SELECT * FROM tasks WHERE status IN ("todo","failed")') + triggered = [] + for t in rows: + deps = _deps(t) + if task['id'] not in deps: + continue + ok, blockers = check_dependencies(t) + if ok: + if runner.submit(t['id']): + triggered.append(t) + else: + _log(t['id'], 'info', f'⏳ 等待前置任务完成:' + '、'.join(blockers)) + return triggered + + def run_task(task_id): """在后台线程中执行任务""" task = db.q('SELECT * FROM tasks WHERE id=?', (task_id,), one=True) @@ -82,6 +161,17 @@ def run_task(task_id): if task['status'] == 'running': return + # DAG 依赖检查 + ok, blockers = check_dependencies(task) + if not ok: + _set_task(task_id, status='failed', error='前置任务未完成:' + '、'.join(blockers), + finished_at=db.now()) + _log(task_id, 'error', '❌ 依赖未满足,无法执行:' + '、'.join(blockers)) + notify.notify('task_failed', f'任务失败:{task["title"]}', + f'项目 #{task["project_id"]} 任务「{task["title"]}」因依赖未完成被拒绝执行:' + + '、'.join(blockers), save_alert=True, level='warn', atype='task_failed') + return + # 确定 Worker worker = None if task['worker_id']: @@ -90,6 +180,9 @@ def run_task(task_id): _set_task(task_id, status='failed', error='指定 Worker 不存在或已停用', finished_at=db.now()) _log(task_id, 'error', '指定 Worker 不存在或已停用') + notify.notify('worker_alert', f'Worker 异常:任务「{task["title"]}」', + f'指定 Worker #{task["worker_id"]} 不存在或已停用', save_alert=True, + level='warn', atype='worker_alert') return else: worker = pick_worker_auto(task) @@ -97,6 +190,9 @@ def run_task(task_id): _set_task(task_id, status='failed', error='无可用 Worker(自动路由失败)', finished_at=db.now()) _log(task_id, 'error', '自动路由失败:无可用 Worker') + notify.notify('worker_alert', f'Worker 异常:任务「{task["title"]}」', + '自动路由失败:没有可用的 Worker', save_alert=True, + level='warn', atype='worker_alert') return _set_task(task_id, worker_id=worker['id']) _log(task_id, 'info', f'自动路由 → Worker「{worker["name"]}」({worker["provider"]}/{worker["model"]})') @@ -106,29 +202,31 @@ def run_task(task_id): if not ok: _set_task(task_id, status='failed', error=reason, finished_at=db.now()) _log(task_id, 'error', reason) + notify.notify('budget_alert', f'成本上限拦截:任务「{task["title"]}」', reason, + save_alert=True, level='warn', atype='budget') return ok, reason = _check_project_budget(task) if not ok: _set_task(task_id, status='failed', error=reason, finished_at=db.now()) _log(task_id, 'error', reason) + notify.notify('budget_alert', f'预算拦截:任务「{task["title"]}」', reason, + save_alert=True, level='warn', atype='budget') return _set_task(task_id, status='running', started_at=db.now(), error='') _log(task_id, 'info', f'开始执行:Worker「{worker["name"]}」 模型 {worker["provider"]}/{worker["model"]}') - messages = [] - if worker['system_prompt']: - messages.append({'role': 'system', 'content': worker['system_prompt']}) - messages.append({'role': 'user', 'content': task['description'] or task['title']}) - try: usage = llm_gateway.chat( - worker['provider'], worker['model'], messages, + worker['provider'], worker['model'], _build_messages(task, worker), temperature=worker['temperature'], max_tokens=worker['max_tokens'], base_url=worker['base_url'] or None, api_key=worker['api_key'] or None) except Exception as e: _set_task(task_id, status='failed', error=str(e), finished_at=db.now()) _log(task_id, 'error', f'执行失败: {e}') + notify.notify('task_failed', f'任务失败:{task["title"]}', + f'项目 #{task["project_id"]} 任务「{task["title"]}」执行出错:{str(e)[:300]}', + save_alert=True, level='warn', atype='task_failed') return _cost_record(task, worker, usage) @@ -146,8 +244,23 @@ def run_task(task_id): _set_task(task_id, **fields) if new_status == 'review': _log(task_id, 'info', '产出已提交,等待人工审核(HITL)') + notify.notify('task_review', f'任务待审核:{task["title"]}', + f'项目 #{task["project_id"]} 任务「{task["title"]}」已完成,等待人工验收。\n' + f'模型 {worker["provider"]}/{worker["model"]} · {usage["total_tokens"]} tokens · ¥{usage["cost"]:.4f}', + save_alert=True, level='info', atype='task_review') else: _log(task_id, 'success', '任务完成(无需审核)') + notify.notify('task_done', f'任务完成:{task["title"]}', + f'项目 #{task["project_id"]} 任务「{task["title"]}」执行完毕,' + f'成本 ¥{usage["cost"]:.4f},tokens {usage["total_tokens"]}', + save_alert=False) + + _budget_alert(task['project_id']) + + # DAG:触发下游就绪任务 + downstream = _trigger_downstream(task) + for t in downstream: + _log(t['id'], 'info', f'🔗 前置任务「{task["title"]}」已完成,自动触发执行') class TaskRunner: diff --git a/notify.py b/notify.py new file mode 100644 index 0000000..3c2d2eb --- /dev/null +++ b/notify.py @@ -0,0 +1,111 @@ +# -*- coding: utf-8 -*- +""" +通知渠道:飞书/企微群机器人 Webhook + 邮件(可选 SMTP) +事件:task_review(待审核)/ task_done(完成)/ task_failed(失败) + budget_alert(预算告警)/ worker_alert(Worker 异常)/ wbs_ready(规划完成) +""" +import json +import smtplib +import requests +import db + +EVENTS = { + 'task_review': '任务待审核', + 'task_done': '任务完成', + 'task_failed': '任务失败', + 'budget_alert': '预算告警', + 'worker_alert': 'Worker 异常', + 'wbs_ready': '规划完成', +} + + +def _channels_for(event): + rows = db.q('SELECT * FROM notify_channels WHERE enabled=1') + out = [] + for r in rows: + try: + evs = json.loads(r['events'] or '[]') + except Exception: + evs = [] + if event in evs: + out.append(r) + return out + + +def _send_feishu(ch, title, text): + payload = {'msg_type': 'text', 'content': {'text': f'【AI Worker 平台】{title}\n{text}'}} + r = requests.post(ch['webhook'], json=payload, timeout=10) + try: + body = r.json() + if body.get('code', 0) not in (0, None): + raise ValueError(f'飞书返回错误: {body.get("msg", r.text[:100])}') + except ValueError: + raise + except Exception: + pass + return r + + +def _send_wecom(ch, title, text): + payload = {'msgtype': 'text', 'text': {'content': f'【AI Worker 平台】{title}\n{text}'}} + return requests.post(ch['webhook'], json=payload, timeout=10) + + +def _send_email(ch, title, text): + from config import EMAIL # 邮件配置在 config.py + if not EMAIL.get('host'): + return None + msg = f'From: {EMAIL["from_name"]} <{EMAIL["user"]}>\nTo: {ch["email"]}\n' \ + f'Subject: 【AI Worker 平台】{title}\nContent-Type: text/plain; charset=utf-8\n\n{text}' + try: + s = smtplib.SMTP(EMAIL['host'], EMAIL['port']) + if EMAIL.get('starttls'): + s.starttls() + if EMAIL.get('user'): + s.login(EMAIL['user'], EMAIL['password']) + s.sendmail(EMAIL['user'], [ch['email']], msg.encode('utf-8')) + s.quit() + return True + except Exception as e: + return f'邮件发送失败: {e}' + + +def notify(event, title, text, save_alert=True, level='info', atype=None): + """触发事件:写告警记录 + 推送所有订阅渠道""" + if save_alert: + db.w('INSERT INTO alerts (type, level, title, detail, read, created_at) ' + 'VALUES (?,?,?,?,0,?)', + (atype or event, level, title, text[:500], db.now())) + ok, fail = [], [] + for ch in _channels_for(event): + try: + if ch['type'] == 'feishu': + r = _send_feishu(ch, title, text) + (ok if r.status_code == 200 else fail).append(f'{ch["name"]}({r.status_code})') + elif ch['type'] == 'wecom': + r = _send_wecom(ch, title, text) + (ok if r.status_code == 200 else fail).append(f'{ch["name"]}({r.status_code})') + elif ch['type'] == 'email': + r = _send_email(ch, title, text) + (ok if r is True else fail).append(f'{ch["name"]}({r})') + except Exception as e: + fail.append(f'{ch["name"]}({e})') + return {'ok': ok, 'fail': fail} + + +def test_channel(ch): + """渠道连通性测试""" + text = '这是一条来自 AI Worker 项目管理平台的连通性测试消息 ✅' + try: + if ch['type'] == 'feishu': + r = _send_feishu(ch, '连通测试', text) + return r.status_code == 200, f'HTTP {r.status_code}' + if ch['type'] == 'wecom': + r = _send_wecom(ch, '连通测试', text) + return r.status_code == 200, f'HTTP {r.status_code}' + if ch['type'] == 'email': + r = _send_email(ch, '连通测试', text) + return r is True, str(r) + except Exception as e: + return False, str(e) + return False, '未知渠道类型' diff --git a/rag.py b/rag.py new file mode 100644 index 0000000..3af043c --- /dev/null +++ b/rag.py @@ -0,0 +1,119 @@ +# -*- coding: utf-8 -*- +""" +RAG 知识库:文档分块 + jieba 分词 + BM25 检索(本地零依赖) +升级路径:接入商用 embedding + pgvector/Qdrant 做向量检索(见 README) +""" +import math +import re +import jieba +import db + +CHUNK_SIZE = 400 # 每块目标字数 +CHUNK_OVERLAP = 80 # 块间重叠 + + +def segment(text): + """jieba 分词,去停用词/单字/空白""" + text = re.sub(r'[\s,。!?、;:""''()【】《》·…—0-9a-zA-Z-]+', ' ', text) + words = [w for w in jieba.cut(text) if len(w.strip()) > 1 and not w.isspace()] + return words + +def split_chunks(content, size=CHUNK_SIZE, overlap=CHUNK_OVERLAP): + """按段落聚合切块,避免从句子中间切断""" + paras = [p.strip() for p in re.split(r'\n+', content) if p.strip()] + chunks, buf, buf_len = [], '', 0 + for p in paras: + if buf_len + len(p) > size and buf: + chunks.append(buf) + tail = buf[-overlap:] if overlap else '' + buf, buf_len = tail + p, len(tail) + len(p) + else: + buf += ('\n' if buf else '') + p + buf_len += len(p) + if buf: + chunks.append(buf) + return chunks or [content[:size]] + + +def rebuild_document(doc_id): + """重新分块索引文档""" + doc = db.q('SELECT * FROM documents WHERE id=?', (doc_id,), one=True) + if not doc: + return 0 + db.w('DELETE FROM doc_chunks WHERE document_id=?', (doc_id,)) + chunks = split_chunks(doc['content']) + for i, c in enumerate(chunks): + db.w('INSERT INTO doc_chunks (document_id, idx, content, tokens) VALUES (?,?,?,?)', + (doc_id, i, c, len(segment(c)))) + db.w('UPDATE documents SET chunk_size=?, updated_at=? WHERE id=?', (len(chunks), db.now(), doc_id)) + return len(chunks) + + +class BM25Index: + """轻量 BM25:按查询词 IDF 加权打分""" + + def __init__(self, chunks): + # chunks: [{id, document_id, content, tokens}] + self.chunks = chunks + self.doc_len = [max(c['tokens'], 1) for c in chunks] + self.avg_len = sum(self.doc_len) / max(len(self.doc_len), 1) + self.N = len(chunks) + self.k1, self.b = 1.5, 0.75 + # 倒排:term -> set of chunk idx + self.inv = {} + self.df = {} + for i, c in enumerate(chunks): + seen = set() + for w in segment(c['content']): + if w in seen: + continue + seen.add(w) + self.inv.setdefault(w, []).append(i) + for w, lst in self.inv.items(): + self.df[w] = len(lst) + + def search(self, query, top_k=5): + q_terms = [w for w in segment(query)] + if not q_terms or not self.N: + return [] + scores = {} + for w in q_terms: + postings = self.inv.get(w, []) + if not postings: + continue + idf = math.log(1 + (self.N - self.df[w] + 0.5) / (self.df[w] + 0.5)) + for idx in postings: + tf = sum(1 for x in segment(self.chunks[idx]['content']) if x == w) + denom = tf + self.k1 * (1 - self.b + self.b * self.doc_len[idx] / self.avg_len) + scores[idx] = scores.get(idx, 0) + idf * (tf * (self.k1 + 1)) / denom + ranked = sorted(scores.items(), key=lambda x: -x[1])[:top_k] + return [{'chunk_id': self.chunks[i]['id'], 'document_id': self.chunks[i]['document_id'], + 'content': self.chunks[i]['content'], 'score': round(s, 4)} + for i, s in ranked] + + +def search_project(project_id, query, top_k=5): + """在项目知识库中检索,返回 (片段列表, 是否命中)""" + chunks = db.q( + 'SELECT c.id, c.document_id, c.content, c.tokens FROM doc_chunks c ' + 'JOIN documents d ON d.id=c.document_id WHERE d.project_id=?', + (project_id,)) + if not chunks: + return [], False + idx = BM25Index(chunks) + return idx.search(query, top_k), True + + +def build_context(project_id, query, top_k=4): + """生成注入提示词的检索上下文(含来源标注)""" + hits, hit = search_project(project_id, query, top_k) + if not hit or not hits: + return '', [] + parts, refs = [], [] + for h in hits: + doc = db.q('SELECT name FROM documents WHERE id=?', (h['document_id'],), one=True) + name = doc['name'] if doc else f'文档#{h["document_id"]}' + parts.append(f'【来源:{name}】\n{h["content"]}') + refs.append(f'{name}#块{h["chunk_id"]}') + ctx = '以下是项目知识库中的相关资料,回答时请优先参考:\n\n' + '\n\n---\n\n'.join(parts) + return ctx, refs diff --git a/static/app.js b/static/app.js index de065d5..5ba2bcd 100644 --- a/static/app.js +++ b/static/app.js @@ -1,4 +1,4 @@ -/* AI Worker 项目管理平台 - 前端 SPA */ +/* AI Worker 项目管理平台 - 前端 SPA(V1) */ 'use strict'; const $ = (sel, el = document) => el.querySelector(sel); @@ -15,6 +15,7 @@ const STATUS = { }; const PRI = {high:'高', medium:'中', low:'低'}; const KANBAN_ORDER = ['todo','running','review','done','rejected','failed']; +const ALERT_LEVEL = {info:'ℹ️', warn:'⚠️', critical:'🚨'}; /* ---------- API ---------- */ async function api(path, opts = {}) { @@ -62,18 +63,32 @@ function openModal(html) { function closeModal() { $('#modal-mask').style.display = 'none'; $('#modal-box').innerHTML = ''; } $('#modal-mask').addEventListener('click', e => { if (e.target.id === 'modal-mask') closeModal(); }); +/* ---------- 告警角标 ---------- */ +async function refreshAlertBadge() { + try { + const r = await api('/api/alerts/unread_count'); + const b = $('#alert-badge'); + if (r.count > 0) { b.textContent = r.count; b.style.display = 'inline'; } + else b.style.display = 'none'; + } catch (e) {} +} +setInterval(refreshAlertBadge, 30000); + /* ---------- Router ---------- */ const routes = { 'dashboard': pageDashboard, 'projects': pageProjects, 'project': pageProject, - 'workers': pageWorkers, 'reports': pageReports, 'logs': pageLogs + 'workers': pageWorkers, 'reports': pageReports, 'logs': pageLogs, + 'alerts': pageAlerts, 'api': pageApiTokens, 'settings': pageSettings }; function router() { const hash = location.hash.replace(/^#\//, '') || 'dashboard'; - const [name, ...rest] = hash.split('/'); + const parts = hash.split('/'); + const name = parts[0]; const fn = routes[name] || pageDashboard; - $$('#sidebar nav a').forEach(a => a.classList.toggle('active', a.dataset.route === (name === 'project' ? 'projects' : name))); + const navMap = {project:'projects', api:'api', settings:'settings', alerts:'alerts'}; + $$('#sidebar nav a').forEach(a => a.classList.toggle('active', a.dataset.route === (navMap[name] || name))); $('#main').innerHTML = '
加载中…
'; - fn(rest).catch(e => { $('#main').innerHTML = `
${esc(e.message)}
`; }); + fn(parts.slice(1)).catch(e => { $('#main').innerHTML = `
${esc(e.message)}
`; }); } window.addEventListener('hashchange', router); @@ -172,58 +187,85 @@ function openProjectModal(p = {}) { }); } -/* ---------- 项目详情(看板) ---------- */ -let kanbanData = null; -async function pageProject([pid]) { +/* ---------- 项目详情(看板 / DAG / 知识库) ---------- */ +let projCtx = null; // {pid, p, tasks, workers, tab} + +async function pageProject([pid, tab]) { const p = (await api(`/api/projects/${pid}`)).data; const tasks = (await api(`/api/tasks?project_id=${pid}`)).data; const workers = (await api('/api/workers')).data; - kanbanData = {pid, tasks, workers}; + projCtx = {pid, p, tasks, workers, tab: tab || 'kanban'}; + renderProjectShell(); + if (projCtx.tab === 'dag') renderDag(); + else if (projCtx.tab === 'kb') renderKb(); + else renderKanban(); +} + +function renderProjectShell() { + const {p, pid, tab} = projCtx; $('#main').innerHTML = `
← 返回

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

- + +
-
${KANBAN_ORDER.map(s => ` -
-
${STATUS[s].label} - 0
-
-
`).join('')}
`; - renderKanban(); - $$('.kb-body').forEach(bd => { - bd.addEventListener('dragover', e => { e.preventDefault(); bd.classList.add('drag-over'); }); - bd.addEventListener('dragleave', () => bd.classList.remove('drag-over')); - bd.addEventListener('drop', e => { - e.preventDefault(); bd.classList.remove('drag-over'); - const tid = e.dataTransfer.getData('text/plain'); - moveTask(parseInt(tid), bd.dataset.drop); - }); - }); + +
`; } +async function refreshProj() { + projCtx.tasks = (await api(`/api/tasks?project_id=${projCtx.pid}`)).data; + if (projCtx.tab === 'dag') renderDag(); + else renderKanban(); +} + +/* --- 看板 --- */ function renderKanban() { - const {tasks} = kanbanData; + const {tasks} = projCtx; const by = {}; tasks.forEach(t => { (by[t.status] = by[t.status] || []).push(t); }); + $('#tab-body').innerHTML = `
${KANBAN_ORDER.map(s => ` +
+
${STATUS[s].label} + 0
+
+
`).join('')}
`; KANBAN_ORDER.forEach(s => { const list = by[s] || []; $('#cnt-' + s).textContent = list.length; $('#kanban .kb-body[data-drop="' + s + '"]').innerHTML = list.map(t => ` -
${esc(t.title)}
●${PRI[t.priority]} ${t.worker_id ? `Worker#${t.worker_id}` : '自动路由'} + ${t.depends_on && t.depends_on.length ? `依赖×${t.depends_on.length}` : ''} ${t.rejection_count ? `打回×${t.rejection_count}` : ''} ${t.status === 'running' ? '' : ''}
${t.status === 'done' ? '
✅ 已交付 v' + t.output_version + '
' : ''}
`).join('') || '
'; }); + bindKanbanDnD(); +} + +function bindKanbanDnD() { + $$('.kb-body').forEach(bd => { + bd.addEventListener('dragover', e => { e.preventDefault(); bd.classList.add('drag-over'); }); + bd.addEventListener('dragleave', () => bd.classList.remove('drag-over')); + bd.addEventListener('drop', e => { + e.preventDefault(); bd.classList.remove('drag-over'); + const tid = parseInt(e.dataTransfer.getData('text/plain')); + moveTask(tid, bd.dataset.drop); + }); + }); $$('.kb-card[draggable="true"]').forEach(c => { c.addEventListener('dragstart', e => { e.dataTransfer.setData('text/plain', c.dataset.id); @@ -234,76 +276,228 @@ function renderKanban() { } async function moveTask(tid, status) { - const t = kanbanData.tasks.find(x => x.id === tid); + const t = projCtx.tasks.find(x => x.id === tid); if (!t || t.status === status) return; if (['running', 'review'].includes(status)) return toast('不能直接拖到该状态', 'err'); if (status === 'todo' && !['rejected', 'failed', 'done'].includes(t.status)) return; if (status === 'done') return toast('完成请走审核流(审核通过)', 'err'); await api(`/api/tasks/${tid}`, {method:'PUT', body:{status}}); - toast('状态已更新', 'ok'); refreshKanban(); + toast('状态已更新', 'ok'); refreshProj(); } -async function refreshKanban() { - const {pid} = kanbanData; - kanbanData.tasks = (await api(`/api/tasks?project_id=${pid}`)).data; - renderKanban(); +/* --- DAG 工作流 --- */ +function renderDag() { + const {tasks} = projCtx; + const idMap = {}; + tasks.forEach(t => { idMap[t.id] = t; }); + const edges = []; + tasks.forEach(t => (t.depends_on || []).forEach(d => { if (idMap[d]) edges.push({from:d, to:t.id}); })); + + // 拓扑分层:最长路径 + const layer = {}; + const visit = id => { + if (layer[id] !== undefined) return layer[id]; + const t = idMap[id]; + let l = 0; + (t.depends_on || []).forEach(d => { if (idMap[d]) l = Math.max(l, visit(d) + 1); }); + layer[id] = l; + return l; + }; + tasks.forEach(t => visit(t.id)); + const layers = {}; + tasks.forEach(t => { (layers[layer[t.id]] = layers[layer[t.id]] || []).push(t.id); }); + + const W = 260, H = 110, PAD = 60; + const width = Math.max(2, Object.keys(layers).length) * W + PAD; + const height = Math.max(...tasks.map(t => layer[t.id] === undefined ? 0 : layer[t.id]), 0) * H + PAD + 60; + const pos = {}; + Object.entries(layers).forEach(([l, ids]) => { + const x = PAD + parseInt(l) * W; + const step = Math.min(H - 60, Math.max(70, (height - 2*PAD) / (ids.length + 1))); + ids.forEach((id, i) => { pos[id] = {x, y: PAD + (i + 1) * step}; }); + }); + const colors = {todo:'#2a3550', running:'#173b6b', review:'#4a3a10', done:'#0f3d2e', rejected:'#4a1c1c', failed:'#3d1c2e', cancelled:'#333'}; + const ready = tasks.filter(t => (t.status === 'todo' || t.status === 'failed') && (!t.depends_on || t.depends_on.every(d => idMap[d]?.status === 'done'))); + $('#tab-body').innerHTML = ` +
+ 共 ${tasks.length} 个节点 · ${edges.length} 条依赖边${ready.length ? ` · ${ready.length} 个任务就绪` : ''} + +
+
+ + ${edges.map(e => { + const a = pos[e.from], b = pos[e.to]; + if (!a || !b) return ''; + const mx = (a.x + b.x) / 2; + const act = idMap[e.from].status === 'done'; + return ``; + }).join('')} + ${tasks.map(t => { + const p = pos[t.id]; + const col = colors[t.status]; + return ` + + #${t.id} ${esc(t.title).slice(0,10)} + ${STATUS[t.status].label}${t.depends_on?.length ? ' · 前置'+t.depends_on.length : ''} + `; + }).join('')} +
+
💡 节点点击查看详情;前置任务完成后下游自动触发。执行整个工作流 = 派发所有就绪任务(含失败重试)。
`; } -function openTaskModal(pid) { - const workers = kanbanData ? kanbanData.workers : []; +async function runWorkflow() { + const r = await api(`/api/projects/${projCtx.pid}/workflow/run`, {method:'POST'}); + toast(`🚀 已派发 ${r.started.length} 个任务${r.blocked.length ? `,${r.blocked.length} 个因依赖阻塞` : ''}`, 'ok'); + setTimeout(refreshProj, 1500); +} + +/* --- 知识库 RAG --- */ +async function renderKb() { + const docs = (await api(`/api/projects/${projCtx.pid}/documents`)).data; + $('#tab-body').innerHTML = ` +
+ 文档共 ${docs.length} 篇 · 任务执行时自动检索注入上下文 + +
+
+ +
+
+
文档列表
+
${docs.map(d => ` +
+
${esc(d.name)} ${d.chunks} 块 ${d.content.length} 字 · ${fmtTime(d.created_at)}
+
+
+
`).join('') || '
暂无文档,添加产品资料/知识手册后 AI Worker 干活时会自动参考
'} +
+
`; +} + +async function kbSearch() { + const q = $('#kb-q').value.trim(); + if (!q) return; + const r = await api(`/api/projects/${projCtx.pid}/search?q=${encodeURIComponent(q)}`); + $('#kb-results').innerHTML = r.data.length ? r.data.map(h => ` +
+
📄 命中 文档#${h.document_id} · 块${h.chunk_id} · 得分 ${h.score}
+
${esc(h.content.slice(0, 160))}${h.content.length > 160 ? '…' : ''}
+
`).join('') : '
无命中(知识库为空或关键词不匹配)
'; +} + +function openDocModal(docId) { + (docId ? api(`/api/documents/${docId}`) : Promise.resolve({data: {}})).then(async r => { + const doc = docId ? r.data : {}; + openModal(` +

${docId ? '编辑文档' : '添加文档'}

+ + + + + `); + $('#dm-file').addEventListener('change', e => { + const f = e.target.files[0]; + if (!f) return; + const rd = new FileReader(); + rd.onload = () => { $('#dm-content').value = rd.result; if (!$('#dm-name').value) $('#dm-name').value = f.name; }; + rd.readAsText(f); + }); + $('#dm-save').addEventListener('click', async () => { + const name = $('#dm-name').value.trim(), content = $('#dm-content').value; + if (!name || !content.trim()) return toast('名称和内容必填', 'err'); + const rr = await api(docId ? `/api/documents/${docId}` : `/api/projects/${projCtx.pid}/documents`, + {method: docId ? 'PUT' : 'POST', body: docId ? {name, content} : {name, content}}); + toast(docId ? '已保存' : `已添加,分块 ${rr.chunks} 块`, 'ok'); closeModal(); renderKb(); + }); + }); +} + +async function delDoc(docId) { + if (!confirm('确认删除该文档及其分块?')) return; + await api(`/api/documents/${docId}`, {method:'DELETE'}); + toast('已删除', 'ok'); renderKb(); +} + +/* --- 新建任务(含依赖) --- */ +function openTaskModal(editTask) { + const t = editTask || {}; + const {tasks, workers, pid} = projCtx || {tasks: [], workers: []}; + const candidates = tasks.filter(x => x.status !== 'running' && x.id !== t.id); + const deps = t.depends_on || []; openModal(` -

新建任务

- - +

${t.id ? '编辑任务' : '新建任务'}

+ +
-
-
+
+
+
+ +
+ ${candidates.length ? candidates.map(c => ` + `).join('') : '暂无其他任务可选'}
`); $('#tm-save').addEventListener('click', async () => { const title = $('#tm-title').value.trim(); if (!title) return toast('请填写任务标题', 'err'); + const body = { + title, description: $('#tm-desc').value, + worker_id: $('#tm-worker').value ? parseInt($('#tm-worker').value) : null, + priority: $('#tm-pri').value, deadline: $('#tm-deadline').value, + review_required: $('#tm-review').value === '1', + depends_on: $$('#tm-deps input:checked').map(x => parseInt(x.value)) + }; try { - const r = await api('/api/tasks', {method:'POST', body:{ - project_id: pid, title, description: $('#tm-desc').value, - worker_id: $('#tm-worker').value ? parseInt($('#tm-worker').value) : null, - priority: $('#tm-pri').value, deadline: $('#tm-deadline').value, - review_required: $('#tm-review').value === '1' - }}); - toast('任务已创建', 'ok'); closeModal(); refreshKanban(); - if (confirm('任务已创建,立即派给 AI Worker 执行?')) runTask(r.id); + if (t.id) { + await api(`/api/tasks/${t.id}`, {method:'PUT', body}); + toast('已保存', 'ok'); closeModal(); refreshProj(); + } else { + const r = await api('/api/tasks', {method:'POST', body:{...body, project_id: projCtx.pid}}); + toast('任务已创建', 'ok'); closeModal(); refreshProj(); + if (confirm('任务已创建,立即派给 AI Worker 执行?')) runTask(r.id); + } } catch (e) { toast(e.message, 'err'); } }); } async function runTask(tid) { - try { await api(`/api/tasks/${tid}/run`, {method:'POST'}); toast('🚀 已派活,AI Worker 开始执行', 'ok'); refreshKanban(); } + try { await api(`/api/tasks/${tid}/run`, {method:'POST'}); toast('🚀 已派活,AI Worker 开始执行', 'ok'); refreshProj(); } catch (e) { toast(e.message, 'err'); } } /* ---------- 任务详情(含审核) ---------- */ async function openTaskDetail(tid) { const t = (await api(`/api/tasks/${tid}`)).data; - const workers = kanbanData ? kanbanData.workers : (await api('/api/workers')).data; + const workers = projCtx?.workers || (await api('/api/workers')).data; const w = t.worker || (t.worker_id ? {id: t.worker_id, name: `#${t.worker_id}`} : null); + const deps = (t.depends_on || []).map(d => projCtx?.tasks.find(x => x.id === d)).filter(Boolean); openModal(`

${esc(t.title)}${STATUS[t.status].label}

执行 Worker
${w ? `${esc(w.name)} ${esc(w.provider)}/${esc(w.model)}` : '未指派(自动路由)'}
优先级
●${PRI[t.priority]} · 返工 ${t.rejection_count} 次
+
前置任务
${deps.length ? deps.map(d => `#${d.id} ${esc(d.title).slice(0,12)}(${STATUS[d.status].label})`).join('') : '无'}
创建/更新
${fmtTime(t.created_at)} / ${fmtTime(t.updated_at)}
截止
${t.deadline || '—'}
@@ -314,14 +508,16 @@ async function openTaskDetail(tid) { ${t.error ? `
错误信息
${esc(t.error)}
` : ''}
执行日志
-
${t.logs.map(l => logLine(l)).join('') || '
暂无日志
'}
+
${t.logs.map(logLine).join('') || '
暂无日志
'}
`); } @@ -341,22 +537,70 @@ async function reviewTask(tid, action) { await api(`/api/tasks/${tid}/review`, {method:'POST', body:{action}}); } toast(action === 'approve' ? '✅ 验收通过,任务完成' : '任务已打回,等待返工', 'ok'); - closeModal(); refreshKanban(); + closeModal(); refreshProj(); } async function cancelTask(tid) { if (!confirm('确认取消该任务?')) return; await api(`/api/tasks/${tid}/cancel`, {method:'POST'}); - toast('已取消', 'ok'); closeModal(); refreshKanban(); + toast('已取消', 'ok'); closeModal(); refreshProj(); } async function reopenTask(tid) { await api(`/api/tasks/${tid}`, {method:'PUT', body:{status:'todo'}}); - toast('任务已重新打开', 'ok'); closeModal(); refreshKanban(); + toast('任务已重新打开', 'ok'); closeModal(); refreshProj(); } async function delTask(tid) { if (!confirm('确认删除该任务及其日志/成本记录?')) return; await api(`/api/tasks/${tid}`, {method:'DELETE'}); - toast('已删除', 'ok'); closeModal(); refreshKanban(); + toast('已删除', 'ok'); closeModal(); refreshProj(); +} + +/* ---------- AI 规划 WBS ---------- */ +function openWbsModal() { + const {p} = projCtx; + openModal(` +

🧠 AI 辅助规划(WBS 任务分解)

+ + + +
`); + $('#wb-gen').addEventListener('click', async () => { + const goal = $('#wb-goal').value.trim(); + if (!goal) return toast('请填写项目目标', 'err'); + $('#wb-result').innerHTML = '
AI 正在拆解任务(约 10-30 秒)…
'; + try { + const r = await api(`/api/projects/${projCtx.pid}/wbs/generate`, {method:'POST', body:{goal}}); + let wbs = r.data; + $('#wb-result').innerHTML = ` +
生成 ${wbs.length} 个任务(可编辑后导入)· 花费 ¥${r.cost.toFixed(4)} / ${r.usage} tokens
+ ${wbs.map((t, i) => ` +
+
+ ${i+1} + +
+ +
依赖任务:${wbs.map((_, j) => ``).join('')}
+
`).join('')} + `; + $('#wb-import').addEventListener('click', async () => { + const tasks = wbs.map((t, i) => ({ + title: $(`[data-f="title"][data-i="${i}"]`).value.trim() || t.title, + description: $(`[data-f="desc"][data-i="${i}"]`).value, + depends_on: $$(`[data-dep="${i}"]:checked`).map(x => parseInt(x.value)) + })); + const rr = await api(`/api/projects/${projCtx.pid}/wbs/import`, {method:'POST', body:{tasks, worker_id: projCtx.workers[0]?.id || null}}); + toast(`✅ 已导入 ${rr.created} 个任务(含依赖关系)`, 'ok'); + closeModal(); refreshProj(); router(); + }); + } catch (e) { $('#wb-result').innerHTML = `
❌ ${esc(e.message)}
`; } + }); } /* ---------- Worker 管理 ---------- */ @@ -493,6 +737,170 @@ async function pageLogs() { [任务#${l.task_id} ${esc(l.task_title || '')}] ${esc(l.message)}`).join('') || '
暂无日志
'}`; } +/* ---------- 告警中心 ---------- */ +async function pageAlerts() { + const alerts = (await api('/api/alerts?limit=200')).data; + $('#main').innerHTML = ` +

告警中心任务失败 · 预算超标 · Worker 异常 · 成本拦截

+
+
${alerts.map(a => ` +
+
${ALERT_LEVEL[a.level] || 'ℹ️'} ${esc(a.title)} ${a.read ? '' : '未读'}
+
${esc(a.detail)}
+
${fmtTime(a.created_at)} · ${esc(a.type)}
+
`).join('') || '
暂无告警 🎉
'}
`; + refreshAlertBadge(); +} +async function readAlert(id) { + await api(`/api/alerts/${id}/read`, {method:'POST'}); + pageAlerts(); +} +async function readAllAlerts() { + await api('/api/alerts/read_all', {method:'POST'}); + toast('已全部标记为已读', 'ok'); pageAlerts(); +} + +/* ---------- 开放 API ---------- */ +async function pageApiTokens() { + const tokens = (await api('/api/tokens')).data; + $('#main').innerHTML = ` +

开放 API外部系统通过 Bearer Token 调平台能力

+
+ +
+
+
+
已生成的 Token
+ + ${tokens.map(t => ` + + `).join('') || ''} +
名称创建时间最近使用操作
${esc(t.name)}${fmtTime(t.created_at)}${t.last_used_at ? fmtTime(t.last_used_at) : '从未使用'}
暂无 Token
+
+
+
调用示例
+
+
POST /api/tasks 创建任务
+
curl -X POST http://<host>:16071/api/tasks \\
+  -H "Authorization: Bearer <TOKEN>" \\
+  -H "Content-Type: application/json" \\
+  -d '{"project_id":1,"title":"生成周报","description":"…"}'
+# 响应: {"ok":true,"id":42}
+
POST /api/tasks/42/run 派活执行
+
curl -X POST http://<host>:16071/api/tasks/42/run \\
+  -H "Authorization: Bearer <TOKEN>"
+
GET /api/tasks/42 查结果(status=review 时人工审核)
+
POST /api/tasks/42/review 审核
+
curl -X POST http://<host>:16071/api/tasks/42/review \\
+  -H "Authorization: Bearer <TOKEN>" -H "Content-Type: application/json" \\
+  -d '{"action":"approve"}'   # 或 {"action":"reject","reason":"…"}
+
+
+
`; +} + +async function createToken() { + const name = prompt('Token 名称(用途说明):') || 'default'; + const r = await api('/api/tokens', {method:'POST', body:{name}}); + alert('✅ 生成成功(只显示一次,请立即复制保存):\n\n' + r.token); + pageApiTokens(); +} +async function delToken(id) { + if (!confirm('删除后使用该 Token 的调用将失效,确认?')) return; + await api(`/api/tokens/${id}`, {method:'DELETE'}); + toast('已删除', 'ok'); pageApiTokens(); +} + +/* ---------- 通知设置 ---------- */ +async function pageSettings() { + const [chans, evs, st] = await Promise.all([api('/api/channels'), api('/api/events'), api('/api/settings')]); + $('#main').innerHTML = ` +

通知与设置飞书 / 企微群机器人 · 邮件 · 预算告警阈值

+
+
+
+
通知渠道
+ + ${chans.data.map(c => ` + + + + + + `).join('') || ''} +
名称类型订阅事件状态操作
${esc(c.name)}${c.type === 'feishu' ? '飞书' : c.type === 'wecom' ? '企微' : '邮件'}${(c.events || []).map(e => `${(evs.data.find(x => x.id === e) || {}).name || e}`).join('') || '—'}${c.enabled ? '启用' : '停用'} + +
暂无渠道,添加飞书/企微机器人后任务事件自动推送
+
💡 飞书群机器人:群设置 → 群机器人 → 添加自定义机器人,复制 Webhook 地址(以 https://open.feishu.cn/open-apis/bot/v2/hook/ 开头)。企微同理(https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=…)。
+
+
+
预算告警
+ +
+
+
告警会写入告警中心并推送到订阅了「预算告警」事件的渠道。
+
邮件通知
+
${st.data.email_configured ? '✅ 已配置 SMTP(config.py EMAIL)' : '⚠️ 未配置 SMTP(config.py 的 EMAIL 填写 host/user/password 后重启生效)'}
+
+
`; +} +async function saveSettings() { + await api('/api/settings', {method:'PUT', body:{budget_alert_ratio: parseFloat($('#st-ratio').value)}}); + toast('已保存', 'ok'); +} +async function testChannel(id) { + toast('正在发送测试消息…'); + try { const r = await api(`/api/channels/${id}/test`, {method:'POST'}); + r.ok ? toast('✅ 测试消息已送达', 'ok') : toast('❌ ' + r.msg, 'err'); } + catch (e) { toast('❌ ' + e.message, 'err'); } +} +async function delChannel(id) { + if (!confirm('确认删除该渠道?')) return; + await api(`/api/channels/${id}`, {method:'DELETE'}); + toast('已删除', 'ok'); pageSettings(); +} +function openChannelModal(c, events) { + c = c || {}; + openModal(` +

${c.id ? '编辑通知渠道' : '添加通知渠道'}

+ + +
+ + +
+ ${events.map(e => ``).join('')} +
+
+ `); + const toggle = () => { + const isMail = $('#cm-type').value === 'email'; + $('#cm-webhook-wrap').style.display = isMail ? 'none' : ''; + $('#cm-email-wrap').style.display = isMail ? '' : 'none'; + }; + $('#cm-type').addEventListener('change', toggle); + toggle(); + $('#cm-save').addEventListener('click', async () => { + const body = { + name: $('#cm-name').value.trim(), type: $('#cm-type').value, + webhook: $('#cm-webhook').value.trim(), email: $('#cm-email').value.trim(), + events: $$('#modal-box .check-group input:checked').map(x => x.value), + enabled: $('#cm-enabled').checked ? 1 : 0 + }; + if (!body.name || (body.type !== 'email' && !body.webhook) || (body.type === 'email' && !body.email)) + return toast('请填写完整', 'err'); + await api(c.id ? `/api/channels/${c.id}` : '/api/channels', {method: c.id ? 'PUT' : 'POST', body}); + toast('已保存', 'ok'); closeModal(); pageSettings(); + }); +} + /* ---------- 启动 ---------- */ (async function init() { try { @@ -501,5 +909,6 @@ async function pageLogs() { $('#logout-btn').style.display = 'block'; } catch (e) { showLogin(); return; } try { await api('/api/health'); } catch (e) { $('#conn-state').textContent = '● 服务异常'; $('#conn-state').className = 'conn err'; } + refreshAlertBadge(); router(); })(); diff --git a/static/index.html b/static/index.html index 4e3d2d8..bdfd599 100644 --- a/static/index.html +++ b/static/index.html @@ -16,6 +16,9 @@ 🧑‍💻 AI Worker 💰 成本报表 📜 运行日志 + 🚨 告警中心 + 🔌 开放 API + ⚙️ 通知设置