Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
c56745a8ab | ||
|
|
339dd89ba6 | ||
|
|
b6b9e9148c | ||
|
|
affa335e8a | ||
|
|
e5d664fb98 |
@@ -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 小工具 / 市场调研简报),一键复制其目标与任务列表生成新项目 |
|
||||
|
||||
---
|
||||
|
||||
|
||||
@@ -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')
|
||||
|
||||
+4
-8
@@ -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:
|
||||
|
||||
@@ -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
|
||||
@@ -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():
|
||||
"""存量 Worker:provider 匹配的接口库自动关联 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
@@ -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('/')
|
||||
|
||||
@@ -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 Worker(V3.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:
|
||||
|
||||
@@ -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
@@ -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
File diff suppressed because it is too large
Load Diff
+1
-1
@@ -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>
|
||||
|
||||
@@ -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}
|
||||
|
||||
/* V1:tabs / 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)}
|
||||
Reference in New Issue
Block a user