From eb92a4e591cc5fd02dfcd0532d55a493fcd0e85a Mon Sep 17 00:00:00 2001 From: hz4th_coder Date: Thu, 13 Aug 2026 12:38:38 +0800 Subject: [PATCH] =?UTF-8?q?V3.0=20=E4=BA=A4=E4=BB=98=E4=BD=93=E7=B3=BB=20+?= =?UTF-8?q?=20=E7=B2=BE=E5=87=86=E6=9D=83=E9=99=90=E7=B3=BB=E7=BB=9F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 交付环节优化: - 新建项目必填送达者(人)邮箱,项目完成/遇到无法绕开的难关时自动邮件及时通知 - 每个项目独立工作目录 data/workspace/project_/,交付物分类存放互不污染(上传/下载/删除) - 网页交付物一键部署到 /demo// 免登录公开 Demo,送达者直接打开链接查看 - 非网页交付物 zip 打包 data/packages/,随邮件附件发送送达者(含手动交付/完成交付/自动交付) - 任务全部完成自动收尾交付;失败任务触发难关通知(30分钟限频去重) 用户(人)管理 + 精准权限: - 管理员增删改用户,管理用户项目所属与 Worker 权限(授权弹窗 + 授权总览矩阵) - 角色体系:admin 全部 / auditor 全量只读 / member 按授权 - 项目授权 view/manage/admin;Worker 授权 view/use/manage;创建者自动成为项目管理员 - 仪表盘/成本报表/日志/协作/评估全部按权限过滤,越权访问 403 新表:project_deliverables / user_projects / user_workers;新模块 delivery.py --- README.md | 42 +++- app.py | 647 +++++++++++++++++++++++++++++++++++++++++++++++--- config.py | 3 + db.py | 60 +++++ delivery.py | 369 ++++++++++++++++++++++++++++ engine.py | 43 ++-- enterprise.py | 92 ++++++- static/app.js | 261 +++++++++++++++++++- 8 files changed, 1451 insertions(+), 66 deletions(-) create mode 100644 delivery.py diff --git a/README.md b/README.md index 15eed5b..dd010de 100644 --- a/README.md +++ b/README.md @@ -3,7 +3,7 @@ > 以「项目」为中心、以「AI Worker」为执行单元的项目管理平台。 > 把大模型团队变成一支可指挥、可审计、可控成本的"虚拟团队"。 -**当前版本:V2.0**(多 Agent 协作 / 自动评估 / 模板市场 / 企业版) +**当前版本:V3.0**(交付体系 / 精准权限 / 多 Agent 协作 / 自动评估 / 模板市场 / 企业版) --- @@ -49,6 +49,24 @@ | ⚖️ 合规 | 全量数据导出(JSON)、审计 CSV 导出、数据保留期清理、PII 脱敏(邮箱/手机/身份证)、数据使用同意凭证 | | 🔒 私有化 | 单机 SQLite + 无 CDN 前端,完全离线可用;Docker 一键部署(见下) | +### 📦 交付体系(V3 新增) +| 能力 | 说明 | +|---|---| +| 📬 送达者邮箱 | 新建项目**必填**送达者(人)邮箱;项目**完成**或遇到**无法绕开的难关**时自动邮件及时通知 | +| 🗂️ 项目工作目录 | 每个项目独立工作目录 `data/workspace/project_/`,中间产物/交付物分类存放,互不污染;支持网页上传/下载/删除 | +| 🌐 网页交付物 | 一键部署到 `data/demo//`,经 `/demo//` **免登录公开访问**,送达者直接打开链接查看 | +| 🗜️ 打包交付 | 工作目录一键 zip 打包(`data/packages/`),随邮件附件发送给送达者 | +| ✉️ 手动通知 | 难关说明/进展可随时手动邮件通知送达者;交付全流程留痕(交付记录) | + +### 🔑 精准权限(V3 新增) +| 能力 | 说明 | +|---|---| +| 👥 用户管理 | 管理员增删改用户,管理用户的项目所属与 Worker 权限(企业版 → 用户与权限 → 🔑 授权) | +| 📁 项目授权 | view 查看 / manage 管理(建任务/执行/上传交付物/发送)/ admin 管理员;创建者自动成为项目管理员 | +| 🤖 Worker 授权 | view 查看档案 / use 使用(可指派任务)/ manage 管理;Worker 的注册/删除仍仅管理员 | +| 🎭 角色体系 | 管理员=全部;审计员=全量**只读**;成员=仅可见被授权内容,仪表盘/报表/日志/协作/评估全部按权限过滤 | +| 🗺️ 授权总览 | 一键查看所有用户的 项目×Worker 授权矩阵,杜绝越权 | + ## 快速开始 ```bash @@ -84,22 +102,23 @@ docker run -d --name aiworker -p 16071:16071 \ ``` ai-worker-platform/ -├── app.py # Flask 应用 + REST API(V1 + V2 路由) +├── app.py # Flask 应用 + REST API(V1 + V2 + V3 路由) ├── config.py # 供应商/定价/鉴权/告警阈值/邮件 SMTP -├── db.py # SQLite 数据层 + V2 迁移(users/audit/agent/eval/templates) -├── engine.py # V1 执行引擎:DAG/自动触发/RAG 注入/告警 +├── db.py # SQLite 数据层 + V2/V3 迁移(users/audit/agent/eval/templates/deliverables/grants) +├── 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/审计/合规 +├── enterprise.py # V2/V3 企业版:用户/RBAC/SSO/审计/合规 + 项目/Worker 精准授权 +├── delivery.py # V3 交付体系:工作目录/Demo 部署/打包/邮件送达/难关通知 ├── llm_gateway.py # 统一模型网关 ├── rag.py # 知识库:分块 + BM25 检索 ├── notify.py # 飞书/企微/邮件通知 ├── dag_verify.py # DAG 全链路验证脚本 ├── seed.py # 演示数据 ├── start.sh # 启停脚本 -├── static/ # 前端 SPA(含 V2 四页) -└── data/ # SQLite 库;logs/ 运行日志 +├── static/ # 前端 SPA(含 V3 交付页/授权管理) +└── data/ # SQLite 库;workspace/ 项目工作目录;demo/ 网页Demo;packages/ 打包件;logs/ 运行日志 ``` ## V2 开放 API 摘要 @@ -120,12 +139,19 @@ ai-worker-platform/ 企业版: - `GET/POST /api/enterprise/users`(管理员)| `PUT/DELETE /api/enterprise/users/` +- `GET/PUT /api/enterprise/users//grants` 精准授权(项目/Worker 权限)| `GET /api/enterprise/grants/overview` 授权总览 - `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` 保留期/脱敏设置 +交付体系(V3): +- `GET /api/projects//workspace` 工作目录文件列表 | `POST .../workspace/upload` 上传(multipart)| `GET .../workspace/download?path=` 下载 | `DELETE .../workspace?path=` 删除 +- `POST /api/projects//deploy` 网页交付物部署 Demo | `POST .../package` zip 打包 | `POST .../deliver` 邮件交付(Demo 链接 + 附件) +- `POST /api/projects//complete` 完成项目并交付 | `POST .../notify_deliverer` 手动通知送达者 | `GET /api/projects//deliverables` 交付记录 +- `GET /demo//` 公开 Demo 地址(送达者免登录访问) + ## 路线图 - **V2.1**:Temporal 持久执行、多 Agent 协作接入项目任务、eval 回归对比视图 -- **V3**:多租户 SaaS 化、工作流画布(拖拽编排)、插件市场 +- **V3**:多租户 SaaS 化、工作流画布(拖拽编排)、插件市场 ✅ 交付体系 + 精准权限(V3.0) diff --git a/app.py b/app.py index 3ede7d9..0e5546d 100644 --- a/app.py +++ b/app.py @@ -19,6 +19,7 @@ import agents import eval as evalmod import templates as tplmod import enterprise +import delivery app = Flask(__name__, static_folder='static', static_url_path='') app.secret_key = config.SECRET_KEY @@ -194,7 +195,106 @@ def require_admin(fn): @app.route('/api/health') def health(): - return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v2.0.0'}) + return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v3.0.0'}) + + +# --------------------------------------------------------------------------- +# V3 精准权限:角色(admin/auditor/member)+ 项目授权 + Worker 授权 +# --------------------------------------------------------------------------- +def _session_user(): + """当前会话用户(dict)或 None""" + if not session.get('authed'): + return None + return enterprise.get_user(session.get('username')) + + +def _is_super(): + """admin / auditor / API Token 视为超管可见范围(auditor 只读)""" + actor = enterprise.current_actor() + if actor.startswith('token:'): + return True + u = _session_user() + return bool(u and u['role'] in ('admin', 'auditor')) + + +def _role(): + u = _session_user() + return u['role'] if u else None + + +def _project_perm(pid): + if _is_super(): + return 'admin' + u = _session_user() + if not u: + return None + return enterprise.user_project_perm(u['id'], pid) + + +def _worker_perm(wid): + if _is_super(): + return 'admin' + u = _session_user() + if not u: + return None + return enterprise.user_worker_perm(u['id'], wid) + + +def _visible_project_ids(): + """当前用户可见项目 id 列表;None = 全部""" + u = _session_user() + if not u or _is_super(): + return None + return enterprise.visible_project_ids(u['id']) + + +def _visible_worker_ids(): + u = _session_user() + if not u or _is_super(): + return None + return enterprise.visible_worker_ids(u['id']) + + +def _check_project_perm(pid, need='view'): + """返回 None=通过,否则为错误响应 (jsonify, code)""" + if not _auth_ok(): + return jsonify({'ok': False, 'error': '未登录或 Token 无效'}), 401 + if not db.q('SELECT id FROM projects WHERE id=?', (pid,), one=True): + return jsonify({'ok': False, 'error': '项目不存在'}), 404 + perm = _project_perm(pid) + if not perm: + return jsonify({'ok': False, 'error': '无权访问该项目'}), 403 + if not enterprise.perm_ok(perm, need): + return jsonify({'ok': False, 'error': f'需要更高级别的项目权限({need})'}), 403 + if need != 'view' and _role() == 'auditor': + return jsonify({'ok': False, 'error': '审计员为只读角色,禁止写操作'}), 403 + return None + + +def _check_worker_perm(wid, need='view'): + if not _auth_ok(): + return jsonify({'ok': False, 'error': '未登录或 Token 无效'}), 401 + if not db.q('SELECT id FROM workers WHERE id=?', (wid,), one=True): + return jsonify({'ok': False, 'error': 'Worker 不存在'}), 404 + perm = _worker_perm(wid) + if not perm: + return jsonify({'ok': False, 'error': '无权访问该 Worker'}), 403 + if not enterprise.perm_ok(perm, need): + return jsonify({'ok': False, 'error': f'需要更高级别的 Worker 权限({need})'}), 403 + if need != 'view' and _role() == 'auditor': + return jsonify({'ok': False, 'error': '审计员为只读角色,禁止写操作'}), 403 + return None + + +def _check_task_perm(tid, need='view'): + """按任务所在项目校验权限""" + t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) + if not t: + return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_project_perm(t['project_id'], need) + if err: + return err + return None # --------------------------------------------------------------------------- @@ -205,18 +305,40 @@ def health(): def projects(): if request.method == 'POST': d = request.get_json(force=True) + deliver_email = (d.get('deliver_email') or '').strip() + if not deliver_email: + return jsonify({'ok': False, 'error': '必填:项目送达者(人)邮箱,用于项目完成/难关通知'}), 400 + if '@' not in deliver_email: + return jsonify({'ok': False, 'error': '送达者邮箱格式不正确'}), 400 pid = db.w( 'INSERT INTO projects (name, description, objective, acceptance_criteria, ' - 'status, budget_limit, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?)', + 'status, budget_limit, deliver_email, deliver_type, deliver_note, workspace_dir, ' + 'created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?)', (d.get('name', '').strip(), d.get('description', ''), d.get('objective', ''), d.get('acceptance_criteria', ''), d.get('status', 'active'), - float(d.get('budget_limit') or 0), db.now(), db.now())) + float(d.get('budget_limit') or 0), deliver_email, + d.get('deliver_type', 'web'), d.get('deliver_note', ''), + '', db.now(), db.now())) + # 工作目录名依赖自增 id:先拿到 id 再补写 + db.w('UPDATE projects SET workspace_dir=? WHERE id=?', (f'project_{pid}', pid)) + delivery.workspace_path(pid) + # 创建者自动成为项目管理员 + u = _session_user() + if u and not _is_super(): + db.w('INSERT INTO user_projects (user_id, project_id, perm, created_at) VALUES (?,?,?,?)', + (u['id'], pid, 'admin', db.now())) + enterprise.audit(enterprise.current_actor(), 'project.create', f'project#{pid}', + f'「{d.get("name","")}」送达者 {deliver_email}', request.remote_addr or '') return jsonify({'ok': True, 'id': pid}) rows = db.q('SELECT * FROM projects ORDER BY id DESC') + visible = _visible_project_ids() + if visible is not None: + rows = [r for r in rows if r['id'] in visible] for r in rows: r['task_count'] = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0', (r['id'],))[0]['c'] r['done_count'] = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND status="done" AND deleted=0', (r['id'],))[0]['c'] + r['perm'] = _project_perm(r['id']) return jsonify({'ok': True, 'data': rows}) @@ -224,22 +346,46 @@ def projects(): @require_auth def project_detail(pid): if request.method == 'GET': + err = _check_project_perm(pid, 'view') + if err: + return err p = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) if not p: return jsonify({'ok': False, 'error': '项目不存在'}), 404 + p['perm'] = _project_perm(pid) return jsonify({'ok': True, 'data': p}) if request.method == 'DELETE': + err = _check_project_perm(pid, 'admin') + if err: + return err + db.w('DELETE FROM user_projects WHERE project_id=?', (pid,)) + db.w('DELETE FROM project_deliverables WHERE project_id=?', (pid,)) db.w('DELETE FROM tasks WHERE project_id=?', (pid,)) db.w('DELETE FROM cost_records WHERE project_id=?', (pid,)) db.w('DELETE FROM projects WHERE id=?', (pid,)) + import shutil + for d in (delivery.workspace_path(pid), delivery.demo_path(pid)): + if os.path.isdir(d): + shutil.rmtree(d, ignore_errors=True) return jsonify({'ok': True}) + err = _check_project_perm(pid, 'manage') + if err: + return err d = request.get_json(force=True) - fields = ['name', 'description', 'objective', 'acceptance_criteria', 'status', 'budget_limit'] + fields = ['name', 'description', 'objective', 'acceptance_criteria', 'status', 'budget_limit', + 'deliver_email', 'deliver_type', 'deliver_note'] sets, args = [], [] for f in fields: if f in d: - sets.append(f'{f}=?') - args.append(d[f]) + if f == 'deliver_email': + v = (d[f] or '').strip() + if v and '@' not in v: + return jsonify({'ok': False, 'error': '送达者邮箱格式不正确'}), 400 + sets.append(f'{f}=?') + args.append(v) + else: + sets.append(f'{f}=?') + args.append(d[f]) if sets: args.append(db.now()) db.w(f'UPDATE projects SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, pid)) @@ -253,6 +399,8 @@ def project_detail(pid): @require_auth def workers(): if request.method == 'POST': + if not (enterprise.current_actor().startswith('token:') or _role() == 'admin'): + return jsonify({'ok': False, 'error': '仅管理员可创建 Worker'}), 403 d = request.get_json(force=True) wid = db.w( 'INSERT INTO workers (name, description, provider, model, base_url, api_key, ' @@ -266,8 +414,12 @@ def workers(): db.now(), db.now())) return jsonify({'ok': True, 'id': wid}) rows = db.q('SELECT * FROM workers ORDER BY id DESC') + visible = _visible_worker_ids() + if visible is not None: + rows = [r for r in rows if r['id'] in visible] for r in rows: r['month_cost'] = round(db.monthly_worker_cost(r['id']), 6) + r['perm'] = _worker_perm(r['id']) return jsonify({'ok': True, 'data': rows}) @@ -275,12 +427,27 @@ def workers(): @require_auth def worker_detail(wid): if request.method == 'GET': + err = _check_worker_perm(wid, 'view') + if err: + return err w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True) + w['perm'] = _worker_perm(wid) return jsonify({'ok': True, 'data': w}) if w else (jsonify({'ok': False, 'error': '不存在'}), 404) if request.method == 'DELETE': + err = _check_worker_perm(wid, 'manage') + if err: + return err + if not (enterprise.current_actor().startswith('token:') or _role() == 'admin'): + return jsonify({'ok': False, 'error': '仅管理员可删除 Worker'}), 403 db.w('UPDATE tasks SET worker_id=NULL WHERE worker_id=?', (wid,)) + db.w('DELETE FROM user_workers WHERE worker_id=?', (wid,)) db.w('DELETE FROM workers WHERE id=?', (wid,)) return jsonify({'ok': True}) + err = _check_worker_perm(wid, 'manage') + if err: + return err + if not (enterprise.current_actor().startswith('token:') or _role() == 'admin'): + return jsonify({'ok': False, 'error': '仅管理员可修改 Worker 配置'}), 403 d = request.get_json(force=True) fields = ['name', 'description', 'provider', 'model', 'base_url', 'api_key', 'system_prompt', 'temperature', 'max_tokens', 'task_cost_limit', @@ -299,6 +466,9 @@ def worker_detail(wid): @app.route('/api/workers//test', methods=['POST']) @require_auth def worker_test(wid): + err = _check_worker_perm(wid, 'use') + if err: + return err w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True) if not w: return jsonify({'ok': False, 'error': '不存在'}), 404 @@ -315,6 +485,9 @@ def worker_test(wid): @require_auth def worker_vision_test(wid): """视觉智能体能力测试:传入图片(URL 或 base64)与问题,验证多模态理解""" + err = _check_worker_perm(wid, 'use') + if err: + return err w = db.q('SELECT * FROM workers WHERE id=?', (wid,), one=True) if not w: return jsonify({'ok': False, 'error': '不存在'}), 404 @@ -363,10 +536,17 @@ def providers(): def tasks(): if request.method == 'POST': d = request.get_json(force=True) + err = _check_project_perm(d.get('project_id'), 'manage') + if err: + return err + if d.get('worker_id'): + werr = _check_worker_perm(d['worker_id'], 'use') + if werr: + return werr if d.get('depends_on'): - ok, err = _validate_depends_on(d.get('project_id'), None, d['depends_on']) + ok, errmsg = _validate_depends_on(d.get('project_id'), None, d['depends_on']) if not ok: - return jsonify({'ok': False, 'error': err}), 400 + return jsonify({'ok': False, 'error': errmsg}), 400 tid = db.w( 'INSERT INTO tasks (project_id, worker_id, title, description, priority, ' 'review_required, deadline, depends_on, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?)', @@ -379,9 +559,15 @@ def tasks(): return jsonify({'ok': True, 'id': tid}) pid = request.args.get('project_id') if pid: + err = _check_project_perm(int(pid), 'view') + if err: + return err rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id DESC', (int(pid),)) else: rows = db.q('SELECT * FROM tasks WHERE deleted=0 ORDER BY id DESC LIMIT 200') + visible = _visible_project_ids() + if visible is not None: + rows = [r for r in rows if r['project_id'] in visible] return jsonify({'ok': True, 'data': [db.serialize_task(t) for t in rows]}) @@ -392,6 +578,9 @@ def task_detail(tid): if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 if request.method == 'GET': + err = _check_task_perm(tid, 'view') + if err: + return err t = db.serialize_task(t) t['logs'] = db.q('SELECT * FROM task_logs WHERE task_id=? ORDER BY id', (tid,)) t['costs'] = db.q('SELECT * FROM cost_records WHERE task_id=? ORDER BY id', (tid,)) @@ -399,6 +588,9 @@ def task_detail(tid): t['worker'] = db.q('SELECT id,name,provider,model FROM workers WHERE id=?', (t['worker_id'],), one=True) return jsonify({'ok': True, 'data': t}) + err = _check_task_perm(tid, 'manage') + if err: + return err if request.method == 'DELETE': # 软删除(进回收站);?hard=1 物理删除 if request.args.get('hard') == '1': @@ -437,10 +629,16 @@ def tasks_trash(): """回收站:已软删除的任务列表""" pid = request.args.get('project_id') if pid: + err = _check_project_perm(int(pid), 'view') + if err: + return err rows = db.q('SELECT * FROM tasks WHERE deleted=1 AND project_id=? ORDER BY deleted_at DESC', (int(pid),)) else: rows = db.q('SELECT * FROM tasks WHERE deleted=1 ORDER BY deleted_at DESC LIMIT 200') + visible = _visible_project_ids() + if visible is not None: + rows = [r for r in rows if r['project_id'] in visible] return jsonify({'ok': True, 'data': [db.serialize_task(t) for t in rows]}) @@ -451,6 +649,9 @@ def task_trash(tid): t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_task_perm(tid, 'manage') + if err: + return err if t['status'] == 'running': return jsonify({'ok': False, 'error': '任务执行中,禁止删除'}), 400 db.w('UPDATE tasks SET deleted=1, deleted_at=? WHERE id=?', (db.now(), tid)) @@ -468,6 +669,9 @@ def task_restore(tid): t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_task_perm(tid, 'manage') + if err: + return err db.w('UPDATE tasks SET deleted=0, deleted_at=NULL WHERE id=?', (tid,)) db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)', (tid, 'info', '♻️ 任务已从回收站恢复', db.now())) @@ -482,6 +686,9 @@ def task_run(tid): t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_task_perm(tid, 'manage') + if err: + return err if t['deleted']: return jsonify({'ok': False, 'error': '任务已在回收站,请先恢复'}), 400 if t['status'] == 'running': @@ -502,6 +709,9 @@ def task_cancel(tid): t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_task_perm(tid, 'manage') + if err: + return err if t['status'] != 'running': return jsonify({'ok': False, 'error': '仅执行中的任务可取消'}), 400 # MVP:标记取消(线程无法强杀,完成后会回到 review;这里直接置 cancelled 并忽略结果) @@ -518,6 +728,9 @@ def task_review(tid): t = db.q('SELECT * FROM tasks WHERE id=?', (tid,), one=True) if not t: return jsonify({'ok': False, 'error': '任务不存在'}), 404 + err = _check_task_perm(tid, 'manage') + if err: + return err if t['status'] != 'review': return jsonify({'ok': False, 'error': '仅待审核状态可审核'}), 400 d = request.get_json(force=True) @@ -534,6 +747,14 @@ def task_review(tid): 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())) + # V3:项目全部任务完成后自动收尾 → 通知送达者 + try: + r = delivery.auto_complete_if_ready(t['project_id'], request.host_url.rstrip('/')) + if r: + db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)', + (0, 'success', f'🎉 项目全部任务完成,已自动交付并通知送达者({r["msg"]})', db.now())) + except Exception: + pass return jsonify({'ok': True}) if action == 'reject': if not reason: @@ -554,22 +775,32 @@ def task_review(tid): @app.route('/api/reports/cost') @require_auth def report_cost(): + visible = _visible_project_ids() + scope_sql, scope_args = '', [] + if visible is not None: + if not visible: + return jsonify({'ok': True, 'data': []}) + scope_sql = 'WHERE project_id IN (%s)' % ','.join('?' * len(visible)) + scope_args = list(visible) group = request.args.get('group', 'project') if group == 'worker': rows = db.q( 'SELECT worker_id, provider, model, COUNT(*) runs, SUM(total_tokens) tokens, ' - 'SUM(cost) cost FROM cost_records GROUP BY worker_id ORDER BY cost DESC') + f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY worker_id ORDER BY cost DESC', + scope_args) for r in rows: w = db.q('SELECT name FROM workers WHERE id=?', (r['worker_id'],), one=True) r['worker_name'] = w['name'] if w else f'#{r["worker_id"]}' elif group == 'model': rows = db.q( 'SELECT provider, model, COUNT(*) runs, SUM(total_tokens) tokens, ' - 'SUM(cost) cost FROM cost_records GROUP BY model ORDER BY cost DESC') + f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY model ORDER BY cost DESC', + scope_args) else: rows = db.q( 'SELECT project_id, COUNT(*) runs, SUM(total_tokens) tokens, ' - 'SUM(cost) cost FROM cost_records GROUP BY project_id ORDER BY cost DESC') + f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY project_id ORDER BY cost DESC', + scope_args) for r in rows: p = db.q('SELECT name FROM projects WHERE id=?', (r['project_id'],), one=True) r['project_name'] = p['name'] if p else f'#{r["project_id"]}' @@ -579,31 +810,60 @@ def report_cost(): @app.route('/api/stats') @require_auth def stats(): + visible = _visible_project_ids() + scope_sql, scope_args = '', [] + if visible is not None: + if not visible: + return jsonify({'ok': True, 'data': { + 'projects': 0, 'tasks': 0, 'workers': 0, 'total_cost': 0, 'total_tokens': 0, + 'by_status': {}, 'recent': [], 'daily_cost': [], 'one_pass_rate': None, + 'rework_count': 0, 'visible_workers': 0}}) + scope_sql = 'WHERE project_id IN (%s)' % ','.join('?' * len(visible)) + scope_args = list(visible) out = {'projects': 0, 'tasks': 0, 'workers': 0, 'total_cost': 0, 'total_tokens': 0, - 'by_status': {}, 'recent': [], 'daily_cost': []} - out['projects'] = db.q('SELECT COUNT(*) c FROM projects')[0]['c'] - out['workers'] = db.q('SELECT COUNT(*) c FROM workers')[0]['c'] - out['tasks'] = db.q('SELECT COUNT(*) c FROM tasks WHERE deleted=0')[0]['c'] - for r in db.q('SELECT status, COUNT(*) c FROM tasks WHERE deleted=0 GROUP BY status'): - out['by_status'][r['status']] = r['c'] - c = db.q('SELECT COALESCE(SUM(cost),0) cost, COALESCE(SUM(total_tokens),0) tokens ' - 'FROM cost_records')[0] + 'by_status': {}, 'recent': [], 'daily_cost': [], 'one_pass_rate': None, + 'rework_count': 0, 'visible_workers': 0} + if visible is None: + out['projects'] = db.q('SELECT COUNT(*) c FROM projects')[0]['c'] + out['workers'] = db.q('SELECT COUNT(*) c FROM workers')[0]['c'] + out['tasks'] = db.q('SELECT COUNT(*) c FROM tasks WHERE deleted=0')[0]['c'] + for r in db.q('SELECT status, COUNT(*) c FROM tasks WHERE deleted=0 GROUP BY status'): + out['by_status'][r['status']] = r['c'] + else: + out['projects'] = len(visible) + vw = _visible_worker_ids() + out['workers'] = len(vw) if vw is not None else db.q('SELECT COUNT(*) c FROM workers')[0]['c'] + out['visible_workers'] = out['workers'] + for r in db.q(f'SELECT project_id, status, COUNT(*) c FROM tasks WHERE deleted=0 ' + f'AND project_id IN ({scope_args_ph(visible)}) GROUP BY project_id, status', visible): + out['by_status'][r['status']] = out['by_status'].get(r['status'], 0) + r['c'] + out['tasks'] = sum(out['by_status'].values()) + c = db.q(f'SELECT COALESCE(SUM(cost),0) cost, COALESCE(SUM(total_tokens),0) tokens ' + f'FROM cost_records {scope_sql}', scope_args)[0] out['total_cost'], out['total_tokens'] = round(c['cost'], 4), c['tokens'] # 一次通过率:done 且 rejection_count=0 done = out['by_status'].get('done', 0) - clean = db.q('SELECT COUNT(*) c FROM tasks WHERE status="done" AND rejection_count=0 AND deleted=0')[0]['c'] + if visible is None: + clean = db.q('SELECT COUNT(*) c FROM tasks WHERE status="done" AND rejection_count=0 AND deleted=0')[0]['c'] + out['rework_count'] = db.q('SELECT COALESCE(SUM(rejection_count),0) c FROM tasks WHERE deleted=0')[0]['c'] + else: + clean = db.q(f'SELECT COUNT(*) c FROM tasks WHERE status="done" AND rejection_count=0 AND deleted=0 ' + f'AND project_id IN ({scope_args_ph(visible)})', visible)[0]['c'] + out['rework_count'] = db.q(f'SELECT COALESCE(SUM(rejection_count),0) c FROM tasks WHERE deleted=0 ' + f'AND project_id IN ({scope_args_ph(visible)})', visible)[0]['c'] out['one_pass_rate'] = round(clean / done * 100, 1) if done else None - out['rework_count'] = db.q('SELECT COALESCE(SUM(rejection_count),0) c FROM tasks WHERE deleted=0')[0]['c'] out['recent'] = db.q( 'SELECT t.id, t.title, t.status, t.project_id, p.name AS project_name, t.updated_at ' 'FROM tasks t LEFT JOIN projects p ON p.id=t.project_id ' - 'WHERE t.deleted=0 ORDER BY t.updated_at DESC LIMIT 10') + 'WHERE t.deleted=0 ' + + (f'AND t.project_id IN ({scope_args_ph(visible)}) ' if visible is not None else '') + + 'ORDER BY t.updated_at DESC LIMIT 10', visible or ()) # 近 7 天成本 import datetime - rows = db.q('SELECT created_at, cost FROM cost_records ORDER BY created_at') + rows = db.q(f'SELECT created_at, cost FROM cost_records {scope_sql} ORDER BY created_at', scope_args) day_map = {} for r in rows: d = datetime.datetime.fromtimestamp(r['created_at']).strftime('%m-%d') @@ -614,13 +874,28 @@ def stats(): return jsonify({'ok': True, 'data': out}) +def scope_args_ph(ids): + """生成 ? 占位符串""" + return ','.join('?' * len(ids)) if ids else 'NULL' + + @app.route('/api/logs') @require_auth def logs(): limit = min(int(request.args.get('limit', 100)), 500) - rows = db.q( - 'SELECT l.*, t.title AS task_title FROM task_logs l ' - 'LEFT JOIN tasks t ON t.id=l.task_id ORDER BY l.id DESC LIMIT ?', (limit,)) + visible = _visible_project_ids() + if visible is None: + rows = db.q( + 'SELECT l.*, t.title AS task_title FROM task_logs l ' + 'LEFT JOIN tasks t ON t.id=l.task_id ORDER BY l.id DESC LIMIT ?', (limit,)) + else: + if not visible: + return jsonify({'ok': True, 'data': []}) + rows = db.q( + f'SELECT l.*, t.title AS task_title FROM task_logs l ' + f'LEFT JOIN tasks t ON t.id=l.task_id ' + f'WHERE (t.project_id IN ({scope_args_ph(visible)}) OR t.id IS NULL) ' + f'ORDER BY l.id DESC LIMIT ?', (*visible, limit)) return jsonify({'ok': True, 'data': rows}) @@ -628,6 +903,9 @@ def logs(): @require_auth def workflow_run(pid): """执行整个工作流:跑所有就绪(无未完成前置)任务""" + err = _check_project_perm(pid, 'manage') + if err: + return err rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 AND status IN ("todo","failed")', (pid,)) started, blocked = [], [] for t in rows: @@ -644,6 +922,9 @@ def workflow_run(pid): @require_auth def project_dag(pid): """DAG 图数据:节点 + 边""" + err = _check_project_perm(pid, 'view') + if err: + return err rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id', (pid,)) nodes, edges, id_map = [], [], {} for t in rows: @@ -690,6 +971,9 @@ def _parse_wbs(text): @app.route('/api/projects//wbs/generate', methods=['POST']) @require_auth def wbs_generate(pid): + err = _check_project_perm(pid, 'manage') + if err: + return err d = request.get_json(force=True) or {} goal = d.get('goal') or '' proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) @@ -699,9 +983,15 @@ def wbs_generate(pid): 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) + worker = None + vw = _visible_worker_ids() + if vw is None: + worker = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id', one=True) + else: + candidates = [r for r in db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id') if r['id'] in vw] + worker = candidates[0] if candidates else None if not worker: - return jsonify({'ok': False, 'error': '请先注册至少一个 Worker 用于规划'}), 400 + return jsonify({'ok': False, 'error': '请先注册至少一个 Worker 用于规划(或联系管理员授权)'}), 400 try: r = llm_gateway.chat(worker['provider'], worker['model'], [ {'role': 'system', 'content': '你只输出 JSON,不输出任何解释文字。'}, @@ -716,6 +1006,9 @@ def wbs_generate(pid): @app.route('/api/projects//wbs/import', methods=['POST']) @require_auth def wbs_import(pid): + err = _check_project_perm(pid, 'manage') + if err: + return err d = request.get_json(force=True) tasks = d.get('tasks') or [] worker_id = d.get('worker_id') @@ -744,6 +1037,9 @@ def wbs_import(pid): @require_auth def documents(pid): if request.method == 'POST': + err = _check_project_perm(pid, 'manage') + if err: + return err d = request.get_json(force=True) doc_id = db.w( 'INSERT INTO documents (project_id, name, content, source, created_at, updated_at) ' @@ -752,6 +1048,9 @@ def documents(pid): d.get('source', 'manual'), db.now(), db.now())) n = rag.rebuild_document(doc_id) return jsonify({'ok': True, 'id': doc_id, 'chunks': n}) + err = _check_project_perm(pid, 'view') + if err: + return err 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'] @@ -764,6 +1063,9 @@ 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 + err = _check_project_perm(doc['project_id'], 'view' if request.method == 'GET' else 'manage') + if err: + return err if request.method == 'GET': return jsonify({'ok': True, 'data': doc}) if request.method == 'DELETE': @@ -782,6 +1084,9 @@ def document_detail(doc_id): @app.route('/api/projects//search') @require_auth def kb_search(pid): + err = _check_project_perm(pid, 'view') + if err: + return err q = request.args.get('q', '') if not q: return jsonify({'ok': True, 'data': []}) @@ -789,6 +1094,228 @@ def kb_search(pid): return jsonify({'ok': True, 'data': hits if hit else [], 'hit': hit}) +# --------------------------------------------------------------------------- +# V3 · 交付体系:工作目录 / Demo 部署 / 打包 / 邮件送达 +# --------------------------------------------------------------------------- +@app.route('/api/projects//deliverables') +@require_auth +def project_deliverables(pid): + """交付记录列表""" + err = _check_project_perm(pid, 'view') + if err: + return err + rows = db.q('SELECT * FROM project_deliverables WHERE project_id=? ORDER BY id DESC', (pid,)) + return jsonify({'ok': True, 'data': rows}) + + +@app.route('/api/projects//workspace') +@require_auth +def project_workspace_list(pid): + """项目工作目录文件列表""" + err = _check_project_perm(pid, 'view') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + files = delivery.list_workspace(proj) + return jsonify({'ok': True, 'data': files, 'dir': f'data/workspace/project_{pid}/'}) + + +@app.route('/api/projects//workspace/upload', methods=['POST']) +@require_auth +def project_workspace_upload(pid): + """上传交付物文件到项目工作目录(multipart,支持多文件)""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if 'files' not in request.files: + return jsonify({'ok': False, 'error': '缺少文件字段 files'}), 400 + saved = [] + try: + for fs in request.files.getlist('files'): + if fs and fs.filename: + rel = delivery.save_upload(proj, fs) + saved.append(rel) + except Exception as e: + return jsonify({'ok': False, 'error': f'上传失败:{e}'}), 400 + enterprise.audit(enterprise.current_actor(), 'deliver.upload', f'project#{pid}', + f'上传 {len(saved)} 个文件:{", ".join(saved)[:200]}', request.remote_addr or '') + return jsonify({'ok': True, 'saved': saved}) + + +@app.route('/api/projects//workspace/download') +@require_auth +def project_workspace_download(pid): + """下载工作目录文件""" + err = _check_project_perm(pid, 'view') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + try: + rel = delivery._safe_relpath(request.args.get('path', '')) + except ValueError as e: + return jsonify({'ok': False, 'error': str(e)}), 400 + return send_from_directory(delivery.workspace_path(proj), rel, as_attachment=True) + + +@app.route('/api/projects//workspace', methods=['DELETE']) +@require_auth +def project_workspace_delete(pid): + """删除工作目录文件""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + try: + rel = delivery.delete_workspace_file(proj, request.args.get('path', '')) + except Exception as e: + return jsonify({'ok': False, 'error': str(e)}), 400 + return jsonify({'ok': True, 'deleted': rel}) + + +@app.route('/api/projects//deploy', methods=['POST']) +@require_auth +def project_deploy(pid): + """网页交付物:部署到 Demo 地址(送达者可直接访问)""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if proj['deliver_type'] != 'web': + return jsonify({'ok': False, 'error': '该项目交付类型不是网页(web),无需部署 Demo'}), 400 + base = request.host_url.rstrip('/') + try: + r = delivery.deploy_demo(proj, base) + except Exception as e: + return jsonify({'ok': False, 'error': f'部署失败:{e}'}), 500 + enterprise.audit(enterprise.current_actor(), 'deliver.deploy', f'project#{pid}', + f'Demo 部署 {r["demo_url"]}({r["copied"]} 个文件)', request.remote_addr or '') + return jsonify({'ok': True, **r}) + + +@app.route('/api/projects//package', methods=['POST']) +@require_auth +def project_package(pid): + """把项目工作目录打包为 zip""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + try: + pkg = delivery.package_project(proj) + except Exception as e: + return jsonify({'ok': False, 'error': f'打包失败:{e}'}), 500 + enterprise.audit(enterprise.current_actor(), 'deliver.package', f'project#{pid}', + f'打包 {pkg["name"]}({pkg["size"]} B)', request.remote_addr or '') + return jsonify({'ok': True, **pkg}) + + +@app.route('/api/projects//deliver', methods=['POST']) +@require_auth +def project_deliver(pid): + """把项目交付物发送给送达者:网页→Demo 链接;文件→打包附件;全部走邮件""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if not (proj.get('deliver_email') or '').strip(): + return jsonify({'ok': False, 'error': '项目未配置送达者邮箱,请先在编辑项目中填写'}), 400 + base = request.host_url.rstrip('/') + # 网页交付物:确保 Demo 已部署 + demo = proj.get('demo_url') or '' + if proj['deliver_type'] == 'web' and not demo: + try: + r = delivery.deploy_demo(proj, base) + demo = r['demo_url'] + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + except Exception as e: + return jsonify({'ok': False, 'error': f'Demo 部署失败:{e}'}), 500 + # 打包附件 + attach = None + try: + pkg = delivery.package_project(proj) + attach = pkg['path'] + except Exception as e: + return jsonify({'ok': False, 'error': f'打包失败:{e}'}), 500 + s = delivery.project_summary(proj) + kind = '网页 Demo' if proj['deliver_type'] == 'web' else '文件包' + extra = (f'本次交付:{kind}\n' + f'任务完成:{s["done"]}/{s["total"]} · 累计成本 ¥{s["cost"]:.4f}\n' + f'交付物已打包为附件({pkg["size"]} B),请查收。') + subject = f'📦 项目交付:{proj["name"]}' + body = delivery._email_body(proj, extra=extra) + ok, msg = delivery.send_mail(proj['deliver_email'], subject, body, + attachments=[attach] if attach else None) + if not ok: + return jsonify({'ok': False, 'error': msg}), 500 + db.w('UPDATE projects SET delivered_at=? WHERE id=?', (db.now(), pid)) + db.w('INSERT INTO project_deliverables (project_id, name, kind, path, demo_url, size, note, created_at) ' + 'VALUES (?,?,?,?,?,?,?,?)', + (pid, f'交付邮件 → {proj["deliver_email"]}', 'email', pkg['name'], demo, + pkg['size'], '手动交付(含附件)', db.now())) + enterprise.audit(enterprise.current_actor(), 'deliver.send', f'project#{pid}', + f'邮件交付 → {proj["deliver_email"]}({pkg["name"]})', request.remote_addr or '') + return jsonify({'ok': True, 'msg': msg, 'demo_url': demo, 'package': pkg['name']}) + + +@app.route('/api/projects//complete', methods=['POST']) +@require_auth +def project_complete(pid): + """完成项目并交付:置 done + 部署 Demo + 打包 + 邮件通知送达者""" + err = _check_project_perm(pid, 'manage') + if err: + return err + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if not (proj.get('deliver_email') or '').strip(): + return jsonify({'ok': False, 'error': '项目未配置送达者邮箱,请先填写'}), 400 + if proj['status'] == 'done': + return jsonify({'ok': False, 'error': '项目已完成,无需重复交付'}), 400 + base = request.host_url.rstrip('/') + ok, msg = delivery.notify_project_complete(proj, base) + if not ok: + return jsonify({'ok': False, 'error': msg}), 500 + db.w('UPDATE projects SET status="done", updated_at=? WHERE id=?', (db.now(), pid)) + enterprise.audit(enterprise.current_actor(), 'project.complete', f'project#{pid}', + f'「{proj["name"]}」完成交付 → {proj["deliver_email"]}', request.remote_addr or '') + return jsonify({'ok': True, 'msg': msg}) + + +@app.route('/api/projects//notify_deliverer', methods=['POST']) +@require_auth +def project_notify_deliverer(pid): + """手动向送达者发送一条消息(如说明难关/进展)""" + err = _check_project_perm(pid, 'manage') + if err: + return err + d = request.get_json(force=True) or {} + proj = db.q('SELECT * FROM projects WHERE id=?', (pid,), one=True) + if not (proj.get('deliver_email') or '').strip(): + return jsonify({'ok': False, 'error': '项目未配置送达者邮箱'}), 400 + subject = (d.get('subject') or f'项目动态:{proj["name"]}').strip() + text = (d.get('message') or '').strip() + if not text: + return jsonify({'ok': False, 'error': '请填写消息内容'}), 400 + ok, msg = delivery.notify_deliverer(proj, subject, delivery._email_body(proj, extra=text)) + if not ok: + return jsonify({'ok': False, 'error': msg}), 500 + return jsonify({'ok': True, 'msg': msg}) + + +@app.route('/demo//') +@app.route('/demo//') +def demo_serve(pid, subpath=''): + """公开 Demo 地址(送达者免登录直接查看网页交付物)""" + proj = db.q('SELECT id FROM projects WHERE id=?', (pid,), one=True) + if not proj: + return '项目不存在', 404 + return send_from_directory(delivery.demo_path(proj), subpath or 'index.html') + + +@app.route('/demo/') +def demo_serve_root(pid): + return redirect(f'/demo/{pid}/') + + # --------------------------------------------------------------------------- # 告警中心 # --------------------------------------------------------------------------- @@ -974,6 +1501,15 @@ def agent_runs(): worker_ids = [r['id'] for r in rows] if not worker_ids: return jsonify({'ok': False, 'error': '请先注册至少一个 Worker'}), 400 + # 成员只能使用自己有权限的 Worker + vw = _visible_worker_ids() + if vw is not None: + allowed = set(vw) + bad = [w for w in worker_ids if w not in allowed] + if bad: + return jsonify({'ok': False, 'error': f'无权使用 Worker:{bad}'}), 403 + if not worker_ids: + return jsonify({'ok': False, 'error': '没有可用的 Worker(请管理员授权)'}), 403 rid = db.w( 'INSERT INTO agent_runs (mode, title, topic, context, worker_ids, params, status, ' 'created_at) VALUES (?,?,?,?,?,?,?,?)', @@ -984,15 +1520,30 @@ def agent_runs(): 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') + vw = _visible_worker_ids() + if vw is not None: + rows = [r for r in rows if _run_visible(r, vw)] return jsonify({'ok': True, 'data': rows}) +def _run_visible(run, visible_worker_ids): + """协作/评估运行是否对当前用户可见:任一参与 Worker 可见即可""" + try: + wids = _json.loads(run.get('worker_ids') or '[]') + except Exception: + wids = [] + return any(wid in visible_worker_ids for wid in wids) + + @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 + vw = _visible_worker_ids() + if vw is not None and not _run_visible(run, vw): + return jsonify({'ok': False, 'error': '无权访问该协作运行'}), 403 if request.method == 'DELETE': db.w('DELETE FROM agent_steps WHERE run_id=?', (rid,)) db.w('DELETE FROM agent_runs WHERE id=?', (rid,)) @@ -1125,6 +1676,9 @@ def eval_runs(): 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') + vw = _visible_worker_ids() + if vw is not None: + rows = [r for r in rows if r['worker_id'] in vw] 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: @@ -1325,6 +1879,43 @@ def ent_user_detail(uid): return jsonify({'ok': True}) +@app.route('/api/enterprise/users//grants', methods=['GET', 'PUT']) +@require_admin +def ent_user_grants(uid): + """用户授权管理:项目权限(view/manage/admin)+ Worker 权限(view/use/manage)""" + u = db.q('SELECT * FROM users WHERE id=?', (uid,), one=True) + if not u: + return jsonify({'ok': False, 'error': '用户不存在'}), 404 + if request.method == 'PUT': + d = request.get_json(force=True) or {} + out = enterprise.set_user_grants(uid, + projects=d.get('projects'), + workers=d.get('workers')) + enterprise.audit(enterprise.current_actor(), 'user.grants', f'user#{uid}', + f'{u["username"]} 授权更新:项目 {out["projects"]} 项 / Worker {out["workers"]} 项', + request.remote_addr or '') + return jsonify({'ok': True, **out}) + data = enterprise.user_grants(uid) + return jsonify({'ok': True, 'data': data}) + + +@app.route('/api/enterprise/grants/overview') +@require_admin +def ent_grants_overview(): + """授权总览:所有用户的 项目/Worker 授权矩阵""" + users = db.q('SELECT id, username, display_name, role FROM users ORDER BY id') + up = db.q('SELECT * FROM user_projects ORDER BY user_id, project_id') + uw = db.q('SELECT * FROM user_workers ORDER BY user_id, worker_id') + pname = {r['id']: r['name'] for r in db.q('SELECT id,name FROM projects')} + wname = {r['id']: r['name'] for r in db.q('SELECT id,name FROM workers')} + for u in users: + u['projects'] = [{'project_id': r['project_id'], 'project_name': pname.get(r['project_id'], f'#{r["project_id"]}'), + 'perm': r['perm']} for r in up if r['user_id'] == u['id']] + u['workers'] = [{'worker_id': r['worker_id'], 'worker_name': wname.get(r['worker_id'], f'#{r["worker_id"]}'), + 'perm': r['perm']} for r in uw if r['user_id'] == u['id']] + return jsonify({'ok': True, 'data': users}) + + @app.route('/api/enterprise/audit') @require_auth def ent_audit(): diff --git a/config.py b/config.py index 0baed97..2094d31 100644 --- a/config.py +++ b/config.py @@ -98,5 +98,8 @@ EMAIL = { 'from_name': 'AI Worker 平台', } +# 公网访问地址(Demo 链接/邮件中的回链基准;留空则用请求 host) +PUBLIC_BASE_URL = os.environ.get('PUBLIC_BASE_URL', 'http://121.40.164.32:16071') + # 自动路由:按模型单价升序挑选可用 Worker AUTO_ROUTE_POOL = 'enabled' # enabled | all diff --git a/db.py b/db.py index a4c49d3..b22b206 100644 --- a/db.py +++ b/db.py @@ -16,6 +16,12 @@ CREATE TABLE IF NOT EXISTS projects ( acceptance_criteria TEXT DEFAULT '', status TEXT DEFAULT 'active', -- planning/active/done/archived budget_limit REAL DEFAULT 0, -- 项目预算上限(元),0=不限 + deliver_email TEXT DEFAULT '', -- 送达者(人)邮箱,新建项目必填 + deliver_type TEXT DEFAULT 'web', -- 交付物类型 web=网页 / file=文件包 + deliver_note TEXT DEFAULT '', -- 交付说明 + workspace_dir TEXT DEFAULT '', -- 项目工作目录(相对 data/ 的目录名) + demo_url TEXT DEFAULT '', -- 网页交付物 Demo 访问地址 + delivered_at INTEGER, -- 最近一次交付/送达时间 created_at INTEGER, updated_at INTEGER ); @@ -258,6 +264,44 @@ CREATE TABLE IF NOT EXISTS enterprise_settings ( value TEXT DEFAULT '' ); +-- =================================================================== +-- V3 表结构:交付体系(工作目录/交付物/Demo/邮件送达) + 用户授权(项目/Worker 权限) +-- =================================================================== +CREATE TABLE IF NOT EXISTS project_deliverables ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + project_id INTEGER NOT NULL, + name TEXT NOT NULL, + kind TEXT DEFAULT 'file', -- file/dir/webpage/package + path TEXT DEFAULT '', -- 相对项目工作目录路径 / 打包文件名 + demo_url TEXT DEFAULT '', -- 网页交付物的 Demo 访问地址 + size INTEGER DEFAULT 0, + note TEXT DEFAULT '', + created_at INTEGER +); + +CREATE TABLE IF NOT EXISTS user_projects ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER NOT NULL, + project_id INTEGER NOT NULL, + perm TEXT DEFAULT 'view', -- view 查看 / manage 管理 / admin 管理员 + created_at INTEGER, + UNIQUE(user_id, project_id) +); + +CREATE TABLE IF NOT EXISTS user_workers ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + user_id INTEGER NOT NULL, + worker_id INTEGER NOT NULL, + perm TEXT DEFAULT 'view', -- view 查看 / use 使用(可指派任务)/ manage 管理(可改配置) + created_at INTEGER, + UNIQUE(user_id, worker_id) +); + +CREATE INDEX IF NOT EXISTS idx_deliverables_project ON project_deliverables(project_id); +CREATE INDEX IF NOT EXISTS idx_user_projects_user ON user_projects(user_id); +CREATE INDEX IF NOT EXISTS idx_user_projects_project ON user_projects(project_id); +CREATE INDEX IF NOT EXISTS idx_user_workers_user ON user_workers(user_id); +CREATE INDEX IF NOT EXISTS idx_user_workers_worker ON user_workers(worker_id); 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); @@ -286,6 +330,22 @@ def _migrate(): conn.execute('ALTER TABLE tasks ADD COLUMN deleted INTEGER DEFAULT 0') conn.execute('ALTER TABLE tasks ADD COLUMN deleted_at INTEGER') conn.execute('CREATE INDEX IF NOT EXISTS idx_tasks_deleted ON tasks(deleted)') + # V3:projects 交付字段 + pcols = {r['name'] for r in conn.execute('PRAGMA table_info(projects)')} + for col, ddl in ( + ('deliver_email', "ALTER TABLE projects ADD COLUMN deliver_email TEXT DEFAULT ''"), + ('deliver_type', "ALTER TABLE projects ADD COLUMN deliver_type TEXT DEFAULT 'web'"), + ('deliver_note', "ALTER TABLE projects ADD COLUMN deliver_note TEXT DEFAULT ''"), + ('workspace_dir', "ALTER TABLE projects ADD COLUMN workspace_dir TEXT DEFAULT ''"), + ('demo_url', "ALTER TABLE projects ADD COLUMN demo_url TEXT DEFAULT ''"), + ('delivered_at', 'ALTER TABLE projects ADD COLUMN delivered_at INTEGER'), + ): + if col not in pcols: + conn.execute(ddl) + # V3:老项目补齐工作目录名 + for r in conn.execute("SELECT id, workspace_dir FROM projects WHERE workspace_dir IS NULL OR workspace_dir=''"): + conn.execute('UPDATE projects SET workspace_dir=? WHERE id=?', + ('project_%d' % r['id'], r['id'])) conn.commit() conn.close() diff --git a/delivery.py b/delivery.py new file mode 100644 index 0000000..64add1f --- /dev/null +++ b/delivery.py @@ -0,0 +1,369 @@ +# -*- coding: utf-8 -*- +""" +V3 交付体系 +- 每个项目独立工作目录:data/workspace/project_/(中间产物与交付物隔离存放) +- 网页交付物 → data/demo// 部署,经 /demo// 公开访问(送达者无需登录) +- 文件包交付物 → zip 打包到 data/packages/,随邮件附件发送 +- 送达者通知:项目完成 / 遇到无法绕开的难关时,邮件及时通知 deliver_email +""" +import os +import re +import time +import shutil +import zipfile +import smtplib +from email.mime.multipart import MIMEMultipart +from email.mime.text import MIMEText +from email.mime.application import MIMEApplication +from email.utils import formataddr + +import db +from config import DATA_DIR, EMAIL, PUBLIC_BASE_URL + +WORKSPACE_ROOT = os.path.join(DATA_DIR, 'workspace') +DEMO_ROOT = os.path.join(DATA_DIR, 'demo') +PACKAGE_ROOT = os.path.join(DATA_DIR, 'packages') + +# 项目维度 blocker 通知去重窗口(秒):同一项目短时间内不重复打扰送达者 +BLOCKER_DEDUP_SECONDS = 1800 + + +def ensure_dirs(): + for d in (WORKSPACE_ROOT, DEMO_ROOT, PACKAGE_ROOT): + os.makedirs(d, exist_ok=True) + + +def workspace_path(project): + """项目工作目录绝对路径(不存在则创建)""" + pid = project['id'] if isinstance(project, dict) else project + d = os.path.join(WORKSPACE_ROOT, f'project_{pid}') + os.makedirs(d, exist_ok=True) + return d + + +def demo_path(project): + """Demo 部署目录绝对路径""" + pid = project['id'] if isinstance(project, dict) else project + return os.path.join(DEMO_ROOT, f'project_{pid}') + + +def package_dir(): + os.makedirs(PACKAGE_ROOT, exist_ok=True) + return PACKAGE_ROOT + + +def _safe_relpath(relpath): + """路径穿越防护:仅允许工作目录内的相对路径""" + relpath = (relpath or '').replace('\\', '/').strip('/') + if not relpath: + return '' + if '..' in relpath.split('/') or relpath.startswith('/'): + raise ValueError('非法路径') + return relpath + + +def list_workspace(project): + """递归列出工作目录文件:相对路径 + 类型 + 大小 + 修改时间""" + root = workspace_path(project) + out = [] + for dirpath, dirnames, filenames in os.walk(root): + # 忽略临时目录 + dirnames[:] = [d for d in dirnames if not d.startswith('.')] + for fn in sorted(filenames): + if fn.startswith('.'): + continue + full = os.path.join(dirpath, fn) + rel = os.path.relpath(full, root).replace(os.sep, '/') + try: + size = os.path.getsize(full) + mtime = int(os.path.getmtime(full)) + except OSError: + size, mtime = 0, 0 + out.append({'path': rel, 'name': fn, 'size': size, + 'ext': os.path.splitext(fn)[1].lstrip('.').lower(), + 'mtime': mtime}) + out.sort(key=lambda x: x['path']) + return out + + +def save_upload(project, file_storage, subdir=''): + """保存上传文件到工作目录,返回相对路径""" + fn = os.path.basename(file_storage.filename or '') + fn = re.sub(r'[\\/:*?"<>|]', '_', fn).strip() + if not fn: + raise ValueError('文件名为空') + rel = _safe_relpath(subdir) + target_dir = os.path.join(workspace_path(project), rel) if rel else workspace_path(project) + os.makedirs(target_dir, exist_ok=True) + target = os.path.join(target_dir, fn) + file_storage.save(target) + return (rel + '/' if rel else '') + fn + + +def delete_workspace_file(project, relpath): + rel = _safe_relpath(relpath) + if not rel: + raise ValueError('请指定要删除的文件') + full = os.path.join(workspace_path(project), rel) + if not os.path.isfile(full): + raise ValueError('文件不存在') + os.remove(full) + return rel + + +def demo_url_of(project, base_url=''): + """生成 Demo 访问地址""" + base = (base_url or PUBLIC_BASE_URL).rstrip('/') + return f'{base}/demo/{project["id"]}/' + + +def deploy_demo(project, base_url=''): + """把项目工作目录部署为可公开访问的 Demo(网页交付物) + - 将工作目录文件复制到 data/demo/project_/ + - 无 index.html 时生成一个简易索引页 + - 记录 demo_url 到项目 + """ + src = workspace_path(project) + dst = demo_path(project) + os.makedirs(dst, exist_ok=True) + # 清空旧内容,避免残留文件污染 + for item in os.listdir(dst): + p = os.path.join(dst, item) + if os.path.isdir(p): + shutil.rmtree(p, ignore_errors=True) + else: + os.remove(p) + copied = 0 + for dirpath, dirnames, filenames in os.walk(src): + dirnames[:] = [d for d in dirnames if not d.startswith('.')] + rel = os.path.relpath(dirpath, src) + if rel == '.': + rel = '' + for fn in filenames: + if fn.startswith('.') or fn.endswith('.zip'): + continue + sub = os.path.join(dst, rel) if rel else dst + os.makedirs(sub, exist_ok=True) + shutil.copy2(os.path.join(dirpath, fn), os.path.join(sub, fn)) + copied += 1 + index = os.path.join(dst, 'index.html') + if not os.path.isfile(index): + files = sorted(list_workspace(project), key=lambda x: x['path']) + links = '\n'.join( + f'
  • {os.path.basename(f["path"])}' + f' ({f["size"]} B)
  • ' + for f in files if f['ext'] in ('html', 'htm') or '/' not in f['path']) + if not links: + links = '
  • (工作目录中暂无网页文件)
  • ' + with open(index, 'w', encoding='utf-8') as fh: + fh.write(f''' + +{project['name']} · Demo + +

    📦 {project['name']} · 交付 Demo

    +

    本页面由 AI Worker 平台自动生成,展示项目工作目录中的交付文件:

    +
      {links}
    ''') + url = demo_url_of(project, base_url) + db.w('UPDATE projects SET demo_url=?, updated_at=? WHERE id=?', (url, db.now(), project['id'])) + return {'copied': copied, 'demo_url': url} + + +def package_project(project, name=''): + """把项目工作目录打包为 zip,落盘到 data/packages/,返回 {path, size, relname}""" + ensure_dirs() + src = workspace_path(project) + ts = time.strftime('%Y%m%d_%H%M%S') + base = name or f'project_{project["id"]}_deliverable' + zip_name = f'{base}_{ts}.zip' + zip_path = os.path.join(PACKAGE_ROOT, zip_name) + with zipfile.ZipFile(zip_path, 'w', zipfile.ZIP_DEFLATED) as zf: + for dirpath, dirnames, filenames in os.walk(src): + dirnames[:] = [d for d in dirnames if not d.startswith('.')] + for fn in filenames: + if fn.startswith('.'): + continue + full = os.path.join(dirpath, fn) + rel = os.path.relpath(full, src) + zf.write(full, os.path.join(os.path.basename(src), rel)) + size = os.path.getsize(zip_path) + db.w('INSERT INTO project_deliverables (project_id, name, kind, path, size, note, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (project['id'], zip_name, 'package', zip_name, size, + '交付物打包(zip)', db.now())) + return {'path': zip_path, 'size': size, 'name': zip_name} + + +def record_file_deliverables(project, files): + """把工作目录文件登记为交付物记录""" + for f in files: + db.w('INSERT INTO project_deliverables (project_id, name, kind, path, size, note, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (project['id'], f['name'], 'file', f['path'], f['size'], '工作目录交付物', db.now())) + + +def auto_complete_if_ready(project_id, base_url=''): + """项目全部任务完成后自动收尾:置 done + 打包 + 通知送达者。 + 返回 {'ok','msg'} 或 None(未满足条件/异常)。引擎线程与审核接口共用。""" + try: + proj = db.q('SELECT * FROM projects WHERE id=?', (project_id,), one=True) + if not proj or proj['status'] == 'done': + return None + total = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0', (project_id,))[0]['c'] + done = db.q('SELECT COUNT(*) c FROM tasks WHERE project_id=? AND status="done" AND deleted=0', + (project_id,))[0]['c'] + if total == 0 or done < total: + return None + db.w('UPDATE projects SET status="done", updated_at=? WHERE id=?', (db.now(), project_id)) + ok, msg = notify_project_complete(proj, base_url) + return {'ok': ok, 'msg': msg} + except Exception: + return None + + +def project_summary(project): + """项目交付摘要:任务统计 + 成本""" + rows = db.q('SELECT status, COUNT(*) c FROM tasks WHERE project_id=? AND deleted=0 GROUP BY status', + (project['id'],)) + by = {r['status']: r['c'] for r in rows} + total = sum(by.values()) + done = by.get('done', 0) + cost = db.q('SELECT COALESCE(SUM(cost),0) t FROM cost_records WHERE project_id=?', + (project['id'],))[0]['t'] + return {'total': total, 'done': done, 'failed': by.get('failed', 0), + 'review': by.get('review', 0), 'cost': round(cost, 4)} + + +# --------------------------------------------------------------------------- +# 邮件发送(支持附件) +# --------------------------------------------------------------------------- +def send_mail(to_addr, subject, text, attachments=None, html=None): + """发送邮件到任意收件人(送达者),支持附件。返回 (ok, msg)""" + if not EMAIL.get('host'): + return False, '邮件服务未配置(config.EMAIL.host 为空)' + if not to_addr: + return False, '收件邮箱为空' + msg = MIMEMultipart() + msg['From'] = formataddr((EMAIL.get('from_name', 'AI Worker 平台'), EMAIL['user'])) + msg['To'] = to_addr + msg['Subject'] = subject + if html: + msg.attach(MIMEText(html, 'html', 'utf-8')) + else: + msg.attach(MIMEText(text, 'plain', 'utf-8')) + for f in (attachments or []): + if not f or not os.path.isfile(f): + continue + with open(f, 'rb') as fh: + subtype = os.path.splitext(f)[1].lstrip('.').lower() or 'octet-stream' + part = MIMEApplication(fh.read(), _subtype=subtype) + part.add_header('Content-Disposition', 'attachment', + filename=('utf-8', '', os.path.basename(f))) + msg.attach(part) + try: + s = smtplib.SMTP(EMAIL['host'], EMAIL['port'], timeout=30) + if EMAIL.get('starttls'): + s.starttls() + if EMAIL.get('user'): + s.login(EMAIL['user'], EMAIL['password']) + s.sendmail(EMAIL['user'], [to_addr], msg.as_string()) + s.quit() + return True, '已发送' + except Exception as e: + return False, f'邮件发送失败: {e}' + + +def _email_body(project, extra=''): + p = project + lines = [ + f'项目名称:{p["name"]}', + f'项目状态:{ {"planning":"规划中","active":"进行中","done":"已完成","archived":"已归档"}.get(p["status"], p["status"]) }', + f'项目目标:{p.get("objective") or "—"}', + ] + if p.get('demo_url'): + lines.append(f'在线 Demo(可直接打开查看):{p["demo_url"]}') + if extra: + lines.append('') + lines.append(extra) + lines.append('') + lines.append('—— 来自 AI Worker 项目管理平台') + return '\n'.join(lines) + + +def notify_deliverer(project, subject, text, attach=None, base_url=''): + """发邮件给送达者,写交付记录。返回 (ok, msg)""" + email = (project.get('deliver_email') or '').strip() + if not email: + return False, '项目未配置送达者邮箱' + ok, msg = send_mail(email, subject, text, attachments=[attach] if attach else None) + if ok: + db.w('UPDATE projects SET delivered_at=? WHERE id=?', (db.now(), project['id'])) + return ok, msg + + +def notify_blocker(project, task_title, detail): + """项目遇到无法绕开的难关 → 及时邮件通知送达者(同项目限频防打扰)""" + email = (project.get('deliver_email') or '').strip() + if not email: + return + # 去重:同项目 30 分钟内只提醒一次(detail 带项目标记) + dup = db.q("SELECT COUNT(*) c FROM alerts WHERE type='blocker' AND detail LIKE ? AND created_at>?", + (f'[project:{project["id"]}]%', db.now() - BLOCKER_DEDUP_SECONDS)) + if dup and dup[0]['c'] > 0: + return + db.w("INSERT INTO alerts (type, level, title, detail, read, created_at) " + "VALUES ('blocker','warn',?,?,0,?)", + (f'项目难关:{project["name"]} · {task_title}', + f'[project:{project["id"]}] {detail[:500]}', db.now())) + subject = f'⚠️ 项目遇到难关:{project["name"]}' + body = _email_body(project, extra=f'任务「{task_title}」遇到无法绕开的难关:\n{detail[:800]}') + ok, msg = send_mail(email, subject, body) + if ok: + db.w('UPDATE projects SET delivered_at=? WHERE id=?', (db.now(), project['id'])) + else: + # 邮件失败也留痕 + db.w("INSERT INTO alerts (type, level, title, detail, read, created_at) " + "VALUES ('notify','warn',?,?,0,?)", + (f'难关通知邮件发送失败:{project["name"]}', msg[:300], db.now())) + + +def notify_project_complete(project, base_url=''): + """项目完成 → 打包 + 邮件送达(含 Demo 链接与附件)。返回 (ok, msg)""" + email = (project.get('deliver_email') or '').strip() + if not email: + return False, '项目未配置送达者邮箱' + # 网页交付物:确保已部署 Demo + if project.get('deliver_type') == 'web' and not project.get('demo_url'): + try: + deploy_demo(project, base_url) + project = db.q('SELECT * FROM projects WHERE id=?', (project['id'],), one=True) + except Exception as e: + pass + # 打包工作目录 + attach = None + try: + pkg = package_project(project) + attach = pkg['path'] + except Exception as e: + pkg = None + s = project_summary(project) + extra = (f'项目已完成 ✅\n任务完成情况:{s["done"]}/{s["total"]}(失败 {s["failed"]})\n' + f'累计成本:¥{s["cost"]:.4f}\n交付物打包:{"已生成附件(见邮件附件)" if attach else "无工作目录文件"}') + subject = f'✅ 项目完成交付:{project["name"]}' + body = _email_body(project, extra=extra) + ok, msg = send_mail(email, subject, body, attachments=[attach] if attach else None) + if ok: + db.w('UPDATE projects SET delivered_at=? WHERE id=?', (db.now(), project['id'])) + db.w('INSERT INTO project_deliverables (project_id, name, kind, path, demo_url, size, note, created_at) ' + 'VALUES (?,?,?,?,?,?,?,?)', + (project['id'], f'完成交付邮件 → {email}', 'email', + project.get('demo_url') or '', project.get('demo_url') or '', + pkg['size'] if pkg else 0, '项目完成通知(含附件)', db.now())) + else: + db.w("INSERT INTO alerts (type, level, title, detail, read, created_at) " + "VALUES ('notify','warn',?,?,0,?)", + (f'完成交付邮件失败:{project["name"]}', msg[:300], db.now())) + return ok, msg + + +ensure_dirs() diff --git a/engine.py b/engine.py index aec4db9..f22f8c7 100644 --- a/engine.py +++ b/engine.py @@ -14,6 +14,20 @@ import llm_gateway import config import rag import notify +import delivery + + +def _notify_failed(task, message): + """任务失败:写告警 + 推送渠道 + 邮件通知送达者(项目难关)""" + notify.notify('task_failed', f'任务失败:{task["title"]}', + f'项目 #{task["project_id"]} 任务「{task["title"]}」{message}', + save_alert=True, level='warn', atype='task_failed') + try: + proj = db.q('SELECT * FROM projects WHERE id=?', (task['project_id'],), one=True) + if proj and (proj.get('deliver_email') or '').strip(): + delivery.notify_blocker(proj, task['title'], message) + except Exception: + pass def _log(task_id, level, message): @@ -182,9 +196,7 @@ def run_task(task_id): _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') + _notify_failed(task, '因依赖未完成被拒绝执行:' + '、'.join(blockers)) return # 确定 Worker @@ -195,9 +207,7 @@ 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') + _notify_failed(task, '指定 Worker 不存在或已停用,无法执行') return else: worker = pick_worker_auto(task) @@ -205,9 +215,7 @@ 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') + _notify_failed(task, '自动路由失败:没有可用的 Worker') return _set_task(task_id, worker_id=worker['id']) _log(task_id, 'info', f'自动路由 → Worker「{worker["name"]}」({worker["provider"]}/{worker["model"]})') @@ -217,15 +225,13 @@ 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') + _notify_failed(task, reason) 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') + _notify_failed(task, reason) return _set_task(task_id, status='running', started_at=db.now(), error='') @@ -239,9 +245,7 @@ def run_task(task_id): 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') + _notify_failed(task, f'执行出错:{str(e)[:300]}') return _cost_record(task, worker, usage) @@ -272,6 +276,13 @@ def run_task(task_id): _budget_alert(task['project_id']) + if new_status == 'done': + # V3:无需审核的任务直接完成后,检查项目是否全部完成 → 自动交付并通知送达者 + try: + delivery.auto_complete_if_ready(task['project_id']) + except Exception: + pass + # DAG:触发下游就绪任务 downstream = _trigger_downstream(task) for t in downstream: diff --git a/enterprise.py b/enterprise.py index f258f7e..37ea739 100644 --- a/enterprise.py +++ b/enterprise.py @@ -293,7 +293,8 @@ 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'] + 'eval_results', 'templates', 'users', 'audit_logs', + 'project_deliverables', 'user_projects', 'user_workers'] out = {'exported_at': time.strftime('%Y-%m-%d %H:%M:%S'), 'platform': 'ai-worker-platform', 'version': 'v2.0.0'} for t in tables: @@ -337,3 +338,92 @@ def generate_consent_token(): tok = uuid.uuid4().hex[:12] audit('system', 'compliance.consent', '数据使用同意', f'consent_token={tok}') return tok + + +# --------------------------------------------------------------------------- +# V3 授权:用户 ↔ 项目 / 用户 ↔ Worker(精准权限) +# --------------------------------------------------------------------------- +PERM_LEVEL = {'view': 0, 'use': 1, 'manage': 2, 'admin': 3} +PROJECT_PERMS = ('view', 'manage', 'admin') +WORKER_PERMS = ('view', 'use', 'manage') + + +def perm_ok(have, need): + """have 权限是否满足 need 权限(None 视为无权限)""" + if not have: + return False + return PERM_LEVEL.get(have, -1) >= PERM_LEVEL.get(need, 99) + + +def user_project_perm(user_id, project_id): + """用户在项目上的权限:None / view / manage / admin""" + r = db.q('SELECT perm FROM user_projects WHERE user_id=? AND project_id=?', + (user_id, project_id), one=True) + return r['perm'] if r else None + + +def user_worker_perm(user_id, worker_id): + """用户在 Worker 上的权限:None / view / use / manage""" + r = db.q('SELECT perm FROM user_workers WHERE user_id=? AND worker_id=?', + (user_id, worker_id), one=True) + return r['perm'] if r else None + + +def visible_project_ids(user_id): + """用户可见的项目 id 列表(admin/auditor 返回 None 表示全部)""" + u = db.q('SELECT role FROM users WHERE id=?', (user_id,), one=True) + if u and u['role'] in ('admin', 'auditor'): + return None + rows = db.q('SELECT project_id FROM user_projects WHERE user_id=?', (user_id,)) + return [r['project_id'] for r in rows] + + +def visible_worker_ids(user_id): + """用户可见的 Worker id 列表(admin/auditor 返回 None 表示全部)""" + u = db.q('SELECT role FROM users WHERE id=?', (user_id,), one=True) + if u and u['role'] in ('admin', 'auditor'): + return None + rows = db.q('SELECT worker_id FROM user_workers WHERE user_id=?', (user_id,)) + return [r['worker_id'] for r in rows] + + +def set_user_grants(user_id, projects=None, workers=None): + """批量覆盖用户授权。projects=[{project_id, perm}], workers=[{worker_id, perm}] + perm 传空/None 表示收回该授权。返回 {'projects': n, 'workers': m}。""" + out = {'projects': 0, 'workers': 0} + if projects is not None: + db.w('DELETE FROM user_projects WHERE user_id=?', (user_id,)) + for g in projects: + perm = (g.get('perm') or '').strip() + if perm not in PROJECT_PERMS: + continue + pid = int(g.get('project_id') or 0) + if not db.q('SELECT id FROM projects WHERE id=?', (pid,), one=True): + continue + db.w('INSERT INTO user_projects (user_id, project_id, perm, created_at) VALUES (?,?,?,?)', + (user_id, pid, perm, db.now())) + out['projects'] += 1 + if workers is not None: + db.w('DELETE FROM user_workers WHERE user_id=?', (user_id,)) + for g in workers: + perm = (g.get('perm') or '').strip() + if perm not in WORKER_PERMS: + continue + wid = int(g.get('worker_id') or 0) + if not db.q('SELECT id FROM workers WHERE id=?', (wid,), one=True): + continue + db.w('INSERT INTO user_workers (user_id, worker_id, perm, created_at) VALUES (?,?,?,?)', + (user_id, wid, perm, db.now())) + out['workers'] += 1 + return out + + +def user_grants(user_id): + """用户现有授权 + 全部可选项目/Worker,供管理界面展示""" + projects = db.q('SELECT p.id, p.name, p.status FROM projects p ORDER BY p.id DESC') + workers = db.q('SELECT id, name, provider, model, status FROM workers ORDER BY id DESC') + for p in projects: + p['perm'] = user_project_perm(user_id, p['id']) + for w in workers: + w['perm'] = user_worker_perm(user_id, w['id']) + return {'projects': projects, 'workers': workers} diff --git a/static/app.js b/static/app.js index 60c251b..252ab3c 100644 --- a/static/app.js +++ b/static/app.js @@ -43,10 +43,17 @@ function toast(msg, type = '') { /* ---------- Login ---------- */ function showLogin() { $('#login-mask').style.display = 'flex'; } function hideLogin() { $('#login-mask').style.display = 'none'; } +function applyRoleUI(role) { + CURRENT_ROLE = role || 'admin'; + $$('#sidebar nav a').forEach(a => { a.style.display = (a.dataset.route === 'enterprise' && CURRENT_ROLE !== 'admin') ? 'none' : ''; }); +} $('#login-btn').addEventListener('click', async () => { try { await api('/api/login', {method:'POST', body:{username: $('#login-user').value, password: $('#login-pwd').value}}); - hideLogin(); $('#logout-btn').style.display = 'block'; router(); + hideLogin(); $('#logout-btn').style.display = 'block'; + const me = await api('/api/me'); + applyRoleUI(me.user?.role || 'admin'); + router(); } catch (e) { $('#login-err').textContent = e.message; } }); $('#login-pwd').addEventListener('keydown', e => { if (e.key === 'Enter') $('#login-btn').click(); }); @@ -154,6 +161,10 @@ async function pageProjects() {
    ${esc(p.objective || p.description || '暂无描述')}
    任务 ${p.task_count} · 完成 ${p.done_count} ${p.budget_limit ? '预算 ¥' + p.budget_limit : ''}
    +
    📧 ${esc(p.deliver_email || '未配置送达者')} + ${p.perm ? `${{view:'查看',manage:'管理',admin:'管理员'}[p.perm] || p.perm}` : ''} + ${p.demo_url ? `🌐 Demo` : ''} +
    `).join('') || '
    还没有项目,点右上角新建
    '} `; } @@ -171,6 +182,15 @@ function openProjectModal(p = {}) {
    +
    +
    +
    +
    +
    +