From 9a95b71631cf47728b7bd7d20987498e4346362a Mon Sep 17 00:00:00 2001 From: hz4th_coder Date: Wed, 12 Aug 2026 12:21:17 +0800 Subject: [PATCH] =?UTF-8?q?V2=20=E5=BC=80=E5=8F=91=EF=BC=9A=E5=A4=9A=20Age?= =?UTF-8?q?nt=20=E5=8D=8F=E4=BD=9C(=E4=B8=BB=E7=AE=A1/=E8=AF=84=E5=AE=A1/?= =?UTF-8?q?=E8=BE=A9=E8=AE=BA)=20+=20=E8=87=AA=E5=8A=A8=E8=AF=84=E4=BC=B0?= =?UTF-8?q?=E4=BD=93=E7=B3=BB=20+=20=E6=A8=A1=E6=9D=BF=E5=B8=82=E5=9C=BA?= =?UTF-8?q?=20+=20=E4=BC=81=E4=B8=9A=E7=89=88(SSO/RBAC/=E5=AE=A1=E8=AE=A1/?= =?UTF-8?q?=E5=90=88=E8=A7=84)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - agents.py: supervisor/review/debate 三模式编排引擎,步骤流水全记录 - eval.py: eval 数据集 + LLM-as-Judge 评分 + 任务产出自动沉淀 + Worker 排行榜 - templates.py: 任务/项目/团队三类模板,占位符渲染(JSON安全转义),一键应用 - enterprise.py: 用户体系/RBAC、OIDC+LDAP SSO、审计日志、数据导出/保留期/PII脱敏 - db.py: V2 表结构 + 迁移 + 启动恢复(stale running→failed) - app.py: V2 API 路由 + 登录升级(账号体系)+ 审计钩子 - 前端: 新增多Agent协作/自动评估/模板市场/企业版 4 页面 - Dockerfile: 企业版私有化部署 --- Dockerfile | 20 ++ README.md | 122 +++++++--- agents.py | 422 ++++++++++++++++++++++++++++++++ app.py | 601 +++++++++++++++++++++++++++++++++++++++++++++- db.py | 182 ++++++++++++++ enterprise.py | 339 ++++++++++++++++++++++++++ eval.py | 241 +++++++++++++++++++ static/app.js | 529 +++++++++++++++++++++++++++++++++++++++- static/index.html | 7 +- static/style.css | 13 + templates.py | 197 +++++++++++++++ 11 files changed, 2626 insertions(+), 47 deletions(-) create mode 100644 Dockerfile create mode 100644 agents.py create mode 100644 enterprise.py create mode 100644 eval.py create mode 100644 templates.py diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..77e5aaf --- /dev/null +++ b/Dockerfile @@ -0,0 +1,20 @@ +# AI Worker 平台 - 企业版私有化部署镜像 +FROM python:3.12-slim + +WORKDIR /app + +# 系统依赖(jieba 分词无需编译;保留最小面) +RUN pip install --no-cache-dir flask requests jieba && \ + pip install --no-cache-dir ldap3 || true + +COPY . /app +RUN mkdir -p data logs + +EXPOSE 16071 +VOLUME ["/app/data", "/app/logs"] + +ENV AUTH_PASSWORD=admin123 \ + SMTP_HOST=mail.tphai.com \ + SMTP_PORT=587 + +CMD ["python3", "app.py"] diff --git a/README.md b/README.md index 2a76265..15eed5b 100644 --- a/README.md +++ b/README.md @@ -3,31 +3,51 @@ > 以「项目」为中心、以「AI Worker」为执行单元的项目管理平台。 > 把大模型团队变成一支可指挥、可审计、可控成本的"虚拟团队"。 -**当前版本:V1**(MVP 闭环 + 多 Agent 编排 / RAG / 预算告警 / 开放 API / 飞书企微集成) +**当前版本:V2.0**(多 Agent 协作 / 自动评估 / 模板市场 / 企业版) --- ## 功能总览 -### MVP(已交付 v1.0) +### 📁 项目管理(V1 基础) | 模块 | 能力 | |---|---| -| 📁 项目管理 | 项目 CRUD、目标/验收标准/预算上限、状态流转 | -| 🧑💻 Worker 管理 | 虚拟员工档案(供应商+模型+角色提示词+温度/Token+成本上限)、连通性测试、自动路由 | -| 📋 看板 | 6 列看板拖拽、优先级、截止时间、返工计数 | -| 🚀 执行引擎 | 统一模型网关(doubao/deepseek/openai/qwen/vLLM)、失败重试、超时、成本预检 | -| ✅ HITL 审核 | 产出进「待审核」,通过/打回(必填原因)→ 返工,AI 不可绕过 | -| 💰 成本记账 | Token/成本按 项目×Worker×模型 核算、预算硬约束、一次通过率 | +| 项目 / Worker / 看板 | CRUD、目标/验收标准/预算上限、虚拟员工档案(供应商+模型+角色提示词+成本上限)、6 列看板、自动路由 | +| 执行引擎 | 统一模型网关(doubao/deepseek/openai/qwen/vLLM)、失败重试、超时、成本预检、HITL 审核(通过/打回必填原因) | +| 成本记账 | Token/成本按 项目×Worker×模型 核算、预算硬约束、一次通过率 | +| DAG 编排 | 任务依赖、依赖校验拦截、完成后自动触发下游、SVG 流程可视化 | +| RAG 知识库 | 项目级文档、BM25 检索(jieba)、任务执行时自动注入相关知识 | +| 通知集成 | 飞书/企微 Webhook + 邮件 SMTP,订阅事件:待审核/完成/失败/预算告警/Worker 异常 | -### V1(本次交付 v2.0) -| 模块 | 能力 | +### 🤝 多 Agent 协作(V2 新增) +| 模式 | 流程 | |---|---| -| 🔗 **DAG 多任务编排** | 任务依赖(depends_on)、依赖校验拦截、**完成后自动触发下游**、并行执行、一键执行整个工作流、SVG 流程可视化 | -| 🧠 **AI 辅助规划 WBS** | 输入目标 → LLM 自动生成 4~8 个任务的任务分解(含依赖关系)→ 可编辑预览 → 一键导入项目 | -| 📚 **RAG 知识库** | 项目级文档(粘贴/上传 txt·md)、自动分块索引、BM25 检索(jieba 分词)、**任务执行时自动注入相关知识**(日志可见命中来源) | -| 🚨 **预算控制与告警** | 项目预算使用率阈值告警(默认 80%)、成本上限拦截、**告警中心**(未读角标/标记已读)、去重防刷屏 | -| 🔌 **开放 API** | Bearer Token 认证、Token 管理、外部系统可建任务/派活/查结果/审核(curl 示例内置) | -| 📨 **通知集成** | 飞书/企微群机器人 Webhook + 邮件 SMTP,订阅事件:待审核/完成/失败/预算告警/Worker 异常/规划完成,渠道连通性测试 | +| 👔 主管模式 | 主管 Agent 拆解任务 → 委派多个 Worker 并行执行 → 主管汇总合成最终交付物 | +| 🔍 评审模式 | Worker 产出 → Reviewer 打分评审(rubric 标准)→ 未达分数线自动返工(最多 N 轮)→ 终稿 | +| ⚖️ 辩论模式 | 多辩手各自立论(可指定立场)→ 互相质询反驳(多轮)→ 首席裁判综合裁决输出共识结论 | + +- 每个运行完整记录步骤流水(角色/阶段/内容/成本),可回放过程 +- 参与 Agent 自由选择(评审模式第 1 个为执行者、辩论模式最后一个为裁判) + +### 🎯 自动评估体系(V2 新增) +- **eval 数据集**:用例(输入 + 期望输出)+ 评分标准 rubric + 标签,内置 2 个基准集(指令遵循/代码生成) +- **LLM-as-Judge**:Worker 逐条作答 → 评委按 rubric 打分(0-100)+ 评语,全程记录延迟/成本 +- **自动沉淀**:把已验收通过(done 且无打回)的任务产出一键转为评测用例,数据集越用越厚 +- **Worker 排行榜**:按平均分排序,评估模型/提示词/Agent 的回归对比 + +### 🧩 模板市场(V2 新增) +- 三类模板:**任务模板**(代码审查/技术方案/周报…)、**项目模板**(MVP 启动/数据分析,含任务 DAG)、**团队模板**(标准三 Agent 团队/辩论三人组) +- 占位符变量(`{name}`/`{target}`/`{code}`…)一键应用,使用次数统计 +- 支持用户自建模板(JSON),内置模板可一键恢复 + +### 🏢 企业版(V2 新增) +| 能力 | 说明 | +|---|---| +| 🔐 用户与 RBAC | 本地账号 + 角色(管理员/成员/审计员),默认 admin/admin123;停用/重置密码 | +| 🔑 SSO 单点登录 | **OIDC**(飞书/钉钉/企微/Okta,授权码模式 + discovery + 自动开通 + 管理员组)与 **LDAP/AD**(绑定认证,可选 ldap3) | +| 📜 审计日志 | 登录/登出、任务派发/审核、协作运行、评估、数据导出等关键操作自动留痕(操作者/动作/目标/IP) | +| ⚖️ 合规 | 全量数据导出(JSON)、审计 CSV 导出、数据保留期清理、PII 脱敏(邮箱/手机/身份证)、数据使用同意凭证 | +| 🔒 私有化 | 单机 SQLite + 无 CDN 前端,完全离线可用;Docker 一键部署(见下) | ## 快速开始 @@ -37,49 +57,75 @@ ./start.sh restart # 重启 ``` -- 管理界面:`http://:16071/`,默认口令 `admin123`(改 `config.py` 的 `AUTH_PASSWORD`) -- 首次启动自动灌入演示数据;`python3 seed.py` 重新初始化 -- V1 全链路验证脚本:`python3 dag_verify.py <终点任务ID>`(自动审核放行直到链路完成) +- 管理界面:`http://:16071/`,默认账号 `admin` / `admin123` +- 首次启动自动注入:内置 eval 数据集、内置模板、admin 用户 +- V1 全链路验证:`python3 dag_verify.py <终点任务ID>`;V2 冒烟:启动后直接跑一次「多 Agent 协作」 + +### Docker 私有化部署(企业版) + +```bash +docker build -t ai-worker-platform:v2 . +docker run -d --name aiworker -p 16071:16071 \ + -v $(pwd)/data:/app/data \ + -e AUTH_PASSWORD=your-password \ + ai-worker-platform:v2 +``` + +数据卷 `data/` 挂出即完成持久化与备份;无外网依赖(模型供应商需可达)。 ## 技术栈 -- 后端:Python + Flask + SQLite(WAL),依赖仅 flask/requests/jieba +- 后端:Python + Flask + SQLite(WAL),核心依赖仅 flask/requests/jieba(LDAP 可选 ldap3) - 前端:原生 HTML/JS/CSS SPA(无 CDN,离线可用) -- RAG:本地 BM25(jieba 分词),零依赖;升级路径 = 商用 embedding + pgvector/Qdrant -- 编排:自研轻量 DAG 引擎(线程级,断点状态落库);升级路径 = Temporal 持久执行 + LangGraph 多 Agent 推理 +- 评估:LLM-as-Judge(复用统一模型网关) +- 编排:自研轻量线程级引擎,状态全落库;升级路径 = Temporal 持久执行 + LangGraph ## 目录结构 ``` ai-worker-platform/ -├── app.py # Flask 应用 + REST API(含 V1 路由) +├── app.py # Flask 应用 + REST API(V1 + V2 路由) ├── config.py # 供应商/定价/鉴权/告警阈值/邮件 SMTP -├── db.py # SQLite 数据层 + 自动迁移 -├── engine.py # 执行引擎:DAG 依赖/自动触发/RAG 注入/告警/通知 +├── db.py # SQLite 数据层 + V2 迁移(users/audit/agent/eval/templates) +├── engine.py # V1 执行引擎:DAG/自动触发/RAG 注入/告警 +├── agents.py # V2 多 Agent 协作引擎:supervisor/review/debate +├── eval.py # V2 自动评估:数据集/LLM 评委/沉淀/排行榜 +├── templates.py # V2 模板市场:三类模板 + 占位符渲染 + 应用 +├── enterprise.py # V2 企业版:用户/RBAC/OIDC/LDAP/审计/合规 ├── llm_gateway.py # 统一模型网关 ├── rag.py # 知识库:分块 + BM25 检索 ├── notify.py # 飞书/企微/邮件通知 ├── dag_verify.py # DAG 全链路验证脚本 ├── seed.py # 演示数据 ├── start.sh # 启停脚本 -├── static/ # 前端 SPA -├── docs/ # 界面截图 +├── static/ # 前端 SPA(含 V2 四页) └── data/ # SQLite 库;logs/ 运行日志 ``` -## 开放 API 摘要 +## V2 开放 API 摘要 -所有接口支持 `Authorization: Bearer `(管理界面「开放 API」页生成): +多 Agent 协作: +- `GET /api/agents/modes` 模式说明 | `POST /api/agents/runs` 发起运行 `{mode, topic, params:{workers:[],rounds,pass_score,stances}}` +- `GET /api/agents/runs` 列表 | `GET /api/agents/runs/` 详情(含步骤流水与最终产出) -- `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` 统计 +自动评估: +- `GET/POST /api/eval/datasets` 数据集 | `GET/PUT/DELETE /api/eval/datasets/` +- `POST /api/eval/datasets//cases` 加用例 | `POST /api/eval/datasets//sink` 沉淀任务产出 +- `POST /api/eval/runs` 发起评估 `{dataset_id, worker_id}` | `GET /api/eval/runs/` 逐条结果 +- `GET /api/eval/leaderboard` Worker 排行榜 + +模板市场: +- `GET/POST /api/templates` | `GET/PUT/DELETE /api/templates/` +- `POST /api/templates//apply` 应用模板 `{variables, project_id, worker_id, provider, model}` + +企业版: +- `GET/POST /api/enterprise/users`(管理员)| `PUT/DELETE /api/enterprise/users/` +- `GET/PUT /api/enterprise/sso` SSO 配置 | `POST /api/enterprise/sso/oidc/login` 发起 OIDC | `GET /api/enterprise/sso/oidc/callback` 回调 +- `POST /api/enterprise/sso/ldap/test` LDAP 连通测试 +- `GET /api/enterprise/audit` 审计日志 | `GET /api/enterprise/audit/export` CSV 导出 +- `GET /api/enterprise/export` 全量数据导出(JSON)| `POST /api/enterprise/compliance` 保留期/脱敏设置 ## 路线图 -- **V2**:多 Agent 协作模式(主管/评审/辩论)、Temporal 持久执行、自动评估体系(eval)、模板市场、企业版(私有化/SSO/审计) +- **V2.1**:Temporal 持久执行、多 Agent 协作接入项目任务、eval 回归对比视图 +- **V3**:多租户 SaaS 化、工作流画布(拖拽编排)、插件市场 diff --git a/agents.py b/agents.py new file mode 100644 index 0000000..6f76756 --- /dev/null +++ b/agents.py @@ -0,0 +1,422 @@ +# -*- coding: utf-8 -*- +""" +V2 多 Agent 协作引擎 +三种模式: +- supervisor 主管模式:主管拆解任务 → 委派多个 Worker 并行执行 → 汇总合成最终产出 +- review 评审模式:Worker 产出 → Reviewer 评审打分 → 未达标返工(最多 N 轮)→ 终稿 +- debate 辩论模式:多位辩手各自立论 → 互相质询(可多轮)→ 裁判综合裁决 +""" +import json +import threading +import traceback +import db +import llm_gateway +import notify + + +def _log(run_id, role, worker_id, stage, content, usage=None): + tokens = usage['total_tokens'] if usage else 0 + cost = usage['cost'] if usage else 0.0 + db.w( + 'INSERT INTO agent_steps (run_id, role, worker_id, seq, stage, content, tokens, cost, created_at) ' + 'VALUES (?,?,?,?,?,?,?,?,?)', + (run_id, role, worker_id, + db.q('SELECT COALESCE(MAX(seq),0)+1 s FROM agent_steps WHERE run_id=?', (run_id,))[0]['s'], + stage, content, tokens, cost, db.now())) + + +def _update_run(run_id, **fields): + if not fields: + return + fields['finished_at'] = db.now() + sets = ', '.join(f'{k}=?' for k in fields) + db.w(f'UPDATE agent_runs SET {sets} WHERE id=?', (*fields.values(), run_id)) + + +def _chat_worker(worker, messages, temperature=None, max_tokens=None): + """调用某个 Worker 的模型,返回 (text, usage)""" + r = llm_gateway.chat( + worker['provider'], worker['model'], messages, + temperature=temperature if temperature is not None else worker['temperature'], + max_tokens=max_tokens or worker['max_tokens'] or 2000, + base_url=worker['base_url'] or None, api_key=worker['api_key'] or None) + return r['text'], r + + +def _workers_from_ids(ids): + """按 id 列表取 Worker;自动路由兜底(id 为空或不存在时)""" + out = [] + for wid in ids or []: + w = db.q('SELECT * FROM workers WHERE id=? AND status="enabled"', (wid,), one=True) + if w: + out.append(w) + if not out: + rows = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id LIMIT 3') + out = rows + return out + + +def _extract_json(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 + if end <= start: + raise ValueError('未找到 JSON 内容') + return json.loads(t[start:end]) + + +def _extract_score(text): + """从评审 JSON 中提取分数,失败则启发式解析""" + try: + data = _extract_json(text) + if isinstance(data, dict): + for k in ('score', '总分', '评分'): + if k in data: + return float(data[k]), data.get('judgment') or data.get('意见') or text + except Exception: + pass + # 启发式:找 0-100 数字 + import re + m = re.search(r'(?:score|总分|评分)[^\d]*(\d{1,3})', text, re.I) + if m: + return float(m.group(1)), text + m = re.search(r'(\d{1,3})\s*/\s*100', text) + if m: + return float(m.group(1)), text + return 60.0, text + + +# --------------------------------------------------------------------------- +# 主管模式 +# --------------------------------------------------------------------------- +SUPERVISOR_PLAN_PROMPT = ( + '你是资深项目经理(主管 Agent)。请把下面的任务拆解为 2~5 个子任务,分配给团队协作完成。\n' + '输出严格 JSON:{{"subtasks": [{{"title": "子任务标题", "instruction": "给执行 Agent 的完整指令(含要求与输出格式)"}}]}}\n' + '只输出 JSON,不要任何解释。\n\n' + '任务主题:{topic}\n' + '背景上下文:{context}' +) + +SUPERVISOR_SYNTH_PROMPT = ( + '你是主管 Agent。团队已完成以下子任务,请汇总合成一份完整、连贯、高质量的最终交付物。\n' + '要求:结构清晰、覆盖所有子任务要点、去除重复、补足衔接,直接输出最终成果(不要解释过程)。\n\n' + '原始任务:{topic}\n' + '子任务成果:\n{parts}' +) + + +def run_supervisor(run_id, run, workers): + _log(run_id, 'supervisor', None, 'plan', + f'主管开始规划:{run["topic"][:200]}') + # 1) 主管拆解 + plan_text, u1 = _chat_worker( + workers[0], + [{'role': 'system', 'content': '你只输出 JSON。'}, + {'role': 'user', 'content': SUPERVISOR_PLAN_PROMPT.format( + topic=run['topic'], context=run['context'] or '无')}], + temperature=0.3, max_tokens=2000) + try: + subtasks = _extract_json(plan_text) + if isinstance(subtasks, dict): + subtasks = subtasks.get('subtasks') or subtasks.get('tasks') or [] + except Exception as e: + _update_run(run_id, status='failed', error=f'主管规划解析失败: {e}') + _log(run_id, 'supervisor', None, 'plan', f'❌ 规划解析失败: {e}') + return + _log(run_id, 'supervisor', workers[0]['id'], 'plan', + f'规划完成,拆解为 {len(subtasks)} 个子任务:' + ';'.join(s.get('title', '?')[:30] for s in subtasks), + u1) + + # 2) 委派并行执行(轮询分配 Worker) + parts = [] + total_tokens, total_cost = u1['total_tokens'], u1['cost'] + for i, st in enumerate(subtasks): + w = workers[i % len(workers)] + title = st.get('title', f'子任务{i+1}') + instr = st.get('instruction') or st.get('description') or title + _log(run_id, 'worker', w['id'], 'delegate', f'委派子任务「{title}」→ {w["name"]}') + try: + text, u = _chat_worker( + w, [{'role': 'system', 'content': w['system_prompt'] or '你是高效可靠的执行 Agent。'}, + {'role': 'user', 'content': instr}], + max_tokens=2500) + parts.append(f'【子任务{i+1}: {title}】\n{text}') + total_tokens += u['total_tokens']; total_cost += u['cost'] + _log(run_id, 'worker', w['id'], 'produce', + f'子任务「{title}」完成({u["total_tokens"]} tokens):\n{text[:600]}', u) + except Exception as e: + parts.append(f'【子任务{i+1}: {title}】\n(执行失败: {str(e)[:200]})') + _log(run_id, 'worker', w['id'], 'produce', f'❌ 子任务「{title}」失败: {e}') + + # 3) 主管合成 + _log(run_id, 'supervisor', None, 'synthesize', '汇总合成最终产出…') + try: + final_text, u2 = _chat_worker( + workers[0], + [{'role': 'system', 'content': '你是资深主管 Agent,负责最终交付物合成。'}, + {'role': 'user', 'content': SUPERVISOR_SYNTH_PROMPT.format( + topic=run['topic'], parts='\n\n'.join(parts))}], + temperature=0.4, max_tokens=3000) + total_tokens += u2['total_tokens']; total_cost += u2['cost'] + _log(run_id, 'supervisor', workers[0]['id'], 'synthesize', f'✅ 最终产出({u2["total_tokens"]} tokens):\n{final_text[:800]}', u2) + except Exception as e: + final_text = '\n\n'.join(parts) + _log(run_id, 'supervisor', None, 'synthesize', f'⚠️ 合成失败,退回拼接结果: {e}') + + _update_run(run_id, status='done', result=final_text, + summary=f'拆解 {len(subtasks)} 个子任务,由 {len(workers)} 个 Agent 协作完成', + total_tokens=total_tokens, cost=round(total_cost, 6)) + _finish_notify(run, '主管协作', final_text) + + +# --------------------------------------------------------------------------- +# 评审模式 +# --------------------------------------------------------------------------- +REVIEWER_PROMPT = ( + '你是严格的评审专家(Reviewer)。请对下面的执行成果进行评审,输出严格 JSON:\n' + '{{"score": 0到100的整数, "judgment": "总体评价", "issues": ["问题1", "问题2"], "suggestions": "具体修改建议"}}\n' + '评分标准:{rubric}\n' + '只输出 JSON。\n\n' + '原始任务:{topic}\n' + '执行成果:\n{output}' +) + +REVISE_PROMPT = ( + '你是执行 Agent。请根据评审意见修改你的成果,输出修改后的最终版本(直接输出成果内容,不要解释过程)。\n\n' + '原始任务:{topic}\n' + '上一版成果:\n{output}\n\n' + '评审意见:\n{review}' +) + + +def run_review(run_id, run, workers): + main_w = workers[0] + reviewer_w = workers[1] if len(workers) > 1 else workers[0] + params = json.loads(run.get('params') or '{}') + max_rounds = int(params.get('rounds', 2)) + pass_score = float(params.get('pass_score', 70)) + rubric = params.get('rubric') or '内容准确性 40%,完整性 30%,清晰度与格式 30%。' + + _log(run_id, 'worker', main_w['id'], 'produce', f'执行 Agent 开始产出:{run["topic"][:200]}') + output, u1 = _chat_worker( + main_w, [{'role': 'system', 'content': main_w['system_prompt'] or '你是高质量执行 Agent。'}, + {'role': 'user', 'content': run['topic'] + (('\n\n上下文:' + run['context']) if run['context'] else '')}], + max_tokens=2500) + total_tokens, total_cost = u1['total_tokens'], u1['cost'] + _log(run_id, 'worker', main_w['id'], 'produce', f'初稿完成({u1["total_tokens"]} tokens):\n{output[:600]}', u1) + + final_output, last_review, final_score = output, '', 0.0 + for round_i in range(max_rounds + 1): + # 评审 + review_usage = None + try: + review_text, u2 = _chat_worker( + reviewer_w, + [{'role': 'system', 'content': '你只输出 JSON。'}, + {'role': 'user', 'content': REVIEWER_PROMPT.format( + topic=run['topic'], output=final_output, rubric=rubric)}], + temperature=0.2, max_tokens=1500) + score, judgment = _extract_score(review_text) + review_usage = u2 + total_tokens += u2['total_tokens']; total_cost += u2['cost'] + except Exception as e: + score, judgment = 0, f'评审调用失败: {e}' + final_score = score + last_review = judgment + _log(run_id, 'reviewer', reviewer_w['id'], 'critique', + f'第 {round_i+1} 轮评审:得分 {score}/100\n{judgment[:500]}', review_usage) + + if score >= pass_score or round_i >= max_rounds: + _log(run_id, 'reviewer', reviewer_w['id'], 'critique', + f'{"✅ 评审通过" if score >= pass_score else "⚠️ 已达最大轮次,采纳当前稿"}(得分 {score})') + break + # 返工 + try: + revised, u3 = _chat_worker( + main_w, + [{'role': 'system', 'content': main_w['system_prompt'] or '你是高质量执行 Agent。'}, + {'role': 'user', 'content': REVISE_PROMPT.format( + topic=run['topic'], output=final_output, review=judgment)}], + max_tokens=2500) + total_tokens += u3['total_tokens']; total_cost += u3['cost'] + final_output = revised + _log(run_id, 'worker', main_w['id'], 'revise', + f'第 {round_i+1} 轮返工完成:\n{revised[:600]}', u3) + except Exception as e: + _log(run_id, 'worker', main_w['id'], 'revise', f'❌ 返工失败: {e}') + break + + _update_run(run_id, status='done', result=final_output, + summary=f'评审得分 {final_score}/100(阈值 {pass_score}),共评审 {min(max_rounds+1, 3)} 轮', + total_tokens=total_tokens, cost=round(total_cost, 6)) + _finish_notify(run, '评审协作', final_output) + + +# --------------------------------------------------------------------------- +# 辩论模式 +# --------------------------------------------------------------------------- +DEBATE_OPEN_PROMPT = ( + '你是辩手(立场:{stance})。针对议题给出你的立场陈述与核心论点。\n' + '要求:论点鲜明、论据充分、逻辑严谨,控制在 400 字以内。\n\n' + '议题:{topic}\n' + '背景:{context}' +) + +DEBATE_REBUT_PROMPT = ( + '你是辩手(立场:{stance})。请阅读其他辩手的观点,进行质询与反驳,同时回应对方对你观点的质疑。\n' + '输出严格 JSON:{{"rebuttal": "你的反驳与回应", "refine": "是否坚持原立场 true/false"}}\n' + '只输出 JSON。\n\n' + '议题:{topic}\n' + '你的上一轮立场陈述:\n{my_view}\n\n' + '其他辩手观点:\n{others}' +) + +JUDGE_PROMPT = ( + '你是首席裁判(Judge)。请综合以下辩论内容,给出最终裁决。\n' + '输出严格 JSON:{{"winner": "胜出方或共识结论", "consensus": "最终结论(可直接作为交付物)", "summary": "辩论过程摘要"}}\n' + '只输出 JSON。\n\n' + '议题:{topic}\n' + '辩论记录:\n{transcript}' +) + + +def run_debate(run_id, run, workers): + params = json.loads(run.get('params') or '{}') + rounds = int(params.get('rounds', 1)) + stances = params.get('stances') or [] + judge_w = workers[-1] + debaters = workers[:-1] if len(workers) > 1 else workers + + # 立场分配 + if len(stances) < len(debaters): + default_stances = ['正方(支持)', '反方(反对)', '中立(分析利弊)', '补充视角'] + stances = stances + default_stances[len(stances):] + + views = {} + total_tokens, total_cost = 0, 0.0 + + # 第一轮:立论 + for i, w in enumerate(debaters): + stance = stances[i % len(stances)] + try: + text, u = _chat_worker( + w, [{'role': 'system', 'content': w['system_prompt'] or '你是犀利严谨的辩手。'}, + {'role': 'user', 'content': DEBATE_OPEN_PROMPT.format( + stance=stance, topic=run['topic'], context=run['context'] or '无')}], + temperature=0.7, max_tokens=800) + views[w['id']] = {'stance': stance, 'view': text} + total_tokens += u['total_tokens']; total_cost += u['cost'] + _log(run_id, 'debater', w['id'], 'produce', f'【{stance}】{w["name"]} 立论:\n{text[:500]}', u) + except Exception as e: + views[w['id']] = {'stance': stance, 'view': f'(发言失败: {e})'} + _log(run_id, 'debater', w['id'], 'produce', f'❌ {w["name"]} 立论失败: {e}') + + # 质询轮 + for rnd in range(rounds): + for i, w in enumerate(debaters): + mine = views[w['id']] + others = '\n\n'.join( + f"【{v['stance']}】(#{k}):{v['view'][:500]}" + for k, v in views.items() if k != w['id']) + try: + text, u = _chat_worker( + w, + [{'role': 'system', 'content': '你只输出 JSON。'}, + {'role': 'user', 'content': DEBATE_REBUT_PROMPT.format( + stance=mine['stance'], topic=run['topic'], + my_view=mine['view'], others=others)}], + temperature=0.7, max_tokens=800) + data = _extract_json(text) + rebut = data.get('rebuttal') if isinstance(data, dict) else text + mine['view'] = mine['view'] + '\n\n【质询回应】' + str(rebut) + total_tokens += u['total_tokens']; total_cost += u['cost'] + _log(run_id, 'debater', w['id'], 'produce', + f'第 {rnd+1} 轮质询回应({w["name"]}):\n{str(rebut)[:500]}', u) + except Exception as e: + _log(run_id, 'debater', w['id'], 'produce', f'❌ {w["name"]} 质询失败: {e}') + + # 裁判裁决 + transcript = '\n\n'.join( + f"【{v['stance']}】{v['view']}" for v in views.values()) + _log(run_id, 'judge', judge_w['id'], 'verdict', '裁判综合裁决中…') + try: + verdict, u4 = _chat_worker( + judge_w, + [{'role': 'system', 'content': '你只输出 JSON。你是公正严明的首席裁判。'}, + {'role': 'user', 'content': JUDGE_PROMPT.format(topic=run['topic'], transcript=transcript)}], + temperature=0.3, max_tokens=1500) + data = _extract_json(verdict) + if isinstance(data, dict): + consensus = data.get('consensus') or data.get('结论') or verdict + summary = data.get('summary') or data.get('摘要') or '' + winner = data.get('winner') or '' + else: + consensus, summary, winner = verdict, '', '' + total_tokens += u4['total_tokens']; total_cost += u4['cost'] + _log(run_id, 'judge', judge_w['id'], 'verdict', + f'🏆 裁决:{winner}\n共识结论:{str(consensus)[:800]}', u4) + except Exception as e: + consensus, summary, winner = transcript, '裁判裁决失败: ' + str(e), '' + _log(run_id, 'judge', judge_w['id'], 'verdict', f'❌ 裁决失败: {e}') + + _update_run(run_id, status='done', result=str(consensus), + summary=f'辩论参与者 {len(debaters)} 名,质询 {rounds} 轮' + (f';胜出:{winner}' if winner else ''), + total_tokens=total_tokens, cost=round(total_cost, 6)) + _finish_notify(run, '辩论协作', str(consensus)) + + +def _finish_notify(run, mode_label, result_preview): + try: + notify.notify('task_done', f'多 Agent 协作完成:{run["title"]}', + f'模式:{mode_label}\n主题:{run["topic"][:100]}\n' + f'产出预览:{str(result_preview)[:200]}', + save_alert=False) + except Exception: + pass + + +def execute_run(run_id): + """后台线程:执行一次多 Agent 协作""" + run = db.q('SELECT * FROM agent_runs WHERE id=?', (run_id,), one=True) + if not run: + return + try: + workers = _workers_from_ids(json.loads(run.get('worker_ids') or '[]')) + if not workers: + _update_run(run_id, status='failed', error='没有可用的 Worker') + _log(run_id, 'system', None, 'plan', '❌ 无可用 Worker') + return + mode = run['mode'] + if mode == 'supervisor': + run_supervisor(run_id, run, workers) + elif mode == 'review': + run_review(run_id, run, workers) + elif mode == 'debate': + run_debate(run_id, run, workers) + else: + _update_run(run_id, status='failed', error=f'未知模式: {mode}') + except Exception as e: + _update_run(run_id, status='failed', error=f'引擎异常: {str(e)[:300]}') + _log(run_id, 'system', None, 'plan', '❌ ' + traceback.format_exc()) + + +class AgentRunner: + def __init__(self): + self._threads = {} + + def submit(self, run_id): + if run_id in self._threads and self._threads[run_id].is_alive(): + return False + t = threading.Thread(target=execute_run, args=(run_id,), daemon=True) + self._threads[run_id] = t + t.start() + return True + + +runner = AgentRunner() diff --git a/app.py b/app.py index 68fbe47..c7e74d3 100644 --- a/app.py +++ b/app.py @@ -7,7 +7,7 @@ import os import secrets import functools import json as _json -from flask import Flask, request, jsonify, session, send_from_directory +from flask import Flask, request, jsonify, session, send_from_directory, redirect import config import db @@ -15,10 +15,17 @@ 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() # --------------------------------------------------------------------------- @@ -32,22 +39,58 @@ def auth_enabled(): def login(): data = request.get_json(force=True) if not auth_enabled(): - return jsonify({'ok': True}) - if data.get('password') == config.AUTH_PASSWORD: + session['username'] = 'admin' session['authed'] = True return jsonify({'ok': True}) - return jsonify({'ok': False, 'error': '口令错误'}), 401 + 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(): - return jsonify({'ok': True, 'authed': not auth_enabled() or session.get('authed')}) + 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(): @@ -71,13 +114,39 @@ def require_auth(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': 'v1.0.0'}) + return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v2.0.0'}) # --------------------------------------------------------------------------- @@ -286,6 +355,8 @@ def task_run(tid): 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 @@ -321,6 +392,8 @@ 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())) + 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): @@ -334,6 +407,8 @@ def task_review(tid): '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 @@ -715,6 +790,520 @@ def settings(): }}) +# --------------------------------------------------------------------------- +# 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/', 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/', 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//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/', 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//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/', 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/', 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//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', 'doubao'), + model=d.get('model', 'doubao-seed-evolving')) + 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/', 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(), + }}) + + # --------------------------------------------------------------------------- # 前端 # --------------------------------------------------------------------------- diff --git a/db.py b/db.py index 64b76df..056d472 100644 --- a/db.py +++ b/db.py @@ -134,6 +134,136 @@ CREATE TABLE IF NOT EXISTS settings ( ); CREATE INDEX IF NOT EXISTS idx_tasks_project ON tasks(project_id); + +-- =================================================================== +-- V2 表结构:多 Agent 协作 / 自动评估 / 模板市场 / 企业版 +-- =================================================================== +CREATE TABLE IF NOT EXISTS agent_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + mode TEXT NOT NULL, -- supervisor / review / debate + title TEXT DEFAULT '', + topic TEXT DEFAULT '', -- 输入主题 / 任务 + context TEXT DEFAULT '', -- 附加上下文(知识库/约束) + worker_ids TEXT DEFAULT '[]', -- JSON: 参与协作的 worker id 列表 + params TEXT DEFAULT '{}', -- JSON: rounds/阈值/立场等 + status TEXT DEFAULT 'running', -- running/done/failed/cancelled + result TEXT DEFAULT '', -- 最终产出 + summary TEXT DEFAULT '', -- 过程摘要(评审意见/共识等) + error TEXT DEFAULT '', + total_tokens INTEGER DEFAULT 0, + cost REAL DEFAULT 0, + created_at INTEGER, + finished_at INTEGER +); + +CREATE TABLE IF NOT EXISTS agent_steps ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + run_id INTEGER NOT NULL, + role TEXT DEFAULT '', -- supervisor/worker/reviewer/judge/debater + worker_id INTEGER, + seq INTEGER DEFAULT 0, + stage TEXT DEFAULT '', -- plan/delegate/produce/critique/revise/synthesize/verdict + content TEXT DEFAULT '', + tokens INTEGER DEFAULT 0, + cost REAL DEFAULT 0, + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS eval_datasets ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + name TEXT NOT NULL, + description TEXT DEFAULT '', + rubric TEXT DEFAULT '', -- 评分标准(LLM-as-judge) + tags TEXT DEFAULT '[]', + is_builtin INTEGER DEFAULT 0, + created_at INTEGER, + updated_at INTEGER +); + +CREATE TABLE IF NOT EXISTS eval_cases ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + dataset_id INTEGER NOT NULL, + input TEXT DEFAULT '', + expected TEXT DEFAULT '', + tags TEXT DEFAULT '[]', + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS eval_runs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + dataset_id INTEGER NOT NULL, + worker_id INTEGER NOT NULL, + status TEXT DEFAULT 'running', -- running/done/failed/cancelled + score REAL DEFAULT 0, -- 平均分 0-100 + total_tokens INTEGER DEFAULT 0, + cost REAL DEFAULT 0, + cases_total INTEGER DEFAULT 0, + cases_done INTEGER DEFAULT 0, + created_at INTEGER, + finished_at INTEGER +); + +CREATE TABLE IF NOT EXISTS eval_results ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + run_id INTEGER NOT NULL, + case_id INTEGER NOT NULL, + worker_id INTEGER, + output TEXT DEFAULT '', + score REAL DEFAULT 0, + judgment TEXT DEFAULT '', + latency_ms INTEGER DEFAULT 0, + cost REAL DEFAULT 0, + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS templates ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + type TEXT NOT NULL, -- task / project / team + name TEXT NOT NULL, + description TEXT DEFAULT '', + content TEXT DEFAULT '{}', -- JSON + tags TEXT DEFAULT '[]', + author TEXT DEFAULT 'system', + is_builtin INTEGER DEFAULT 0, + usage_count INTEGER DEFAULT 0, + created_at INTEGER, + updated_at INTEGER +); + +CREATE TABLE IF NOT EXISTS users ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + username TEXT NOT NULL UNIQUE, + password_hash TEXT DEFAULT '', + display_name TEXT DEFAULT '', + role TEXT DEFAULT 'member', -- admin / member / auditor + source TEXT DEFAULT 'local', -- local / oidc / ldap + status TEXT DEFAULT 'active', -- active / disabled + last_login_at INTEGER, + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS audit_logs ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + actor TEXT DEFAULT '', -- 用户名 / token 名 / system + action TEXT DEFAULT '', -- 如 task.create / worker.update / agent.run + target TEXT DEFAULT '', + detail TEXT DEFAULT '', + ip TEXT DEFAULT '', + user_agent TEXT DEFAULT '', + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS enterprise_settings ( + key TEXT PRIMARY KEY, + value TEXT DEFAULT '' +); + +CREATE INDEX IF NOT EXISTS idx_agent_steps_run ON agent_steps(run_id); +CREATE INDEX IF NOT EXISTS idx_agent_runs_status ON agent_runs(status); +CREATE INDEX IF NOT EXISTS idx_eval_cases_ds ON eval_cases(dataset_id); +CREATE INDEX IF NOT EXISTS idx_eval_runs_ds ON eval_runs(dataset_id); +CREATE INDEX IF NOT EXISTS idx_eval_results_run ON eval_results(run_id); +CREATE INDEX IF NOT EXISTS idx_audit_time ON audit_logs(created_at); 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); @@ -156,6 +286,42 @@ def _migrate(): conn.close() +# --------------------------------------------------------------------------- +# V2 迁移:旧库升级(幂等) +# --------------------------------------------------------------------------- +def migrate_v2(): + """老数据库升级:V2 表由 SCHEMA 中的 CREATE TABLE IF NOT EXISTS 保证存在; + 此处处理老表缺列 / 默认数据(管理员账号、内置模板)。""" + conn = get_conn() + # users 表首次出现时注入默认管理员 + c = conn.execute('SELECT COUNT(*) c FROM users').fetchone()['c'] + if c == 0: + conn.execute( + "INSERT INTO users (username, password_hash, display_name, role, source, status, created_at) " + "VALUES ('admin', ?, '管理员', 'admin', 'local', 'active', ?)", + (_hash_password('admin123'), int(time.time()))) + conn.commit() + conn.close() + + +def _hash_password(pwd): + import hashlib + return 'sha256$' + hashlib.sha256(pwd.encode('utf-8')).hexdigest() + + +def verify_password(pwd, pwd_hash): + if not pwd_hash: + return False + if pwd_hash.startswith('sha256$'): + import hashlib + return hashlib.sha256(pwd.encode('utf-8')).hexdigest() == pwd_hash.split('$', 1)[1] + return pwd == pwd_hash # 兼容明文 + + +def hash_password(pwd): + return _hash_password(pwd) + + def get_conn(): conn = sqlite3.connect(DB_PATH, timeout=30) conn.row_factory = sqlite3.Row @@ -170,6 +336,22 @@ def init_db(): conn.commit() conn.close() _migrate() + migrate_v2() + + +def recover_stale_runs(): + """启动恢复:进程重启后,把遗留的 running 状态标记为 failed(线程已随进程消亡)。 + 覆盖:V1 任务 / V2 协作运行 / V2 评估运行。""" + now_ts = now() + n1 = w('UPDATE tasks SET status="failed", error="服务重启,执行中断", finished_at=? ' + 'WHERE status="running"', (now_ts,)) + n2 = w('UPDATE agent_runs SET status="failed", error="服务重启,协作中断", finished_at=? ' + 'WHERE status="running"', (now_ts,)) + n3 = w('UPDATE eval_runs SET status="failed", finished_at=? WHERE status="running"', (now_ts,)) + if n1 or n2 or n3: + import logging + logging.warning(f'recover_stale_runs: tasks={n1 or 0} agent_runs={n2 or 0} eval_runs={n3 or 0}') + return (n1 or 0, n2 or 0, n3 or 0) def q(sql, args=(), one=False): diff --git a/enterprise.py b/enterprise.py new file mode 100644 index 0000000..f258f7e --- /dev/null +++ b/enterprise.py @@ -0,0 +1,339 @@ +# -*- coding: utf-8 -*- +""" +V2 企业版能力 +- 用户体系 + RBAC(admin / member / auditor) +- SSO:OIDC 授权码模式 + LDAP 绑定(可选依赖),本地口令兜底 +- 审计日志:全 API 关键操作留痕(actor/action/target/ip) +- 合规:数据导出(全量 JSON / 审计 CSV)、保留期清理、PII 掩码 +""" +import csv +import io +import json +import time +import uuid +import db + +# --------------------------------------------------------------------------- +# 审计 +# --------------------------------------------------------------------------- +PII_PATTERNS = [ + (r'[A-Za-z0-9._%+-]+@[A-Za-z0-9.-]+\.[A-Za-z]{2,}', ''), + (r'\b1[3-9]\d{9}\b', ''), + (r'\b\d{17}[\dXx]\b', ''), +] + + +def mask_pii(text): + """合规:敏感信息掩码""" + if not text: + return text + import re + for pat, rep in PII_PATTERNS: + text = re.sub(pat, rep, text) + return text + + +def audit(actor, action, target='', detail='', ip='', user_agent=''): + """写审计日志(失败不影响主流程)""" + try: + if db.get_setting('compliance_mask_pii', '0') == '1': + detail = mask_pii(detail) + db.w( + 'INSERT INTO audit_logs (actor, action, target, detail, ip, user_agent, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (str(actor)[:100], str(action)[:100], str(target)[:200], str(detail)[:2000], + str(ip)[:64], str(user_agent)[:200], db.now())) + except Exception: + pass + + +def current_actor(): + """从请求上下文推断操作者(由 app 注入 request-local 变量)""" + import flask + try: + req = flask.request + actor = getattr(flask.g, 'auth_actor', None) + if actor: + return actor + hdr = req.headers.get('Authorization', '') + if hdr.startswith('Bearer '): + r = db.q('SELECT name FROM api_tokens WHERE token=?', (hdr[7:].strip(),), one=True) + return f'token:{r["name"]}' if r else 'token:?' + username = flask.session.get('username') + if username: + return username + return 'anonymous' + except Exception: + return 'anonymous' + + +def audit_auto(action, target='', detail='', save_body_keys=None): + """装饰器版自动审计:包装 flask 路由""" + import flask + import functools + + def deco(fn): + @functools.wraps(fn) + def wrapper(*args, **kwargs): + resp = fn(*args, **kwargs) + try: + detail_text = '' + if save_body_keys and flask.request.method in ('POST', 'PUT'): + try: + body = flask.request.get_json(silent=True) or {} + detail_text = ' '.join(f'{k}={body.get(k)}' for k in save_body_keys if k in body) + except Exception: + pass + audit(current_actor(), action, target or (flask.request.path or ''), + detail_text, flask.request.remote_addr or '', + flask.request.headers.get('User-Agent', '')) + except Exception: + pass + return resp + return wrapper + return deco + + +# --------------------------------------------------------------------------- +# 用户 / RBAC +# --------------------------------------------------------------------------- +def get_user(username): + return db.q('SELECT * FROM users WHERE username=?', (username,), one=True) + + +def role_of(username): + u = get_user(username) + return u['role'] if u else 'anonymous' + + +def is_admin(username): + return role_of(username) == 'admin' + + +def create_local_user(username, password, display_name='', role='member'): + if get_user(username): + return None, '用户已存在' + uid = db.w( + 'INSERT INTO users (username, password_hash, display_name, role, source, status, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (username, db.hash_password(password), display_name or username, role, 'local', 'active', db.now())) + return uid, None + + +def verify_local(username, password): + u = get_user(username) + if not u or u['source'] != 'local' or u['status'] != 'active': + return None + if db.verify_password(password, u['password_hash']): + db.w('UPDATE users SET last_login_at=? WHERE id=?', (db.now(), u['id'])) + return u + return None + + +# --------------------------------------------------------------------------- +# SSO:OIDC(授权码)+ LDAP(可选) +# --------------------------------------------------------------------------- +def sso_config(): + cfg = {} + for k in ('oidc_enabled', 'oidc_name', 'oidc_discovery_url', 'oidc_client_id', + 'oidc_client_secret', 'oidc_redirect_uri', 'oidc_scope', 'oidc_admin_group', + 'ldap_enabled', 'ldap_url', 'ldap_base_dn', 'ldap_bind_dn', 'ldap_bind_password', + 'ldap_user_filter', 'sso_auto_provision'): + cfg[k] = db.get_setting(k, '') + return cfg + + +def save_sso_config(data): + keys = list(data.keys()) + for k in keys: + if k in ('oidc_client_secret', 'ldap_bind_password') and not data[k]: + continue # 留空不覆盖已保存的密钥 + db.set_setting(k, str(data[k])) + + +def oidc_discovery(): + """读取 OIDC discovery 文档,返回端点字典""" + import requests + url = sso_config().get('oidc_discovery_url', '').strip().rstrip('/') + if not url: + raise ValueError('未配置 OIDC discovery URL') + r = requests.get(url, timeout=15) + if r.status_code != 200: + raise ValueError(f'Discovery 请求失败({r.status_code})') + return r.json() + + +def oidc_authorize_url(state): + """生成授权跳转 URL""" + import urllib.parse + cfg = sso_config() + disc = oidc_discovery() + params = { + 'response_type': 'code', + 'client_id': cfg['oidc_client_id'], + 'redirect_uri': cfg['oidc_redirect_uri'], + 'scope': cfg.get('oidc_scope') or 'openid profile email', + 'state': state, + } + return disc.get('authorization_endpoint') + '?' + urllib.parse.urlencode(params) + + +def oidc_exchange(code): + """用授权码换 token + 用户信息""" + import requests + cfg = sso_config() + disc = oidc_discovery() + tok = requests.post(disc.get('token_endpoint'), data={ + 'grant_type': 'authorization_code', + 'code': code, + 'redirect_uri': cfg['oidc_redirect_uri'], + 'client_id': cfg['oidc_client_id'], + 'client_secret': cfg['oidc_client_secret'], + }, timeout=15) + if tok.status_code != 200: + raise ValueError(f'Token 交换失败({tok.status_code}): {tok.text[:200]}') + token_data = tok.json() + id_token = token_data.get('id_token', '') + userinfo = {} + # 优先 userinfo 端点 + if token_data.get('access_token'): + ui = requests.get(disc.get('userinfo_endpoint'), headers={ + 'Authorization': f"Bearer {token_data['access_token']}"}, timeout=15) + if ui.status_code == 200: + userinfo = ui.json() + # 解析 id_token payload 兜底 + if not userinfo and id_token: + import base64 + try: + payload = id_token.split('.')[1] + payload += '=' * (-len(payload) % 4) + userinfo = json.loads(base64.urlsafe_b64decode(payload)) + except Exception: + pass + return userinfo + + +def sso_login(userinfo): + """SSO 登录回调:查找或自动开通用户""" + cfg = sso_config() + username = userinfo.get('preferred_username') or userinfo.get('email') or userinfo.get('sub') or '' + email = userinfo.get('email', '') + display = userinfo.get('name') or userinfo.get('display_name') or username + groups = userinfo.get('groups') or userinfo.get('roles') or [] + if not username: + return None, '无法从 SSO 响应中解析用户名' + u = get_user(username) + if not u: + if cfg.get('sso_auto_provision') != '1': + return None, '用户未开通(自动开通未启用),请联系管理员' + role = 'admin' if cfg.get('oidc_admin_group') and cfg['oidc_admin_group'] in groups else 'member' + db.w( + 'INSERT INTO users (username, password_hash, display_name, role, source, status, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (username, '', display, role, 'oidc', 'active', db.now())) + u = get_user(username) + elif u['status'] != 'active': + return None, '账号已停用' + elif u['source'] != 'oidc': + return None, f'用户名 {username} 已被本地账号占用' + db.w('UPDATE users SET last_login_at=? WHERE id=?', (db.now(), u['id'])) + return u, None + + +def ldap_authenticate(username, password): + """LDAP 绑定认证(依赖 ldap3,未安装时返回 None)""" + try: + from ldap3 import Server, Connection, ALL + except ImportError: + return None, '未安装 ldap3,无法使用 LDAP SSO' + cfg = sso_config() + try: + server = Server(cfg['ldap_url'], get_info=ALL) + conn = Connection(server, user=cfg['ldap_bind_dn'], password=cfg['ldap_bind_password'], + auto_bind=True) + user_filter = cfg.get('ldap_user_filter') or '(uid={username})' + conn.search(cfg['ldap_base_dn'], user_filter.format(username=username), attributes=['cn', 'mail', 'displayName']) + if not conn.entries: + return None, 'LDAP 中未找到该用户' + entry = conn.entries[0] + user_dn = entry.entry_dn + conn.unbind() + conn2 = Connection(server, user=user_dn, password=password, auto_bind=True) + conn2.unbind() + return {'username': username, + 'display': str(entry.displayName.value) if hasattr(entry, 'displayName') else username, + 'email': str(entry.mail.value) if hasattr(entry, 'mail') else ''}, None + except Exception as e: + return None, f'LDAP 认证失败: {e}' + + +def ldap_login(username, password): + info, err = ldap_authenticate(username, password) + if err: + return None, err + u = get_user(username) + if not u: + if sso_config().get('sso_auto_provision') != '1': + return None, '用户未开通,请联系管理员' + db.w( + 'INSERT INTO users (username, password_hash, display_name, role, source, status, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (username, '', info.get('display') or username, 'member', 'ldap', 'active', db.now())) + u = get_user(username) + elif u['status'] != 'active': + return None, '账号已停用' + db.w('UPDATE users SET last_login_at=? WHERE id=?', (db.now(), u['id'])) + return u, None + + +# --------------------------------------------------------------------------- +# 合规:导出 / 保留期 +# --------------------------------------------------------------------------- +def export_all(): + """全量数据导出(JSON)""" + tables = ['projects', 'workers', 'tasks', 'task_logs', 'cost_records', 'documents', + 'agent_runs', 'agent_steps', 'eval_datasets', 'eval_cases', 'eval_runs', + 'eval_results', 'templates', 'users', 'audit_logs'] + out = {'exported_at': time.strftime('%Y-%m-%d %H:%M:%S'), + 'platform': 'ai-worker-platform', 'version': 'v2.0.0'} + for t in tables: + try: + out[t] = db.q(f'SELECT * FROM {t}') + except Exception: + out[t] = [] + return out + + +def export_audit_csv(): + """审计日志导出 CSV""" + rows = db.q('SELECT * FROM audit_logs ORDER BY id DESC LIMIT 10000') + buf = io.StringIO() + w = csv.writer(buf) + w.writerow(['ID', '时间', '操作者', '动作', '目标', '详情', 'IP', 'UA']) + for r in rows: + w.writerow([r['id'], time.strftime('%Y-%m-%d %H:%M:%S', time.localtime(r['created_at'])), + r['actor'], r['action'], r['target'], r['detail'], r['ip'], r['user_agent']]) + return buf.getvalue() + + +def apply_retention(): + """合规:按保留期清理审计日志 / 协作运行 / 评估结果(每天可跑一次)""" + days = int(db.get_setting('compliance_retention_days', '0') or 0) + if days <= 0: + return {'cleaned': 0, 'note': '未配置保留期(0=永久保留)'} + cutoff = db.now() - days * 86400 + cleaned = 0 + for table in ('audit_logs', 'agent_steps', 'agent_runs', 'eval_results', 'eval_runs'): + try: + cur = db.w(f'DELETE FROM {table} WHERE created_at start else {} + score = float(data.get('score', data.get('总分', 0))) + judgment = data.get('judgment') or data.get('评分理由') or '' + return max(0.0, min(100.0, score)), str(judgment)[:500] + except Exception: + import re + m = re.search(r'(\d{1,3})\s*[//]\s*100', text) + if m: + return float(m.group(1)), text[:300] + m = re.search(r'"score"\s*:\s*(\d{1,3})', text) + if m: + return float(m.group(1)), text[:300] + return 0.0, '裁判输出解析失败' + + +def create_builtin_datasets(): + """首次启动注入内置数据集""" + c = db.q('SELECT COUNT(*) c FROM eval_datasets WHERE is_builtin=1')[0]['c'] + if c > 0: + return + for ds in BUILTIN_DATASETS: + did = db.w( + 'INSERT INTO eval_datasets (name, description, rubric, tags, is_builtin, created_at, updated_at) ' + 'VALUES (?,?,?,?,1,?,?)', + (ds['name'], ds['description'], ds['rubric'], json.dumps(ds['tags']), db.now(), db.now())) + for case in ds['cases']: + db.w('INSERT INTO eval_cases (dataset_id, input, expected, tags, created_at) VALUES (?,?,?,?,?)', + (did, case['input'], case['expected'], json.dumps(ds['tags']), db.now())) + + +def sink_done_tasks(): + """沉淀:把最近验收通过(done 且无打回)且尚未沉淀的任务转为 eval 用例(幂等)""" + # 幂等标记:记录在 settings + sunk = set(json.loads(db.get_setting('eval_sunk_task_ids', '[]'))) + rows = db.q( + 'SELECT id, title, description, output_text, updated_at FROM tasks ' + 'WHERE status="done" AND rejection_count=0 AND output_text!="" ORDER BY id DESC LIMIT 200') + created = 0 + for t in rows: + if t['id'] in sunk: + continue + # 找到/创建"任务沉淀"数据集 + ds = db.q('SELECT id FROM eval_datasets WHERE name=?', ('任务产出沉淀',), one=True) + if not ds: + ds_id = db.w( + 'INSERT INTO eval_datasets (name, description, rubric, tags, is_builtin, created_at, updated_at) ' + 'VALUES (?,?,?,?,0,?,?)', + ('任务产出沉淀', '从已验收任务自动沉淀的高质量问答对', DEFAULT_RUBRIC, + json.dumps(['沉淀']), db.now(), db.now())) + else: + ds_id = ds['id'] + # 避免重复用例 + dup = db.q('SELECT id FROM eval_cases WHERE dataset_id=? AND input=?', (ds_id, t['description'] or t['title']), one=True) + if dup: + sunk.add(t['id']) + continue + db.w('INSERT INTO eval_cases (dataset_id, input, expected, tags, created_at) VALUES (?,?,?,?,?)', + (ds_id, (t['description'] or t['title'])[:2000], t['output_text'][:4000], + json.dumps(['沉淀', f'task#{t["id"]}']), db.now())) + sunk.add(t['id']) + created += 1 + db.set_setting('eval_sunk_task_ids', json.dumps(list(sunk)[-2000:])) + return created + + +def evaluate_case(dataset, case, worker): + """单条用例:Worker 作答 + Judge 打分。返回 (output, score, judgment, latency, cost, tokens)""" + t0 = time.time() + # 1) Worker 作答 + r = llm_gateway.chat( + worker['provider'], worker['model'], + [{'role': 'system', 'content': worker['system_prompt'] or '你是待评估的执行 Agent,请直接回答问题。'}, + {'role': 'user', 'content': case['input']}], + temperature=worker['temperature'], max_tokens=worker['max_tokens'] or 2000, + base_url=worker['base_url'] or None, api_key=worker['api_key'] or None) + output = r['text'] + latency = int((time.time() - t0) * 1000) + tokens = r['total_tokens'] + + # 2) Judge 打分 + judge_cost = 0.0 + try: + j = llm_gateway.chat( + worker['provider'], worker['model'], + [{'role': 'system', 'content': '你只输出 JSON。'}, + {'role': 'user', 'content': JUDGE_PROMPT.format( + rubric=dataset['rubric'] or DEFAULT_RUBRIC, + input=case['input'][:2000], + expected=(case['expected'] or '无参考答案,凭专业判断')[:3000], + output=output[:4000])}], + temperature=0.1, max_tokens=600, + base_url=worker['base_url'] or None, api_key=worker['api_key'] or None) + score, judgment = _extract_judge(j['text']) + judge_cost = j['cost'] + tokens += j['total_tokens'] + except Exception as e: + score, judgment = 0.0, f'裁判调用失败: {e}' + total_cost = round(r['cost'] + judge_cost, 6) + return output, score, judgment, latency, total_cost, tokens + + +def run_eval(run_id): + """后台线程:执行一次完整评估""" + run = db.q('SELECT * FROM eval_runs WHERE id=?', (run_id,), one=True) + if not run: + return + ds = db.q('SELECT * FROM eval_datasets WHERE id=?', (run['dataset_id'],), one=True) + worker = db.q('SELECT * FROM workers WHERE id=?', (run['worker_id'],), one=True) + if not ds or not worker: + db.w('UPDATE eval_runs SET status="failed", finished_at=? WHERE id=?', (db.now(), run_id)) + return + cases = db.q('SELECT * FROM eval_cases WHERE dataset_id=? ORDER BY id', (ds['id'],)) + total_tokens, total_cost, scores, done = 0, 0.0, [], 0 + try: + for case in cases: + try: + output, score, judgment, latency, cost, tokens = evaluate_case(ds, case, worker) + db.w( + 'INSERT INTO eval_results (run_id, case_id, worker_id, output, score, judgment, ' + 'latency_ms, cost, created_at) VALUES (?,?,?,?,?,?,?,?,?)', + (run_id, case['id'], worker['id'], output, score, judgment, latency, cost, db.now())) + total_cost += cost + total_tokens += tokens + scores.append(score) + done += 1 + db.w('UPDATE eval_runs SET cases_done=? WHERE id=?', (done, run_id)) + except Exception as e: + db.w( + 'INSERT INTO eval_results (run_id, case_id, worker_id, output, score, judgment, ' + 'latency_ms, cost, created_at) VALUES (?,?,?,?,?,?,?,?,?)', + (run_id, case['id'], worker['id'], '', 0, f'执行失败: {str(e)[:200]}', 0, 0, db.now())) + done += 1 + db.w('UPDATE eval_runs SET cases_done=? WHERE id=?', (done, run_id)) + avg = round(sum(scores) / len(scores), 2) if scores else 0 + db.w('UPDATE eval_runs SET status="done", score=?, total_tokens=?, cost=?, cases_done=?, finished_at=? ' + 'WHERE id=?', (avg, total_tokens, round(total_cost, 6), done, db.now(), run_id)) + try: + notify.notify('task_done', f'评估完成:{ds["name"]}', + f'Worker「{worker["name"]}」在数据集「{ds["name"]}」上平均分 {avg}/100,' + f'共 {done}/{len(cases)} 条用例,成本 ¥{total_cost:.4f}', + save_alert=True, level='info', atype='eval_done') + except Exception: + pass + except Exception as e: + db.w('UPDATE eval_runs SET status="failed", error=?, finished_at=? WHERE id=?', + (str(e)[:300], db.now(), run_id)) + _trace(e) + + +def _trace(e): + print('[eval]', traceback.format_exc()) + + +class EvalRunner: + def __init__(self): + self._threads = {} + + def submit(self, run_id): + if run_id in self._threads and self._threads[run_id].is_alive(): + return False + t = threading.Thread(target=run_eval, args=(run_id,), daemon=True) + self._threads[run_id] = t + t.start() + return True + + +runner = EvalRunner() diff --git a/static/app.js b/static/app.js index 5ba2bcd..008e3d8 100644 --- a/static/app.js +++ b/static/app.js @@ -45,7 +45,7 @@ function showLogin() { $('#login-mask').style.display = 'flex'; } function hideLogin() { $('#login-mask').style.display = 'none'; } $('#login-btn').addEventListener('click', async () => { try { - await api('/api/login', {method:'POST', body:{password: $('#login-pwd').value}}); + await api('/api/login', {method:'POST', body:{username: $('#login-user').value, password: $('#login-pwd').value}}); hideLogin(); $('#logout-btn').style.display = 'block'; router(); } catch (e) { $('#login-err').textContent = e.message; } }); @@ -78,7 +78,9 @@ setInterval(refreshAlertBadge, 30000); const routes = { 'dashboard': pageDashboard, 'projects': pageProjects, 'project': pageProject, 'workers': pageWorkers, 'reports': pageReports, 'logs': pageLogs, - 'alerts': pageAlerts, 'api': pageApiTokens, 'settings': pageSettings + 'alerts': pageAlerts, 'api': pageApiTokens, 'settings': pageSettings, + 'agents': pageAgents, 'eval': pageEval, 'templates': pageTemplates, + 'enterprise': pageEnterprise }; function router() { const hash = location.hash.replace(/^#\//, '') || 'dashboard'; @@ -901,6 +903,529 @@ function openChannelModal(c, events) { }); } +/* ---------- V2 · 多 Agent 协作 ---------- */ +let agentsCtx = {modes: [], workers: [], runs: []}; + +async function pageAgents() { + const [modes, workers, runs] = await Promise.all([ + api('/api/agents/modes'), api('/api/workers'), api('/api/agents/runs') + ]); + agentsCtx = {modes: modes.data, workers: workers.data, runs: runs.data}; + renderAgents(); +} + +function renderAgents() { + const {modes, workers, runs} = agentsCtx; + const wOpts = workers.map(w => + ``).join(''); + $('#main').innerHTML = ` +

