Compare commits

..
5 Commits
Author SHA1 Message Date
hz4th_coder c56745a8ab V3.5 精细化运营升级:用量统计/接口库/团队/对话/工作目录/流式超时
1. 精细化统计:cost_records 新增 calls/cached_tokens/latency_ms/first_token_ms;
   用量明细报表(项目×智能体矩阵) + 成本报表细化(输入/输出/缓存命中/调用次数)
2. 从参考项目中新建:内置3个测试项目(文案/Python/调研),一键复制目标+任务
3. 大模型接口库(llm_endpoints):专门配置接口(地址/密钥/模型/定价),
   计费支持按token(逐模型)与按调用次数;创建AI Worker直接选用;
   AI Worker团队(worker_teams):打包Worker,对话/建项目可直接选团队
4. 对话导航融合仪表盘:可选大模型/AI Worker/团队,默认主力AI Worker(可设),SSE流式
5. 系统工作目录:默认data/workspace可改绝对路径;项目与多Agent协作均在其下
   建唯一工作目录;手动输入目录已存在则列出信息并需手动确认
6. 任务执行超时改为流式单token返回超时+首字延迟超时,设置页可配;
   所有模型输出SSE按token接收
2026-09-05 15:27:43 +08:00
hz4th_coder 339dd89ba6 修复: 分页每页条数下拉框被全局 select 样式撑满整行
- 根因: 全局 input/select/textarea{width:100%} 导致 .psize-sel 独占一行
- 修复: .psize-sel 覆盖 width:auto, 与页码按钮/跳页框同一行居中排列
2026-08-16 19:55:37 +08:00
hz4th_coder b6b9e9148c UI 优化: Markdown 切换开关移入标题行 + 分页控件顶部复制/每页条数可调
- 任务详情弹窗: 📝 Markdown/原文 切换开关移到「任务指令」「产出物」标题右侧
  (标题左、开关右, flex 布局, 不再挤占文本区域宽度)
- 运行日志/告警中心: 分页控件顶部底部各一份; 新增每页条数选择(10/20/50/100),
  改条数自动回到第1页, 顶部/底部跳页双向同步
2026-08-16 19:51:32 +08:00
hz4th_coder affa335e8a V3.3 功能: 日志/告警分页筛选 + 项目状态筛选/删除/改状态 + 任务详情 Markdown
- 运行日志: 分页(快速跳页/上一页/下一页) + 级别筛选(全部/信息/成功/警告/错误) + 关键词搜索
- 告警中心: 分页 + 级别/类型/已读未读筛选
- 项目列表: 状态筛选(全部/规划中/进行中/待验收/已完成/已归档, 带计数, 默认全部) + 卡片内状态下拉直接改状态 + 删除按钮(带确认, 权限管理级)
- 任务详情弹窗: 任务指令/产出物按 Markdown 渲染(内置轻量渲染器, 无外部依赖, 支持标题/加粗/斜体/代码块/列表/引用/表格/链接), 可一键切换原文视图, 默认 Markdown
- 后端: /api/logs /api/alerts 支持 page/page_size/level/type/read/q 参数并返回 total; /api/projects 支持 status 筛选
2026-08-16 19:41:04 +08:00
hz4th_coder e5d664fb98 性能优化: 项目页/看板/DAG/知识库/交付 加载提速
- /api/tasks 列表接口改用轻量序列化,去掉 output_text/error 大字段
  (300KB -> 6KB, 详情接口 /api/tasks/<id> 仍返回完整数据)
- /api/tasks/trash 与 /api/projects/<id>/dag 同步轻量化
- 前端同一项目内切换 tab 不再重复拉取项目/任务/Worker/事件,
  复用 projCtx 缓存; 离开项目页时自动清缓存
