Compare commits

..
9 Commits
Author SHA1 Message Date
hz4th_coder eb92a4e591 V3.0 交付体系 + 精准权限系统
交付环节优化:
- 新建项目必填送达者(人)邮箱,项目完成/遇到无法绕开的难关时自动邮件及时通知
- 每个项目独立工作目录 data/workspace/project_<id>/,交付物分类存放互不污染(上传/下载/删除)
- 网页交付物一键部署到 /demo/<id>/ 免登录公开 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
2026-08-13 12:38:38 +08:00
hz4th_coder 78ccca61c3 DAG画布连线编辑:新建/重连/删除 + 环检测
- 连线操作三件套:
  * 新建:拖节点右侧输出口到目标节点(幽灵线预览 + 合法/非法高亮)
  * 重连:单击连线选中 → 操作条「🔗重连」→ 点击新目标节点(自动移除旧边)
  * 删除:双击连线 / 悬停中点✕按钮 / 选中后 Delete 键 / 操作条删除按钮
- 单击连线选中高亮 + 顶部操作条(重连/删除/取消),点击空白取消选中
- 前端防环校验(DFS)+ 非法目标标红,后端 PUT/POST depends_on 全面校验:
  依赖存在性、非自依赖、全图环检测(返回环路径)
- 节点新增左右连接端口(悬停/连线模式显示);拖节点时连线与端口同步更新
- 修复:控件条点击不再误触发取消选中/视图复位
2026-08-13 11:31:40 +08:00
hz4th_coder 3cbe3a5996 修复:DAG 画布 tab 切换后节点拖不动(多层根因)
根因链:
1. mousemove handler 更新 rect 用闭包 s.pos(旧 dagState 对象),切回 tab 后 renderDag 重建 dagState,s 指向旧对象 → 节点被设回旧布局值 → 视觉上拖不动 → 改为 handler 内现查的 st.pos
2. overflow:visible 让节点绘制到 svg 元素盒外,盒外内容看得见但命中测试不到 → svg 盒加 ±3000 边距(DAG_M)+viewBox 同步扩大,节点永远可命中
3. 保存的越界节点位置(y=620>画布560)使 fit 看不到 → dagFitView 改为基于节点实际包围盒计算
4. 打开时仅'任一节点可见'就保留旧视图 → 改为'全部节点可见',任一越界自动适配

实测:刷新进入/SPA切看板切回/连续切换3次后,真实鼠标拖动均精确跟随;缩放显示/100%按钮/删除区正常
2026-08-13 00:25:01 +08:00
hz4th_coder dfdcd78f1a DAG画布:节点拖不动根因修复 + 缩放比例显示/一键100%
- 根因:dagLayout 画布宽(1100-1360px)超出视口可用宽度(约1060px),加上保存的偏移视图可能把节点推到视口外 → 节点不可见/不可点,表现为'拖不动'
- 修复:打开画布时若无保存视图、或节点全部在视口外,自动适配缩放(≤100%)并居中,全部节点始终可见可拖
- 新增右上角缩放控件:实时显示缩放百分比(滚轮/适配/拖动联动) + ⤢适配按钮 + 100%一键回原始视图
- 实测:真实鼠标拖动在 0.82x 缩放下精确跟随;滚轮 112%→100% 显示联动
2026-08-12 23:55:57 +08:00
hz4th_coder 9c45a143b9 DAG画布:文字清晰 + 无边画布 + 删除区回收站
1. 文字看不清修复:.dag-node hover stroke 继承到文字导致白色小字被描边盖糊 → label/sub 加 stroke:none,字号加大(13/11px)加粗
2. 无边画布:bg rect 扩大至±4000且与画布同色 + svg overflow:visible,大幅平移缩放不再露出边缘/裁剪节点
3. 删除区与回收站:
- 画布左下角🗑️删除区,节点拖入即软删除(deleted=1)
- 点击删除区打开回收站:恢复/彻底删除
- 后端 tasks 加 deleted/deleted_at 列,全查询过滤,trash/restore/hard-delete API,看板删除改软删
- 修复画布重建后事件失效:box监听每次重绘重绑、window监听只绑一次且拖拽状态提升为模块级(dagDrag)
2026-08-12 23:44:20 +08:00
hz4th_coder cf092da69a 优化:多Agent协作表单按模式动态展示 + DAG画布文字可读性
1. 多Agent协作:
- 选择不同模式时表单只显示该模式专属参数(主管:参与Agent+上下文;评审:+返工轮次/分数线;辩论:+质询轮次/立场)
- 参数改为下拉选择(轮次/分数线预设档位),提示文案按模式动态切换,点模式卡片可直接切换