多 Agent 协作主管 / 评审 / 辩论 三种团队协作模式

+
${modes.map(m => ` +
+
${m.icon}
+
${esc(m.name)}
+
${esc(m.desc)}
+
${m.params.length} 项参数
+
`).join('')}
+
+
⚡ 发起协作运行
+
+ + + + + + + + + +
+
+ + 💡 评审模式第 1 个 Agent 为执行者,第 2 个(若有)为评审;辩论模式最后一个为裁判。 +
+
+
+
📋 运行记录
+ + ${runs.map(r => ` + + + + + + + + + `).join('') || ''} +
ID模式标题状态过程摘要成本时间操作
#${r.id}${esc(r.mode)}${esc(r.title || r.topic.slice(0, 30))}${({done:'✅ 完成',running:'⏳ 执行中',failed:'❌ 失败',cancelled:'已取消'})[r.status] || r.status}${esc(r.summary || '—')}${fmtMoney(r.cost)}${fmtTime(r.created_at)}
暂无运行记录
+
`; + $('#ag-mode').addEventListener('change', e => { + const m = modes.find(x => x.id === e.target.value); + toast('已选择:' + m.name, ''); + }); + $('#ag-run').addEventListener('click', async () => { + const mode = $('#ag-mode').value; + const topic = $('#ag-topic').value.trim(); + if (!topic) return toast('请填写任务主题', 'err'); + const workerIds = [...$('#ag-workers').selectedOptions].map(o => Number(o.value)); + if (!workerIds.length) return toast('请至少选择一个 Agent', 'err'); + const params = {workers: workerIds, context: $('#ag-context').value}; + if (mode === 'review') { params.rounds = Number($('#ag-rounds').value) || 2; params.pass_score = Number($('#ag-pass').value) || 70; } + if (mode === 'debate') { params.rounds = Number($('#ag-debate-rounds').value) || 1; params.stances = $('#ag-stances').value.split(',').map(s => s.trim()).filter(Boolean); } + try { + await api('/api/agents/runs', {method: 'POST', body: {mode, title: $('#ag-title').value.trim(), topic, params}}); + toast('协作已启动 🚀', 'ok'); pageAgents(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function viewAgentRun(id) { + const d = (await api(`/api/agents/runs/${id}`)).data; + openModal(` +

