diff --git a/README.md b/README.md index 6aa77b7..c0703de 100644 --- a/README.md +++ b/README.md @@ -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_ / agent_run_);支持**手动输入新的系统工作目录与项目目录**——不存在自动创建,**已存在则列出目录信息(文件数/大小/样例)并需手动勾选确认** | +| ⏱️ 流式单 token 超时 | 所有模型输出改为 **SSE 按 token 流式接收**;超时按「**单 token 返回超时**」(相邻 token 间隔)与「**首字延迟超时**」判定,均在设置页可配(默认 60s / 120s,另有整体兑底 600s);超时中断保留已产出的部分内容 | +| ⭐ 主力 AI Worker | Worker 列表可一键「⭐ 设为主力」(或设置页选择),对话默认使用;仪表盘标注主力 | +| 📋 从参考项目中新建 | 项目列表新增入口:内置 3 个不同维度的简单测试项目(产品文案速写 / Python 小工具 / 市场调研简报),一键复制其目标与任务列表生成新项目 | --- diff --git a/agents.py b/agents.py index 3bd7528..68e0610 100644 --- a/agents.py +++ b/agents.py @@ -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_) + 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') diff --git a/app.py b/app.py index 760cb54..611a3d2 100644 --- a/app.py +++ b/app.py @@ -6,8 +6,9 @@ AI Worker 项目管理平台 - MVP import os import secrets import functools +import time import json as _json -from flask import Flask, request, jsonify, session, send_from_directory, redirect +from flask import Flask, request, jsonify, session, send_from_directory, redirect, Response import config import db @@ -231,7 +232,75 @@ def require_admin(fn): @app.route('/api/health') def health(): - return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v3.4.0'}) + return jsonify({'ok': True, 'service': 'ai-worker-platform', 'version': 'v3.5.0'}) + + +# --------------------------------------------------------------------------- +# V3.5 内置参考测试项目(从参考项目新建用,不进普通项目列表) +# --------------------------------------------------------------------------- +def seed_reference_projects(): + """幂等:内置 3 个不同维度的简单参考测试项目,供「从参考项目中新建」快速复制。""" + if db.q('SELECT COUNT(*) c FROM projects WHERE is_reference=1')[0]['c'] > 0: + return + ts = db.now() + refs = [ + { + 'name': '产品文案速写(参考)', + 'objective': '为一款智能保温杯快速生成 3 条电商卖点文案(标题 + 正文 + 话题标签)', + 'acceptance_criteria': '输出 3 条文案,每条包含:≤20字标题、80字左右正文、5 个话题标签', + 'tasks': [ + {'title': '生成 3 条产品标题', + 'description': '为一款智能保温杯生成 3 条电商标题,每条不超过 20 字,突出卖点(保温/智能/便携),直接输出编号列表。'}, + {'title': '为每条标题写正文', + 'description': '为上面生成的 3 条标题各写一段约 80 字的小红书风格正文,口语化、有情绪、含使用场景。'}, + {'title': '提炼话题标签并汇总', + 'description': '提炼 10 个与产品相关的话题标签(如 #智能保温杯 #秋冬好物),并把 标题+正文+标签 汇总为一份 Markdown 文档输出。'}, + ], + }, + { + 'name': 'Python 小工具(参考)', + 'objective': '用 Python 实现一个命令行单词统计工具 wc.py(支持 -l/-w/-c 参数)', + 'acceptance_criteria': '代码可运行;-l 行数 / -w 单词数 / -c 字符数;附 3 组测试与 README 说明', + 'tasks': [ + {'title': '编写 wc.py 主程序', + 'description': '用 Python 实现命令行工具 wc.py,支持 -l(行数)、-w(单词数)、-c(字符数)三个参数,从标准输入或文件读取,直接输出完整代码。'}, + {'title': '编写单元测试', + 'description': '为 wc.py 编写 3 组单元测试(空文件/单行/多行英文文本),用 pytest 或 unittest,直接输出测试代码。'}, + {'title': '编写 README 说明', + 'description': '编写 README.md,说明 wc.py 的用法、参数含义、示例命令与运行环境,直接输出 Markdown 内容。'}, + ], + }, + { + 'name': '市场调研简报(参考)', + 'objective': '输出一份关于「宠物智能用品」市场的调研简报(约 3 页 A4)', + 'acceptance_criteria': '含市场规模/竞品/用户洞察/机会建议 四个章节,结构清晰可直接阅读', + 'tasks': [ + {'title': '整理市场规模数据要点', + 'description': '整理宠物智能用品市场的规模与增速要点(可基于常识估算并注明口径),输出 200 字以内的数据摘要。'}, + {'title': '分析 5 个主要竞品', + 'description': '分析 5 个主要的宠物智能用品品牌/产品特点(定位、价格带、卖点),以列表形式输出。'}, + {'title': '总结用户痛点与机会建议', + 'description': '总结宠物主人的核心痛点,并给出 3-5 条市场切入机会建议。'}, + {'title': '汇总为调研简报', + 'description': '把以上内容汇总为一份结构化 Markdown 调研简报,包含:市场规模、竞品分析、用户洞察、机会建议 四节,并加引言与结论。'}, + ], + }, + ] + for ref in refs: + pid = db.w( + 'INSERT INTO projects (name, description, objective, acceptance_criteria, status, ' + 'budget_limit, deliver_type, workspace_dir, is_reference, auto_status, created_at, updated_at) ' + 'VALUES (?,?,?,?,?,?,?,?,1,?,?,?)', + (ref['name'], '内置参考测试项目,可「从参考项目新建」快速复制', ref['objective'], + ref['acceptance_criteria'], 'active', 0, 'web', '', 'none', ts, ts)) + db.w('UPDATE projects SET workspace_dir=? WHERE id=?', (f'ref_{pid}', pid)) + for i, t in enumerate(ref['tasks']): + db.w('INSERT INTO tasks (project_id, worker_id, title, description, priority, ' + 'review_required, deadline, depends_on, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?)', + (pid, None, t['title'], t['description'], 'medium', 0, '', '[]', ts, ts)) + + +seed_reference_projects() # --------------------------------------------------------------------------- @@ -357,6 +426,42 @@ def _attach_ai_team(project): return project +def _copy_tasks_from_reference(ref_project, new_pid, worker_id=None): + """从参考项目复制任务到新项目:depends_on 存的是参考项目旧任务 id,需映射到新任务 id; + worker 沿用参考任务的 worker(若仍存在且启用),否则自动路由。返回复制数量。""" + rows = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id', (ref_project['id'],)) + old_ids = [r['id'] for r in rows] + created = [] + for t in rows: + wid = t['worker_id'] + if wid: + w = db.q('SELECT id FROM workers WHERE id=? AND status="enabled"', (wid,), one=True) + if not w: + wid = None + elif worker_id: + wid = worker_id + dep_ids = [] + try: + dep_ids = _json.loads(t['depends_on'] or '[]') + except Exception: + dep_ids = [] + mapped = [] + for d in dep_ids: + if isinstance(d, int) and d in old_ids: + idx = old_ids.index(d) + if idx < len(created): + mapped.append(created[idx]) + tid = db.w( + 'INSERT INTO tasks (project_id, worker_id, title, description, priority, ' + 'review_required, deadline, depends_on, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?)', + (new_pid, wid, t['title'], t['description'], t['priority'], t['review_required'], + t['deadline'] or '', _json.dumps(mapped), db.now(), db.now())) + created.append(tid) + db.w('INSERT INTO task_logs (task_id, level, message, created_at) VALUES (?,?,?,?)', + (tid, 'info', f'从参考项目复制:{t["title"]}', db.now())) + return len(created) + + # --------------------------------------------------------------------------- # 项目 # --------------------------------------------------------------------------- @@ -391,7 +496,23 @@ def projects(): tuple(team_ids)) if len(team) != len(set(team_ids)): return jsonify({'ok': False, 'error': '存在不可用的干活 Worker(不存在或已停用)'}), 400 + # 团队快捷选择:team_id 提供时校验存在性(前端会自动填充到 team_worker_ids) + team_id = d.get('team_id') + if team_id: + tg = db.q('SELECT id FROM worker_teams WHERE id=?', (int(team_id),), one=True) + if not tg: + return jsonify({'ok': False, 'error': '所选团队不存在'}), 400 review_required = 1 if d.get('auto_review') else 0 + # 工作目录:手动输入(相对系统工作目录名或绝对路径);已存在需手动确认 + workspace_dir = (d.get('workspace_dir') or '').strip() + confirm_existing = bool(d.get('confirm_existing_dir')) + if workspace_dir: + target = workspace_dir if os.path.isabs(workspace_dir) \ + else os.path.join(delivery.get_workspace_root(), workspace_dir.lstrip('/')) + info = delivery.dir_info(target) + if info and not confirm_existing: + return jsonify({'ok': False, 'need_confirm': True, 'error': '工作目录已存在,请确认后使用', + 'dir_info': info}), 400 pid = db.w( 'INSERT INTO projects (name, description, objective, acceptance_criteria, ' 'status, budget_limit, deliver_user_id, deliver_email, deliver_type, deliver_note, workspace_dir, ' @@ -400,9 +521,10 @@ def projects(): d.get('acceptance_criteria', ''), d.get('status', 'active'), float(d.get('budget_limit') or 0), int(duid), du['email'], d.get('deliver_type', 'web'), d.get('deliver_note', ''), - '', int(mgr_id), 'none', db.now(), db.now())) - # 工作目录名依赖自增 id:先拿到 id 再补写 - db.w('UPDATE projects SET workspace_dir=? WHERE id=?', (f'project_{pid}', pid)) + workspace_dir, int(mgr_id), 'none', db.now(), db.now())) + # 工作目录名依赖自增 id:未手动指定时用 project_(天然无重复) + if not workspace_dir: + db.w('UPDATE projects SET workspace_dir=? WHERE id=?', (f'project_{pid}', pid)) delivery.workspace_path(pid) # 干活团队 autostart.set_team(pid, team_ids) @@ -411,17 +533,29 @@ def projects(): if u and not _is_super(): db.w('INSERT INTO user_projects (user_id, project_id, perm, created_at) VALUES (?,?,?,?)', (u['id'], pid, 'admin', db.now())) + # 从参考项目复制任务(from_reference 传参考项目 id) + copied = 0 + ref_id = d.get('from_reference') + if ref_id: + ref = db.q('SELECT * FROM projects WHERE id=? AND is_reference=1', (int(ref_id),), one=True) + if ref: + copied = _copy_tasks_from_reference(ref, pid, d.get('from_reference_worker_id')) enterprise.audit(enterprise.current_actor(), 'project.create', f'project#{pid}', - f'「{d.get("name","")}」送达者用户#{duid},AI主管#{mgr_id},团队{len(team_ids)}人', + f'「{d.get("name","")}」送达者用户#{duid},AI主管#{mgr_id},团队{len(team_ids)}人' + + (f',参考项目#{ref_id}复制{copied}任务' if copied else ''), request.remote_addr or '') # 🚀 创建后立即自动开工:AI 主管规划 → 派活 → 执行 → 监控 - autostart.launch(pid, review_required=review_required) - return jsonify({'ok': True, 'id': pid, 'auto_started': True}) + if copied: + # 已有任务:跳过规划直接派活执行 + autostart.launch(pid, review_required=review_required) + else: + autostart.launch(pid, review_required=review_required) + return jsonify({'ok': True, 'id': pid, 'auto_started': True, 'copied_tasks': copied}) err = _check_perm_point('project.view') if err: return err status_f = (request.args.get('status') or '').strip() - rows = db.q('SELECT * FROM projects ORDER BY id DESC') + rows = db.q('SELECT * FROM projects WHERE is_reference=0 ORDER BY id DESC') visible = _visible_project_ids() if visible is not None: rows = [r for r in rows if r['id'] in visible] @@ -453,6 +587,7 @@ def project_detail(pid): p['perm'] = _project_perm(pid) _attach_deliver_user(p) _attach_ai_team(p) + p['workspace_info'] = delivery.dir_info(delivery.workspace_path(p)) if p.get('workspace_dir') else None return jsonify({'ok': True, 'data': p}) if request.method == 'DELETE': err = _check_project_perm(pid, 'admin') @@ -517,6 +652,18 @@ def project_detail(pid): if len(team) != len(set(team_ids)): return jsonify({'ok': False, 'error': '存在不可用的干活 Worker(不存在或已停用)'}), 400 autostart.set_team(pid, team_ids) + # V3.5 工作目录:手动指定(相对系统工作目录名或绝对路径);已存在需手动确认 + if 'workspace_dir' in d: + wd = (d.get('workspace_dir') or '').strip() + confirm_existing = bool(d.get('confirm_existing_dir')) + if wd: + target = wd if os.path.isabs(wd) else os.path.join(delivery.get_workspace_root(), wd.lstrip('/')) + info = delivery.dir_info(target) + if info and not confirm_existing: + return jsonify({'ok': False, 'need_confirm': True, 'error': '工作目录已存在,请确认后使用', + 'dir_info': info}), 400 + db.w('UPDATE projects SET workspace_dir=?, updated_at=? WHERE id=?', (wd, db.now(), pid)) + delivery.workspace_path(pid) return jsonify({'ok': True}) @@ -536,12 +683,13 @@ def workers(): wid = db.w( 'INSERT INTO workers (name, description, provider, model, base_url, api_key, ' 'system_prompt, temperature, max_tokens, task_cost_limit, monthly_cost_limit, ' - 'status, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?)', + 'endpoint_id, status, created_at, updated_at) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)', (d.get('name', '').strip(), d.get('description', ''), d.get('provider', ''), d.get('model', ''), d.get('base_url', ''), d.get('api_key', ''), d.get('system_prompt', ''), float(d.get('temperature', 0.7)), int(d.get('max_tokens', 2000)), float(d.get('task_cost_limit') or 0), - float(d.get('monthly_cost_limit') or 0), d.get('status', 'enabled'), + float(d.get('monthly_cost_limit') or 0), + d.get('endpoint_id'), d.get('status', 'enabled'), db.now(), db.now())) return jsonify({'ok': True, 'id': wid}) err = _check_perm_point('worker.view') @@ -551,9 +699,14 @@ def workers(): visible = _visible_worker_ids() if visible is not None: rows = [r for r in rows if r['id'] in visible] + eps = {e['id']: e for e in db.q('SELECT id, name, provider, base_url, billing, input_price, output_price, price_per_call FROM llm_endpoints')} + mid = db.main_worker_id() for r in rows: r['month_cost'] = round(db.monthly_worker_cost(r['id']), 6) r['perm'] = _worker_perm(r['id']) + ep = eps.get(r.get('endpoint_id')) + r['endpoint'] = dict(ep) if ep else None + r['is_main'] = (r['id'] == mid) return jsonify({'ok': True, 'data': rows}) @@ -588,7 +741,7 @@ def worker_detail(wid): d = request.get_json(force=True) fields = ['name', 'description', 'provider', 'model', 'base_url', 'api_key', 'system_prompt', 'temperature', 'max_tokens', 'task_cost_limit', - 'monthly_cost_limit', 'status'] + 'monthly_cost_limit', 'endpoint_id', 'status'] sets, args = [], [] for f in fields: if f in d: @@ -600,6 +753,22 @@ def worker_detail(wid): return jsonify({'ok': True}) +@app.route('/api/workers//main', methods=['POST']) +@require_auth +def worker_set_main(wid): + """V3.5 设置主力 AI Worker(对话默认使用)""" + err = _check_worker_perm(wid, 'use') + if err: + return err + w = db.q('SELECT id FROM workers WHERE id=?', (wid,), one=True) + if not w: + return jsonify({'ok': False, 'error': 'Worker 不存在'}), 404 + db.set_main_worker(wid) + enterprise.audit(enterprise.current_actor(), 'worker.main', f'worker#{wid}', + f'设置为主力 AI Worker', request.remote_addr or '') + return jsonify({'ok': True, 'main_worker_id': wid}) + + @app.route('/api/workers//test', methods=['POST']) @require_auth def worker_test(wid): @@ -613,9 +782,11 @@ def worker_test(wid): if not w: return jsonify({'ok': False, 'error': '不存在'}), 404 try: - r = llm_gateway.test_connection(w['provider'], w['model'], - base_url=w['base_url'] or None, - api_key=w['api_key'] or None) + cfg = llm_gateway.worker_llm_cfg(w) + r = llm_gateway.test_connection(cfg['provider'], cfg['model'], + base_url=cfg['base_url'] or None, + api_key=cfg['api_key'] or None, + endpoint=cfg['endpoint']) return jsonify({'ok': True, 'data': r}) except Exception as e: return jsonify({'ok': False, 'error': str(e)}) @@ -641,16 +812,17 @@ def worker_vision_test(wid): if not image_url and not image_b64: return jsonify({'ok': False, 'error': '请提供图片 URL 或 base64 数据'}), 400 try: + cfg = llm_gateway.worker_llm_cfg(w) if image_b64: r = llm_gateway.chat_vision( - w['provider'], w['model'], question, image_url=f'data:image/png;base64,{image_b64}', + cfg['provider'], cfg['model'], question, image_url=f'data:image/png;base64,{image_b64}', temperature=0.3, max_tokens=w['max_tokens'] or 2000, - base_url=w['base_url'] or None, api_key=w['api_key'] or None) + base_url=cfg['base_url'] or None, api_key=cfg['api_key'] or None) else: r = llm_gateway.chat_vision( - w['provider'], w['model'], question, image_url=image_url, + cfg['provider'], cfg['model'], question, image_url=image_url, temperature=0.3, max_tokens=w['max_tokens'] or 2000, - base_url=w['base_url'] or None, api_key=w['api_key'] or None) + base_url=cfg['base_url'] or None, api_key=cfg['api_key'] or None) return jsonify({'ok': True, 'data': r}) except Exception as e: return jsonify({'ok': False, 'error': str(e)}) @@ -671,6 +843,480 @@ def providers(): return jsonify({'ok': True, 'data': data}) +# --------------------------------------------------------------------------- +# V3.5 · 大模型接口库(专门配置大模型接口,创建 AI Worker 时直接选用) +# --------------------------------------------------------------------------- +@app.route('/api/endpoints', methods=['GET', 'POST']) +@require_auth +def endpoints_api(): + if request.method == 'POST': + err = _check_perm_point('worker.create') + if err: + return err + d = request.get_json(force=True) + if not (d.get('name') or '').strip(): + return jsonify({'ok': False, 'error': '接口名称必填'}), 400 + if not (d.get('base_url') or '').strip(): + return jsonify({'ok': False, 'error': 'Base URL 必填'}), 400 + models = [m.strip() for m in (d.get('models') or []) if m and m.strip()] + eid = db.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 (?,?,?,?,?,?,?,?,?,?,?,?,?,?)', + (d.get('name', '').strip(), d.get('provider', 'custom') or 'custom', + d.get('base_url', '').strip(), d.get('api_key', ''), _json.dumps(models), + _json.dumps(d.get('pricing') or {}), float(d.get('input_price') or 0), + float(d.get('output_price') or 0), float(d.get('price_per_call') or 0), + d.get('billing', 'token'), d.get('description', ''), d.get('status', 'enabled'), + db.now(), db.now())) + return jsonify({'ok': True, 'id': eid}) + err = _check_perm_point('worker.view') + if err: + return err + rows = db.q('SELECT * FROM llm_endpoints ORDER BY id DESC') + wc = {w['id']: w for w in db.q('SELECT id, name, endpoint_id FROM workers WHERE endpoint_id>0')} + for r in rows: + r['models'] = _json.loads(r.get('models') or '[]') + r['pricing'] = _json.loads(r.get('pricing') or '{}') + r['worker_count'] = sum(1 for w in wc.values() if w.get('endpoint_id') == r['id']) + r['workers'] = [w['name'] for w in wc.values() if w.get('endpoint_id') == r['id']][:20] + return jsonify({'ok': True, 'data': rows}) + + +@app.route('/api/endpoints/', methods=['GET', 'PUT', 'DELETE']) +@require_auth +def endpoint_detail(eid): + e = db.q('SELECT * FROM llm_endpoints WHERE id=?', (eid,), one=True) + if not e: + return jsonify({'ok': False, 'error': '接口不存在'}), 404 + if request.method == 'GET': + e['models'] = _json.loads(e.get('models') or '[]') + e['pricing'] = _json.loads(e.get('pricing') or '{}') + return jsonify({'ok': True, 'data': e}) + if request.method == 'DELETE': + used = db.q('SELECT COUNT(*) c FROM workers WHERE endpoint_id=?', (eid,))[0]['c'] + if used: + return jsonify({'ok': False, 'error': f'该接口正被 {used} 个 AI Worker 使用,请先解除绑定'}), 400 + db.w('DELETE FROM llm_endpoints WHERE id=?', (eid,)) + return jsonify({'ok': True}) + d = request.get_json(force=True) + fields = ['name', 'provider', 'base_url', 'api_key', 'description', 'billing', + 'input_price', 'output_price', 'price_per_call', 'status'] + sets, args = [], [] + for f in fields: + if f in d: + sets.append(f'{f}=?') + args.append(d[f]) + if 'models' in d: + sets.append('models=?') + args.append(_json.dumps([m for m in (d['models'] or []) if m and str(m).strip()])) + if 'pricing' in d: + sets.append('pricing=?') + args.append(_json.dumps(d['pricing'] or {})) + if sets: + args.append(db.now()) + db.w(f'UPDATE llm_endpoints SET {", ".join(sets)}, updated_at=? WHERE id=?', (*args, eid)) + return jsonify({'ok': True}) + + +@app.route('/api/endpoints//test', methods=['POST']) +@require_auth +def endpoint_test(eid): + e = db.q('SELECT * FROM llm_endpoints WHERE id=?', (eid,), one=True) + if not e: + return jsonify({'ok': False, 'error': '接口不存在'}), 404 + d = request.get_json(silent=True) or {} + model = (d.get('model') or '').strip() or (_json.loads(e.get('models') or '[]') or [''])[0] + base_url = d.get('base_url') or e.get('base_url') or '' + api_key = d.get('api_key') if d.get('api_key') is not None else e.get('api_key') or '' + try: + r = llm_gateway.test_connection(e.get('provider') or 'custom', model, + base_url=base_url or None, api_key=api_key or None, + endpoint=e) + return jsonify({'ok': True, 'data': r}) + except Exception as ex: + return jsonify({'ok': False, 'error': str(ex)}) + + +# --------------------------------------------------------------------------- +# V3.5 · AI Worker 团队(从 AI Worker 中创建/管理,创建项目与对话可直接选团队) +# --------------------------------------------------------------------------- +@app.route('/api/teams', methods=['GET', 'POST']) +@require_auth +def teams_api(): + if request.method == 'POST': + err = _check_perm_point('worker.create') + if err: + return err + d = request.get_json(force=True) + name = (d.get('name') or '').strip() + if not name: + return jsonify({'ok': False, 'error': '团队名称必填'}), 400 + worker_ids = [int(x) for x in (d.get('worker_ids') or []) if x] + if not worker_ids: + return jsonify({'ok': False, 'error': '至少选择 1 个 AI Worker 成员'}), 400 + tid = db.w( + 'INSERT INTO worker_teams (name, description, worker_ids, created_at, updated_at) VALUES (?,?,?,?,?)', + (name, d.get('description', ''), _json.dumps(worker_ids), db.now(), db.now())) + return jsonify({'ok': True, 'id': tid}) + err = _check_perm_point('worker.view') + if err: + return err + rows = db.q('SELECT * FROM worker_teams ORDER BY id DESC') + wmap = {w['id']: w for w in db.q('SELECT id, name, provider, model, status FROM workers')} + for r in rows: + r['worker_ids'] = _json.loads(r.get('worker_ids') or '[]') + r['workers'] = [dict(wmap[i]) for i in r['worker_ids'] if i in wmap] + r['enabled_workers'] = [w for w in r['workers'] if w.get('status') == 'enabled'] + return jsonify({'ok': True, 'data': rows}) + + +@app.route('/api/teams/', methods=['GET', 'PUT', 'DELETE']) +@require_auth +def team_detail(tid): + t = db.q('SELECT * FROM worker_teams WHERE id=?', (tid,), one=True) + if not t: + return jsonify({'ok': False, 'error': '团队不存在'}), 404 + if request.method == 'GET': + t['worker_ids'] = _json.loads(t.get('worker_ids') or '[]') + wmap = {w['id']: w for w in db.q('SELECT id, name, provider, model, status FROM workers')} + t['workers'] = [dict(wmap[i]) for i in t['worker_ids'] if i in wmap] + return jsonify({'ok': True, 'data': t}) + if request.method == 'DELETE': + db.w('DELETE FROM worker_teams WHERE id=?', (tid,)) + return jsonify({'ok': True}) + d = request.get_json(force=True) + if 'name' in d and (d.get('name') or '').strip(): + db.w('UPDATE worker_teams SET name=?, updated_at=? WHERE id=?', (d['name'].strip(), db.now(), tid)) + if 'description' in d: + db.w('UPDATE worker_teams SET description=? WHERE id=?', (d.get('description', ''), tid)) + if 'worker_ids' in d: + worker_ids = [int(x) for x in (d.get('worker_ids') or []) if x] + if not worker_ids: + return jsonify({'ok': False, 'error': '至少选择 1 个 AI Worker 成员'}), 400 + db.w('UPDATE worker_teams SET worker_ids=?, updated_at=? WHERE id=?', + (_json.dumps(worker_ids), db.now(), tid)) + return jsonify({'ok': True}) + + +# --------------------------------------------------------------------------- +# V3.5 · 从参考项目新建 / 工作目录探测 +# --------------------------------------------------------------------------- +@app.route('/api/reference_projects') +@require_auth +def reference_projects(): + """内置参考测试项目(3 个,含任务列表),供「从参考项目中新建」快速复制""" + err = _check_perm_point('project.view') + if err: + return err + rows = db.q('SELECT * FROM projects WHERE is_reference=1 ORDER BY id') + out = [] + for r in rows: + r['tasks'] = db.q('SELECT * FROM tasks WHERE project_id=? AND deleted=0 ORDER BY id', (r['id'],)) + for t in r['tasks']: + t['depends_on'] = _json.loads(t.get('depends_on') or '[]') + out.append(r) + return jsonify({'ok': True, 'data': out}) + + +@app.route('/api/workspace/probe') +@require_auth +def workspace_probe(): + """探测工作目录:相对系统工作目录名或绝对路径;已存在则返回相关信息供用户确认""" + err = _check_perm_point('project.view') + if err: + return err + path = (request.args.get('path') or '').strip() + if not path: + return jsonify({'ok': True, 'data': {'exists': False}}) + target = path if os.path.isabs(path) else os.path.join(delivery.get_workspace_root(), path.lstrip('/')) + info = delivery.dir_info(target) + return jsonify({'ok': True, 'data': info or {'exists': False, 'path': target}, + 'root': delivery.get_workspace_root()}) + + +@app.route('/api/settings/workspace', methods=['GET', 'POST']) +@require_auth +def settings_workspace(): + """系统工作目录:默认 data/workspace;可改成任意绝对路径(无则创建,已存在需确认)""" + err = _check_perm_point('setting.manage') + if err: + return err + if request.method == 'POST': + d = request.get_json(force=True) or {} + path = (d.get('path') or '').strip() + if not path: + return jsonify({'ok': False, 'error': '请输入工作目录路径'}), 400 + if not os.path.isabs(path): + path = os.path.join(config.DATA_DIR, path.lstrip('/')) + info = delivery.dir_info(path) + if info and not d.get('confirm_existing'): + return jsonify({'ok': False, 'need_confirm': True, 'error': '该目录已存在,请确认后使用', 'dir_info': info}), 400 + os.makedirs(path, exist_ok=True) + db.set_setting('sys_workspace_root', path) + db.set_setting('sys_workspace_root_confirmed', '1' if (info and d.get('confirm_existing')) else '1') + return jsonify({'ok': True, 'root': path, 'info': delivery.dir_info(path)}) + cur = delivery.get_workspace_root() + return jsonify({'ok': True, 'data': { + 'root': cur, + 'default_root': config.DEFAULT_WORKSPACE_ROOT, + 'info': delivery.dir_info(cur), + }}) + + +# --------------------------------------------------------------------------- +# V3.5 · 对话(融合在仪表盘顶部):可选 大模型接口 / AI Worker / 团队,默认主力 AI Worker +# --------------------------------------------------------------------------- +def _chat_target_cfg(session_row): + """解析对话目标,返回 (label, worker_or_none, llm_cfg, system_prompt, model)""" + ttype = session_row.get('target_type') or 'worker' + tid = session_row.get('target_id') or 0 + model = session_row.get('model') or '' + if ttype == 'model': + ep = db.q('SELECT * FROM llm_endpoints WHERE id=? AND status="enabled"', (tid,), one=True) + if not ep: + raise ValueError('对话目标:大模型接口不存在或已停用') + models = _json.loads(ep.get('models') or '[]') or [''] + if not model: + model = models[0] + cfg = {'provider': ep.get('provider') or 'custom', 'model': model, + 'base_url': ep.get('base_url') or '', 'api_key': ep.get('api_key') or '', + 'endpoint': ep} + return f'接口「{ep["name"]}」/{model}', None, cfg, '', model + if ttype == 'team': + team = db.q('SELECT * FROM worker_teams WHERE id=?', (tid,), one=True) + if not team: + raise ValueError('对话目标:团队不存在') + wids = [int(x) for x in _json.loads(team.get('worker_ids') or '[]')] + rows = db.q(f'SELECT * FROM workers WHERE id IN ({scope_args_ph(wids)}) AND status="enabled"', + tuple(wids)) if wids else [] + if not rows: + raise ValueError(f'团队「{team["name"]}」没有可用的 AI Worker 成员') + # 轮询分配成员:按该会话消息数取模 + n = db.q('SELECT COUNT(*) c FROM chat_messages WHERE session_id=? AND role="assistant"', + (session_row['id'],))[0]['c'] + worker = rows[n % len(rows)] + return f'团队「{team["name"]}」成员「{worker["name"]}」', worker, None, worker.get('system_prompt') or '', '' + # worker + worker = db.q('SELECT * FROM workers WHERE id=? AND status="enabled"', (tid,), one=True) + if not worker: + raise ValueError('对话目标:AI Worker 不存在或已停用') + return f'Worker「{worker["name"]}」', worker, None, worker.get('system_prompt') or '', '' + + +def _chat_history(sid, limit=24): + rows = db.q('SELECT role, content FROM chat_messages WHERE session_id=? AND error="" ' + 'ORDER BY id DESC LIMIT ?', (sid, limit)) + rows.reverse() + out = [] + for r in rows: + if (r['content'] or '').strip(): + out.append({'role': r['role'], 'content': r['content'][:8000]}) + return out + + +@app.route('/api/chat/options') +@require_auth +def chat_options(): + """对话配置:可选 大模型接口/Worker/团队 + 主力 AI Worker""" + eps = db.q('SELECT id, name, provider, base_url, models, billing, input_price, output_price, ' + 'price_per_call, status FROM llm_endpoints WHERE status="enabled" ORDER BY id') + for e in eps: + e['models'] = _json.loads(e.get('models') or '[]') + workers = db.q('SELECT id, name, provider, model, status, endpoint_id FROM workers WHERE status="enabled" ORDER BY id') + teams = db.q('SELECT * FROM worker_teams ORDER BY id DESC') + for t in teams: + t['worker_ids'] = _json.loads(t.get('worker_ids') or '[]') + t['worker_count'] = len(t['worker_ids']) + return jsonify({'ok': True, 'data': { + 'endpoints': eps, 'workers': workers, 'teams': teams, + 'main_worker_id': db.main_worker_id(), + }}) + + +@app.route('/api/chat/sessions', methods=['GET', 'POST']) +@require_auth +def chat_sessions(): + if request.method == 'POST': + d = request.get_json(force=True) or {} + ttype = d.get('target_type') or 'worker' + tid = int(d.get('target_id') or 0) + model = d.get('model') or '' + if not tid: + return jsonify({'ok': False, 'error': '请选择对话目标'}), 400 + sid = db.w( + 'INSERT INTO chat_sessions (title, target_type, target_id, model, created_at, updated_at) ' + 'VALUES (?,?,?,?,?,?)', + ((d.get('title') or '新对话').strip(), ttype, tid, model, db.now(), db.now())) + return jsonify({'ok': True, 'id': sid}) + rows = db.q('SELECT * FROM chat_sessions ORDER BY id DESC LIMIT 50') + for r in rows: + last = db.q('SELECT role, content, created_at FROM chat_messages WHERE session_id=? ' + 'ORDER BY id DESC LIMIT 1', (r['id'],), one=True) + r['last_message'] = (last['content'][:80] if last and last.get('content') else '') or '' + r['last_at'] = last['created_at'] if last else r['created_at'] + return jsonify({'ok': True, 'data': rows}) + + +@app.route('/api/chat/sessions/', methods=['GET', 'DELETE']) +@require_auth +def chat_session_detail(sid): + s = db.q('SELECT * FROM chat_sessions WHERE id=?', (sid,), one=True) + if not s: + return jsonify({'ok': False, 'error': '会话不存在'}), 404 + if request.method == 'DELETE': + db.w('DELETE FROM chat_messages WHERE session_id=?', (sid,)) + db.w('DELETE FROM chat_sessions WHERE id=?', (sid,)) + return jsonify({'ok': True}) + msgs = db.q('SELECT * FROM chat_messages WHERE session_id=? ORDER BY id', (sid,)) + try: + label, worker, cfg, sys_prompt, model = _chat_target_cfg(s) + s['target_label'] = label + except Exception as e: + s['target_label'] = f'({str(e)})' + return jsonify({'ok': True, 'data': {'session': s, 'messages': msgs}}) + + +def _chat_gen(sid, user_content): + """SSE 生成器:流式转发模型输出(单 token 返回超时 / 首字延迟超时 由 llm_gateway 处理)""" + import json as _j + def sse(obj): + return f'data: {_j.dumps(obj, ensure_ascii=False)}\n\n' + s = db.q('SELECT * FROM chat_sessions WHERE id=?', (sid,), one=True) + if not s: + yield sse({'type': 'error', 'message': '会话不存在'}); return + db.w('UPDATE chat_sessions SET updated_at=? WHERE id=?', (db.now(), sid)) + worker = None + cfg = None + sys_prompt = '' + model = '' + chunks = [] + worker_id = None + try: + label, worker, cfg, sys_prompt, model = _chat_target_cfg(s) + except Exception as e: + yield sse({'type': 'error', 'message': str(e)}); return + # 组装消息 + history = _chat_history(sid, 24) + history.append({'role': 'user', 'content': user_content}) + if sys_prompt and not any(m['role'] == 'system' for m in history): + history.insert(0, {'role': 'system', 'content': sys_prompt}) + chunks = [] + t0 = time.time() + try: + if worker is not None: + def on_chunk(piece): + chunks.append(piece) + r = llm_gateway.chat_worker(worker, history, on_chunk=on_chunk) + for piece in chunks: + yield sse({'type': 'delta', 'content': piece}) + worker_id = worker['id'] + model = r['model'] + else: + # 大模型接口:流式 + r = llm_gateway.chat_stream(cfg['provider'], cfg['model'], history, + base_url=cfg['base_url'] or None, api_key=cfg['api_key'] or None, + on_chunk=lambda p: chunks.append(p)) + for piece in chunks: + yield sse({'type': 'delta', 'content': piece}) + r['cost'] = llm_gateway.calc_cost_ex(cfg['model'], r['prompt_tokens'], + r['completion_tokens'], 1, cfg) + worker_id = None + model = r['model'] + db.w('INSERT INTO chat_messages (session_id, role, content, model, worker_id, prompt_tokens, ' + 'completion_tokens, cached_tokens, cost, latency_ms, first_token_ms, created_at) ' + 'VALUES (?,?,?,?,?,?,?,?,?,?,?,?)', + (sid, 'assistant', ''.join(chunks), model, worker_id, + r['prompt_tokens'], r['completion_tokens'], r.get('cached_tokens', 0), + r['cost'], r.get('elapsed_ms', 0), r.get('first_token_ms') or 0, db.now())) + yield sse({'type': 'done', 'usage': { + 'prompt_tokens': r['prompt_tokens'], 'completion_tokens': r['completion_tokens'], + 'cached_tokens': r.get('cached_tokens', 0), 'total_tokens': r['total_tokens'], + 'cost': r['cost'], 'first_token_ms': r.get('first_token_ms'), + 'elapsed_ms': r.get('elapsed_ms', 0), 'model': model, 'label': label}}) + except llm_gateway.LLMError as e: + partial = ''.join(chunks) + db.w('INSERT INTO chat_messages (session_id, role, content, model, worker_id, error, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (sid, 'assistant', partial, model or '', worker_id, + str(e)[:500], db.now())) + yield sse({'type': 'error', 'message': str(e), 'partial': partial}) + except Exception as e: + db.w('INSERT INTO chat_messages (session_id, role, content, model, worker_id, error, created_at) ' + 'VALUES (?,?,?,?,?,?,?)', + (sid, 'assistant', ''.join(chunks), model or '', worker_id, str(e)[:500], db.now())) + yield sse({'type': 'error', 'message': str(e)}) + + +@app.route('/api/chat/sessions//messages', methods=['POST']) +@require_auth +def chat_send(sid): + """发送对话消息,服务端按 token 流式返回(SSE)""" + s = db.q('SELECT id FROM chat_sessions WHERE id=?', (sid,), one=True) + if not s: + return jsonify({'ok': False, 'error': '会话不存在'}), 404 + d = request.get_json(force=True) + content = (d.get('content') or '').strip() + if not content: + return jsonify({'ok': False, 'error': '消息内容为空'}), 400 + if len(content) > 60000: + content = content[:60000] + db.w('INSERT INTO chat_messages (session_id, role, content, created_at) VALUES (?,?,?,?)', + (sid, 'user', content, db.now())) + db.w('UPDATE chat_sessions SET updated_at=? WHERE id=?', (db.now(), sid)) + return Response(_chat_gen(sid, content), mimetype='text/event-stream', + headers={'Cache-Control': 'no-cache', 'X-Accel-Buffering': 'no'}) + + +# --------------------------------------------------------------------------- +# V3.5 · 用量明细报表(每个项目 × 每个智能体:输入/输出 token、调用次数、缓存命中、成本) +# --------------------------------------------------------------------------- +@app.route('/api/reports/usage') +@require_auth +@require_perm('report.view') +def report_usage(): + visible = _visible_project_ids() + scope_sql, scope_args = '', [] + if visible is not None: + if not visible: + return jsonify({'ok': True, 'data': {'by_project': [], 'by_worker': [], 'by_model': [], 'matrix': []}}) + scope_sql = 'WHERE project_id IN (%s)' % ','.join('?' * len(visible)) + scope_args = list(visible) + # 项目 × Worker 明细矩阵 + matrix = db.q( + f'SELECT project_id, worker_id, provider, model, ' + f'COUNT(*) calls, SUM(prompt_tokens) prompt_tokens, SUM(completion_tokens) completion_tokens, ' + f'SUM(cached_tokens) cached_tokens, SUM(total_tokens) total_tokens, SUM(cost) cost ' + f'FROM cost_records {scope_sql} GROUP BY project_id, worker_id ORDER BY cost DESC', scope_args) + pname = {r['id']: r['name'] for r in db.q('SELECT id,name FROM projects')} + wname = {r['id']: r['name'] for r in db.q('SELECT id,name FROM workers')} + for r in matrix: + r['project_name'] = pname.get(r['project_id'], f'#{r["project_id"]}') + r['worker_name'] = wname.get(r['worker_id'], f'#{r["worker_id"]}') + r['cost'] = round(r['cost'] or 0, 6) + def rows_for(col): + return db.q( + f'SELECT {col}, COUNT(*) calls, SUM(prompt_tokens) prompt_tokens, ' + f'SUM(completion_tokens) completion_tokens, SUM(cached_tokens) cached_tokens, ' + f'SUM(total_tokens) total_tokens, SUM(cost) cost ' + f'FROM cost_records {scope_sql} GROUP BY {col} ORDER BY cost DESC', scope_args) + by_project = rows_for('project_id') + for r in by_project: + r['project_name'] = pname.get(r['project_id'], f'#{r["project_id"]}') + r['cost'] = round(r['cost'] or 0, 6) + by_worker = rows_for('worker_id') + for r in by_worker: + r['worker_name'] = wname.get(r['worker_id'], f'#{r["worker_id"]}') + r['cost'] = round(r['cost'] or 0, 6) + by_model = rows_for('model') + for r in by_model: + r['cost'] = round(r['cost'] or 0, 6) + return jsonify({'ok': True, 'data': { + 'by_project': by_project, 'by_worker': by_worker, 'by_model': by_model, + 'matrix': matrix, + }}) + + # --------------------------------------------------------------------------- # 任务 # --------------------------------------------------------------------------- @@ -929,25 +1575,31 @@ def report_cost(): group = request.args.get('group', 'project') if group == 'worker': rows = db.q( - 'SELECT worker_id, provider, model, COUNT(*) runs, SUM(total_tokens) tokens, ' - f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY worker_id ORDER BY cost DESC', - scope_args) + 'SELECT worker_id, provider, model, COUNT(*) runs, COUNT(DISTINCT task_id) task_calls, ' + 'SUM(prompt_tokens) prompt_tokens, SUM(completion_tokens) completion_tokens, ' + 'SUM(cached_tokens) cached_tokens, SUM(total_tokens) tokens, SUM(cost) cost ' + f'FROM cost_records {scope_sql} GROUP BY worker_id ORDER BY cost DESC', scope_args) for r in rows: w = db.q('SELECT name FROM workers WHERE id=?', (r['worker_id'],), one=True) r['worker_name'] = w['name'] if w else f'#{r["worker_id"]}' elif group == 'model': rows = db.q( - 'SELECT provider, model, COUNT(*) runs, SUM(total_tokens) tokens, ' - f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY model ORDER BY cost DESC', - scope_args) + 'SELECT provider, model, COUNT(*) runs, COUNT(DISTINCT task_id) task_calls, ' + 'SUM(prompt_tokens) prompt_tokens, SUM(completion_tokens) completion_tokens, ' + 'SUM(cached_tokens) cached_tokens, SUM(total_tokens) tokens, SUM(cost) cost ' + f'FROM cost_records {scope_sql} GROUP BY model ORDER BY cost DESC', scope_args) else: rows = db.q( - 'SELECT project_id, COUNT(*) runs, SUM(total_tokens) tokens, ' - f'SUM(cost) cost FROM cost_records {scope_sql} GROUP BY project_id ORDER BY cost DESC', - scope_args) + 'SELECT project_id, COUNT(*) runs, COUNT(DISTINCT task_id) task_calls, ' + 'SUM(prompt_tokens) prompt_tokens, SUM(completion_tokens) completion_tokens, ' + 'SUM(cached_tokens) cached_tokens, SUM(total_tokens) tokens, SUM(cost) cost ' + f'FROM cost_records {scope_sql} GROUP BY project_id ORDER BY cost DESC', scope_args) for r in rows: p = db.q('SELECT name FROM projects WHERE id=?', (r['project_id'],), one=True) r['project_name'] = p['name'] if p else f'#{r["project_id"]}' + for r in rows: + for k in ('prompt_tokens', 'completion_tokens', 'cached_tokens', 'tokens', 'cost'): + r[k] = r[k] or 0 return jsonify({'ok': True, 'data': rows}) @@ -962,12 +1614,14 @@ def stats(): return jsonify({'ok': True, 'data': { 'projects': 0, 'tasks': 0, 'workers': 0, 'total_cost': 0, 'total_tokens': 0, 'by_status': {}, 'recent': [], 'daily_cost': [], 'one_pass_rate': None, - 'rework_count': 0, 'visible_workers': 0}}) + 'rework_count': 0, 'visible_workers': 0, 'total_calls': 0, 'total_cached': 0, + 'avg_first_token_ms': 0}}) scope_sql = 'WHERE project_id IN (%s)' % ','.join('?' * len(visible)) scope_args = list(visible) out = {'projects': 0, 'tasks': 0, 'workers': 0, 'total_cost': 0, 'total_tokens': 0, 'by_status': {}, 'recent': [], 'daily_cost': [], 'one_pass_rate': None, - 'rework_count': 0, 'visible_workers': 0} + 'rework_count': 0, 'visible_workers': 0, 'total_calls': 0, 'total_cached': 0, + 'avg_first_token_ms': 0} if visible is None: out['projects'] = db.q('SELECT COUNT(*) c FROM projects')[0]['c'] out['workers'] = db.q('SELECT COUNT(*) c FROM workers')[0]['c'] @@ -983,9 +1637,14 @@ def stats(): f'AND project_id IN ({scope_args_ph(visible)}) GROUP BY project_id, status', visible): out['by_status'][r['status']] = out['by_status'].get(r['status'], 0) + r['c'] out['tasks'] = sum(out['by_status'].values()) - c = db.q(f'SELECT COALESCE(SUM(cost),0) cost, COALESCE(SUM(total_tokens),0) tokens ' + c = db.q(f'SELECT COALESCE(SUM(cost),0) cost, COALESCE(SUM(total_tokens),0) tokens, ' + f'COUNT(*) calls, COALESCE(SUM(cached_tokens),0) cached, ' + f'COALESCE(AVG(first_token_ms),0) avg_first_ms ' f'FROM cost_records {scope_sql}', scope_args)[0] out['total_cost'], out['total_tokens'] = round(c['cost'], 4), c['tokens'] + out['total_calls'] = c['calls'] + out['total_cached'] = c['cached'] + out['avg_first_token_ms'] = int(c['avg_first_ms'] or 0) # 一次通过率:done 且 rejection_count=0 done = out['by_status'].get('done', 0) @@ -1816,11 +2475,19 @@ def settings(): db.set_setting(k, v) if 'budget_alert_ratio' in d: config.BUDGET_ALERT_RATIO = float(d['budget_alert_ratio']) + if 'main_worker_id' in d and d.get('main_worker_id'): + db.set_main_worker(int(d['main_worker_id'])) return jsonify({'ok': True}) + tk, fk, rk = db.get_llm_timeouts() return jsonify({'ok': True, 'data': { 'budget_alert_ratio': config.BUDGET_ALERT_RATIO, 'auth_enabled': auth_enabled(), 'email_configured': bool(config.EMAIL.get('host')), + 'token_timeout': tk, + 'first_token_timeout': fk, + 'request_timeout': rk, + 'main_worker_id': db.main_worker_id(), + 'sys_workspace_root': delivery.get_workspace_root(), }}) diff --git a/autostart.py b/autostart.py index d01145b..b07e0a1 100644 --- a/autostart.py +++ b/autostart.py @@ -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: diff --git a/config.py b/config.py index 2094d31..a963515 100644 --- a/config.py +++ b/config.py @@ -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 diff --git a/db.py b/db.py index 2743204..79adef4 100644 --- a/db.py +++ b/db.py @@ -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() diff --git a/delivery.py b/delivery.py index f3ce695..0f3eccf 100644 --- a/delivery.py +++ b/delivery.py @@ -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_(唯一,不存在则创建)""" + 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('/') diff --git a/engine.py b/engine.py index 21f7734..d70f434 100644 --- a/engine.py +++ b/engine.py @@ -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: diff --git a/eval.py b/eval.py index 604b741..1914758 100644 --- a/eval.py +++ b/eval.py @@ -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'] diff --git a/llm_gateway.py b/llm_gateway.py index 680112b..7f614ef 100644 --- a/llm_gateway.py +++ b/llm_gateway.py @@ -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']} diff --git a/static/app.js b/static/app.js index c85d6be..a8406a2 100644 --- a/static/app.js +++ b/static/app.js @@ -214,18 +214,280 @@ function router() { } window.addEventListener('hashchange', router); -/* ---------- 仪表盘 ---------- */ +/* ---------- 💬 对话(融合在仪表盘顶部)V3.5 ---------- */ +let chatCtx = {options: null, sessions: [], sid: null, ttype: 'worker', tid: null, model: '', busy: false, msgs: []}; + +function chatTargetLabel() { + const c = chatCtx; + if (c.ttype === 'model') { + const e = (c.options?.endpoints || []).find(x => String(x.id) === String(c.tid)); + return '接口「' + (e?.name || '#' + c.tid) + '」/' + (c.model || '模型'); + } + if (c.ttype === 'team') { + const t = (c.options?.teams || []).find(x => String(x.id) === String(c.tid)); + return '团队「' + (t?.name || '#' + c.tid) + '」'; + } + const w = (c.options?.workers || []).find(x => String(x.id) === String(c.tid)); + return (w?.name || 'AI Worker') + (c.tid === c.options?.main_worker_id ? ' ⭐主力' : ''); +} + +function chatPanelHtml() { + const c = chatCtx; + const eps = (c.options?.endpoints || []).map(e => ``).join('') || ''; + const ws = (c.options?.workers || []).map(w => ``).join('') || ''; + const ts = (c.options?.teams || []).map(t => ``).join('') || ''; + const sessOpts = (c.sessions || []).map(s => ``).join(''); + return ` +
+
+
💬 对话 可选 大模型 / AI Worker / 团队 · 输出按 token 流式 · 默认主力 AI Worker
+
+ + + ${c.sid ? `` : ''} +
+
+
+ + + + + + ${esc(chatTargetLabel())} +
+
👆 选择对话目标,开始输入吧(新会话需先发送第一条消息)
+
+ + +
+
`; +} + +function chatRenderTarget() { + const c = chatCtx; + const lbl = $('#chat-target-label'); + if (lbl) lbl.textContent = chatTargetLabel(); +} + +async function chatInit() { + try { + const [opt, sess] = await Promise.all([api('/api/chat/options'), api('/api/chat/sessions')]); + chatCtx.options = opt.data; + chatCtx.sessions = sess.data; + // 默认:主力 AI Worker + if (!chatCtx.tid) { + chatCtx.ttype = 'worker'; + chatCtx.tid = chatCtx.options.main_worker_id || (chatCtx.options.workers[0]?.id) || null; + } + } catch (e) {} +} + +async function chatNewSession() { + if (!chatCtx.tid) return toast('请先选择对话目标', 'err'); + try { + const r = await api('/api/chat/sessions', {method: 'POST', body: { + target_type: chatCtx.ttype, target_id: chatCtx.tid, + model: chatCtx.ttype === 'model' ? chatCtx.model : '' + }}); + chatCtx.sid = r.id; + chatCtx.msgs = []; + const sess = (await api('/api/chat/sessions')).data; + chatCtx.sessions = sess; + chatRenderMessages(); + } catch (e) { toast(e.message, 'err'); } +} + +async function chatSwitchSession() { + const v = Number($('#chat-session').value || 0); + if (!v) { chatCtx.sid = null; chatCtx.msgs = []; chatRenderMessages(); return; } + try { + const d = (await api(`/api/chat/sessions/${v}`)).data; + chatCtx.sid = v; + chatCtx.ttype = d.session.target_type; + chatCtx.tid = d.session.target_id; + chatCtx.model = d.session.model || ''; + chatCtx.msgs = d.messages; + // 同步选择器 + const tt = $('#chat-ttype'); if (tt) tt.value = d.session.target_type; + chatSyncSelectors(); + chatRenderMessages(); + } catch (e) { toast(e.message, 'err'); } +} + +async function chatDelSession() { + if (!chatCtx.sid || !confirm('删除该会话及全部消息?')) return; + await api(`/api/chat/sessions/${chatCtx.sid}`, {method: 'DELETE'}); + chatCtx.sid = null; chatCtx.msgs = []; + chatCtx.sessions = (await api('/api/chat/sessions')).data; + chatRenderMessages(); +} + +function chatChangeType() { + chatCtx.ttype = $('#chat-ttype').value; + chatSyncSelectors(); +} + +function chatSyncSelectors() { + const c = chatCtx; + const tt = $('#chat-ttype'); + $('#chat-worker').style.display = c.ttype === 'worker' ? '' : 'none'; + $('#chat-endpoint').style.display = c.ttype === 'model' ? '' : 'none'; + $('#chat-model').style.display = c.ttype === 'model' ? '' : 'none'; + $('#chat-team').style.display = c.ttype === 'team' ? '' : 'none'; + if (c.ttype === 'worker') { + const w = $('#chat-worker'); + w.value = c.tid || ''; + if (!w.value && c.options?.workers.length) { w.value = c.options.main_worker_id || c.options.workers[0].id; c.tid = Number(w.value); } + } else if (c.ttype === 'team') { + const t = $('#chat-team'); + if (!t.value && c.options?.teams.length) { t.value = c.options.teams[0].id; c.tid = Number(t.value); } + } else if (c.ttype === 'model') { + const e = $('#chat-endpoint'); + if (!e.value && c.options?.endpoints.length) { e.value = c.options.endpoints[0].id; } + chatEndpointChange(); + } + chatRenderTarget(); +} + +function chatEndpointChange() { + const sel = $('#chat-endpoint'); + const opt = sel.selectedOptions[0]; + chatCtx.tid = Number(sel.value || 0); + let models = []; + try { models = JSON.parse(opt?.dataset?.models || '[]'); } catch (e) {} + const m = $('#chat-model'); + m.innerHTML = models.map(x => ``).join(''); + chatCtx.model = m.value || models[0] || ''; + chatRenderTarget(); +} + +function chatRenderMessages() { + const body = $('#chat-body'); + if (!body) return; + const c = chatCtx; + if (!c.sid) { body.innerHTML = '
👆 选择对话目标,点「发送」自动创建新会话
'; return; } + if (!c.msgs.length) { body.innerHTML = '
新会话已建立,开始提问吧 ✍️
'; return; } + body.innerHTML = c.msgs.map(m => { + if (m.role === 'user') return `
${esc(m.content)}
`; + if (m.error) return `
❌ 出错了:${esc(m.error)}
${m.content ? `
${esc(m.content)}
` : ''}
`; + return `
${esc(m.content)}
` + + `
⚡ ${m.total_tokens ?? (m.prompt_tokens + m.completion_tokens)} tokens · 输入 ${m.prompt_tokens} / 输出 ${m.completion_tokens}${m.cached_tokens ? ' / 缓存命中 ' + m.cached_tokens : ''} · ${fmtMoney(m.cost)} · 首字 ${m.first_token_ms || '—'}ms · ${m.latency_ms || '—'}ms · ${esc(m.model || '')}
`; + }).join(''); + body.scrollTop = body.scrollHeight; +} + +function chatKeydown(e) { + if (e.key === 'Enter' && !e.shiftKey) { e.preventDefault(); chatSend(); } +} + +async function chatSend() { + const input = $('#chat-input'); + const content = (input.value || '').trim(); + if (!content || chatCtx.busy) return; + if (!chatCtx.tid) return toast('请先选择对话目标', 'err'); + chatCtx.busy = true; + const sendBtn = $('#chat-send'); if (sendBtn) sendBtn.disabled = true; + // 无会话则先创建 + if (!chatCtx.sid) { + try { + const r = await api('/api/chat/sessions', {method: 'POST', body: { + target_type: chatCtx.ttype, target_id: chatCtx.tid, + model: chatCtx.ttype === 'model' ? chatCtx.model : '' + }}); + chatCtx.sid = r.id; + chatCtx.msgs = []; + chatCtx.sessions = (await api('/api/chat/sessions')).data; + } catch (e) { chatCtx.busy = false; if (sendBtn) sendBtn.disabled = false; return toast(e.message, 'err'); } + } + // 渲染用户消息 + AI 占位 + const body = $('#chat-body'); + const bubble = document.createElement('div'); + const msgWrap = document.createElement('div'); + msgWrap.className = 'chat-msg ai'; + bubble.className = 'chat-bubble ai'; + bubble.textContent = ''; + msgWrap.appendChild(bubble); + const usage = document.createElement('div'); + usage.className = 'chat-usage'; + usage.textContent = '⏳ 等待首字…'; + msgWrap.appendChild(usage); + const userWrap = document.createElement('div'); + userWrap.className = 'chat-msg user'; + const userBubble = document.createElement('div'); + userBubble.className = 'chat-bubble user'; + userBubble.textContent = content; + userWrap.appendChild(userBubble); + if (chatCtx.msgs.length === 0 && body) body.innerHTML = ''; + body.appendChild(userWrap); + body.appendChild(msgWrap); + body.scrollTop = body.scrollHeight; + input.value = ''; + let acc = ''; + try { + const res = await fetch(`/api/chat/sessions/${chatCtx.sid}/messages`, { + method: 'POST', headers: {'Content-Type': 'application/json'}, + body: JSON.stringify({content}) + }); + if (!res.ok) throw new Error('HTTP ' + res.status); + const reader = res.body.getReader(); + const dec = new TextDecoder(); + let buf = ''; + let done_usage = null; + while (true) { + const {done, value} = await reader.read(); + if (done) break; + buf += dec.decode(value, {stream: true}); + let idx; + while ((idx = buf.indexOf('\n\n')) >= 0) { + const chunk = buf.slice(0, idx); + buf = buf.slice(idx + 2); + const line = chunk.split('\n').find(l => l.startsWith('data: ')); + if (!line) continue; + let evt; + try { evt = JSON.parse(line.slice(6)); } catch (e) { continue; } + if (evt.type === 'delta') { + acc += evt.content; + bubble.textContent = acc; + body.scrollTop = body.scrollHeight; + } else if (evt.type === 'done') { + done_usage = evt.usage; + } else if (evt.type === 'error') { + throw new Error(evt.message); + } + } + } + if (done_usage) { + usage.textContent = `⚡ ${done_usage.total_tokens} tokens · 输入 ${done_usage.prompt_tokens} / 输出 ${done_usage.completion_tokens}${done_usage.cached_tokens ? ' / 缓存命中 ' + done_usage.cached_tokens : ''} · ${fmtMoney(done_usage.cost)} · 首字 ${done_usage.first_token_ms || '—'}ms · 总 ${done_usage.elapsed_ms}ms · ${esc(done_usage.model)}`; + } + // 刷新会话消息列表(含用户消息 + 助手消息) + const d = (await api(`/api/chat/sessions/${chatCtx.sid}`)).data; + chatCtx.msgs = d.messages; + } catch (e) { + usage.textContent = '❌ ' + e.message; + } finally { + chatCtx.busy = false; + if (sendBtn) sendBtn.disabled = false; + } +} + +/* ---------- 仪表盘(含对话) ---------- */ async function pageDashboard() { + await chatInit(); const s = (await api('/api/stats')).data; const total = Object.values(s.by_status).reduce((a, b) => a + b, 0); const pct = n => total ? Math.round(n / total * 100) : 0; $('#main').innerHTML = ` -