2. DAG画布:
- 修复文字看不清:CSS .dag-node fill 覆盖节点内联状态色 → 移除,文字改纯白+半透明白
- 节点加宽 110→140px,标题超 13 字符自动省略号,不再溢出
- 节点描边加亮(#5a6b8f)提升状态区分度
2026-08-12 23:23:36 +08:00
hz4th_coder 5deb447c8e DAG 工作流画布升级:节点可拖动 + 滚轮缩放 + 画布平移
- 节点拖拽调整布局(位移按缩放比例换算,拖动不误触详情点击)
- 滚轮以鼠标为锚点缩放(0.25x-3x),双击空白一键复位
- 拖拽空白区域平移画布(cursor grab/grabbing)
- 布局与视图状态 localStorage 按项目记忆,刷新/重开自动恢复
- 实测:节点精确拖动、锚点缩放、平移、复位、点击详情全部通过
2026-08-12 17:26:30 +08:00
hz4th_coder f28d73f1bb fix: 多Agent协作 JSON 截断导致主管规划解析失败
- 根因: deepseek-v4-flash 规划输出超 max_tokens=2000 被截断,JSON 不完整解析失败
- agents.py: 新增 _chat_json() 统一处理 JSON 类调用,解析失败自动加大 max_tokens 重试(×1/×2/×3)
- 主管规划/辩论质询/辩论裁决/评审 全部改用 _chat_json 并提高初始上限(3000/1200/2500/2000)
- 实测: 主管(计算器网页版4子任务) 评审(92分) 辩论(指定阵容)
2026-08-12 17:08:16 +08:00
hz4th_coder dc49864af4 智能体大模型接口全面切换 + 新增视觉智能体
- 全部 AI Worker 切换到 deepseek/deepseek-v4-flash(api.deepseek.com,实测757ms/次,成本降40倍)
- 团队模板默认 provider/model 同步改为 deepseek/deepseek-v4-flash
- 新增视觉智能体「视觉分析师」:autodl/qwen3.6-plus 多模态
- llm_gateway: chat_vision 多模态调用(URL/base64) + 空content自动重试 + 聚合后端路由容错
- engine: 任务描述支持 ![图](url) 图片注入(视觉任务直接派活)
- API: POST /api/workers/<id>/vision_test;前端 Worker 页新增🖼️视觉测试按钮
- 实测: 截图结构化分析  雪羊图片识别  辩论模式DeepSeek 55秒完成
2026-08-12 13:22:56 +08:00
15 changed files with 2423 additions and 158 deletions
+34 -8
View File
@@ -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_<id>/`,中间产物/交付物分类存放,互不污染;支持网页上传/下载/删除 |
| 🌐 网页交付物 | 一键部署到 `data/demo/<id>/`,经 `/demo/<id>/` **免登录公开访问**,送达者直接打开链接查看 |
| 🗜️ 打包交付 | 工作目录一键 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 APIV1 + V2 路由)
├── app.py # Flask 应用 + REST APIV1 + 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/ 网页Demopackages/ 打包件;logs/ 运行日志
```
## V2 开放 API 摘要
@@ -120,12 +139,19 @@ ai-worker-platform/
企业版:
- `GET/POST /api/enterprise/users`(管理员)| `PUT/DELETE /api/enterprise/users/<id>`
- `GET/PUT /api/enterprise/users/<id>/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/<pid>/workspace` 工作目录文件列表 `POST .../workspace/upload` 上传(multipart)| `GET .../workspace/download?path=` 下载 `DELETE .../workspace?path=` 删除
- `POST /api/projects/<pid>/deploy` 网页交付物部署 Demo `POST .../package` zip 打包 `POST .../deliver` 邮件交付(Demo 链接 + 附件)
- `POST /api/projects/<pid>/complete` 完成项目并交付 `POST .../notify_deliverer` 手动通知送达者 `GET /api/projects/<pid>/deliverables` 交付记录
- `GET /demo/<pid>/` 公开 Demo 地址(送达者免登录访问)
## 路线图
- **V2.1**Temporal 持久执行、多 Agent 协作接入项目任务、eval 回归对比视图
- **V3**:多租户 SaaS 化、工作流画布(拖拽编排)、插件市场
- **V3**:多租户 SaaS 化、工作流画布(拖拽编排)、插件市场 ✅ 交付体系 + 精准权限(V3.0)
+38 -19
View File
@@ -92,6 +92,23 @@ def _extract_score(text):
return 60.0, text
def _chat_json(worker, messages, max_tokens=2000, temperature=None):
"""调用 LLM 并解析 JSON;解析失败自动加大 max_tokens 重试(防止长文截断)。
返回 (data, raw_text, usage);全部失败抛最后异常。"""
last_err = None
for attempt in range(3):
mt = max_tokens * (attempt + 1) # 2000 → 4000 → 6000
text, usage = _chat_worker(worker, messages, temperature=temperature, max_tokens=mt)
try:
data = _extract_json(text)
if isinstance(data, dict) and not data:
raise ValueError('空 JSON')
return data, text, usage
except Exception as e:
last_err = e
raise last_err or ValueError('JSON 解析失败')
# ---------------------------------------------------------------------------
# 主管模式
# ---------------------------------------------------------------------------
@@ -114,17 +131,17 @@ SUPERVISOR_SYNTH_PROMPT = (
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)
# 1) 主管拆解(JSON 解析失败自动加大 max_tokens 重试)
try:
subtasks = _extract_json(plan_text)
if isinstance(subtasks, dict):
subtasks = subtasks.get('subtasks') or subtasks.get('tasks') or []
data, _, u1 = _chat_json(
workers[0],
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': SUPERVISOR_PLAN_PROMPT.format(
topic=run['topic'], context=run['context'] or '')}],
max_tokens=3000, temperature=0.3)
subtasks = data.get('subtasks') or data.get('tasks') or []
if isinstance(data, list):
subtasks = data
except Exception as e:
_update_run(run_id, status='failed', error=f'主管规划解析失败: {e}')
_log(run_id, 'supervisor', None, 'plan', f'❌ 规划解析失败: {e}')
@@ -216,15 +233,19 @@ def run_review(run_id, run, workers):
# 评审
review_usage = None
try:
review_text, u2 = _chat_worker(
review_data, review_text, u2 = _chat_json(
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)
max_tokens=2000, temperature=0.2)
review_usage = u2
total_tokens += u2['total_tokens']; total_cost += u2['cost']
if isinstance(review_data, dict):
score = float(review_data.get('score', review_data.get('总分', 60)))
judgment = review_data.get('judgment') or review_data.get('意见') or review_text
else:
score, judgment = _extract_score(review_text)
except Exception as e:
score, judgment = 0, f'评审调用失败: {e}'
final_score = score
@@ -325,14 +346,13 @@ def run_debate(run_id, run, workers):
f"{v['stance']}】(#{k}){v['view'][:500]}"
for k, v in views.items() if k != w['id'])
try:
text, u = _chat_worker(
data, _, u = _chat_json(
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)
max_tokens=1200, temperature=0.7)
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']
@@ -346,12 +366,11 @@ def run_debate(run_id, run, workers):
f"{v['stance']}{v['view']}" for v in views.values())
_log(run_id, 'judge', judge_w['id'], 'verdict', '裁判综合裁决中…')
try:
verdict, u4 = _chat_worker(
data, verdict, u4 = _chat_json(
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)
max_tokens=2500, temperature=0.3)
if isinstance(data, dict):
consensus = data.get('consensus') or data.get('结论') or verdict
summary = data.get('summary') or data.get('摘要') or ''
+765 -39
View File
File diff suppressed because it is too large Load Diff
+13 -2
View File
@@ -33,8 +33,14 @@ PROVIDERS = {
},
'deepseek': {
'name': 'DeepSeek',
'base_url': 'https://api.deepseek.com/v1',
'api_key': os.environ.get('DEEPSEEK_API_KEY', ''),
'base_url': 'https://api.deepseek.com',
'api_key': os.environ.get('DEEPSEEK_API_KEY', 'sk-edb9df58ff574f8c98df1cd6a425e97c'),
'timeout': 600,
},
'autodl': {
'name': 'AutoDL 多模态(视觉)',
'base_url': 'https://www.autodl.art/api/v1',
'api_key': os.environ.get('AUTODL_API_KEY', 'F9MBfolzuapqTsD4KmUf9qen720rXvUZ3Sp3IrWiCTukqonx'),
'timeout': 600,
},
'openai': {
@@ -66,10 +72,12 @@ MODEL_PRICING = {
'doubao-pro-32k': {'input': 2.0, 'output': 8.0},
'deepseek-chat': {'input': 2.0, 'output': 8.0},
'deepseek-reasoner': {'input': 4.0, 'output': 16.0},
'deepseek-v4-flash': {'input': 0.2, 'output': 0.6},
'gpt-4o': {'input': 17.5, 'output': 70.0},
'gpt-4o-mini': {'input': 1.1, 'output': 4.4},
'qwen-plus': {'input': 0.8, 'output': 2.0},
'qwen-max': {'input': 4.0, 'output': 12.0},
'qwen3.6-plus': {'input': 2.0, 'output': 8.0},
}
DEFAULT_PRICE = {'input': 2.0, 'output': 8.0}
@@ -90,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
+64
View File
@@ -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);
@@ -282,6 +326,26 @@ def _migrate():
if 'depends_on' not in cols:
conn.execute("ALTER TABLE tasks ADD COLUMN depends_on TEXT DEFAULT '[]'")
conn.execute('CREATE INDEX IF NOT EXISTS idx_tasks_depends ON tasks(depends_on)')
if 'deleted' not in cols:
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)')
# V3projects 交付字段
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()
+369
View File
@@ -0,0 +1,369 @@
# -*- coding: utf-8 -*-
"""
V3 交付体系
- 每个项目独立工作目录:data/workspace/project_<id>/(中间产物与交付物隔离存放)
- 网页交付物 → data/demo/<id>/ 部署,经 /demo/<id>/ 公开访问(送达者无需登录)
- 文件包交付物 → 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_<id>/
- 无 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'<li><a href="{os.path.basename(f["path"])}">{os.path.basename(f["path"])}</a>'
f' <small>({f["size"]} B)</small></li>'
for f in files if f['ext'] in ('html', 'htm') or '/' not in f['path'])
if not links:
links = '<li>(工作目录中暂无网页文件)</li>'
with open(index, 'w', encoding='utf-8') as fh:
fh.write(f'''<!DOCTYPE html>
<html lang="zh-CN"><head><meta charset="UTF-8">
<title>{project['name']} · Demo</title>
<style>body{{font-family:system-ui;max-width:720px;margin:40px auto;padding:0 16px;color:#333}}
h1{{font-size:20px}} li{{margin:8px 0}} a{{color:#2f6fed}}</style></head>
<body><h1>📦 {project['name']} · 交付 Demo</h1>
<p>本页面由 AI Worker 平台自动生成,展示项目工作目录中的交付文件:</p>
<ul>{links}</ul></body></html>''')
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()
Binary file not shown.

After

Width:  |  Height:  |  Size: 93 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 72 KiB

Binary file not shown.

After

Width:  |  Height:  |  Size: 88 KiB

+45 -19
View File
@@ -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):
@@ -123,7 +137,9 @@ def _budget_alert(project_id):
def _build_messages(task, worker):
"""构造提示词:任务指令 + RAG 知识库上下文"""
"""构造提示词:任务指令 + RAG 知识库上下文
支持图片注入:任务描述中的 ![说明](图片URL) 或 图片:URL 会转成多模态消息(视觉 Worker)。"""
import re as _re
messages = []
if worker['system_prompt']:
messages.append({'role': 'system', 'content': worker['system_prompt']})
@@ -132,13 +148,26 @@ def _build_messages(task, worker):
if ctx:
user_text = f'{ctx}\n\n----\n\n任务指令:{user_text}'
_log(task['id'], 'info', f'📚 RAG 知识库命中 {len(refs)} 个片段:' + ''.join(refs[:5]))
messages.append({'role': 'user', 'content': user_text})
# 提取图片 URL![alt](url) 或 图片:url 或 image:url
img_urls = []
for m in _re.finditer(r'!\[[^\]]*\]\(([^)\s]+)\)', user_text):
img_urls.append(m.group(1))
for m in _re.finditer(r'(?:图片|image)\s*[:]\s*(https?://\S+)', user_text, _re.I):
img_urls.append(m.group(1))
if img_urls:
content = [{'type': 'text', 'text': user_text}]
for u in img_urls:
content.append({'type': 'image_url', 'image_url': {'url': u}})
messages.append({'role': 'user', 'content': content})
_log(task['id'], 'info', f'🖼️ 检测到 {len(img_urls)} 张图片,已注入多模态消息')
else:
messages.append({'role': 'user', 'content': user_text})
return messages
def _trigger_downstream(task):
"""DAG:任务完成后自动触发所有就绪的下游任务"""
rows = db.q('SELECT * FROM tasks WHERE status IN ("todo","failed")')
rows = db.q('SELECT * FROM tasks WHERE deleted=0 AND status IN ("todo","failed")')
triggered = []
for t in rows:
deps = _deps(t)
@@ -167,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
@@ -180,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)
@@ -190,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"]}')
@@ -202,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='')
@@ -224,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)
@@ -257,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:
+91 -1
View File
@@ -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}
+56 -2
View File
@@ -30,7 +30,12 @@ def calc_cost(model, prompt_tokens, completion_tokens):
def chat(provider, model, messages, temperature=0.7, max_tokens=None,
base_url=None, api_key=None, timeout=None, retries=None):
"""调用 OpenAI 兼容 chat/completions,返回 {text, usage, cost, model}"""
"""调用 OpenAI 兼容 chat/completions,返回 {text, usage, cost, model}
messages 支持两种格式:
- 纯文本:[{'role':'user','content':'...'}]
- 多模态:[{'role':'user','content':[{'type':'text','text':'...'},
{'type':'image_url','image_url':{'url':'...'}}]}]
"""
cfg = get_provider_cfg(provider)
url = (base_url or cfg['base_url']).rstrip('/') + '/chat/completions'
key = api_key or cfg['api_key']
@@ -56,7 +61,14 @@ def chat(provider, model, messages, temperature=0.7, max_tokens=None,
resp = requests.post(url, json=payload, headers=headers, timeout=timeout)
if resp.status_code == 200:
data = resp.json()
text = data['choices'][0]['message']['content'] or ''
msg = data['choices'][0]['message']
text = msg.get('content') or ''
if not text:
# 推理模型偶发 content 为空:用 reasoning_content 兜底
text = msg.get('reasoning_content') or ''
if not text:
last_err = LLMError('模型返回空内容,重试中…')
continue
usage = data.get('usage', {})
pt = usage.get('prompt_tokens', 0)
ct = usage.get('completion_tokens', 0)
@@ -84,6 +96,48 @@ def chat(provider, model, messages, temperature=0.7, max_tokens=None,
raise last_err or LLMError('未知错误')
def chat_vision(provider, model, text, image_url=None, image_path=None,
temperature=0.4, max_tokens=2000, base_url=None, api_key=None,
retries=3):
"""多模态视觉调用:文本 + 图片(URL 或本地路径/base64)。
返回与 chat() 相同结构。
容错:聚合 API 偶发路由到纯文本后端(不认识 image_url),自动重试。"""
import base64 as _b64
import time as _time
content = [{'type': 'text', 'text': text}]
img_url = image_url
if image_path:
with open(image_path, 'rb') as f:
raw = f.read()
mime = 'image/png'
if image_path.lower().endswith(('.jpg', '.jpeg')):
mime = 'image/jpeg'
elif image_path.lower().endswith('.gif'):
mime = 'image/gif'
elif image_path.lower().endswith('.webp'):
mime = 'image/webp'
img_url = f'data:{mime};base64,{_b64.b64encode(raw).decode()}'
if img_url:
content.append({'type': 'image_url', 'image_url': {'url': img_url}})
messages = [{'role': 'user', 'content': content}]
last_err = None
for attempt in range(max(1, retries)):
try:
return chat(provider, model, messages,
temperature=temperature, max_tokens=max_tokens,
base_url=base_url, api_key=api_key)
except LLMError as e:
last_err = e
msg = str(e)
# 仅对“多模态格式不被支持/图片无效”类错误重试(聚合后端路由问题)
if any(k in msg for k in ('image_url', 'InvalidParameter', 'invalid_parameter',
'does not appear to be valid', 'image')):
_time.sleep(2 * (attempt + 1))
continue
raise
raise last_err or LLMError('视觉调用失败')
def test_connection(provider, model, base_url=None, api_key=None):
"""连通性测试:发一条最小请求"""
t0 = time.time()
+913 -62
View File
File diff suppressed because it is too large Load Diff
+34 -5
View File
@@ -131,13 +131,42 @@ tr:hover td{background:var(--panel2)}
.tabs{display:flex;gap:4px;margin-bottom:16px;border-bottom:1px solid var(--border)}
.tabs a{padding:9px 16px;color:var(--muted);text-decoration:none;border-bottom:2px solid transparent;font-size:13px}
.tabs a.active{color:var(--accent);border-bottom-color:var(--accent)}
.dag-box{background:var(--panel);border:1px solid var(--border);border-radius:12px;padding:14px;overflow:auto}
.dag-node{fill:var(--panel2);stroke:var(--border);stroke-width:1.5;rx:10;cursor:pointer}
.dag-box{background:var(--panel);border:1px solid var(--border);border-radius:12px;overflow:hidden;height:520px;position:relative;touch-action:none;user-select:none;cursor:grab}
.dag-box svg{display:block;transition:transform .05s linear;overflow:visible}
.dag-node{stroke:var(--border);stroke-width:1.5;rx:10;cursor:grab}
.dag-node:hover{stroke:var(--accent)}
.dag-label{font-size:12px;fill:var(--text);text-anchor:middle}
.dag-sub{font-size:10px;fill:var(--muted);text-anchor:middle}
.dag-edge{stroke:var(--border);stroke-width:1.5;fill:none;marker-end:url(#arrow)}
.dag-node rect{stroke-width:1.5;rx:10}
.dag-label{font-size:13px;font-weight:600;fill:#ffffff;stroke:none;text-anchor:middle;pointer-events:none}
.dag-sub{font-size:11px;fill:rgba(255,255,255,.82);stroke:none;text-anchor:middle;pointer-events:none}
.dag-edge{stroke:var(--border);stroke-width:1.5;fill:none;marker-end:url(#arrow);pointer-events:none}
.dag-edge.act{stroke:var(--accent)}
.dag-edge-g{cursor:pointer}
.dag-edge-hit{fill:none;stroke:transparent;stroke-width:14;pointer-events:stroke;cursor:pointer}
.dag-edge-g:hover .dag-edge{stroke:var(--accent)}
.dag-edge.sel{stroke:var(--accent);stroke-width:2.5}
.dag-edge-x{opacity:0;fill:var(--danger);stroke:#fff;stroke-width:1.5;pointer-events:none;cursor:pointer;transition:opacity .12s}
.dag-edge-xt{opacity:0;font-size:10px;fill:#fff;text-anchor:middle;pointer-events:none;transition:opacity .12s}
.dag-edge-g:hover .dag-edge-x,.dag-edge-g:hover .dag-edge-xt,.dag-edge-g.sel .dag-edge-x,.dag-edge-g.sel .dag-edge-xt{opacity:1}
.dag-edge-g:hover .dag-edge-x,.dag-edge-g.sel .dag-edge-x{pointer-events:all}
.dag-port{fill:var(--panel2);stroke:var(--accent);stroke-width:2;opacity:0;pointer-events:all;transition:opacity .15s}
.dag-node:hover .dag-port{opacity:1}
.dag-connecting .dag-port{opacity:1}
.dag-port.out{cursor:crosshair}
.dag-port.in{cursor:crosshair}
.dag-node.drop-ok rect{stroke:var(--accent2);stroke-width:2.5}
.dag-node.drop-bad rect{stroke:var(--danger);stroke-width:2.5}
.dag-ghost{stroke:var(--accent);stroke-width:2;stroke-dasharray:6 4;fill:none;marker-end:url(#arrow);pointer-events:none}
.dag-ghost.bad{stroke:var(--danger)}
.dag-edge-actions{position:absolute;top:14px;left:50%;transform:translateX(-50%);background:rgba(23,30,46,.92);border:1px solid var(--border);border-radius:8px;padding:4px 8px;font-size:12px;color:var(--text);display:flex;align-items:center;gap:6px;z-index:7;backdrop-filter:blur(2px);box-shadow:0 4px 14px rgba(0,0,0,.35)}
.dag-edge-actions span{color:var(--muted);white-space:nowrap}
.dag-edge-actions b{color:var(--accent)}
.dag-trash{position:absolute;left:14px;bottom:14px;background:rgba(255,107,107,.1);border:1.5px dashed var(--danger);color:var(--danger);border-radius:10px;padding:8px 14px;font-size:12px;cursor:pointer;z-index:5;display:flex;align-items:center;gap:6px;backdrop-filter:blur(2px)}
.dag-trash:hover,.dag-trash.hover{background:rgba(255,107,107,.25);border-style:solid}
.dag-trash .cnt{background:var(--danger);color:#fff;border-radius:8px;padding:0 6px;font-size:11px}
.dag-zoom{position:absolute;right:14px;top:14px;background:rgba(23,30,46,.88);border:1px solid var(--border);border-radius:8px;padding:4px 10px;font-size:12px;color:var(--muted);display:flex;align-items:center;gap:6px;z-index:6}
.dag-zoom b{color:var(--text);font-weight:600;min-width:42px;text-align:center}
.dag-zoom button{background:var(--panel2);border:1px solid var(--border);color:var(--text);border-radius:5px;padding:2px 8px;font-size:11px;cursor:pointer}
.dag-zoom button:hover{border-color:var(--accent);color:var(--accent)}
.alert-row{border-left:3px solid var(--border);padding:12px 14px;margin-bottom:8px;background:var(--panel2);border-radius:0 10px 10px 0;cursor:pointer}
.alert-row.unread{border-left-color:var(--warn)}
.alert-row.critical{border-left-color:var(--danger)}
+1 -1
View File
@@ -180,7 +180,7 @@ def apply_project_template(tpl, variables, worker_id=None):
return pid, created
def apply_team_template(tpl, variables, provider='doubao', model='doubao-seed-evolving'):
def apply_team_template(tpl, variables, provider='deepseek', model='deepseek-v4-flash'):
"""应用团队模板 → 批量注册 Worker"""
content = render(json.loads(tpl['content']), variables)
created = []