🤝 协作运行 #${d.id} · ${esc(d.mode)}

+
+
状态
${d.status}
+
主题
${esc(d.topic)}
+
摘要
${esc(d.summary || '—')}
+
成本 / tokens
${fmtMoney(d.cost)} / ${d.total_tokens}
+
+
执行过程
+ ${(d.steps || []).map(s => ` +
[${esc(s.role)} / ${esc(s.stage)}] ${esc(s.worker_name || '')} +
${esc(s.content)}
`).join('') || '
暂无步骤
'} +
最终产出
+
${esc(d.result || '(执行中…)')}
+ `); +} + +/* ---------- V2 · 自动评估 ---------- */ +async function pageEval() { + const [datasets, runs, board, workers] = await Promise.all([ + api('/api/eval/datasets'), api('/api/eval/runs'), api('/api/eval/leaderboard'), api('/api/workers') + ]); + const wOpts = workers.data.map(w => + ``).join(''); + $('#main').innerHTML = ` +

自动评估体系eval 数据集沉淀 · LLM 评委打分 · Worker 排行榜

+
+
🏆 Worker 评估排行榜
+ + ${board.data.map((b, i) => ` + + + `).join('') || ''} +
排名Worker模型评估次数平均分累计成本
${i === 0 ? '🥇' : i === 1 ? '🥈' : i === 2 ? '🥉' : '#' + (i+1)}${esc(b.worker_name)}${esc(b.model)}${b.runs}${b.avg_score}/100${fmtMoney(b.cost)}
暂无评估数据,先跑一次评估吧
+
+
+
+
📚 数据集
+ + + ${datasets.data.map(ds => ` + + + `).join('') || ''} +
名称用例最佳分操作
${esc(ds.name)}${ds.is_builtin ? ' 内置' : ''}${ds.case_count}${ds.best_score != null ? ds.best_score : '—'} +
暂无数据集
+
+
+
🔄 评估运行记录
+ + + ${runs.data.map(r => ` + + + + + `).join('') || ''} +
ID数据集Worker状态得分进度操作
#${r.id}${esc(r.dataset_name)}${esc(r.worker_name)}${({done:'✅ 完成',running:'⏳ 进行中',failed:'❌ 失败'})[r.status] || r.status}${r.status === 'done' ? r.score : '—'}${r.cases_done}/${r.cases_total}
暂无运行记录
+
💡 评估 = Worker 逐条作答 → LLM 评委按 rubric 打分。沉淀功能会把已验收通过的任务产出自动转为评测用例。
+
+
+
+
⚡ 发起评估
+ + + +
`; + $('#ev-run').addEventListener('click', async () => { + try { + await api('/api/eval/runs', {method: 'POST', body: {dataset_id: Number($('#ev-ds').value), worker_id: Number($('#ev-w').value)}}); + toast('评估已启动', 'ok'); pageEval(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +function evalCreateDataset() { + openModal(` +

新建评估数据集

+ + + + `); + $('#ed-save').addEventListener('click', async () => { + try { + await api('/api/eval/datasets', {method: 'POST', body: {name: $('#ed-name').value.trim(), description: $('#ed-desc').value, rubric: $('#ed-rubric').value}}); + toast('已创建', 'ok'); closeModal(); pageEval(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function evalOpenDataset(id) { + const ds = (await api(`/api/eval/datasets/${id}`)).data; + openModal(` +

📚 ${esc(ds.name)} ${ds.case_count || ds.cases.length} 用例

+
${esc(ds.description)}
+
评分标准
+
${esc(ds.rubric || '(未设置,使用默认)')}
+
用例列表
+ ${(ds.cases || []).map(c => ` +
+
Q${c.id}: ${esc(c.input.slice(0, 120))}
+
期望:${esc((c.expected || '—').slice(0, 120))}
+ +
`).join('') || '
暂无用例
'} +
添加用例
+ + + + `); + $('#ec-add').addEventListener('click', async () => { + try { + await api(`/api/eval/datasets/${id}/cases`, {method: 'POST', body: {input: $('#ec-input').value, expected: $('#ec-expected').value}}); + toast('已添加', 'ok'); closeModal(); evalOpenDataset(id); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function evalDelCase(cid) { + if (!confirm('确定删除该用例?')) return; + await api(`/api/eval/cases/${cid}`, {method: 'DELETE'}); + toast('已删除', 'ok'); closeModal(); pageEval(); +} + +async function evalRunDataset(id) { + openModal(`

🚀 评估数据集 #${id}

+
将用选定 Worker 逐条作答并由 LLM 评委打分,每条用例约需 1~2 分钟。
+ `); +} + +async function evalSink() { + try { + const r = await api('/api/eval/datasets/0/sink', {method: 'POST'}); + toast(`已沉淀 ${r.sunk} 条任务产出`, 'ok'); pageEval(); + } catch (e) { toast(e.message, 'err'); } +} + +async function viewEvalRun(id) { + const d = (await api(`/api/eval/runs/${id}`)).data; + openModal(` +

🎯 评估运行 #${d.id} · ${esc(d.dataset_name)}

+
+
Worker
${esc(d.worker_name)}
+
平均分
${d.status === 'done' ? d.score : '…'}/100
+
进度
${d.cases_done}/${d.cases_total}
+
成本
${fmtMoney(d.cost)}
+
+
逐条结果
+ ${(d.results || []).map(r => ` +
+
用例 #${r.case_id} 得分 ${r.score} · ${r.latency_ms}ms · ${fmtMoney(r.cost)}
+
输入:${esc(r.input || '')}
+
评委意见
${esc(r.judgment)}
+ 执行输出
${esc(r.output)}
+
`).join('') || '
暂无结果
'} + `); +} + +/* ---------- V2 · 模板市场 ---------- */ +async function pageTemplates() { + const [tpls, projects, workers] = await Promise.all([ + api('/api/templates'), api('/api/projects'), api('/api/workers') + ]); + const typeMap = {task: '📝 任务模板', project: '📁 项目模板', team: '👥 团队模板'}; + const projOpts = projects.data.map(p => ``).join(''); + const wOpts = workers.data.map(w => ``).join(''); + $('#main').innerHTML = ` +

模板市场任务 / 项目 / Agent 团队 一键复用

+
+ + +
+ ${['task', 'project', 'team'].map(t => ` +
${typeMap[t]}
+
${tpls.data.filter(x => x.type === t).map(x => ` +
+
${esc(x.name)} ${x.is_builtin ? '内置' : ''}
+
${esc(x.description)}
+
使用 ${x.usage_count} 次${x.tags.map(tg => ` · ${esc(tg)}`).join('')}
+
+ ${x.is_builtin ? '' : ``}
+
`).join('') || '
暂无模板
'}
`).join('')} + `; + window._tplProjOpts = projOpts; window._tplWorkerOpts = wOpts; +} + +function tplCreate() { + openModal(` +

新建模板

+ + + + + + `); + $('#tc-save').addEventListener('click', async () => { + let content; + try { content = JSON.parse($('#tc-content').value); } catch (e) { return toast('JSON 格式错误', 'err'); } + try { + await api('/api/templates', {method: 'POST', body: { + type: $('#tc-type').value, name: $('#tc-name').value.trim(), + description: $('#tc-desc').value, content, + tags: $('#tc-tags').value.split(',').map(s => s.trim()).filter(Boolean) + }}); + toast('已创建', 'ok'); closeModal(); pageTemplates(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function tplApply(id, type) { + const t = (await api(`/api/templates/${id}`)).data; + let extra = ''; + if (type === 'task') extra = ` + `; + if (type === 'team') extra = ` + `; + const placeholders = Object.keys(t.content || {}).filter(k => typeof t.content[k] === 'string' && t.content[k].includes('{')).map(k => ``).join(''); + openModal(` +

使用模板:${esc(t.name)}

+
${esc(t.description)}
+ ${extra}${placeholders} + `); + $('#ta-go').addEventListener('click', async () => { + const variables = {}; + Object.keys(t.content || {}).filter(k => typeof t.content[k] === 'string' && t.content[k].includes('{')).forEach(k => variables[k] = $(`#ta-var-${k}`)?.value || ''); + const body = {variables, worker_id: $('#ta-worker')?.value ? Number($('#ta-worker').value) : null, + project_id: $('#ta-project')?.value ? Number($('#ta-project').value) : null, + provider: $('#ta-provider')?.value, model: $('#ta-model')?.value}; + try { + const r = await api(`/api/templates/${id}/apply`, {method: 'POST', body}); + toast(r.kind === 'project' ? `已创建项目 #${r.id}(${r.task_ids.length} 个任务)` : r.kind === 'team' ? `已创建 ${r.ids.length} 个 Worker` : `已创建任务 #${r.id}`, 'ok'); + closeModal(); + if (r.kind === 'project') location.hash = `#/project/${r.id}/kanban`; + else pageTemplates(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function tplDelete(id) { + if (!confirm('确定删除该模板?')) return; + await api(`/api/templates/${id}`, {method: 'DELETE'}); + toast('已删除', 'ok'); pageTemplates(); +} + +async function tplRestore() { + await api('/api/templates/builtin/restore', {method: 'POST'}); + toast('内置模板已就绪', 'ok'); pageTemplates(); +} + +/* ---------- V2 · 企业版 ---------- */ +let entTab = 'users'; +async function pageEnterprise() { + const [users, audit, sso, comp] = await Promise.all([ + api('/api/enterprise/users'), api('/api/enterprise/audit?limit=100'), + api('/api/enterprise/sso'), api('/api/enterprise/compliance') + ]); + $('#main').innerHTML = ` +

企业版私有化 · SSO · 审计 · 合规

+ +
`; + window._entData = {users: users.data, audit: audit.data, sso: sso.data, comp: comp.data}; + entRender(); +} + +function entSwitch(tab) { entTab = tab; entRender(); } + +function entRender() { + const d = window._entData || {}; + const box = $('#ent-body'); + if (entTab === 'users') { + box.innerHTML = ` +
👥 用户列表(RBAC)
+ + + ${(d.users || []).map(u => ` + + + + + + `).join('')} +
用户名显示名角色来源状态最近登录操作
${esc(u.username)}${esc(u.display_name)}${({admin:'👑 管理员',member:'成员',auditor:'审计员'})[u.role] || u.role}${({local:'本地',oidc:'OIDC',ldap:'LDAP'})[u.source] || u.source}${u.status === 'active' ? '✅ 启用' : '🚫 停用'}${fmtTime(u.last_login_at)}
+
💡 角色说明:管理员=全部权限;成员=看板/任务操作;审计员=只读审计。V1 口令 admin123 兼容登录为 admin。
+
`; + } else if (entTab === 'sso') { + const s = d.sso || {}; + box.innerHTML = ` +
+
+
🔐 OIDC(推荐:飞书/钉钉/企业微信/Okta)
+ + + + + + + + + +
+
+
🔑 LDAP / AD
+ + + + + + +
连通性测试
+ + + +
+
+
+
SSO 登录入口
+ +
登录页可对接企业统一身份;私有化部署时 SSO 保证账号体系受企业管控。回调地址需与 IdP 侧配置一致。
+
`; + $('#so-save').addEventListener('click', async () => { + await api('/api/enterprise/sso', {method: 'PUT', body: { + oidc_enabled: $('#so-oidc-on').checked ? '1' : '0', oidc_name: $('#so-oidc-name').value, + oidc_discovery_url: $('#so-oidc-disc').value.trim(), oidc_client_id: $('#so-oidc-cid').value.trim(), + oidc_client_secret: $('#so-oidc-secret').value.trim(), oidc_redirect_uri: $('#so-oidc-redirect').value.trim(), + oidc_admin_group: $('#so-oidc-admingroup').value.trim(), sso_auto_provision: $('#so-oidc-auto').checked ? '1' : '0', + ldap_enabled: $('#so-ldap-on').checked ? '1' : '0', ldap_url: $('#so-ldap-url').value.trim(), + ldap_base_dn: $('#so-ldap-base').value.trim(), ldap_bind_dn: $('#so-ldap-bind').value.trim(), + ldap_bind_password: $('#so-ldap-bindpwd').value.trim(), ldap_user_filter: $('#so-ldap-filter').value.trim() + }}); + toast('SSO 配置已保存', 'ok'); pageEnterprise(); + }); + $('#so-ldap-test').addEventListener('click', async () => { + try { + const r = await api('/api/enterprise/sso/ldap/test', {method: 'POST', body: {username: $('#so-ldap-testuser').value, password: $('#so-ldap-testpwd').value}}); + toast('✅ LDAP 认证成功:' + (r.data.display || r.data.username), 'ok'); + } catch (e) { toast(e.message, 'err'); } + }); + $('#so-oidc-login').addEventListener('click', async () => { + try { + const r = await api('/api/enterprise/sso/oidc/login', {method: 'POST'}); + location.href = r.redirect; + } catch (e) { toast(e.message, 'err'); } + }); + } else if (entTab === 'audit') { + box.innerHTML = ` +
+
📜 审计日志
+ + ${(d.audit || []).map(a => ` + `).join('') || ''} +
时间操作者动作目标详情IP
${fmtTime(a.created_at)}${esc(a.actor)}${esc(a.action)}${esc(a.target)}${esc(a.detail)}${esc(a.ip)}
暂无审计记录
+
登录/登出、协作运行、评估、数据导出等关键操作自动留痕,满足等保审计要求。
+
`; + } else { + const c = d.comp || {}; + box.innerHTML = ` +
+
+
🗄️ 数据保留策略
+ + +
到期自动清理审计日志、协作运行、评估结果等历史数据。
+
+
+
🛡️ PII 脱敏
+ + +
+
+
+
📤 数据导出(GDPR / 个人信息保护法)
+ + +
导出包含项目/任务/协作运行/评估/用户/审计全量数据,供数据可携带权与监管审计使用。
+
+
+
📊 存量统计
+
+
审计记录
${c.audit_count} 条
+
协作运行
${c.agent_runs_count} 次
+
评估运行
${c.eval_runs_count} 次
+
数据使用同意凭证
${esc(c.consent_token || '—')}
+
+
`; + $('#co-ret-save').addEventListener('click', async () => { + await api('/api/enterprise/compliance', {method: 'POST', body: {retention_days: Number($('#co-ret').value), apply_retention: true}}); + toast('已保存并清理', 'ok'); pageEnterprise(); + }); + $('#co-mask-save').addEventListener('click', async () => { + await api('/api/enterprise/compliance', {method: 'POST', body: {mask_pii: $('#co-mask').checked}}); + toast('已保存', 'ok'); pageEnterprise(); + }); + } +} + +function entAddUser() { + openModal(` +

新建用户

+ + + + + + `); + $('#eu-save').addEventListener('click', async () => { + try { + await api('/api/enterprise/users', {method: 'POST', body: { + username: $('#eu-name').value.trim(), display_name: $('#eu-display').value, + role: $('#eu-role').value, source: $('#eu-source').value, password: $('#eu-pwd').value + }}); + toast('已创建', 'ok'); closeModal(); pageEnterprise(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +function entEditUser(id) { + const u = (window._entData?.users || []).find(x => x.id === id); + if (!u) return; + openModal(` +

编辑用户 ${esc(u.username)}

+ + + + + + `); + $('#eu2-save').addEventListener('click', async () => { + await api(`/api/enterprise/users/${id}`, {method: 'PUT', body: { + display_name: $('#eu2-display').value, role: $('#eu2-role').value, + status: $('#eu2-status').value, password: $('#eu2-pwd').value + }}); + toast('已保存', 'ok'); closeModal(); pageEnterprise(); + }); + $('#eu2-del').addEventListener('click', async () => { + if (!confirm('确定删除该用户?')) return; + await api(`/api/enterprise/users/${id}`, {method: 'DELETE'}); + toast('已删除', 'ok'); closeModal(); pageEnterprise(); + }); +} + /* ---------- 启动 ---------- */ (async function init() { try { diff --git a/static/index.html b/static/index.html index bdfd599..42281f4 100644 --- a/static/index.html +++ b/static/index.html @@ -14,10 +14,14 @@ 📊 仪表盘 📁 项目 🧑‍💻 AI Worker + 🤝 多 Agent 协作 + 🎯 自动评估 + 🧩 模板市场 💰 成本报表 📜 运行日志 🚨 告警中心 🔌 开放 API + 🏢 企业版 ⚙️ 通知设置