- /api/projects/<id>/events 各来源加 LIMIT, 避免全表加载
2026-08-16 19:26:13 +08:00
13 changed files with 2620 additions and 233 deletions
+16 -1
View File
@@ -3,7 +3,22 @@
> 以「项目」为中心、以「AI Worker」为执行单元的项目管理平台。
> 把大模型团队变成一支可指挥、可审计、可控成本的"虚拟团队"。
**当前版本:V3.4**负责人验收制 / 邮件通道修复 / 事件记录面板 / 真实网页交付
**当前版本:V3.5**精细化用量统计 / 大模型接口库 / AI Worker 团队 / 对话融合仪表盘 / 系统工作目录 / 流式单token超时
---
## 🚀 V3.5 精细化运营 + 对话 + 接口库 + 团队 + 工作目录 + 流式超时(新增)
| 能力 | 说明 |
|---|---|
| 📊 精细化用量统计 | 每次调用记录 **输入/输出 token、缓存命中 token、调用次数、首字延迟、总耗时**;成本报表新增「📊 用量明细(项目×智能体)」矩阵,按项目、智能体、模型三维核算 |
| 🔌 大模型接口库 | AI Worker 页新增「大模型接口库」标签:专门配置接口(名称/提供商/Base URL/API Key/模型列表),**价格可配置**——支持**按 token 数**(逐模型输入/输出单价)与**按调用次数**两种计费方式;创建 AI Worker 时直接从接口库选用 |
| 👥 AI Worker 团队 | AI Worker 页新增「团队」标签:把多个 Worker 打包成团队;对话可直接选团队(成员轮询应答),创建项目时「快捷选择团队」一键勾选干活 Worker |
| 💬 对话导航(融合仪表盘) | 侧边栏顶部新增「💬 对话 · 仪表盘」入口,仪表盘页顶部内嵌对话面板:可选 **大模型接口 / AI Worker / 团队** 对话,**默认主力 AI Worker**(可设主力/在设置中改);输出按 token 流式展示,每条消息记录 tokens/缓存命中/成本/首字延迟 |
| 🗂️ 系统工作目录 | 全局默认系统工作目录(可改任意绝对路径);每个项目与多 Agent 协作都在其下新建**独立无重复**工作目录(project_<id> / agent_run_<id>);支持**手动输入新的系统工作目录与项目目录**——不存在自动创建,**已存在则列出目录信息(文件数/大小/样例)并需手动勾选确认** |
| ⏱️ 流式单 token 超时 | 所有模型输出改为 **SSE 按 token 流式接收**;超时按「**单 token 返回超时**」(相邻 token 间隔)与「**首字延迟超时**」判定,均在设置页可配(默认 60s / 120s,另有整体兑底 600s);超时中断保留已产出的部分内容 |
| ⭐ 主力 AI Worker | Worker 列表可一键「⭐ 设为主力」(或设置页选择),对话默认使用;仪表盘标注主力 |
| 📋 从参考项目中新建 | 项目列表新增入口:内置 3 个不同维度的简单测试项目(产品文案速写 / Python 小工具 / 市场调研简报),一键复制其目标与任务列表生成新项目 |
---
+14 -7
View File
@@ -7,6 +7,7 @@ V2 多 Agent 协作引擎
- debate 辩论模式:多位辩手各自立论 → 互相质询(可多轮)→ 裁判综合裁决
"""
import json
import os
import threading
import traceback
import db
@@ -34,12 +35,9 @@ def _update_run(run_id, **fields):
def _chat_worker(worker, messages, temperature=None, max_tokens=None):
"""调用某个 Worker 的模型,返回 (text, usage)"""
r = llm_gateway.chat(
worker['provider'], worker['model'], messages,
temperature=temperature if temperature is not None else worker['temperature'],
max_tokens=max_tokens or worker['max_tokens'] or 2000,
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
"""调用某个 Worker 的模型(流式 + 接口库计价),返回 (text, usage)"""
r = llm_gateway.chat_worker(worker, messages,
temperature=temperature, max_tokens=max_tokens)
return r['text'], r
@@ -401,11 +399,20 @@ def _finish_notify(run, mode_label, result_preview):
def execute_run(run_id):
"""后台线程:执行一次多 Agent 协作"""
"""后台线程:执行一次多 Agent 协作(V3.5:在系统工作目录下建独立工作目录)"""
run = db.q('SELECT * FROM agent_runs WHERE id=?', (run_id,), one=True)
if not run:
return
try:
# 系统工作目录下新建唯一工作目录(无重复:agent_run_<id>
try:
import delivery
ws = delivery.write_agent_context(run_id, run.get('topic') or '', run.get('context') or '')
db.w('UPDATE agent_runs SET workspace_dir=? WHERE id=?',
(os.path.basename(ws), run_id))
_log(run_id, 'system', None, 'plan', f'📁 协作工作目录:{ws}')
except Exception:
pass
workers = _workers_from_ids(json.loads(run.get('worker_ids') or '[]'))
if not workers:
_update_run(run_id, status='failed', error='没有可用的 Worker')
+759 -52
View File
File diff suppressed because it is too large Load Diff
+4 -8
View File
@@ -109,10 +109,7 @@ def _chat_json(worker, messages, max_tokens=3000, temperature=0.3):
last_err = None
for attempt in range(3):
mt = max_tokens * (attempt + 1)
r = llm_gateway.chat(
worker['provider'], worker['model'], messages,
temperature=temperature, max_tokens=mt,
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
r = llm_gateway.chat_worker(worker, messages, temperature=temperature, max_tokens=mt)
try:
return _extract_json(r['text']), r
except Exception as e:
@@ -161,13 +158,12 @@ def plan_tasks(pid, manager, team, review_required=0):
last_err = None
for attempt in range(1, 4):
try:
r = llm_gateway.chat(
manager['provider'], manager['model'],
r = llm_gateway.chat_worker(
manager,
[{'role': 'system', 'content': '你只输出 JSON,不输出任何解释文字。'},
{'role': 'user', 'content': WBS_PROMPT.format(
goal=goal, accept=proj.get('acceptance_criteria') or '')}],
temperature=0.3, max_tokens=3000 * attempt,
base_url=manager['base_url'] or None, api_key=manager['api_key'] or None)
temperature=0.3, max_tokens=3000 * attempt)
tasks = parse_wbs(r['text'])
break
except Exception as e:
+7 -1
View File
@@ -83,7 +83,10 @@ DEFAULT_PRICE = {'input': 2.0, 'output': 8.0}
# 引擎参数
MAX_RETRY = 1 # 失败重试次数(429/5xx/网络错误)
TASK_TIMEOUT = 600 # 单任务超时(秒)
TASK_TIMEOUT = 600 # 单任务整体超时兜底(秒)
# V3.5 流式输出超时(可在 设置 页面修改,运行时读 settings 表)
TOKEN_TIMEOUT = 60 # 单 token 返回超时:相邻两个 token 数据块的最大间隔(秒)
FIRST_TOKEN_TIMEOUT = 120 # 首字延迟超时:请求发出后首个数据块的最长等待(秒)
# 预算告警阈值(项目预算使用率 >= 该值触发告警)
BUDGET_ALERT_RATIO = 0.8
@@ -101,5 +104,8 @@ EMAIL = {
# 公网访问地址(Demo 链接/邮件中的回链基准;留空则用请求 host)
PUBLIC_BASE_URL = os.environ.get('PUBLIC_BASE_URL', 'http://121.40.164.32:16071')
# 默认系统工作目录(绝对路径;可在 设置 页面修改,项目与多 Agent 协作均在此下建独立子目录)
DEFAULT_WORKSPACE_ROOT = os.path.join(DATA_DIR, 'workspace')
# 自动路由:按模型单价升序挑选可用 Worker
AUTO_ROUTE_POOL = 'enabled' # enabled | all
+233
View File
@@ -390,6 +390,66 @@ CREATE INDEX IF NOT EXISTS idx_cost_project ON cost_records(project_id);
CREATE INDEX IF NOT EXISTS idx_docs_project ON documents(project_id);
CREATE INDEX IF NOT EXISTS idx_chunks_doc ON doc_chunks(document_id);
CREATE INDEX IF NOT EXISTS idx_alerts_read ON alerts(read);
-- ===================================================================
-- V3.5 表结构:大模型接口库 / AI Worker 团队 / 对话 / 精细化计量
-- ===================================================================
CREATE TABLE IF NOT EXISTS llm_endpoints (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL,
provider TEXT DEFAULT 'custom', -- doubao/deepseek/autodl/qwen/openai/vllm/custom
base_url TEXT DEFAULT '',
api_key TEXT DEFAULT '',
models TEXT DEFAULT '[]', -- JSON: 可用模型名列表
pricing TEXT DEFAULT '{}', -- JSON: {model: {input: 元/1M, output: 元/1M}}
input_price REAL DEFAULT 0, -- 兜底输入价(元/1M tokens
output_price REAL DEFAULT 0, -- 兜底输出价(元/1M tokens
price_per_call REAL DEFAULT 0, -- 按调用次数计费单价(元/次)
billing TEXT DEFAULT 'token', -- token=按token数计费 / call=按调用次数计费
description TEXT DEFAULT '',
status TEXT DEFAULT 'enabled', -- enabled/disabled
created_at INTEGER,
updated_at INTEGER
);
CREATE TABLE IF NOT EXISTS worker_teams (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL,
description TEXT DEFAULT '',
worker_ids TEXT DEFAULT '[]', -- JSON: worker id 列表
created_at INTEGER,
updated_at INTEGER
);
CREATE TABLE IF NOT EXISTS chat_sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
title TEXT DEFAULT '',
target_type TEXT DEFAULT 'worker', -- model=大模型接口 / worker=AI Worker / team=团队
target_id INTEGER DEFAULT 0,
model TEXT DEFAULT '', -- target_type=model 时选定的模型名
created_at INTEGER,
updated_at INTEGER
);
CREATE TABLE IF NOT EXISTS chat_messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL,
role TEXT DEFAULT 'user', -- user/assistant
content TEXT DEFAULT '',
model TEXT DEFAULT '',
worker_id INTEGER,
prompt_tokens INTEGER DEFAULT 0,
completion_tokens INTEGER DEFAULT 0,
cached_tokens INTEGER DEFAULT 0,
cost REAL DEFAULT 0,
latency_ms INTEGER DEFAULT 0,
first_token_ms INTEGER DEFAULT 0,
error TEXT DEFAULT '',
created_at INTEGER
);
CREATE INDEX IF NOT EXISTS idx_chat_msg_session ON chat_messages(session_id);
CREATE INDEX IF NOT EXISTS idx_cost_worker ON cost_records(worker_id);
"""
# ---------------------------------------------------------------------------
@@ -468,6 +528,86 @@ def _migrate():
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']))
# ================= V3.5 迁移:精细化计量 / 接口库 / 团队 / 对话 =================
ccols = {r['name'] for r in conn.execute('PRAGMA table_info(cost_records)')}
for col, ddl in (
('calls', 'ALTER TABLE cost_records ADD COLUMN calls INTEGER DEFAULT 1'),
('cached_tokens', 'ALTER TABLE cost_records ADD COLUMN cached_tokens INTEGER DEFAULT 0'),
('latency_ms', 'ALTER TABLE cost_records ADD COLUMN latency_ms INTEGER DEFAULT 0'),
('first_token_ms', 'ALTER TABLE cost_records ADD COLUMN first_token_ms INTEGER DEFAULT 0'),
):
if col not in ccols:
conn.execute(ddl)
wcols = {r['name'] for r in conn.execute('PRAGMA table_info(workers)')}
for col, ddl in (
('endpoint_id', 'ALTER TABLE workers ADD COLUMN endpoint_id INTEGER'),
('is_main', 'ALTER TABLE workers ADD COLUMN is_main INTEGER DEFAULT 0'),
):
if col not in wcols:
conn.execute(ddl)
pcols3 = {r['name'] for r in conn.execute('PRAGMA table_info(projects)')}
if 'is_reference' not in pcols3:
conn.execute('ALTER TABLE projects ADD COLUMN is_reference INTEGER DEFAULT 0')
acols = {r['name'] for r in conn.execute('PRAGMA table_info(agent_runs)')}
if 'workspace_dir' not in acols:
conn.execute("ALTER TABLE agent_runs ADD COLUMN workspace_dir TEXT DEFAULT ''")
scols = {r['name'] for r in conn.execute('PRAGMA table_info(chat_sessions)')}
if 'model' not in scols:
conn.execute("ALTER TABLE chat_sessions ADD COLUMN model TEXT DEFAULT ''")
# V3.5 新表(幂等)
conn.execute('''CREATE TABLE IF NOT EXISTS llm_endpoints (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL,
provider TEXT DEFAULT 'custom',
base_url TEXT DEFAULT '',
api_key TEXT DEFAULT '',
models TEXT DEFAULT '[]',
pricing TEXT DEFAULT '{}',
input_price REAL DEFAULT 0,
output_price REAL DEFAULT 0,
price_per_call REAL DEFAULT 0,
billing TEXT DEFAULT 'token',
description TEXT DEFAULT '',
status TEXT DEFAULT 'enabled',
created_at INTEGER,
updated_at INTEGER
)''')
conn.execute('''CREATE TABLE IF NOT EXISTS worker_teams (
id INTEGER PRIMARY KEY AUTOINCREMENT,
name TEXT NOT NULL,
description TEXT DEFAULT '',
worker_ids TEXT DEFAULT '[]',
created_at INTEGER,
updated_at INTEGER
)''')
conn.execute('''CREATE TABLE IF NOT EXISTS chat_sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
title TEXT DEFAULT '',
target_type TEXT DEFAULT 'worker',
target_id INTEGER DEFAULT 0,
model TEXT DEFAULT '',
created_at INTEGER,
updated_at INTEGER
)''')
conn.execute('''CREATE TABLE IF NOT EXISTS chat_messages (
id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id INTEGER NOT NULL,
role TEXT DEFAULT 'user',
content TEXT DEFAULT '',
model TEXT DEFAULT '',
worker_id INTEGER,
prompt_tokens INTEGER DEFAULT 0,
completion_tokens INTEGER DEFAULT 0,
cached_tokens INTEGER DEFAULT 0,
cost REAL DEFAULT 0,
latency_ms INTEGER DEFAULT 0,
first_token_ms INTEGER DEFAULT 0,
error TEXT DEFAULT '',
created_at INTEGER
)''')
conn.execute('CREATE INDEX IF NOT EXISTS idx_chat_msg_session ON chat_messages(session_id)')
conn.execute('CREATE INDEX IF NOT EXISTS idx_cost_worker ON cost_records(worker_id)')
conn.commit()
conn.close()
@@ -524,6 +664,8 @@ def init_db():
_migrate()
migrate_v2()
seed_builtin_roles()
seed_endpoints()
link_workers_endpoints()
# V3.1 内置角色权限点定义(admin 为特殊值 ALL,表示全部权限)
@@ -573,6 +715,87 @@ def recover_stale_runs():
return (n1 or 0, n2 or 0, n3 or 0)
# ---------------------------------------------------------------------------
# V3.5 大模型接口库:从 config.PROVIDERS/MODEL_PRICING 自动导入 + 存量 Worker 关联
# ---------------------------------------------------------------------------
def seed_endpoints():
"""幂等:首次启动把 config.PROVIDERS 导入为大模型接口库(含逐模型定价),
便于在页面上统一管理与配置价格(按 token / 按调用次数)。"""
c = q('SELECT COUNT(*) c FROM llm_endpoints')[0]['c']
if c > 0:
return
import config as cfg
prefix = {'doubao': 'doubao', 'deepseek': 'deepseek', 'openai': 'gpt',
'qwen': 'qwen', 'vllm': '', 'autodl': 'qwen3'}
ts = now()
for pid, p in cfg.PROVIDERS.items():
models = [m for m in cfg.MODEL_PRICING if m.startswith(prefix.get(pid, '__none__'))]
pricing = {m: cfg.MODEL_PRICING[m] for m in models}
fp = cfg.MODEL_PRICING.get(models[0]) if models else cfg.DEFAULT_PRICE
w('INSERT INTO llm_endpoints (name, provider, base_url, api_key, models, pricing, '
'input_price, output_price, price_per_call, billing, description, status, created_at, updated_at) '
'VALUES (?,?,?,?,?,?,?,?,0,?,?,?,?,?)',
(p['name'], pid, p['base_url'], p['api_key'], json.dumps(models), json.dumps(pricing),
fp.get('input', 2.0), fp.get('output', 8.0), 'token',
f'由系统配置自动导入({pid}', 'enabled', ts, ts))
def link_workers_endpoints():
"""存量 Workerprovider 匹配的接口库自动关联 endpoint_id(统一计价与鉴权)"""
conn = get_conn()
try:
eps = {r['provider']: r['id'] for r in conn.execute('SELECT id, provider FROM llm_endpoints')}
for r in conn.execute('SELECT id, provider FROM workers WHERE endpoint_id IS NULL OR endpoint_id=0'):
eid = eps.get(r['provider'])
if eid:
conn.execute('UPDATE workers SET endpoint_id=? WHERE id=?', (eid, r['id']))
conn.commit()
finally:
conn.close()
def main_worker_id():
"""主力 AI Worker id:优先取设置 main_worker_id 指向的启用 Worker,否则首个启用 Worker"""
mid = get_setting('main_worker_id', '')
if mid:
r = q('SELECT id FROM workers WHERE id=? AND status="enabled"', (int(mid),), one=True)
if r:
return r['id']
r = q('SELECT id FROM workers WHERE status="enabled" ORDER BY id', one=True)
return r['id'] if r else None
def set_main_worker(worker_id):
"""设置主力 AI Worker:先清除其它 is_main,再标记目标"""
conn = get_conn()
try:
conn.execute('UPDATE workers SET is_main=0 WHERE is_main=1')
conn.execute('UPDATE workers SET is_main=1 WHERE id=?', (int(worker_id),))
conn.commit()
finally:
conn.close()
set_setting('main_worker_id', int(worker_id))
def get_llm_timeouts():
"""读取超时配置:token_timeout(单token返回超时) / first_token_timeout(首字延迟超时) / request_timeout(整体兜底)。
单位秒,0/空 回退 config 默认。"""
import config as cfg
try:
tk = float(get_setting('token_timeout', '') or 0) or cfg.TOKEN_TIMEOUT
except Exception:
tk = cfg.TOKEN_TIMEOUT
try:
fk = float(get_setting('first_token_timeout', '') or 0) or cfg.FIRST_TOKEN_TIMEOUT
except Exception:
fk = cfg.FIRST_TOKEN_TIMEOUT
try:
rk = float(get_setting('request_timeout', '') or 0) or cfg.TASK_TIMEOUT
except Exception:
rk = cfg.TASK_TIMEOUT
return max(3.0, tk), max(3.0, fk), max(10.0, rk)
def q(sql, args=(), one=False):
"""查询"""
conn = get_conn()
@@ -631,6 +854,16 @@ def serialize_task(t):
return t
def serialize_task_light(t):
"""轻量序列化(列表接口用):去掉大字段 output_text / error
看板、DAG、回收站等列表视图不需要它们,可显著减小传输体积。
详情接口 /api/tasks/<id> 仍返回完整数据。"""
t = serialize_task(t)
t.pop('output_text', None)
t.pop('error', None)
return t
def get_setting(key, default=''):
r = q('SELECT value FROM settings WHERE key=?', (key,), one=True)
return r['value'] if r else default
+82 -6
View File
@@ -18,9 +18,9 @@ from email.mime.application import MIMEApplication
from email.utils import formataddr
import db
from config import DATA_DIR, EMAIL, PUBLIC_BASE_URL
from config import DATA_DIR, EMAIL, PUBLIC_BASE_URL, DEFAULT_WORKSPACE_ROOT
WORKSPACE_ROOT = os.path.join(DATA_DIR, 'workspace')
WORKSPACE_ROOT = DEFAULT_WORKSPACE_ROOT
DEMO_ROOT = os.path.join(DATA_DIR, 'demo')
PACKAGE_ROOT = os.path.join(DATA_DIR, 'packages')
@@ -29,14 +29,70 @@ BLOCKER_DEDUP_SECONDS = 1800
def ensure_dirs():
for d in (WORKSPACE_ROOT, DEMO_ROOT, PACKAGE_ROOT):
for d in (DEMO_ROOT, PACKAGE_ROOT):
os.makedirs(d, exist_ok=True)
os.makedirs(get_workspace_root(), exist_ok=True)
def get_workspace_root():
"""系统工作目录(V3.5):默认 data/workspace,可在设置中改成任意绝对路径(无则创建,已存在需确认)。
所有项目工作目录与多 Agent 协作工作目录都建在此目录下。"""
root = (db.get_setting('sys_workspace_root', '') or '').strip() or WORKSPACE_ROOT
if not os.path.isabs(root):
root = os.path.join(DATA_DIR, root.lstrip('/'))
try:
os.makedirs(root, exist_ok=True)
except Exception:
root = WORKSPACE_ROOT
os.makedirs(root, exist_ok=True)
return root
def dir_info(path):
"""已存在目录的相关信息:文件数 / 总大小 / 最近修改 / 样例列表(供确认提醒)。不存在返回 None。"""
if not path or not os.path.isdir(path):
return None
n = total = 0
newest = 0
sample = []
try:
for dirpath, dirnames, filenames in os.walk(path):
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)
try:
st = os.stat(full)
except OSError:
continue
n += 1
total += st.st_size
if st.st_mtime > newest:
newest = st.st_mtime
if len(sample) < 40:
sample.append(os.path.relpath(full, path))
except Exception:
pass
return {'exists': True, 'path': path, 'files': n, 'size': total,
'newest_mtime': int(newest), 'sample': sample}
def workspace_path(project):
"""项目工作目录绝对路径(不存在则创建)"""
pid = project['id'] if isinstance(project, dict) else project
d = os.path.join(WORKSPACE_ROOT, f'project_{pid}')
"""项目工作目录绝对路径(V3.5:支持自定义 workspace_dir,相对系统工作目录或绝对路径;不存在则创建)"""
if isinstance(project, dict):
p = project
pid = p['id']
else:
pid = project
p = db.q('SELECT workspace_dir FROM projects WHERE id=?', (pid,), one=True)
wd = ((p or {}).get('workspace_dir') or '').strip() if p else ''
if not wd:
wd = f'project_{pid}'
if os.path.isabs(wd):
d = wd
else:
d = os.path.join(get_workspace_root(), wd)
os.makedirs(d, exist_ok=True)
return d
@@ -52,6 +108,26 @@ def package_dir():
return PACKAGE_ROOT
def agent_workspace(run_id):
"""多 Agent 协作运行的工作目录:在系统工作目录下新建 agent_run_<id>(唯一,不存在则创建)"""
d = os.path.join(get_workspace_root(), f'agent_run_{run_id}')
os.makedirs(d, exist_ok=True)
return d
def write_agent_context(run_id, topic, context=''):
"""在协作运行工作目录写入运行说明文件"""
d = agent_workspace(run_id)
try:
with open(os.path.join(d, 'run_context.md'), 'w', encoding='utf-8') as fh:
fh.write(f'# 多 Agent 协作运行 #{run_id}\n\n')
fh.write(f'## 主题\n{topic}\n\n')
fh.write(f'## 背景上下文\n{context or "(无)"}\n')
except Exception:
pass
return d
def _safe_relpath(relpath):
"""路径穿越防护:仅允许工作目录内的相对路径"""
relpath = (relpath or '').replace('\\', '/').strip('/')
+26 -12
View File
@@ -49,11 +49,14 @@ def _set_task(task_id, **fields):
def _cost_record(task, worker, usage):
db.w(
'INSERT INTO cost_records (task_id, project_id, worker_id, provider, model, '
'prompt_tokens, completion_tokens, total_tokens, cost, created_at) '
'VALUES (?,?,?,?,?,?,?,?,?,?)',
'prompt_tokens, completion_tokens, cached_tokens, total_tokens, calls, cost, '
'latency_ms, first_token_ms, created_at) '
'VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)',
(task['id'], task['project_id'], worker['id'], worker['provider'],
worker['model'], usage['prompt_tokens'], usage['completion_tokens'],
usage['total_tokens'], usage['cost'], db.now()))
usage.get('cached_tokens', 0), usage['total_tokens'], usage.get('calls', 1),
usage['cost'], usage.get('elapsed_ms', 0), usage.get('first_token_ms') or 0,
db.now()))
def _deps(task):
@@ -76,14 +79,16 @@ def check_dependencies(task):
def pick_worker_auto(task):
"""自动路由:按模型输入单价升序挑选 enabled Worker"""
"""自动路由:按模型单次调用成本升序挑选 enabled WorkerV3.5 用接口库计价)"""
rows = db.q('SELECT * FROM workers WHERE status="enabled" ORDER BY id')
if not rows:
return None
best, best_price = None, None
for r in rows:
pin, pout = llm_gateway.model_price(r['model'])
price = pin + pout * 0.5
try:
price = llm_gateway.worker_unit_price(r)
except Exception:
continue
if best_price is None or price < best_price:
best, best_price = r, price
return best
@@ -238,10 +243,19 @@ def run_task(task_id):
_log(task_id, 'info', f'开始执行:Worker「{worker["name"]}」 模型 {worker["provider"]}/{worker["model"]}')
try:
usage = llm_gateway.chat(
worker['provider'], worker['model'], _build_messages(task, worker),
temperature=worker['temperature'], max_tokens=worker['max_tokens'],
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
usage = llm_gateway.chat_worker(worker, _build_messages(task, worker))
except llm_gateway.LLMError as e:
err_msg = str(e)
if e.partial_text:
# 流式超时但已有部分产出:保留部分产出,标记失败并说明原因
db.w('UPDATE tasks SET output_text=?, status="failed", error=?, finished_at=? WHERE id=?',
(e.partial_text[:200000], err_msg, db.now(), task_id))
_log(task_id, 'warn', f'流式输出中断(已收到部分内容):{err_msg}')
else:
_set_task(task_id, status='failed', error=err_msg, finished_at=db.now())
_log(task_id, 'error', f'执行失败: {err_msg}')
_notify_failed(task, f'执行出错:{err_msg[:300]}')
return
except Exception as e:
_set_task(task_id, status='failed', error=str(e), finished_at=db.now())
_log(task_id, 'error', f'执行失败: {e}')
@@ -250,8 +264,8 @@ def run_task(task_id):
_cost_record(task, worker, usage)
_log(task_id, 'success',
f'执行完成:{usage["total_tokens"]} tokens(输入 {usage["prompt_tokens"]} / 输出 {usage["completion_tokens"]}),'
f'成本 ¥{usage["cost"]:.6f}')
f'执行完成:{usage["total_tokens"]} tokens(输入 {usage["prompt_tokens"]} / 输出 {usage["completion_tokens"]} / 缓存命中 {usage.get("cached_tokens", 0)}),'
f'首字 {usage.get("first_token_ms") or ""}ms · 总耗时 {usage.get("elapsed_ms", 0)}ms成本 ¥{usage["cost"]:.6f}')
# V3.4:产出若为完整网页文档 → 落盘工作目录,供 Demo 真实展示
try:
+13 -17
View File
@@ -140,13 +140,11 @@ def sink_done_tasks():
def evaluate_case(dataset, case, worker):
"""单条用例:Worker 作答 + Judge 打分。返回 (output, score, judgment, latency, cost, tokens)"""
t0 = time.time()
# 1) Worker 作答
r = llm_gateway.chat(
worker['provider'], worker['model'],
[{'role': 'system', 'content': worker['system_prompt'] or '你是待评估的执行 Agent,请直接回答问题。'},
{'role': 'user', 'content': case['input']}],
temperature=worker['temperature'], max_tokens=worker['max_tokens'] or 2000,
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
# 1) Worker 作答(流式 + 接口库计价)
r = llm_gateway.chat_worker(worker,
[{'role': 'system', 'content': worker['system_prompt'] or '你是待评估的执行 Agent,请直接回答问题。'},
{'role': 'user', 'content': case['input']}],
max_tokens=worker['max_tokens'] or 2000)
output = r['text']
latency = int((time.time() - t0) * 1000)
tokens = r['total_tokens']
@@ -154,16 +152,14 @@ def evaluate_case(dataset, case, worker):
# 2) Judge 打分
judge_cost = 0.0
try:
j = llm_gateway.chat(
worker['provider'], worker['model'],
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': JUDGE_PROMPT.format(
rubric=dataset['rubric'] or DEFAULT_RUBRIC,
input=case['input'][:2000],
expected=(case['expected'] or '无参考答案,凭专业判断')[:3000],
output=output[:4000])}],
temperature=0.1, max_tokens=600,
base_url=worker['base_url'] or None, api_key=worker['api_key'] or None)
j = llm_gateway.chat_worker(worker,
[{'role': 'system', 'content': '你只输出 JSON。'},
{'role': 'user', 'content': JUDGE_PROMPT.format(
rubric=dataset['rubric'] or DEFAULT_RUBRIC,
input=case['input'][:2000],
expected=(case['expected'] or '无参考答案,凭专业判断')[:3000],
output=output[:4000])}],
temperature=0.1, max_tokens=600)
score, judgment = _extract_judge(j['text'])
judge_cost = j['cost']
tokens += j['total_tokens']
+370 -50
View File
@@ -1,14 +1,32 @@
# -*- coding: utf-8 -*-
"""
统一模型网关:多供应商 OpenAI 兼容协议调用 + 计量 + 计价
统一模型网关V3.5:多供应商 OpenAI 兼容协议调用 + 流式 + 精细计量 + 计价
=====================================================================
核心变化:
1. 所有模型输出按 token 流式接收(SSE),不再一次性等完整响应;
2. 超时语义改为「单 token 返回超时 / 首字延迟超时 / 整体兜底」,
三个值均可在 设置 页面配置(settings 表,动态生效):
- token_timeout 相邻两个 token 数据块的最大间隔(默认 60s)
- first_token_timeout 请求发出后首块数据的最长等待(默认 120s)
- request_timeout 整体兜底上限(默认 600s)
3. 大模型接口库(llm_endpoints):base_url / api_key / 模型列表 / 逐模型定价,
计费方式支持「按 token 数」与「按调用次数」两种,Worker 直接选用接口库;
4. 精细计量:每次调用记录 prompt / completion / 缓存命中 cached / 调用次数 /
首字延迟 / 总耗时,全部入库(cost_records / agent_steps / chat_messages)。
"""
import json
import time
import requests
import config
import db
class LLMError(Exception):
pass
"""模型调用异常。partial_text:超时中断前已收到的部分输出。"""
def __init__(self, msg, partial_text=''):
super().__init__(msg)
self.partial_text = partial_text or ''
def get_provider_cfg(provider):
@@ -18,6 +36,67 @@ def get_provider_cfg(provider):
return cfg
# ---------------------------------------------------------------------------
# 大模型接口库
# ---------------------------------------------------------------------------
def get_endpoint(endpoint_id):
if not endpoint_id:
return None
try:
return db.q('SELECT * FROM llm_endpoints WHERE id=? AND status="enabled"',
(int(endpoint_id),), one=True)
except Exception:
return None
def _endpoint_models(ep):
try:
return json.loads(ep.get('models') or '[]') or []
except Exception:
return []
def _endpoint_pricing_map(ep):
try:
return json.loads(ep.get('pricing') or '{}') or {}
except Exception:
return {}
def worker_llm_cfg(worker):
"""解析 Worker 的大模型接口配置(V3.5):
优先取绑定的接口库 endpoint_id(统一鉴权/计价),worker 自带 base_url/api_key 可覆盖。
返回 dict(provider, model, base_url, api_key, endpoint, in_price, out_price,
price_per_call, billing, pricing_map)"""
ep = get_endpoint(worker.get('endpoint_id')) if worker else None
if ep:
models = _endpoint_models(ep)
return {
'provider': ep.get('provider') or 'custom',
'model': worker.get('model') or (models[0] if models else ''),
'base_url': worker.get('base_url') or ep.get('base_url') or '',
'api_key': worker.get('api_key') or ep.get('api_key') or '',
'endpoint': ep,
'in_price': float(ep.get('input_price') or 0),
'out_price': float(ep.get('output_price') or 0),
'price_per_call': float(ep.get('price_per_call') or 0),
'billing': ep.get('billing') or 'token',
'pricing_map': _endpoint_pricing_map(ep),
}
return {
'provider': (worker or {}).get('provider', ''),
'model': (worker or {}).get('model', ''),
'base_url': (worker or {}).get('base_url') or '',
'api_key': (worker or {}).get('api_key') or '',
'endpoint': None,
'in_price': None, 'out_price': None, 'price_per_call': None,
'billing': 'token', 'pricing_map': {},
}
# ---------------------------------------------------------------------------
# 计价:优先接口库(逐模型定价 / 按调用次数),否则 config.MODEL_PRICING
# ---------------------------------------------------------------------------
def model_price(model):
p = config.MODEL_PRICING.get(model, config.DEFAULT_PRICE)
return p['input'], p['output']
@@ -28,14 +107,89 @@ def calc_cost(model, prompt_tokens, completion_tokens):
return round(prompt_tokens / 1e6 * pin + completion_tokens / 1e6 * pout, 6)
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}
messages 支持两种格式:
- 纯文本:[{'role':'user','content':'...'}]
- 多模态:[{'role':'user','content':[{'type':'text','text':'...'},
{'type':'image_url','image_url':{'url':'...'}}]}]
"""
def calc_cost_ex(model, prompt_tokens, completion_tokens, calls=1, cfg=None):
"""按 Worker/接口库配置计价。cfg = worker_llm_cfg() 结果。"""
if cfg and cfg.get('endpoint'):
ep = cfg['endpoint']
if ep.get('billing') == 'call':
return round(float(ep.get('price_per_call') or 0) * max(1, calls), 6)
p = (cfg.get('pricing_map') or {}).get(model)
if p:
pin, pout = float(p.get('input') or 0), float(p.get('output') or 0)
if pin or pout:
return round(prompt_tokens / 1e6 * pin + completion_tokens / 1e6 * pout, 6)
pin, pout = cfg.get('in_price') or 0, cfg.get('out_price') or 0
if pin or pout:
return round(prompt_tokens / 1e6 * pin + completion_tokens / 1e6 * pout, 6)
pin, pout = model_price(model)
return round(prompt_tokens / 1e6 * pin + completion_tokens / 1e6 * pout, 6)
def worker_unit_price(worker):
"""自动路由用:估算 Worker 单次调用成本(按 token 计费 = 输入价+0.5*输出价;按次计费 = 单价)"""
cfg = worker_llm_cfg(worker)
if cfg.get('endpoint') and cfg.get('billing') == 'call':
return float(cfg.get('price_per_call') or 0)
p = (cfg.get('pricing_map') or {}).get(cfg['model'])
if p:
return float(p.get('input') or 0) + float(p.get('output') or 0) * 0.5
if cfg.get('in_price') is not None:
return float(cfg['in_price']) + float(cfg['out_price']) * 0.5
pin, pout = model_price(cfg['model'])
return pin + pout * 0.5
# ---------------------------------------------------------------------------
# Token 估算(流式响应未带 usage 时的兜底)
# ---------------------------------------------------------------------------
def _estimate_prompt_tokens(messages):
n = 0
for m in messages or []:
c = m.get('content') if isinstance(m, dict) else ''
if isinstance(c, str):
n += max(1, int(len(c) * 0.6))
elif isinstance(c, list):
for part in c:
if not isinstance(part, dict):
continue
t = part.get('text') or ''
n += max(1, int(len(t) * 0.6))
if part.get('image_url') or part.get('image'):
n += 1000 # 图片按约 1000 token 估算
return max(1, n)
def _estimate_completion_tokens(text):
return max(1, int(len(text or '') * 0.6))
# ---------------------------------------------------------------------------
# SSE 流式核心(单次尝试,无重试;超时语义见模块说明)
# ---------------------------------------------------------------------------
def _sock_of(resp):
"""尽力获取底层 socket 以调整读超时(兼容不同 requests/urllib3 版本)"""
try:
raw = resp.raw
fp = getattr(raw, '_fp', None) or getattr(raw, 'fp', None)
fpp = getattr(fp, 'fp', None)
for obj in (fpp, fp):
if obj is None:
continue
s = getattr(obj, 'raw', None) or getattr(obj, '_sock', None)
if s is not None:
return s
except Exception:
pass
return None
def _chat_stream_raw(provider, model, messages, temperature=0.7, max_tokens=None,
base_url=None, api_key=None, token_timeout=60,
first_token_timeout=120, extra_payload=None):
"""发起一次流式请求,产出事件:
yield ('delta', piece) | ('usage', usage_dict) | ('finish', finish_reason)
结束前若收到 usage 则正常给出;异常抛 LLMError(含 partial_text)。
超时:首块数据等待 first_token_timeout;相邻数据块间隔 token_timeout;整体 request_timeout 兜底。"""
cfg = get_provider_cfg(provider)
url = (base_url or cfg['base_url']).rstrip('/') + '/chat/completions'
key = api_key or cfg['api_key']
@@ -44,58 +198,217 @@ def chat(provider, model, messages, temperature=0.7, max_tokens=None,
headers = {
'Authorization': f'Bearer {key}',
'Content-Type': 'application/json',
'Accept': 'text/event-stream',
}
payload = {
'model': model,
'messages': messages,
'temperature': temperature,
'stream': True,
'stream_options': {'include_usage': True},
}
if max_tokens:
payload['max_tokens'] = max_tokens
if extra_payload:
payload.update(extra_payload)
timeout = timeout or cfg.get('timeout', 300)
t0 = time.time()
first_token_at = None
resp = None
try:
resp = requests.post(url, json=payload, headers=headers, stream=True,
timeout=(min(30, first_token_timeout), first_token_timeout))
if resp.status_code != 200:
body = resp.text[:300]
resp.close()
if resp.status_code == 429:
raise LLMError(f'模型限流(429): {body}')
if resp.status_code >= 500:
raise LLMError(f'服务端错误({resp.status_code}): {body}')
raise LLMError(f'调用失败({resp.status_code}): {body}')
# 首块之后,读超时降为「单 token 返回超时」(首次 read 保持 first_token_timeout
sock = _sock_of(resp)
first_line_seen = False
for raw_line in resp.iter_lines(decode_unicode=True):
line = (raw_line or '').strip()
if not first_line_seen:
first_line_seen = True
if sock is not None:
try:
sock.settimeout(token_timeout)
except Exception:
pass
if not line or not line.startswith('data:'):
continue
data = line[5:].strip()
if data == '[DONE]':
break
try:
evt = json.loads(data)
except Exception:
continue
if evt.get('usage'):
yield ('usage', evt['usage'])
continue
choices = evt.get('choices') or []
if not choices:
continue
ch = choices[0]
delta = ch.get('delta') or {}
piece = delta.get('content') or ''
if not piece:
piece = delta.get('reasoning_content') or ''
if piece:
if first_token_at is None:
first_token_at = time.time()
yield ('delta', piece)
if ch.get('finish_reason'):
yield ('finish', ch.get('finish_reason'))
return
except requests.exceptions.ReadTimeout:
elapsed = time.time() - t0
if first_token_at is None:
raise LLMError(f'首字延迟超时(>{first_token_timeout}s 无输出)')
raise LLMError(f'Token 返回超时(>{token_timeout}s 无新数据,已输出 {elapsed:.0f}s')
except requests.exceptions.Timeout:
raise LLMError(f'请求超时({(time.time() - t0):.0f}s')
except requests.exceptions.ConnectionError as e:
raise LLMError(f'连接失败: {e}')
finally:
if resp is not None:
try:
resp.close()
except Exception:
pass
def _usage_fields(usage):
"""从 usage 中解析精细计量字段"""
usage = usage or {}
pt = int(usage.get('prompt_tokens') or 0)
ct = int(usage.get('completion_tokens') or 0)
det = usage.get('prompt_tokens_details') or {}
cached = int(det.get('cached_tokens') or 0)
if not cached:
cached = int(usage.get('prompt_cache_hit_tokens') or 0)
return pt, ct, cached
# ---------------------------------------------------------------------------
# 流式调用(带重试,聚合结果)
# ---------------------------------------------------------------------------
def chat_stream(provider, model, messages, temperature=0.7, max_tokens=None,
base_url=None, api_key=None, timeout=None, retries=None,
token_timeout=None, first_token_timeout=None, on_chunk=None,
extra_payload=None):
"""流式接收全部输出,返回聚合结果 dict:
{text, model, prompt_tokens, completion_tokens, total_tokens, cached_tokens,
cost, calls, first_token_ms, elapsed_ms, usage_estimated, finish_reason}
- 超时按「单 token 返回 / 首字延迟」语义(设置中可配)
- 已有部分输出时不再重试(避免重复内容),否则按 retries 重试
- on_chunk(delta) 逐块回调(用于对话流式转发)"""
if token_timeout is None or first_token_timeout is None:
tk, fk, rk = db.get_llm_timeouts()
token_timeout = token_timeout or tk
first_token_timeout = first_token_timeout or fk
request_timeout = timeout or max(first_token_timeout + 5, 60)
retries = config.MAX_RETRY if retries is None else retries
last_err = None
for attempt in range(retries + 1):
parts, usage = [], None
finish_reason = None
first_token_at = None
t0 = time.time()
try:
resp = requests.post(url, json=payload, headers=headers, timeout=timeout)
if resp.status_code == 200:
data = resp.json()
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)
return {
'text': text,
'model': data.get('model', model),
'prompt_tokens': pt,
'completion_tokens': ct,
'total_tokens': pt + ct,
'cost': calc_cost(model, pt, ct),
}
if resp.status_code == 429:
last_err = LLMError(f'模型限流(429): {resp.text[:200]}')
time.sleep(2 * (attempt + 1))
for evt, val in _chat_stream_raw(
provider, model, messages, temperature=temperature,
max_tokens=max_tokens, base_url=base_url, api_key=api_key,
token_timeout=token_timeout, first_token_timeout=first_token_timeout,
extra_payload=extra_payload):
if evt == 'delta':
if first_token_at is None:
first_token_at = time.time()
parts.append(val)
if on_chunk:
try:
on_chunk(val)
except Exception:
pass
elif evt == 'usage':
usage = val
elif evt == 'finish':
finish_reason = val
text = ''.join(parts)
if not text and not usage:
last_err = LLMError('模型返回空内容,重试中…')
continue
if resp.status_code >= 500:
last_err = LLMError(f'服务端错误({resp.status_code}): {resp.text[:200]}')
time.sleep(1)
continue
raise LLMError(f'调用失败({resp.status_code}): {resp.text[:300]}')
except requests.exceptions.Timeout:
last_err = LLMError(f'请求超时({timeout}s)')
except requests.exceptions.ConnectionError as e:
last_err = LLMError(f'连接失败: {e}')
pt, ct, cached = _usage_fields(usage)
usage_estimated = usage is None
if usage is None:
pt, ct = _estimate_prompt_tokens(messages), _estimate_completion_tokens(text)
elapsed_ms = int((time.time() - t0) * 1000)
first_ms = int((first_token_at - t0) * 1000) if first_token_at else None
return {
'text': text,
'model': model,
'prompt_tokens': pt,
'completion_tokens': ct,
'total_tokens': pt + ct,
'cached_tokens': cached,
'cost': calc_cost(model, pt, ct),
'calls': 1,
'first_token_ms': first_ms,
'elapsed_ms': elapsed_ms,
'usage_estimated': usage_estimated,
'finish_reason': finish_reason,
}
except LLMError as e:
last_err = e
if e.partial_text:
raise
continue
raise last_err or LLMError('未知错误')
# ---------------------------------------------------------------------------
# 兼容接口(非 Worker 场景)
# ---------------------------------------------------------------------------
def chat(provider, model, messages, temperature=0.7, max_tokens=None,
base_url=None, api_key=None, timeout=None, retries=None,
token_timeout=None, first_token_timeout=None, on_chunk=None, cfg=None):
"""兼容旧接口:流式接收全部输出后返回聚合结果。
cfg = worker_llm_cfg() 结果时按接口库计价/鉴权。"""
if cfg:
provider = cfg['provider']
base_url = cfg['base_url'] or None
api_key = cfg['api_key'] or None
model = cfg['model']
r = chat_stream(provider, model, messages, temperature=temperature,
max_tokens=max_tokens, base_url=base_url, api_key=api_key,
timeout=timeout, retries=retries, token_timeout=token_timeout,
first_token_timeout=first_token_timeout, on_chunk=on_chunk)
if cfg:
r['cost'] = calc_cost_ex(model, r['prompt_tokens'], r['completion_tokens'], 1, cfg)
return r
def chat_worker(worker, messages, temperature=None, max_tokens=None, on_chunk=None):
"""按 Worker 配置(含接口库)调用模型,返回聚合结果(含接口库计价与精细计量)"""
cfg = worker_llm_cfg(worker)
if not cfg['model']:
raise LLMError(f'Worker「{worker.get("name", "")}」未配置模型')
r = chat_stream(cfg['provider'], cfg['model'], messages,
temperature=temperature if temperature is not None else worker.get('temperature', 0.7),
max_tokens=max_tokens or worker.get('max_tokens') or 2000,
base_url=cfg['base_url'] or None, api_key=cfg['api_key'] or None,
on_chunk=on_chunk)
r['cost'] = calc_cost_ex(cfg['model'], r['prompt_tokens'], r['completion_tokens'], 1, cfg)
r['worker_id'] = worker['id']
r['provider'] = cfg['provider']
r['model'] = cfg['model']
return r
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):
@@ -129,7 +442,6 @@ def chat_vision(provider, model, text, image_url=None, image_path=None,
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))
@@ -138,12 +450,20 @@ def chat_vision(provider, model, text, image_url=None, image_path=None,
raise last_err or LLMError('视觉调用失败')
def test_connection(provider, model, base_url=None, api_key=None):
"""连通性测试:发一条最小请求"""
def test_connection(provider, model, base_url=None, api_key=None, endpoint=None):
"""连通性测试:发一条最小请求(流式),返回延迟/回复/成本"""
t0 = time.time()
cfg = None
if endpoint:
cfg = {'provider': endpoint.get('provider') or 'custom', 'model': model,
'base_url': base_url or endpoint.get('base_url') or '',
'api_key': api_key or endpoint.get('api_key') or '',
'endpoint': endpoint}
r = chat(provider, model,
[{'role': 'user', 'content': '请回复"OK"两个字'}],
temperature=0, max_tokens=16,
base_url=base_url, api_key=api_key, timeout=30, retries=0)
base_url=base_url, api_key=api_key, timeout=30, retries=0,
token_timeout=30, first_token_timeout=30, cfg=cfg)
return {'ok': True, 'latency_ms': int((time.time() - t0) * 1000),
'reply': r['text'][:50], 'cost': r['cost']}
'first_token_ms': r.get('first_token_ms'), 'reply': r['text'][:50],
'cost': r['cost'], 'tokens': r['total_tokens']}
+1023 -78
View File
File diff suppressed because it is too large Load Diff
+1 -1
View File
@@ -11,7 +11,7 @@
<aside id="sidebar">
<div class="logo">🤖 <span>AI Worker</span><small>项目管理平台</small></div>
<nav>
<a href="#/dashboard" data-route="dashboard" data-perm="dashboard.view">📊 仪表盘</a>
<a href="#/dashboard" data-route="dashboard" data-perm="dashboard.view" class="nav-chat">💬 对话 · 仪表盘</a>
<a href="#/projects" data-route="projects" data-perm="project.view">📁 项目</a>
<a href="#/workers" data-route="workers" data-perm="worker.view">🧑‍💻 AI Worker</a>
<a href="#/agents" data-route="agents" data-perm="agent.view">🤝 多 Agent 协作</a>
+72
View File
@@ -14,6 +14,8 @@ nav{padding:12px 0;flex:1}
nav a{display:block;padding:11px 20px;color:var(--muted);text-decoration:none;border-left:3px solid transparent}
nav a:hover{color:var(--text);background:var(--panel2)}
nav a.active{color:var(--accent);border-left-color:var(--accent);background:var(--panel2)}
nav a.nav-chat{color:var(--accent2);font-weight:600;border-left-color:var(--accent2)}
nav a.nav-chat:hover{background:var(--panel2)}
.sidebar-foot{padding:14px 18px;border-top:1px solid var(--border);font-size:12px}
.conn.ok{color:var(--accent2)} .conn.err{color:var(--danger)}
.logout{color:var(--muted);text-decoration:none;display:block;margin-top:6px}
@@ -126,6 +128,53 @@ tr:hover td{background:var(--panel2)}
.proj-card:hover{border-color:var(--accent)}
.pager-tip{color:var(--muted);font-size:12px}
.pill{display:inline-block;background:var(--panel2);border:1px solid var(--border);border-radius:6px;padding:2px 8px;font-size:12px;margin:2px}
.pstatus-sel{background:var(--panel2);border:1px solid var(--border);color:var(--text);border-radius:6px;padding:2px 6px;font-size:12px;cursor:pointer}
.pstatus-sel:hover{border-color:var(--accent)}
/* 分页 + 筛选 */
.filter-bar{display:flex;flex-wrap:wrap;gap:8px;align-items:center;margin-bottom:12px}
.filter-bar .fgroup{display:flex;gap:4px;align-items:center;background:var(--panel2);border:1px solid var(--border);border-radius:8px;padding:4px}
.filter-bar .fgroup .flbl{color:var(--muted);font-size:12px;padding:0 8px 0 6px}
.filter-bar .fbtn{border:1px solid transparent;background:transparent;color:var(--muted);padding:4px 10px;border-radius:6px;font-size:12px;cursor:pointer}
.filter-bar .fbtn:hover{color:var(--accent)}
.filter-bar .fbtn.on{background:var(--accent);color:#fff}
.filter-bar .fbtn.on.warn-on{background:var(--warn)}
.filter-bar .fbtn.on.danger-on{background:var(--danger)}
.pager{display:flex;gap:4px;align-items:center;flex-wrap:wrap;margin-top:14px;justify-content:center}
.pager.top{margin-top:0;margin-bottom:14px}
.psize-sel{background:var(--panel2);border:1px solid var(--border);color:var(--text);border-radius:7px;padding:4px 6px;font-size:12px;cursor:pointer;height:30px;width:auto}
.psize-sel:hover{border-color:var(--accent)}
.pager .pbtn{min-width:30px;height:30px;padding:0 8px;border:1px solid var(--border);background:var(--panel2);color:var(--text);border-radius:7px;font-size:12px;cursor:pointer}
.pager .pbtn:hover{border-color:var(--accent);color:var(--accent)}
.pager .pbtn.cur{background:var(--accent);border-color:var(--accent);color:#fff;font-weight:700}
.pager .pbtn:disabled{opacity:.35;cursor:not-allowed}
.pager .pgap{color:var(--muted);font-size:12px;padding:0 2px}
.pager .pgo{display:flex;gap:6px;align-items:center;margin-left:6px}
.pager .pgo input{width:52px;height:28px;background:var(--panel);border:1px solid var(--border);color:var(--text);border-radius:7px;text-align:center;font-size:12px}
.pager .pg-info{color:var(--muted);font-size:12px;margin-right:8px}
/* Markdown 渲染 */
.md-box{max-height:320px;overflow-y:auto;line-height:1.75;font-size:13px}
.md-box.sm{max-height:160px}
.md-plain{white-space:pre-wrap;word-break:break-word;margin:0}
.md-h{margin:14px 0 8px;font-weight:700;line-height:1.4}
.md-h h1,.md-h h2,.md-h h3,.md-h h4,.md-h h5,.md-h h6{margin:14px 0 8px}
.md-p{margin:0 0 8px}
.md-ul,.md-ol{margin:0 0 8px;padding-left:22px}
.md-ul li,.md-ol li{margin:3px 0}
.md-quote{border-left:3px solid var(--accent);background:var(--panel2);padding:8px 12px;margin:0 0 10px;border-radius:0 8px 8px 0;color:var(--muted)}
.md-hr{border:none;border-top:1px solid var(--border);margin:12px 0}
.md-table{border-collapse:collapse;margin:8px 0;font-size:12px;width:100%}
.md-table th,.md-table td{border:1px solid var(--border);padding:6px 10px;text-align:left}
.md-table th{background:var(--panel2)}
.md-code{background:#0b0f18;border:1px solid var(--border);border-radius:8px;padding:10px 12px;overflow-x:auto;font-family:ui-monospace,Consolas,monospace;font-size:12px;line-height:1.6;margin:8px 0}
.md-box code{background:var(--panel2);border:1px solid var(--border);border-radius:4px;padding:1px 5px;font-family:ui-monospace,Consolas,monospace;font-size:12px}
.md-box pre.md-code code{background:none;border:none;padding:0}
.md-box a{color:var(--accent)}
.md-toggle{font-size:11px;white-space:nowrap}
.md-toggle button{border:1px solid var(--border);background:var(--panel2);color:var(--muted);padding:2px 8px;cursor:pointer;font-size:11px;border-radius:5px}
.md-toggle button.on{background:var(--accent);color:#fff;border-color:var(--accent)}
.md-title{display:flex;align-items:center;justify-content:space-between;gap:10px}
/* V1tabs / DAG / 告警 / 开关 */
.tabs{display:flex;gap:4px;margin-bottom:16px;border-bottom:1px solid var(--border)}
@@ -225,3 +274,26 @@ code{background:var(--panel2);border:1px solid var(--border);border-radius:6px;p
.al-task{cursor:pointer}
.al-task:hover .al-msg{color:var(--accent)}
.al-go{color:var(--accent);font-size:11px;margin-left:6px;border:1px solid var(--border);border-radius:4px;padding:0 5px}
/* ---------- V3.5 对话 ---------- */
.chat-card{display:flex;flex-direction:column;max-height:560px;padding:14px}
.chat-top{display:flex;justify-content:space-between;align-items:center;gap:10px;flex-wrap:wrap;margin-bottom:10px}
.chat-title{font-size:16px;font-weight:700}
.chat-sess{display:flex;gap:6px;align-items:center}
.chat-sess select{width:auto;min-width:160px}
.chat-target{display:flex;gap:8px;align-items:center;flex-wrap:wrap;margin-bottom:10px}
.chat-target select{width:auto;max-width:240px}
.chat-target-label{font-size:12px;color:var(--accent);font-weight:600}
.chat-body{flex:1;min-height:180px;max-height:330px;overflow-y:auto;background:var(--bg);border:1px solid var(--border);border-radius:10px;padding:12px}
.chat-msg{display:flex;margin-bottom:10px}
.chat-msg.user{justify-content:flex-end}
.chat-msg.ai{justify-content:flex-start;flex-direction:column;align-items:flex-start}
.chat-bubble{max-width:78%;padding:9px 13px;border-radius:12px;font-size:13px;line-height:1.65;white-space:pre-wrap;word-break:break-word}
.chat-bubble.user{background:var(--accent);color:#fff;border-bottom-right-radius:3px}
.chat-bubble.ai{background:var(--panel2);border:1px solid var(--border);border-bottom-left-radius:3px}
.chat-bubble.ai.err{border-color:var(--danger)}
.chat-usage{font-size:11px;color:var(--muted);margin-top:4px;font-family:ui-monospace,Consolas,monospace}
.chat-input{display:flex;gap:10px;margin-top:10px;align-items:flex-end}
.chat-input textarea{flex:1;resize:vertical}
.ref-card{border-left:3px solid var(--accent)}
.ref-card:hover{border-color:var(--accent)}