仪表盘AI 虚拟团队运营总览

+ ${chatPanelHtml()} +

运营总览项目 · 任务 · AI Worker · 用量

项目
${s.projects}
进行中 ${s.by_status.done ?? 0} 个已完成任务
任务总数
${s.tasks}
完成率 ${pct(s.by_status.done ?? 0)}% · 待审核 ${s.by_status.review ?? 0}
AI Worker
${s.workers}
虚拟员工档案数
-
累计成本
${fmtMoney(s.total_cost)}
${(s.total_tokens/1e6).toFixed(2)}M tokens 消耗
+
累计成本
${fmtMoney(s.total_cost)}
${(s.total_tokens/1e6).toFixed(2)}M tokens · ${s.total_calls ?? 0} 次调用 · 缓存命中 ${(s.total_cached ?? 0)/1e6 >= 0.001 ? (s.total_cached/1e6).toFixed(2) + 'M' : (s.total_cached ?? 0)} tokens
@@ -234,6 +496,7 @@ async function pageDashboard() {
一次通过率
${s.one_pass_rate ?? '—'}%(验收一次通过)
返工次数
${s.rework_count} 次打回
执行失败
${s.by_status.failed ?? 0} 个任务
+
平均首字延迟
${s.avg_first_token_ms ? s.avg_first_token_ms + ' ms' : '—'}
@@ -271,7 +534,10 @@ async function pageProjects() { const statusCls = s => s === 'active' ? 'running' : s === 'review' ? 'review' : s === 'done' ? 'done' : s === 'archived' ? 'cancelled' : 'todo'; $('#main').innerHTML = `

项目立项 → 规划 → 执行 → 验收 → 归档

-
+
+ + +
状态 @@ -323,6 +589,7 @@ async function delProject(pid) { function openProjectModal(p = {}) { const usersP = api('/api/users'); const workersP = api('/api/workers'); + const teamsP = api('/api/teams'); openModal(`

${p.id ? '编辑项目' : '新建项目'}

@@ -345,10 +612,21 @@ function openProjectModal(p = {}) {
+
🗂️ 工作目录(V3.5)默认系统工作目录下 project_(自动无重复);可手动指定
+
+
+
+ +
+ +
🤖 AI 虚拟团队(创建后自动开工)
+
+ +
加载 Worker 中…
@@ -376,7 +654,7 @@ function openProjectModal(p = {}) { curMgr = (m || ws[0]).id; } $('#pm-mgr').innerHTML = '' + - ws.map(w => ``).join(''); + ws.map(w => ``).join(''); const curTeam = new Set((p.team_workers || []).map(w => w.id)); let defaults = ws.filter(w => /牛|马/.test(w.name) && String(w.id) !== String(curMgr)); if (!defaults.length) defaults = ws.filter(w => String(w.id) !== String(curMgr)); @@ -384,6 +662,43 @@ function openProjectModal(p = {}) { `).join(''); }).catch(() => {}); + teamsP.then(r => { + const sel = $('#pm-team-preset'); + if (!sel) return; + sel.innerHTML = '' + + r.data.map(t => ``).join(''); + }).catch(() => {}); + window.applyTeamPreset = () => { + const tid = Number($('#pm-team-preset').value || 0); + if (!tid) return; + teamsP.then(r => { + const t = r.data.find(x => x.id === tid); + if (!t) return; + const ids = new Set((t.enabled_workers || t.workers || []).map(w => w.id)); + $$('#pm-team input').forEach(cb => cb.checked = ids.has(Number(cb.value))); + toast(`已按团队「${t.name}」勾选干活 Worker`, 'ok'); + }); + }; + window.probeWsDir = async () => { + const path = $('#pm-wsdir').value.trim(); + const box = $('#pm-wsinfo'); + if (!path) { box.style.display = 'none'; return; } + box.style.display = 'block'; + box.innerHTML = '
检查中…
'; + try { + const r = await api('/api/workspace/probe?path=' + encodeURIComponent(path)); + if (r.data.exists) { + const d = r.data; + box.innerHTML = `
⚠️ 目录已存在:${esc(d.path)}
+
文件 ${d.files} 个 · 大小 ${(d.size/1024/1024).toFixed(2)} MB${d.newest_mtime ? ' · 最近修改 ' + fmtTime(d.newest_mtime) : ''}
+
${(d.sample || []).map(s => '
· ' + esc(s) + '
').join('')}
`; + $('#pm-wsconfirm-wrap').style.display = 'flex'; + } else { + box.innerHTML = `
✅ 目录不存在:${esc(r.data.path)},创建时将自动新建
`; + $('#pm-wsconfirm-wrap').style.display = 'none'; + } + } catch (e) { box.innerHTML = `
${esc(e.message)}
`; } + }; $('#pm-save').addEventListener('click', async () => { const body = { name: $('#pm-name').value.trim(), objective: $('#pm-objective').value, @@ -394,8 +709,11 @@ function openProjectModal(p = {}) { deliver_note: $('#pm-dnote').value.trim(), manager_worker_id: Number($('#pm-mgr').value || 0), team_worker_ids: [...$$('#pm-team input:checked')].map(i => Number(i.value)), + workspace_dir: $('#pm-wsdir').value.trim(), + confirm_existing_dir: !!$('#pm-wsconfirm')?.checked, auto_review: !!$('#pm-auto-review').checked }; + if (p.from_reference) body.from_reference = p.from_reference; if (!body.name) return toast('请填写项目名称', 'err'); if (!body.deliver_user_id) return toast('必填:选择送达者用户', 'err'); if (!body.manager_worker_id) return toast('必填:选择 AI 主管 Worker', 'err'); @@ -406,10 +724,45 @@ function openProjectModal(p = {}) { toast('已保存(新主管/团队将在下次开工时生效)', 'ok'); closeModal(); router(); } else { const r = await api('/api/projects', {method: 'POST', body}); - toast('项目已创建,AI 主管已开工 🚀', 'ok'); closeModal(); + toast(`项目已创建,AI 主管已开工 🚀${r.copied_tasks ? `(从参考项目复制 ${r.copied_tasks} 个任务)` : ''}`, 'ok'); closeModal(); location.hash = `#/project/${r.id}`; } - } catch (e) { toast(e.message, 'err'); } + } catch (e) { + if (e.message && e.message.includes('need_confirm') === false) toast(e.message, 'err'); + else if (e.message) toast(e.message, 'err'); + } + }); +} + +/* 从参考项目中新建(V3.5) */ +async function openRefProjectModal() { + const r = await api('/api/reference_projects'); + const refs = r.data || []; + if (!refs.length) { toast('暂无参考项目', 'err'); return; } + openModal(` +

📋 从参考项目中新建

+
选择下面的内置简单测试项目作为参考,将复制其目标与任务列表,快速生成一个新项目(仍可继续编辑)。
+ ${refs.map(x => ` +
+
${esc(x.name)} ${x.tasks.length} 个任务
+
🎯 ${esc(x.objective)}
+
✅ 验收标准:${esc(x.acceptance_criteria)}
+
点击 → 以此参考新建项目
+
`).join('')} + `); +} + +async function openProjectModalFromRef(refId) { + const d = (await api('/api/reference_projects')).data; + const ref = d.find(x => x.id === refId); + if (!ref) return; + closeModal(); + openProjectModal({ + from_reference: ref.id, + name: ref.name.replace('(参考)', '(副本)'), + objective: ref.objective, + acceptance_criteria: ref.acceptance_criteria, + description: ref.description || '' }); } @@ -1691,37 +2044,62 @@ async function deliverAction(action) { } } -/* ---------- Worker 管理 ---------- */ -async function pageWorkers() { - const [ws, ps] = await Promise.all([api('/api/workers'), api('/api/providers')]); +/* ---------- Worker 管理(V3.5:AI Worker / 大模型接口库 / 团队) ---------- */ +let workersTab = 'list'; + +async function pageWorkers(args = []) { + workersTab = args[0] || 'list'; + const [ws, ps, eps, teams] = await Promise.all([api('/api/workers'), api('/api/providers'), api('/api/endpoints'), api('/api/teams')]); const canManage = CURRENT_ROLE === 'admin'; - $('#main').innerHTML = ` -

AI Worker 管理模型 + 角色提示词 + 工具权限 + 成本上限 = 虚拟员工档案

- ${canManage ? `
` : '
👁️ 仅显示你有权限查看的 Worker(由管理员授权)
'} - - ${ws.data.map(w => ` + const tabBar = ` + `; + $('#main').innerHTML = `

AI Worker模型接口 + 角色提示词 + 工具权限 + 成本上限 = 虚拟员工档案

${tabBar}
`; + if (workersTab === 'endpoints') renderEndpoints(eps.data, canManage); + else if (workersTab === 'teams') renderTeams(teams.data, ws.data, canManage); + else renderWorkersList(ws.data, ps.data, eps.data, canManage); +} + +function renderWorkersList(ws, ps, eps, canManage) { + const epName = eid => { const e = eps.find(x => x.id === eid); return e ? e.name : '—'; }; + $('#workers-body').innerHTML = ` + ${canManage ? `
` : '
👁️ 仅显示你有权限查看的 Worker(由管理员授权)
'} +
ID名称模型角色提示词温度/上限本月成本成本上限状态权限操作
+ ${ws.map(w => ` - + + - + - - `).join('') || ''} + `).join('') || ''}
ID名称接口库模型角色提示词温度/上限本月成本状态权限操作
#${w.id}${esc(w.name)}
${esc(w.description || '')}
${esc(w.name)}${w.is_main ? ' ⭐主力' : ''}
${esc(w.description || '')}
${esc(epName(w.endpoint_id))} ${esc(w.provider)} ${esc(w.model)}${esc((w.system_prompt || '—').slice(0, 40))}${esc((w.system_prompt || '—').slice(0, 36))} ${w.temperature} / ${w.max_tokens} ${fmtMoney(w.month_cost)}${w.task_cost_limit ? '单任务 ¥' + w.task_cost_limit : '—'}
${w.monthly_cost_limit ? '月 ¥' + w.monthly_cost_limit : ''}
${w.status === 'enabled' ? '启用' : '停用'} ${{view:'查看',use:'使用',manage:'管理',admin:'管理员'}[w.perm] || w.perm} - - ${canManage ? ` + + ${canManage && !w.is_main ? `` : ''} + ${canManage ? ` ` : ''}
还没有 Worker,注册第一个虚拟员工吧
还没有 Worker,先在「🔌 大模型接口库」确认接口,再注册第一个虚拟员工吧
`; } +async function setMainWorker(wid) { + try { + await api(`/api/workers/${wid}/main`, {method: 'POST'}); + toast('⭐ 已设置为主力 AI Worker(对话默认使用)', 'ok'); + router(); + } catch (e) { toast(e.message, 'err'); } +} + async function testWorker(wid) { - toast('正在测试连通性…'); + toast('正在测试连通性(流式)…'); try { const r = await api(`/api/workers/${wid}/test`, {method:'POST'}); - toast(`✅ 连通正常:${r.data.latency_ms}ms,回复「${r.data.reply}」,成本 ¥${r.data.cost}`, 'ok'); } + toast(`✅ 连通正常:${r.data.latency_ms}ms(首字 ${r.data.first_token_ms ?? '—'}ms),回复「${r.data.reply}」,成本 ¥${r.data.cost}`, 'ok'); } catch (e) { toast('❌ ' + e.message, 'err'); } } @@ -1752,19 +2130,32 @@ async function delWorker(wid) { toast('已删除', 'ok'); router(); } -function openWorkerModal(w, providers) { +function openWorkerModal(w, providers, eps) { w = w || {}; - const sel = (providerId) => providers.find(p => p.id === providerId); + providers = providers || []; + eps = eps || []; + const curEp = eps.find(e => String(e.id) === String(w.endpoint_id)); openModal(`

${w.id ? '编辑 Worker' : '注册 Worker'}

- + + ${eps.map(e => ``).join('')} + ${!eps.length ? '' : ''} + + + -
-
-
+
+
+
${w.id && w.endpoint_id ? '计价由接口库「' + esc(curEp?.name || '') + '」统一管理' : ''}
+
+
+
@@ -1785,27 +2176,45 @@ function openWorkerModal(w, providers) { ${w.id ? `` : ''}
`); - window.onProviderChange = () => { - const p = providers.find(x => x.id === $('#wm-provider').value); - if (!p) return; - const base = $('#wm-base'); - if (!base.value || base.placeholder === base.value) base.value = p.base_url; - base.placeholder = p.base_url; - const model = $('#wm-model'); - if (p.models.length && !model.value) model.value = p.models[0]; + window.onEndpointChange = () => { + const sel = $('#wm-endpoint'); + const opt = sel.selectedOptions[0]; + let models = []; + try { models = JSON.parse(opt?.dataset?.models || '[]'); } catch (e) {} + const m = $('#wm-model'); + m.innerHTML = models.map(x => ``).join('') || ''; + const info = $('#wm-price-info'); + if (opt) { + let pricing = {}; + try { pricing = JSON.parse(opt.dataset.pricing || '{}'); } catch (e) {} + const first = m.value || models[0] || ''; + const p = pricing[first]; + info.textContent = opt.dataset.billing === 'call' + ? '💰 按调用次数计费:¥' + opt.dataset.price + ' / 次' + : (p ? `💰 计价:输入 ¥${p.input}/1M · 输出 ¥${p.output}/1M tokens` : '💰 由接口库统一计价'); + } else info.textContent = ''; }; $('#wm-save').addEventListener('click', async () => { + const endpoint_id = Number($('#wm-endpoint').value || 0); + const customModel = $('#wm-model-custom')?.value.trim() || ''; const body = { - name: $('#wm-name').value.trim(), provider: $('#wm-provider').value, - model: $('#wm-model').value.trim(), base_url: $('#wm-base').value.trim(), - api_key: $('#wm-key').value.trim(), system_prompt: $('#wm-prompt').value, - description: $('#wm-desc').value.trim(), temperature: parseFloat($('#wm-temp').value || 0.7), + name: $('#wm-name').value.trim(), + endpoint_id: endpoint_id || null, + provider: endpoint_id ? (eps.find(e => e.id === endpoint_id)?.provider || 'custom') : (w.provider || 'custom'), + model: customModel || $('#wm-model').value.trim(), + base_url: $('#wm-base').value.trim(), + api_key: $('#wm-key').value.trim(), + system_prompt: $('#wm-prompt').value, + description: $('#wm-desc').value.trim(), + temperature: parseFloat($('#wm-temp').value || 0.7), max_tokens: parseInt($('#wm-max').value || 2000), task_cost_limit: parseFloat($('#wm-tasklimit').value || 0), monthly_cost_limit: parseFloat($('#wm-monthlimit').value || 0), status: $('#wm-status').value }; - if (!body.name || !body.provider || !body.model) return toast('名称/供应商/模型必填', 'err'); + if (!body.name) return toast('请填写 Worker 名称', 'err'); + if (!endpoint_id) return toast('必选:大模型接口(先在大模型接口库配置)', 'err'); + if (!body.model) return toast('请选择模型', 'err'); try { await api(w.id ? `/api/workers/${w.id}` : '/api/workers', {method: w.id ? 'PUT' : 'POST', body}); toast('已保存', 'ok'); closeModal(); router(); @@ -1813,19 +2222,185 @@ function openWorkerModal(w, providers) { }); } -/* ---------- 成本报表 ---------- */ +/* ---------- 🔌 大模型接口库(V3.5) ---------- */ +function renderEndpoints(eps, canManage) { + $('#workers-body').innerHTML = ` + ${canManage ? `
+ 专门配置大模型接口(地址/密钥/模型/定价),创建 AI Worker 时直接选用;支持按 token 或按调用次数计费
` : ''} + + ${eps.map(e => { + const billingTxt = e.billing === 'call' ? `按次 ¥${e.price_per_call}/次` : '按 token'; + return ` + + + + + + + + + + + `; + }).join('') || ''} +
ID名称提供商Base URL模型计费方式价格状态关联 Worker操作
#${e.id}${esc(e.name)}
${esc(e.description || '')}
${esc(e.provider)}${esc(e.base_url)}${(e.models || []).slice(0, 3).map(m => `${esc(m)}`).join('')}${(e.models || []).length > 3 ? `+${e.models.length - 3}` : ''}${billingTxt}${e.billing === 'call' ? '' : `输入 ¥${e.input_price}/1M · 输出 ¥${e.output_price}/1M`}${e.status === 'enabled' ? '启用' : '停用'}${e.worker_count ? `${e.worker_count} 个` : '—'} + ${canManage ? ` + ` : ''}
还没有大模型接口,点右上角添加(或系统已从 config.py 自动导入)
`; +} + +function openEndpointModal(e = {}) { + const models = e.models || []; + const pricing = e.pricing || {}; + openModal(` +

${e.id ? '编辑大模型接口' : '添加大模型接口'}

+ +
+
+
+
+ + + + +
+ + +
+
+
+
+
+ + + + `); + window.toggleBilling = () => { + const call = $('#ep-billing').value === 'call'; + $('#ep-token-pricing').style.display = call ? 'none' : ''; + $('#ep-call-pricing').style.display = call ? '' : 'none'; + }; + toggleBilling(); + $('#ep-save').addEventListener('click', async () => { + const models = $('#ep-models').value.split(/\n/).map(s => s.trim()).filter(Boolean); + const pricing = {}; + $('#ep-pricing').value.split(/\n/).forEach(l => { + const parts = l.trim().split(/\s+/); + if (parts.length >= 3 && !isNaN(parseFloat(parts[1]))) pricing[parts[0]] = {input: parseFloat(parts[1]), output: parseFloat(parts[2])}; + }); + const body = { + name: $('#ep-name').value.trim(), provider: $('#ep-provider').value, + base_url: $('#ep-base').value.trim(), api_key: $('#ep-key').value.trim(), + models, pricing, billing: $('#ep-billing').value, + input_price: parseFloat($('#ep-in').value || 0), output_price: parseFloat($('#ep-out').value || 0), + price_per_call: parseFloat($('#ep-percall').value || 0), + description: $('#ep-desc').value.trim(), status: $('#ep-status').value + }; + if (!body.name || !body.base_url) return toast('名称 / Base URL 必填', 'err'); + if (!models.length) return toast('请填写至少一个模型', 'err'); + try { + await api(e.id ? `/api/endpoints/${e.id}` : '/api/endpoints', {method: e.id ? 'PUT' : 'POST', body}); + toast('已保存', 'ok'); closeModal(); router(); + } catch (err) { toast(err.message, 'err'); } + }); +} + +async function testEndpoint(eid) { + toast('正在测试连通性…'); + try { + const r = await api(`/api/endpoints/${eid}/test`, {method:'POST'}); + toast(`✅ 连通正常:${r.data.latency_ms}ms(首字 ${r.data.first_token_ms ?? '—'}ms),回复「${r.data.reply}」,成本 ¥${r.data.cost}`, 'ok'); + } catch (e) { toast('❌ ' + e.message, 'err'); } +} + +async function delEndpoint(eid) { + if (!confirm('删除该大模型接口?绑定它的 Worker 仍保留但计价回退默认价。')) return; + try { + await api(`/api/endpoints/${eid}`, {method: 'DELETE'}); + toast('已删除', 'ok'); router(); + } catch (e) { toast(e.message, 'err'); } +} + +/* ---------- 👥 团队(V3.5) ---------- */ +function renderTeams(teams, ws, canManage) { + $('#workers-body').innerHTML = ` + ${canManage ? `
+ 团队 = 一组 AI Worker 的集合,可在「对话」中选择团队对话,或在创建项目时一键勾选为干活团队
` : ''} + + ${teams.map(t => ` + + + + + + + + `).join('') || ''} +
ID团队名称说明成员 Worker可用成员操作
#${t.id}${esc(t.name)}${esc(t.description || '—')}${(t.workers || []).map(w => `${esc(w.name)}${w.status !== 'enabled' ? ' 🚫' : ''}`).join(' ') || '—'}${(t.enabled_workers || []).length} 人可用${canManage ? ` + ` : ``}
还没有团队,创建一个把常用 Worker 打包吧
`; +} + +function openTeamModal(t = {}, allWorkers = []) { + const memberSet = new Set(t.worker_ids || (t.workers || []).map(w => w.id)); + openModal(` +

${t.id ? '编辑团队' : '创建团队'}

+ + + +
+ ${(allWorkers.length ? allWorkers : (t.workers || [])).map(w => ``).join('') || '请先注册 Worker'} +
+ `); + $('#tm-save').addEventListener('click', async () => { + const worker_ids = [...$$('#modal-box .check-group input:checked')].map(x => Number(x.value)); + if (!$('#tm-name').value.trim()) return toast('请填写团队名称', 'err'); + if (!worker_ids.length) return toast('至少选择 1 个成员', 'err'); + const body = {name: $('#tm-name').value.trim(), description: $('#tm-desc').value.trim(), worker_ids}; + try { + await api(t.id ? `/api/teams/${t.id}` : '/api/teams', {method: t.id ? 'PUT' : 'POST', body}); + toast('已保存', 'ok'); closeModal(); router(); + } catch (e) { toast(e.message, 'err'); } + }); +} + +async function delTeam(tid) { + if (!confirm('删除该团队?(成员 Worker 不受影响)')) return; + await api(`/api/teams/${tid}`, {method: 'DELETE'}); + toast('已删除', 'ok'); router(); +} + +/* ---------- 成本报表 / 用量明细 ---------- */ async function pageReports() { const group = location.hash.includes('group=') ? location.hash.split('group=')[1] : 'project'; + if (group === 'usage') return pageUsageReport(); const r = (await api('/api/reports/cost?group=' + group)).data; const s = (await api('/api/stats')).data; - const head = {project:['项目','调用次数','Tokens','成本'], worker:['Worker','模型','调用次数','Tokens','成本'], model:['供应商/模型','调用次数','Tokens','成本']}[group]; + const head = {project:['项目','调用次数','输入tokens','输出tokens','缓存命中','总Tokens','成本'], worker:['Worker','模型','调用次数','输入tokens','输出tokens','缓存命中','总Tokens','成本'], model:['供应商/模型','调用次数','输入tokens','输出tokens','缓存命中','总Tokens','成本']}[group]; + const fmtT = n => (n / 1e6).toFixed(3) + 'M'; $('#main').innerHTML = ` -

成本报表按项目 × 任务 × Worker × 模型 多维核算

+

成本报表按项目 × 任务 × Worker × 模型 多维核算(含输入/输出/缓存命中 token 与调用次数)

按项目 按 Worker 按模型 - 累计成本 ${fmtMoney(s.total_cost)} · ${(s.total_tokens/1e6).toFixed(2)}M tokens + 📊 用量明细(项目×智能体) + 累计成本 ${fmtMoney(s.total_cost)} · ${(s.total_tokens/1e6).toFixed(2)}M tokens · ${s.total_calls ?? 0} 次调用 · 缓存命中 ${(s.total_cached ?? 0)/1e6 >= 0.001 ? (s.total_cached/1e6).toFixed(2)+'M' : (s.total_cached ?? 0)} tokens
${head.map(h => ``).join('')} ${r.map(x => { @@ -1834,12 +2409,58 @@ async function pageReports() { return `${group === 'project' ? `` : group === 'worker' ? `` : ``} - + + + + + `; - }).join('') || ''} + }).join('') || ''}
${h}
${esc(x.project_name || '—')}${esc(x.worker_name || '—')}${esc(x.provider || '')} ${esc(x.model || '')}${esc(x.provider || '')} ${esc(x.model || '')}${x.runs}${(tokens/1e6).toFixed(3)}M${x.runs}${fmtT(x.prompt_tokens || 0)}${fmtT(x.completion_tokens || 0)}${fmtT(x.cached_tokens || 0)}${fmtT(tokens)} ${fmtMoney(cost)}
暂无成本数据,先派几个任务吧
暂无成本数据,先派几个任务吧
`; } +async function pageUsageReport() { + const d = (await api('/api/reports/usage')).data; + const fmtT = n => (n / 1e6).toFixed(3) + 'M'; + const s = (await api('/api/stats')).data; + const cols = (r) => ` + ${r.calls} + ${fmtT(r.prompt_tokens || 0)} + ${fmtT(r.completion_tokens || 0)} + ${fmtT(r.cached_tokens || 0)} + ${fmtT(r.total_tokens || 0)} + ${fmtMoney(r.cost || 0)}`; + $('#main').innerHTML = ` +

用量明细每个项目 × 每个智能体:调用次数 / 输入 / 输出 / 缓存命中 token / 成本

+
+ 按项目 + 按 Worker + 按模型 + 📊 用量明细(项目×智能体) + 累计 ${s.total_calls ?? 0} 次调用 · 缓存命中 ${(s.total_cached ?? 0)/1e6 >= 0.001 ? (s.total_cached/1e6).toFixed(2)+'M' : (s.total_cached ?? 0)} tokens +
+
+
📊 项目 × 智能体 明细矩阵
+ + ${d.matrix.map(r => `${cols(r)}`).join('') || ''} +
项目智能体模型调用次数输入tokens输出tokens缓存命中总Tokens成本
${esc(r.project_name)}${esc(r.worker_name)}${esc(r.model || '')}
暂无用量数据
+
+
+
+
按项目汇总
+ + ${d.by_project.map(r => `${cols(r)}`).join('') || ''} +
项目调用输入输出缓存总Tokens成本
${esc(r.project_name)}
+
+
+
按智能体汇总
+ + ${d.by_worker.map(r => `${cols(r)}`).join('') || ''} +
智能体调用输入输出缓存总Tokens成本
${esc(r.worker_name)}
+
+
`; +} + /* ---------- 运行日志 ---------- */ let logsState = {page: 1, pageSize: 50, level: '', q: ''}; window.logsState = logsState; // 供分页组件通过 window[stateVar] 读写 @@ -2026,11 +2647,14 @@ async function delToken(id) { toast('已删除', 'ok'); pageApiTokens(); } -/* ---------- 通知设置 ---------- */ +/* ---------- 通知设置(V3.5:含执行超时 / 系统工作目录 / 主力 Worker) ---------- */ async function pageSettings() { - const [chans, evs, st] = await Promise.all([api('/api/channels'), api('/api/events'), api('/api/settings')]); + const [chans, evs, st, workers, ws] = await Promise.all([ + api('/api/channels'), api('/api/events'), api('/api/settings'), + api('/api/workers'), api('/api/settings/workspace') + ]); $('#main').innerHTML = ` -

通知与设置飞书 / 企微群机器人 · 邮件 · 预算告警阈值

+

通知与设置飞书 / 企微群机器人 · 邮件 · 预算告警 · 执行超时 · 系统工作目录 · 主力 Worker

@@ -2042,26 +2666,91 @@ async function pageSettings() { ${(c.events || []).map(e => `${(evs.data.find(x => x.id === e) || {}).name || e}`).join('') || '—'} ${c.enabled ? '启用' : '停用'} - + `).join('') || '暂无渠道,添加飞书/企微机器人后任务事件自动推送'} -
💡 飞书群机器人:群设置 → 群机器人 → 添加自定义机器人,复制 Webhook 地址(以 https://open.feishu.cn/open-apis/bot/v2/hook/ 开头)。企微同理(https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=…)。
+
💡 飞书群机器人:群设置 → 群机器人 → 添加自定义机器人,复制 Webhook 地址。企微同理。
-
预算告警
- -
-
-
告警会写入告警中心并推送到订阅了「预算告警」事件的渠道。
-
邮件通知
-
${st.data.email_configured ? '✅ 已配置 SMTP(config.py EMAIL)' : '⚠️ 未配置 SMTP(config.py 的 EMAIL 填写 host/user/password 后重启生效)'}
+
⏱️ 执行超时(V3.5 流式)
+
所有模型输出按 token 流式接收:单 token 返回超时 = 相邻两个 token 数据块的最大间隔;首字延迟超时 = 请求发出后首块数据的最长等待。命中即中断并标记失败。
+
+
+
+
+
+
+
+
+ +
⭐ 主力 AI Worker(对话默认使用)
+ +
+
+
+
🗂️ 系统工作目录(V3.5)
+
所有项目与多 Agent 协作都在此目录下新建独立工作目录(自动无重复)。默认 ${esc(ws.data.default_root)};可手动改成任意绝对路径——不存在则自动创建,已存在会列出目录信息并需手动确认。
+
+
+ + +
+
${ws.data.info ? wsRootInfoHtml(ws.data.info) : '
✅ 目录已就绪(不存在时会自动创建)
'}
+ +
邮件通知
+
${st.data.email_configured ? '✅ 已配置 SMTP(config.py EMAIL)' : '⚠️ 未配置 SMTP(config.py 的 EMAIL 填写 host/user/password 后重启生效)'}
`; + // 主力 Worker + window.saveMainWorker = async () => { + try { + await api('/api/settings', {method:'PUT', body:{main_worker_id: Number($('#st-main').value)}}); + toast('⭐ 主力 AI Worker 已更新', 'ok'); + } catch (e) { toast(e.message, 'err'); } + }; + window.probeWsRoot = async () => { + const path = $('#st-wsroot').value.trim(); + const box = $('#st-wsinfo'); + if (!path) return; + box.innerHTML = '
检查中…
'; + try { + const r = await api('/api/workspace/probe?path=' + encodeURIComponent(path)); + if (r.data.exists) { + box.innerHTML = wsRootInfoHtml(r.data); + $('#st-wsconfirm-wrap').style.display = 'flex'; + } else { + box.innerHTML = `
✅ 目录不存在:${esc(r.data.path)},切换时将自动创建
`; + $('#st-wsconfirm-wrap').style.display = 'none'; + } + } catch (e) { box.innerHTML = `
${esc(e.message)}
`; } + }; + window.saveWsRoot = async () => { + const path = $('#st-wsroot').value.trim(); + if (!path) return toast('请输入目录路径', 'err'); + try { + const r = await api('/api/settings/workspace', {method:'POST', body:{path, confirm_existing: !!$('#st-wsconfirm')?.checked}}); + toast('✅ 系统工作目录已切换:' + r.root, 'ok'); + pageSettings(); + } catch (e) { toast(e.message, 'err'); pageSettings(); } + }; +} + +function wsRootInfoHtml(info) { + return `
📂 ${esc(info.path)}
+
文件 ${info.files} 个 · 大小 ${(info.size/1024/1024).toFixed(2)} MB${info.newest_mtime ? ' · 最近修改 ' + fmtTime(info.newest_mtime) : ''}
+
${(info.sample || []).map(s => '
· ' + esc(s) + '
').join('') || ''}
`; } async function saveSettings() { - await api('/api/settings', {method:'PUT', body:{budget_alert_ratio: parseFloat($('#st-ratio').value)}}); - toast('已保存', 'ok'); + await api('/api/settings', {method:'PUT', body:{ + budget_alert_ratio: parseFloat($('#st-ratio').value), + token_timeout: parseFloat($('#st-token').value), + first_token_timeout: parseFloat($('#st-first').value), + request_timeout: parseFloat($('#st-req').value) + }}); + toast('已保存(新超时设置立即生效)', 'ok'); } async function testChannel(id) { toast('正在发送测试消息…'); diff --git a/static/index.html b/static/index.html index e4b49db..640d032 100644 --- a/static/index.html +++ b/static/index.html @@ -11,7 +11,7 @@