Pamudu13/gemma-3-chat
0
1# -- coding: utf-8 --2 3import inspect4import signal5import struct6import sys7import os8import argparse9import socket10import errno11 12from PIL import Image13 14_is_steam_build = os.environ.get("IS_STEAM_BUILD", "0") == "1"15from py.cli_tool import read_file_tool_local16from py.task_tools import query_task_progress17from py.ws_manager import ws_manager18 19# === Pre-load heavy tool modules at startup to avoid blocking first request ===20from py.web_search import (21 DDGsearch, searxng, Tavily_search, Google_search,22 Brave_search, Exa_search, Serper_search, bochaai_search,23 jina_crawler, Crawl4Ai_search, firecrawl_search, simple_fetch, markdown_new,24 duckduckgo_tool, searxng_tool, tavily_tool, google_tool,25 brave_tool, exa_tool, serper_tool, bochaai_tool,26 jina_crawler_tool, simple_fetch_tool, Crawl4Ai_tool, firecrawl_tool, markdown_new_tool,27)28from py.know_base import kb_tool, query_knowledge_base, rerank_knowledge_base29from py.agent_tool import get_agent_tool, agent_tool_call30from py.a2a_tool import get_a2a_tool, a2a_tool_call31from py.llm_tool import get_llm_tool, custom_llm_tool32from py.pollinations import (33 pollinations_image_tool, openai_image_tool, openai_chat_image_tool,34 pollinations_image, openai_image, openai_chat_image,35)36from py.code_interpreter import e2b_code_tool, local_run_code_tool, e2b_code, local_run_code37from py.custom_http import fetch_custom_http38from py.comfyui_tool import comfyui_tool_call39from py.utility_tools import (40 time_tool, weather_tool, location_tool, timer_weather_tool,41 wikipedia_summary_tool, wikipedia_section_tool, arxiv_tool,42 get_weather, get_location_coordinates, get_weather_by_city,43 get_wikipedia_summary_and_sections, get_wikipedia_section_content, search_arxiv_papers,44)45from py.autoBehavior import auto_behavior_tool, auto_behavior46from py.guard import load_safety_words, check_content_safety47from py.cdp_tool import (48 all_cdp_tools, list_pages, navigate_page, new_page, close_page, select_page,49 take_snapshot, wait_for, click, fill, hover, press_key, evaluate_script,50 take_screenshot, fill_form, drag, handle_dialog,51)52from py.computer_use_tool import (53 computer_use_tools, mouse_use_tools, keyboard_use_tools, desktopVision_use_tools,54 mouse_move, mouse_click, mouse_double_click, mouse_drag, mouse_scroll, mouse_hold,55 copy_to_input_box, keyboard_press, keyboard_sequence, keyboard_hotkey, keyboard_hold,56 logical_type, wait, screenshot, logical_click,57)58from py.mode_change import mode_change_tool, update_workspace_settings59from py.acpx_tools import acp_agent_tool, acpx_agent60 61# Extended CLI tool imports for dispatch_tool62from py.cli_tool import (63 docker_sandbox, list_files_tool, read_file_tool, read_file_range_tool,64 tail_file_tool, search_files_tool, edit_file_tool,65 edit_file_string_tool, glob_files_tool, todo_write_tool, list_processes_tool,66 get_process_logs_tool, kill_process_tool, docker_manage_ports_tool,67 read_skill_tool, shell_tool_local, list_files_tool_local,68 read_file_tool_local, read_file_range_tool_local, tail_file_tool_local,69 search_files_tool_local, edit_file_tool_local,70 edit_file_string_tool_local, glob_files_tool_local, todo_write_tool_local,71 local_net_tool, send_process_input_tool, read_skill_tool_local,72 get_tools_for_mode, get_local_tools_for_mode,73)74from py.task_tools import (75 create_subtask_tool, query_tasks_tool, cancel_subtask_tool, finish_task_tool, finish_main_task_tool,76 create_subtask, cancel_subtask, finish_task, finish_main_task,77)78from py.load_files import get_files_content, file_tool, image_tool79 80import shortuuid81if not os.environ.get("HOME") or os.environ.get("HOME") == "/":82 os.environ["HOME"] = "/tmp"83os.environ["MEM0_DIR"] = os.path.join(os.environ.get("HOME", "/tmp"), ".mem0")84os.environ["MEM0_TELEMETRY"] = "False"85parser = argparse.ArgumentParser(description="Run the ASGI application server.")86parser.add_argument("--host", default="127.0.0.1")87parser.add_argument("--port", type=int, default=3456)88args, _ = parser.parse_known_args()89 90HOST = args.host91PREFERED_PORT = args.port92 93def is_addr_in_use_error(e):94 """跨平台判断是否为地址被占用错误"""95 if hasattr(e, 'errno'):96 if e.errno == errno.EADDRINUSE:97 return True98 # Windows 有时用 WSAEADDRINUSE (10048)99 if sys.platform == 'win32' and e.errno == 10048:100 return True101 # Windows winerror 属性102 if hasattr(e, 'winerror') and e.winerror == 10048:103 return True104 # macOS/Linux 错误消息105 if 'address already in use' in str(e).lower():106 return True107 return False108 109def is_permission_error(e):110 """跨平台判断是否为权限/拒绝访问错误"""111 if isinstance(e, PermissionError):112 return True113 if hasattr(e, 'errno'):114 if e.errno in (errno.EACCES, errno.EPERM):115 return True116 # Windows ERROR_ACCESS_DENIED (5)117 if sys.platform == 'win32' and e.errno == 13:118 return True119 if hasattr(e, 'winerror') and e.winerror in (5, 10013):120 return True121 err_str = str(e).lower()122 if any(x in err_str for x in ['permission', 'denied', 'access', 'not permitted']):123 return True124 return False125 126def force_bind_or_fallback(host, preferred_port):127 """128 跨平台端口绑定:129 1. 尝试强制绑定指定端口(处理TIME_WAIT)130 2. 如果被真正占用/无权限/系统保留,自动降级到随机端口131 3. 绝不抛出异常导致退出132 """133 # 尝试绑定首选端口134 sock = None135 try:136 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)137 # 关键:允许快速复用 TIME_WAIT 状态的端口138 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)139 sock.bind((host, preferred_port))140 sock.close()141 return preferred_port142 143 except (socket.error, OSError, PermissionError) as e:144 # 判断错误类型145 if is_addr_in_use_error(e):146 reason = "in use"147 elif is_permission_error(e):148 reason = "permission denied/system reserved"149 else:150 reason = f"error ({e})"151 152 print(f"Port {preferred_port} unavailable ({reason}), auto-assigning...", 153 file=sys.stderr, flush=True)154 155 # 关闭失败的 socket156 try:157 if sock:158 sock.close()159 except:160 pass161 162 # 降级:让系统分配端口163 return auto_assign_port(host)164 165 except Exception as e:166 # 捕获所有其他异常167 print(f"Unexpected error binding port {preferred_port}: {e}, auto-assigning...", 168 file=sys.stderr, flush=True)169 try:170 if sock:171 sock.close()172 except:173 pass174 return auto_assign_port(host)175 176def auto_assign_port(host):177 """自动分配可用端口,带多重降级"""178 # 尝试 127.0.0.1179 try:180 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)181 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)182 sock.bind((host, 0))183 port = sock.getsockname()[1]184 sock.close()185 print(f"Auto-assigned port: {port}", file=sys.stderr, flush=True)186 return port187 except Exception as e:188 print(f"Failed to bind {host}: {e}", file=sys.stderr, flush=True)189 try:190 sock.close()191 except:192 pass193 194 # 降级 1: 尝试 0.0.0.0 (所有接口)195 if host != "0.0.0.0":196 try:197 print("Trying 0.0.0.0...", file=sys.stderr, flush=True)198 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)199 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)200 sock.bind(("0.0.0.0", 0))201 port = sock.getsockname()[1]202 sock.close()203 print(f"Auto-assigned port on 0.0.0.0: {port}", file=sys.stderr, flush=True)204 return port205 except Exception as e:206 print(f"Failed to bind 0.0.0.0: {e}", file=sys.stderr, flush=True)207 try:208 sock.close()209 except:210 pass211 212 # 降级 2: 尝试 localhost213 if host != "localhost":214 try:215 print("Trying localhost...", file=sys.stderr, flush=True)216 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)217 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)218 sock.bind(("localhost", 0))219 port = sock.getsockname()[1]220 sock.close()221 print(f"Auto-assigned port on localhost: {port}", file=sys.stderr, flush=True)222 return port223 except Exception as e:224 print(f"Failed to bind localhost: {e}", file=sys.stderr, flush=True)225 try:226 sock.close()227 except:228 pass229 230 # 最后手段:硬编码高位端口(极端情况)231 fallback_ports = [45678, 45679, 45680, 0]232 for fp in fallback_ports:233 try:234 sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)235 sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)236 sock.bind((host if host != "0.0.0.0" else "127.0.0.1", fp))237 port = sock.getsockname()[1]238 sock.close()239 print(f"Fallback to hardcoded port: {port}", file=sys.stderr, flush=True)240 return port241 except:242 try:243 sock.close()244 except:245 pass246 continue247 248 # 理论上不会到这里,如果真的到了,返回一个肯定能用的249 return 0250 251# 执行端口查找252FINAL_PORT = force_bind_or_fallback(HOST, PREFERED_PORT)253PORT = FINAL_PORT254os.environ['DYNAMIC_PORT'] = str(FINAL_PORT)255 256# 同时调用 change_port 保持同步257from py.get_setting import change_port, reset_user_data_dir, set_custom_user_data_dir258change_port(FINAL_PORT)259 260# ==========================================261# 第二步:屏蔽掉后面库可能产生的骚扰警告262# ==========================================263import warnings264warnings.filterwarnings("ignore") # 忽略普通警告265os.environ['TF_CPP_MIN_LOG_LEVEL'] = '3' # 如果有 tensorflow 等库,减少其日志输出266 267import hashlib268import importlib269import mimetypes270import pathlib271import sys272import traceback273import platform274import requests275 276from py.agent import add_tool_to_project_config, is_tool_allowed_by_project_config277sys.stdout.reconfigure(encoding='utf-8')278import base64279from datetime import datetime280import glob281from io import BytesIO282import io283import os284from pathlib import Path285import pickle286import socket287import sys288import tempfile289import httpx290import ipaddress291from urllib.parse import urlparse, urlunparse, urljoin292from urllib.robotparser import RobotFileParser293import websockets294from py.load_files import check_robots_txt, get_file_content, is_private_ip, sanitize_url295 296# 修复 sherpa-onnx 在 macOS arm64 上的 onnxruntime dylib 路径问题297import site298try:299 _sp = site.getsitepackages()[0]300 _sherpa_lib = os.path.join(_sp, "sherpa_onnx", "lib")301 _onnx_capi = os.path.join(_sp, "onnxruntime", "capi")302 import glob as _glob303 _dylibs = _glob.glob(os.path.join(_onnx_capi, "libonnxruntime*.dylib"))304 if _dylibs:305 _dylib = _dylibs[0]306 _target = os.path.join(_sherpa_lib, os.path.basename(_dylib))307 if not os.path.exists(_target):308 os.makedirs(_sherpa_lib, exist_ok=True)309 os.symlink(os.path.abspath(_dylib), _target)310except Exception:311 pass312 313def fix_macos_environment():314 """315 专门修复 macOS 下找不到 node (nvm) 和 uv (python framework) 的问题316 """317 if sys.platform != 'darwin':318 return319 320 user_home = Path.home()321 paths_to_add = []322 323 # ---------------------------------------------------------324 # 1. 自动发现 NVM 安装的 Node.js325 # 路径通常是: ~/.nvm/versions/node/vX.X.X/bin326 # ---------------------------------------------------------327 nvm_path = user_home / ".nvm" / "versions" / "node"328 if nvm_path.exists():329 # 获取所有版本文件夹 (如 v20.19.5, v18.0.0)330 # 使用 glob 匹配所有 v 开头的文件夹331 node_versions = sorted(nvm_path.glob("v*"), key=lambda p: p.name, reverse=True)332 333 # 将所有版本的 bin 目录都加入,或者只加最新的334 for version_dir in node_versions:335 bin_path = version_dir / "bin"336 if bin_path.exists():337 paths_to_add.append(str(bin_path))338 # 如果只想用最新的 node,这里可以 break339 # break 340 341 # ---------------------------------------------------------342 # 2. 自动发现 Python Framework 中的 uv343 # 路径通常是: /Library/Frameworks/Python.framework/Versions/X.X/bin344 # ---------------------------------------------------------345 py_framework_path = Path("/Library/Frameworks/Python.framework/Versions")346 if py_framework_path.exists():347 # 查找所有版本,如 3.13, 3.12348 py_versions = py_framework_path.glob("*")349 for ver in py_versions:350 bin_path = ver / "bin"351 if bin_path.exists():352 paths_to_add.append(str(bin_path))353 354 # ---------------------------------------------------------355 # 3. 补充 macOS 常见的其他路径 (Homebrew, Cargo, Local)356 # uv 也经常被安装在 .local/bin 或 .cargo/bin 下357 # ---------------------------------------------------------358 common_extras = [359 "/opt/homebrew/bin", # Apple Silicon Mac Homebrew360 "/usr/local/bin", # Intel Mac Homebrew361 str(user_home / ".local" / "bin"), # 用户级安装通常在这里362 str(user_home / ".cargo" / "bin"), # Rust 工具链 (uv 可能在这里)363 ]364 paths_to_add.extend(common_extras)365 366 # ---------------------------------------------------------367 # 4. 将发现的路径注入到当前进程的环境变量中368 # ---------------------------------------------------------369 current_path = os.environ.get("PATH", "")370 new_path_str = current_path371 372 # 将新路径加到最前面 (优先级最高)373 for p in paths_to_add:374 if p and os.path.isdir(p):375 # 避免重复添加376 if p not in new_path_str:377 new_path_str = p + os.pathsep + new_path_str378 379 # 更新环境变量380 os.environ['PATH'] = new_path_str381 382 # (可选) 打印调试信息383 # print(f"Fixed macOS PATH. Added: {paths_to_add}")384 385# --- 在程序最开始的地方调用这个函数 ---386fix_macos_environment()387 388def _fix_onnx_dll():389 if sys.platform == 'darwin':390 return391 # 1. 找到 uv 虚拟环境里的 onnxruntime392 spec = importlib.util.find_spec("onnxruntime")393 if spec is None or spec.origin is None:394 return # 没装 onnxruntime,随它去395 # DLL 就在 site-packages/onnxruntime/capi 里396 dll_dir = pathlib.Path(spec.origin).with_name("capi")397 if not dll_dir.is_dir():398 return399 400 # 2. 置顶搜索路径401 os.environ["PATH"] = str(dll_dir) + os.pathsep + os.environ["PATH"]402 if hasattr(os, "add_dll_directory"): # Python 3.8+403 os.add_dll_directory(str(dll_dir))404 405 # 3. 如果已经有人 import 过 onnxruntime,清掉缓存406 for mod in list(sys.modules):407 if mod.startswith("onnxruntime"):408 del sys.modules[mod]409 410_fix_onnx_dll()411 412# 在程序最开始设置413if hasattr(sys, '_MEIPASS'):414 # 打包后的程序415 os.environ['PYTHONPATH'] = sys._MEIPASS416 os.environ['PATH'] = sys._MEIPASS + os.pathsep + os.environ.get('PATH', '')417import asyncio418import copy419from functools import partial420import json421import re422import shutil423from fastapi import BackgroundTasks, Body, FastAPI, File, Form, HTTPException, UploadFile, WebSocket, Request, WebSocketDisconnect424from fastapi_mcp import FastApiMCP425import logging426from fastapi.staticfiles import StaticFiles427from fastapi.middleware.cors import CORSMiddleware428from openai import AsyncOpenAI429from pydantic import BaseModel430from fastapi import status431from fastapi.responses import HTMLResponse, JSONResponse, StreamingResponse,Response432import uuid433import time434from typing import Any, AsyncIterator, List, Dict,Optional, Tuple, Union435import shortuuid436from py.mcp_clients import McpClient437from contextlib import asynccontextmanager, suppress438from concurrent.futures import ThreadPoolExecutor439import aiofiles440import argparse441from py.dify_openai import DifyOpenAIAsync442from py.ClaudeAsOpenAI import AsyncClaudeAsOpenAI443from py.GeminiAsOpenAI import AsyncGeminiAsOpenAI444from py.get_setting import EXT_DIR, IS_DOCKER, SKILLS_DIR, _copy_default_skills, convert_to_opus_simple, load_covs, load_settings, save_covs,save_settings,clean_temp_files_task,base_path,configure_host_port,UPLOAD_FILES_DIR,AGENT_DIR,MEMORY_CACHE_DIR,KB_DIR,DEFAULT_VRM_DIR,DEFAULT_THA_DIR,THA_USER_MODELS_DIR,USER_DATA_DIR,LOG_DIR,TOOL_TEMP_DIR,COVS_PATH,DATABASE_PATH445from py.llm_tool import get_image_base64,get_image_media_type446timetamp = time.time()447log_path = os.path.join(LOG_DIR, f"backend_{timetamp}.log")448 449logger = None 450os.environ["no_proxy"] = "localhost,127.0.0.1"451local_timezone = None452settings = None453client = None454fast_client = None 455reasoner_client = None456HA_client = None457ChromeMCP_client = None458sql_client = None459mcp_client_list = {}460node_ext_mcp_clients: Dict[str, McpClient] = {}461node_ext_mcp_tools: Dict[str, List[Dict]] = {} # 存储每个扩展的工具列表462locales = {}463sleep_guard = None464scheduler_task = None465global_http_client = None # 用于共享底层的 TCP 连接池466openai_tts_clients_cache = {} # 缓存 OpenAI TTS Client467tetos_speakers_cache = {} # 缓存 Tetos Speaker 对象468openai_asr_clients_cache = {}469_TOOL_HOOKS = {}470ALLOWED_EXTENSIONS = [471 # 办公文档472 'doc', 'docx', 'ppt', 'pptx', 'xls', 'xlsx', 'pdf', 'pages', 473 'numbers', 'key', 'rtf', 'odt', 'epub',474 475 # 编程开发476 'js', 'ts', 'py', 'java', 'c', 'cpp', 'h', 'hpp', 'go', 'rs',477 'swift', 'kt', 'dart', 'rb', 'php', 'html', 'css', 'scss', 'less',478 'vue', 'svelte', 'jsx', 'tsx', 'json', 'xml', 'yml', 'yaml', 479 'sql', 'sh',480 481 # 数据配置482 'csv', 'tsv', 'txt', 'md', 'log', 'conf', 'ini', 'env', 'toml'483]484ALLOWED_IMAGE_EXTENSIONS = ['png', 'jpg', 'jpeg', 'gif', 'webp', 'bmp']485 486ALLOWED_VIDEO_EXTENSIONS = ['mp4', 'webm', 'ogg', 'mov', 'avi']487 488# 1. 先清空系统可能给错的条目489for ext in ("js", "mjs", "css", "html", "htm", "json", "xml", "map", "svg"):490 mimetypes.add_type("", f".{ext}") # 先删掉491# 2. 再写死我们想要的492mimetypes.add_type("application/javascript", ".js")493mimetypes.add_type("application/javascript", ".mjs")494mimetypes.add_type("text/css", ".css")495mimetypes.add_type("text/html", ".html")496mimetypes.add_type("text/html", ".htm")497mimetypes.add_type("application/json", ".json")498mimetypes.add_type("application/xml", ".xml")499mimetypes.add_type("application/json", ".map")500mimetypes.add_type("image/svg+xml", ".svg")501 502import platform503import ctypes504import io505if platform.system() == "Windows":506 try:507 # 设置 DPI 感知,确保截屏尺寸和 size() 返回的一致508 ctypes.windll.shcore.SetProcessDpiAwareness(1) 509 except Exception:510 ctypes.windll.user32.SetProcessDPIAware()511 512def draw_grid_on_image(image, grid_spacing: int = 10):513 """在图片上绘制网格和千分比坐标标签"""514 from PIL import ImageDraw515 draw = ImageDraw.Draw(image)516 width, height = image.size517 518 # 颜色设置 (半透明红色或亮绿色,视情况而定)519 line_color = (255, 0, 0, 128) # 红色线520 text_color = (255, 0, 0, 255)521 522 # 绘制垂直线 (百分比 0-100,但标签显示为千分比 0-1000‰)523 for x_pc in range(0, 101, grid_spacing):524 x = int(width * (x_pc / 100.0))525 # 确保不超出边界526 x = min(x, width - 1)527 draw.line([(x, 0), (x, height)], fill=line_color, width=1)528 x_permille = x_pc529 draw.text((x + 2, 5), f"{x_permille}%", fill=text_color)530 531 # 绘制水平线532 for y_pc in range(0, 101, grid_spacing):533 y = int(height * (y_pc / 100.0))534 y = min(y, height - 1)535 draw.line([(0, y), (width, y)], fill=line_color, width=1)536 y_permille = y_pc537 draw.text((5, y + 2), f"{y_permille}%", fill=text_color)538 539 return image540 541def draw_action_feedback(image, action_str: str):542 """543 解析返回结果字符串,并在图像上绘制动作反馈轨迹。544 (已针对红色网格优化,全面移除红色,使用高对比度的青/蓝/绿/黄色)545 """546 from PIL import ImageDraw, Image547 image = image.convert("RGBA")548 overlay = Image.new("RGBA", image.size, (255, 255, 255, 0))549 draw = ImageDraw.Draw(overlay)550 w, h = image.size551 552 def to_px(tx, ty):553 return int(float(tx) * w / 1000), int(float(ty) * h / 1000)554 555 # 1. 匹配 MOVE(x,y) -> 画个白色小圆点带黑边556 move_match = re.search(r"\[LAST_ACTION: MOVE\((\d+\.?\d*),(\d+\.?\d*)\)\]", action_str)557 if move_match:558 x, y = move_match.groups()559 px, py = to_px(x, y)560 r = 6561 draw.ellipse([px-r, py-r, px+r, py+r], fill=(255, 255, 255, 200), outline=(0, 0, 0, 255), width=1)562 563 # 2. 匹配 CLICK(x,y) -> 画个青色(Cyan)半透明十字靶心 (对比红色网格极佳)564 click_match = re.search(r"\[LAST_ACTION: CLICK\((\d+\.?\d*),(\d+\.?\d*)\)\]", action_str)565 if click_match:566 x, y = click_match.groups()567 px, py = to_px(x, y)568 r = 12569 # 青色底圈570 draw.ellipse([px-r, py-r, px+r, py+r], fill=(0, 255, 255, 150), outline=(255, 255, 255, 255), width=2)571 # 白色十字572 draw.line([px-r-5, py, px+r+5, py], fill=(255, 255, 255, 255), width=2)573 draw.line([px, py-r-5, px, py+r+5], fill=(255, 255, 255, 255), width=2)574 575 # 3. 匹配 DOUBLE_CLICK(x,y) -> 画个蓝色(Blue)双圈靶心576 dclick_match = re.search(r"\[LAST_ACTION: DOUBLE_CLICK\((\d+\.?\d*),(\d+\.?\d*)\)\]", action_str)577 if dclick_match:578 x, y = dclick_match.groups()579 px, py = to_px(x, y)580 r = 14581 # 蓝色底圈582 draw.ellipse([px-r, py-r, px+r, py+r], fill=(0, 100, 255, 150), outline=(255, 255, 255, 255), width=2)583 # 内层白圈584 draw.ellipse([px-(r-4), py-(r-4), px+(r-4), py+(r-4)], outline=(255, 255, 255, 255), width=1)585 586 # 4. 匹配 DRAG(x1,y1,x2,y2) -> 绿色轨迹线,绿色起点,黄色终点587 drag_match = re.search(r"\[LAST_ACTION: DRAG\((\d+\.?\d*),(\d+\.?\d*),(\d+\.?\d*),(\d+\.?\d*)\)\]", action_str)588 if drag_match:589 x1, y1, x2, y2 = drag_match.groups()590 p1 = to_px(x1, y1)591 p2 = to_px(x2, y2)592 593 # 绿色带透明度的连接线594 draw.line([p1, p2], fill=(0, 255, 0, 200), width=4)595 596 # 绿色起点圆597 draw.ellipse([p1[0]-6, p1[1]-6, p1[0]+6, p1[1]+6], fill=(0, 255, 0, 255), outline=(255,255,255,255), width=1)598 599 # 黄色终点靶心 (黄色在网格上也很显眼)600 r_end = 8601 draw.ellipse([p2[0]-r_end, p2[1]-r_end, p2[0]+r_end, p2[1]+r_end], fill=(255, 215, 0, 180), outline=(255,255,255,255), width=2)602 603 # 合并图层,转回 RGB (防 JPG 格式不支持 Alpha 通道)604 combined = Image.alpha_composite(image, overlay)605 return combined.convert("RGB")606 607def scale_to_fit(width: int, height: int, max_w: int = 1920, max_h: int = 1080) -> tuple[int, int]:608 """计算等比例缩放后的尺寸"""609 # 计算宽和高的缩放比例610 scale_w = max_w / width611 scale_h = max_h / height612 613 # 取较小的那个缩放比例,确保长宽都不超过限制614 scale = min(scale_w, scale_h, 1.0) # 如果原图比 1920x1080 小,则不放大(1.0)615 616 new_width = int(width * scale)617 new_height = int(height * scale)618 return new_width, new_height619 620def _get_target_message(message, role):621 """622 根据角色获取目标消息623 624 参数:625 message (list): 消息列表引用626 role (str): 要操作的角色,可选值: 'user', 'assistant', 'system'627 628 返回:629 dict: 目标消息字典630 """631 # 验证输入参数632 if not isinstance(message, list):633 raise TypeError("message必须是列表类型")634 635 if role not in ['user', 'assistant', 'system']:636 raise ValueError("role必须是'user'或'assistant'或'system'")637 638 target_message = None639 640 # 根据role决定要操作的对象641 if role == 'user':642 # 查找最后一个role为'user'的消息643 for msg in reversed(message):644 if isinstance(msg, dict) and msg['role'] == 'user':645 target_message = msg646 break647 elif role == 'assistant':648 # 检查最后一个消息649 if message and message[-1]['role'] == 'assistant':650 target_message = message[-1]651 else:652 # 如果最后一个消息不是assistant,创建一个新的653 new_assistant_msg = {'role': 'assistant', 'content': '','reasoning_content': ''}654 message.append(new_assistant_msg)655 target_message = new_assistant_msg656 elif role == 'system':657 # 查找第一个role为'system'的消息658 if message and message[0]['role'] == 'system':659 target_message = message[0]660 else:661 # 如果没有找到system消息,创建一个新的662 target_message = {'role': 'system', 'content': ''}663 message.insert(0, target_message)664 665 return target_message666 667def content_append(message, role, content):668 """669 将content添加到指定role消息的末尾670 """671 target_message = _get_target_message(message, role)672 if target_message:673 current_content = target_message.get('content', '')674 target_message['content'] = current_content + content675 676def content_prepend(message, role, content):677 """678 将content添加到指定role消息的前面679 """680 target_message = _get_target_message(message, role)681 if target_message:682 current_content = target_message.get('content', '')683 target_message['content'] = content + current_content684 685def content_replace(message, role, content):686 """687 用content替换指定role消息的内容688 """689 target_message = _get_target_message(message, role)690 if target_message:691 target_message['content'] = content692 693def content_new(message, role, content):694 """695 用content替换指定role消息的内容696 """697 message.append({'role': role, 'content': content})698 699configure_host_port(args.host, args.port)700 701def get_client_class(config, provider_id):702 if not config or 'modelProviders' not in config:703 return AsyncOpenAI704 vendor = 'OpenAI'705 for provider in config['modelProviders']:706 if provider['id'] == provider_id:707 vendor = provider['vendor']708 break709 # 假设你已经导入了 DifyOpenAIAsync 和 AsyncOpenAI710 if vendor == 'Dify':711 return DifyOpenAIAsync 712 elif vendor == 'customAnthropic':713 return AsyncClaudeAsOpenAI714 elif vendor == 'Gemini':715 return AsyncGeminiAsOpenAI716 else: 717 return AsyncOpenAI718 719from py.node_runner import node_mgr720@asynccontextmanager721async def lifespan(app: FastAPI):722 # --- [核心防御] 立即清理系统环境变量中的 SOCKS 代理,防止 httpx 崩溃 ---723 for env_key in ['HTTP_PROXY', 'HTTPS_PROXY', 'ALL_PROXY', 'http_proxy', 'https_proxy', 'all_proxy']:724 val = os.environ.get(env_key, "")725 if val.lower().startswith('socks'):726 # 彻底移除会导致崩溃的 socks 环境变量727 os.environ.pop(env_key, None)728 729 # 1. 准备所有独立的初始化任务730 from py.get_setting import init_db, init_covs_db, load_settings, save_settings731 from tzlocal import get_localzone732 733 asyncio.create_task(clean_temp_files_task())734 735 # 并行执行耗时操作736 init_db_task = init_db()737 init_covs_task = init_covs_db()738 load_locales_task = asyncio.to_thread(lambda: json.load(open(base_path + "/config/locales.json", "r", encoding="utf-8")))739 settings_task = load_settings() 740 timezone_task = asyncio.to_thread(get_localzone)741 copy_skills_task = _copy_default_skills()742 743 results = await asyncio.gather(744 init_db_task, 745 init_covs_task, 746 load_locales_task, 747 settings_task, 748 timezone_task,749 copy_skills_task750 )751 752 # 2. 解包结果753 global settings, client, reasoner_client, fast_client, mcp_client_list, local_timezone, logger, locales, global_http_client,scheduler_task,sleep_guard754 _, _, locales, settings, local_timezone, _ = results755 756 from py.sleep_guard import SleepGuard757 sleep_guard = SleepGuard(verbose=True)758 759 load_safety_words()760 761 if _is_steam_build:762 settings.setdefault("systemSettings", {})763 settings["systemSettings"]["contentSafety"] = True764 765 try:766 await asyncio.to_thread(sleep_guard.start)767 if sleep_guard.is_running():768 print("🛡️ 防休眠保护已启动,系统将不会自动休眠")769 else:770 print("⚠️ 防休眠启动失败,系统可能会在空闲时休眠")771 except Exception as e:772 print(f"防休眠启动异常: {e}")773 774 775 from py.scheduler import AgentScheduler776 # 传入全局 settings 对象的引用777 # 因为 python 字典是引用传递,后续 UI 修改了 settings,这里拿到的也是最新的778 scheduler = AgentScheduler(settings)779 scheduler_task = asyncio.create_task(scheduler.start_loop())780 781 # --- [日志系统初始化] ---782 783 timestamp = time.time()784 log_path = os.path.join(LOG_DIR, f"backend_{timestamp}.log")785 logger = logging.getLogger("app")786 787 if not logger.handlers:788 logger.setLevel(logging.INFO)789 790 # 1. 格式化器(控制台和文件共用一套格式)791 formatter = logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")792 793 # 2. 控制台输出(保留,方便实时看)794 console_handler = logging.StreamHandler()795 console_handler.setFormatter(formatter)796 logger.addHandler(console_handler)797 798 # 3. 【新增】文件输出(这才是真正存盘的)799 # 确保目录存在,否则会报错800 os.makedirs(LOG_DIR, exist_ok=True)801 file_handler = logging.FileHandler(log_path, encoding='utf-8')802 file_handler.setFormatter(formatter)803 logger.addHandler(file_handler)804 805 logger.info("===== 日志系统初始化成功 =====")806 logger.info(f"用户数据目录: {USER_DATA_DIR}")807 logger.info(f"设置数据库路径: {DATABASE_PATH}")808 logger.info(f"日志文件保存至: {log_path}") # 额外加一行,方便确认路径809 810 # --- [代理与 HTTP 客户端初始化] ---811 proxy_url = None812 trust_env = False813 814 if settings:815 sys_set = settings.get("systemSettings", {})816 mode = sys_set.get("proxyMode")817 manual_url = sys_set.get("proxy", "").strip()818 isChinaProxy = sys_set.get("isChinaProxy", False)819 820 if mode == "manual" and manual_url:821 # 手动模式:如果是 socks,由于没安装库,直接跳过并警告822 if manual_url.lower().startswith("socks"):823 logger.error("检测到手动设置了 SOCKS 代理,但当前环境不支持。代理已失效。")824 proxy_url = None825 else:826 proxy_url = manual_url827 elif mode == "system":828 # 系统模式:信任环境(此时环境里已经没有 socks 了,很安全)829 trust_env = True830 if isChinaProxy:831 # 2. 注入 Node.js / NPM 镜像源 (重点)832 # 设置这个环境变量后,所有的 npm install (包括你的 node_runner) 都会默认使用这个源833 os.environ["npm_config_registry"] = "https://registry.npmmirror.com/"834 835 # 3. 注入 UV / Pip 镜像源 (重点)836 # 这样后续如果调用 uv 或 pip,也会自动使用国内镜像837 os.environ["UV_INDEX_URL"] = "https://mirrors.aliyun.com/pypi/simple/"838 839 # 初始化全局带连接池的 HTTP 客户端840 timeout_config = httpx.Timeout(None, connect=10.0)841 global_http_client = httpx.AsyncClient(842 timeout=timeout_config,843 proxy=proxy_url,844 trust_env=trust_env845 )846 847 # --- [模型 Client 初始化] ---848 # 辅助函数:统一注入 global_http_client849 def create_model_client(provider_key, config_node=None):850 if not settings: return AsyncOpenAI(http_client=global_http_client)851 852 target_cfg = config_node if config_node else settings853 p_name = target_cfg.get('selectedProvider', settings.get('selectedProvider'))854 c_cls = get_client_class(settings, p_name)855 856 return c_cls(857 api_key=target_cfg.get('api_key') or settings.get('api_key', ''),858 base_url=target_cfg.get('base_url') or settings.get('base_url') or "https://api.openai.com/v1",859 http_client=global_http_client # 强制使用我们定义的带代理控制的客户端860 )861 862 if settings:863 client = create_model_client('main')864 reasoner_client = create_model_client('reasoner', settings.get('reasoner', {}))865 866 fast_cfg = settings.get('fast', {})867 if fast_cfg.get('enabled'):868 fast_client = create_model_client('fast', fast_cfg)869 else:870 fast_client = None871 else:872 client = AsyncOpenAI(http_client=global_http_client)873 reasoner_client = AsyncOpenAI(http_client=global_http_client)874 fast_client = AsyncOpenAI(http_client=global_http_client)875 876 # --- [其他初始化:ASR / MCP] ---877 try:878 from py.sherpa_asr import _get_recognizer879 asyncio.get_running_loop().run_in_executor(None, _get_recognizer)880 except Exception as e:881 logger.error(f"尝试启动sherpa失败: {e}")882 883 try:884 from py.moss_tts import _get_moss_runtime885 # 将重型的本地 TTS 加载也扔到后台线程池,如果没有下载模型它只会静默返回 None886 asyncio.get_running_loop().run_in_executor(None, _get_moss_runtime)887 except Exception as e:888 logger.error(f"尝试预热 MOSS TTS 失败: {e}")889 890 # MCP 初始化逻辑 (保持你原有的逻辑,但内部会复用 global_http_client)891 mcp_init_tasks = []892 893 async def init_mcp_with_timeout(server_name: str, server_config: dict, timeout=6.0, max_wait_failure=5.0):894 if server_config.get("disabled"):895 return server_name, None, "disabled"896 897 mcp_client = mcp_client_list.get(server_name) or McpClient()898 mcp_client_list[server_name] = mcp_client899 failure_event = asyncio.Event()900 first_error = None901 902 async def on_failure(msg: str):903 nonlocal first_error904 if first_error: return905 first_error = msg906 logger.error(f"MCP {server_name} failure: {msg}")907 settings.setdefault("mcpServers", {}).setdefault(server_name, {})["disabled"] = True908 mcp_client.disabled = True909 await mcp_client.close()910 failure_event.set()911 912 init_task = asyncio.create_task(mcp_client.initialize(server_name, server_config, on_failure_callback=on_failure))913 try:914 await asyncio.wait_for(init_task, timeout=timeout)915 try:916 await asyncio.wait_for(failure_event.wait(), timeout=max_wait_failure)917 except asyncio.TimeoutError:918 pass919 return server_name, (None if first_error else mcp_client), first_error920 except Exception as exc:921 return server_name, None, str(exc)922 finally:923 if not init_task.done(): init_task.cancel()924 925 async def check_results():926 for task in asyncio.as_completed(mcp_init_tasks):927 name, m_client, err = await task928 if err:929 settings['mcpServers'][name]['processingStatus'] = 'server_error'930 elif m_client:931 mcp_client_list[name] = m_client932 await save_settings(settings)933 await ws_manager.broadcast_settings_update(settings)934 935 if settings and settings.get('mcpServers'):936 mcp_init_tasks = [asyncio.create_task(init_mcp_with_timeout(k, v)) for k, v in settings['mcpServers'].items()]937 if mcp_init_tasks: asyncio.create_task(check_results())938 else:939 asyncio.create_task(ws_manager.broadcast_settings_update(settings or {}))940 941 # --- [自动启动 Telegram 机器人] ---942 try:943 tg_config_data = (settings or {}).get("telegramBotConfig", {}) or {}944 if not tg_config_data.get("bot_token"):945 env_token = os.environ.get("TELEGRAM_BOT_TOKEN", "").strip()946 if env_token:947 tg_config_data["bot_token"] = env_token948 if settings is not None:949 settings["telegramBotConfig"] = tg_config_data950 if tg_config_data and tg_config_data.get("bot_token"):951 from py.telegram_bot_manager import TelegramBotConfig952 tg_config = TelegramBotConfig(**tg_config_data)953 BotContainer.get_telegram().start_bot(tg_config)954 logger.info("🤖 自动启动 Telegram 机器人成功")955 except Exception as e:956 logger.warning(f"自动启动 Telegram 机器人跳过或失败: {e}")957 958 # --- [启动完成] ---959 print(f"REAL_PORT_FOUND:{PORT}", flush=True)960 yield961 962 # --- [关闭逻辑] ---963 print("System shutting down, cleaning up...")964 965 try:966 # 注意:此处需要根据您实际的文件结构导入 process_manager967 # 假设上述 ProcessManager 代码保存在 py/agent_tool.py 中968 from py.cli_tool import process_manager 969 970 print("正在清理工具管理的后台进程...")971 await process_manager.kill_all()972 except Exception as e:973 print(f"清理后台进程时发生异常: {e}")974 975 try:976 await asyncio.to_thread(sleep_guard.stop)977 print("🛡️ 防休眠保护已停止,系统将恢复正常休眠策略")978 except Exception as e:979 print(f"防休眠停止异常: {e}")980 981 if scheduler_task:982 scheduler_task.cancel()983 from py.node_runner import node_mgr984 ext_ids = list(node_mgr.exts.keys())985 for ext_id in ext_ids:986 try: await node_mgr.stop(ext_id)987 except: pass988 989 if global_http_client:990 await global_http_client.aclose()991 print("All processes terminated.")992 993 994app = FastAPI(lifespan=lifespan)995 996app.add_middleware(997 CORSMiddleware,998 allow_origins=["*"],999 allow_credentials=True,1000 allow_methods=["*"],1001 allow_headers=["*"],1002)1003 1004@app.middleware("http")1005async def cors_options_workaround(request: Request, call_next):1006 if request.method == "OPTIONS":1007 return Response(1008 status_code=200,1009 headers={1010 "Access-Control-Allow-Origin": "*",1011 "Access-Control-Allow-Methods": "*",1012 "Access-Control-Allow-Headers": "*",1013 "Access-Control-Max-Age": "86400", # 预检缓存 24 h1014 }1015 )1016 return await call_next(request)1017 1018@app.middleware("http")1019async def inject_steam_build_flag(request: Request, call_next):1020 response = await call_next(request)1021 if _is_steam_build and "text/html" in response.headers.get("content-type", ""):1022 body = b""1023 async for chunk in response.body_iterator:1024 body += chunk1025 body_str = body.decode("utf-8")1026 body_str = body_str.replace("</head>", '<script>window.__IS_STEAM_BUILD__=true;</script></head>')1027 headers = dict(response.headers)1028 headers.pop("content-length", None)1029 return HTMLResponse(content=body_str, status_code=response.status_code, headers=headers)1030 return response1031 1032async def t(text: str) -> str:1033 global locales1034 settings = await load_settings()1035 target_language = settings["currentLanguage"]1036 return locales[target_language].get(text, text)1037 1038 1039# 全局存储异步工具状态1040async_tools = {}1041async_tools_lock = asyncio.Lock()1042 1043async def execute_tool(tool_id: str, tool_name: str, args: dict, settings: dict,user_prompt: str):1044 try:1045 results = await dispatch_tool(tool_name, args, settings)1046 if isinstance(results, AsyncIterator):1047 buffer = []1048 async for chunk in results:1049 buffer.append(chunk)1050 results = "".join(buffer)1051 1052 if tool_name in ["query_knowledge_base"] and type(results) == list:1053 from py.know_base import rerank_knowledge_base1054 if settings["KBSettings"]["is_rerank"]:1055 results = await rerank_knowledge_base(user_prompt,results)1056 results = json.dumps(results, ensure_ascii=False, indent=4)1057 async with async_tools_lock:1058 async_tools[tool_id] = {1059 "status": "completed",1060 "result": results,1061 "name": tool_name,1062 "parameters": args,1063 }1064 except Exception as e:1065 async with async_tools_lock:1066 async_tools[tool_id] = {1067 "status": "error",1068 "result": str(e),1069 "name": tool_name,1070 "parameters": args,1071 }1072 1073async def get_image_content(image_url: str) -> str:1074 import hashlib1075 settings = await load_settings()1076 base64_image = await get_image_base64(image_url)1077 media_type = await get_image_media_type(image_url)1078 url= f"data:{media_type};base64,{base64_image}"1079 image_hash = hashlib.md5(image_url.encode()).hexdigest()1080 content = ""1081 if settings['vision']['enabled']:1082 # 如果uploaded_files/{item['image_url']['hash']}.txt存在,则读取文件内容,否则调用vision api1083 if os.path.exists(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt")):1084 with open(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt"), "r", encoding='utf-8') as f:1085 content += f"\n\n图片(URL:{image_url} 哈希值:{image_hash})信息如下:\n\n"+str(f.read())+"\n\n"1086 else:1087 images_content = [{"type": "text", "text": "请仔细描述图片中的内容,包含图片中可能存在的文字、数字、颜色、形状、大小、位置、人物、物体、场景等信息。"},{"type": "image_url", "image_url": {"url": url}}]1088 client = AsyncOpenAI(api_key=settings['vision']['api_key'],base_url=settings['vision']['base_url'])1089 1090 extra = {}1091 1092 if settings['vision']['temperature'] !=1:1093 extra['temperature'] = settings['vision']['temperature']1094 1095 response = await client.chat.completions.create(1096 model=settings['vision']['model'],1097 messages = [{"role": "user", "content": images_content}],1098 **extra1099 )1100 content = f"\n\nn图片(URL:{image_url} 哈希值:{image_hash})信息如下:\n\n"+str(response.choices[0].message.content)+"\n\n"1101 with open(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt"), "w", encoding='utf-8') as f:1102 f.write(str(response.choices[0].message.content))1103 else: 1104 # 如果uploaded_files/{item['image_url']['hash']}.txt存在,则读取文件内容,否则调用vision api1105 if os.path.exists(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt")):1106 with open(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt"), "r", encoding='utf-8') as f:1107 content += f"\n\nn图片(URL:{image_url} 哈希值:{image_hash})信息如下:\n\n"+str(f.read())+"\n\n"1108 else:1109 images_content = [{"type": "text", "text": "请仔细描述图片中的内容,包含图片中可能存在的文字、数字、颜色、形状、大小、位置、人物、物体、场景等信息。"},{"type": "image_url", "image_url": {"url": url}}]1110 client = AsyncOpenAI(api_key=settings['api_key'],base_url=settings['base_url'])1111 1112 extra = {}1113 1114 if settings['temperature'] !=1:1115 extra['temperature'] = settings['temperature']1116 1117 response = await client.chat.completions.create(1118 model=settings['model'],1119 messages = [{"role": "user", "content": images_content}],1120 **extra1121 )1122 content = f"\n\nn图片(URL:{image_url} 哈希值:{image_hash})信息如下:\n\n"+str(response.choices[0].message.content)+"\n\n"1123 with open(os.path.join(UPLOAD_FILES_DIR, f"{image_hash}.txt"), "w", encoding='utf-8') as f:1124 f.write(str(response.choices[0].message.content))1125 return content1126 1127# 存储等待中的MCP调用结果1128mcp_call_results: Dict[str, asyncio.Future] = {}1129 1130async def call_node_extension_tool(ext_id: str, tool_name: str, tool_params: dict) -> str:1131 """通过WebSocket调用Node扩展的工具"""1132 import uuid1133 1134 call_id = str(uuid.uuid4())1135 future = asyncio.Future()1136 mcp_call_results[call_id] = future1137 1138 # 广播给所有连接,找到对应的扩展1139 await ws_manager.broadcast({1140 "type": "call_mcp_tool",1141 "data": {1142 "ext_id": ext_id,1143 "tool_name": tool_name,1144 "tool_params": tool_params,1145 "call_id": call_id1146 }1147 })1148 1149 try:1150 # 等待结果,超时30秒1151 result = await asyncio.wait_for(future, timeout=30.0)1152 return str(result)1153 except asyncio.TimeoutError:1154 return f"调用扩展 {ext_id} 的工具 {tool_name} 超时"1155 finally:1156 if call_id in mcp_call_results:1157 del mcp_call_results[call_id]1158 1159async def dispatch_tool(tool_name: str, tool_params: dict, settings: dict,is_sub_agent:bool=False,force_allow: bool = False) -> str | List | AsyncIterator[str] | None :1160 global mcp_client_list,_TOOL_HOOKS,HA_client,ChromeMCP_client,sql_client, node_ext_mcp_clients, node_ext_mcp_tools1161 print("dispatch_tool",tool_name,tool_params)1162 1163 from py.utility_tools import time1164 # ==================== 1. 定义工具映射表 ====================1165 _TOOL_HOOKS = {1166 "searxng": searxng,1167 "Tavily_search": Tavily_search,1168 "query_knowledge_base": query_knowledge_base,1169 "jina_crawler": jina_crawler,1170 "Crawl4Ai_search": Crawl4Ai_search,1171 "firecrawl_search": firecrawl_search,1172 "simple_fetch":simple_fetch,1173 "markdown_new":markdown_new,1174 "agent_tool_call": agent_tool_call,1175 "a2a_tool_call": a2a_tool_call,1176 "custom_llm_tool": custom_llm_tool,1177 "get_file_content":get_file_content,1178 "get_image_content": get_image_content,1179 "e2b_code": e2b_code,1180 "local_run_code": local_run_code,1181 "openai_image": openai_image,1182 "openai_chat_image":openai_chat_image,1183 "Google_search": Google_search,1184 "Brave_search": Brave_search,1185 "Exa_search": Exa_search,1186 "Serper_search": Serper_search,1187 "bochaai_search": bochaai_search,1188 "comfyui_tool_call": comfyui_tool_call,1189 "time": time,1190 "get_weather": get_weather,1191 "get_location_coordinates": get_location_coordinates,1192 "get_weather_by_city":get_weather_by_city,1193 "get_wikipedia_summary_and_sections": get_wikipedia_summary_and_sections,1194 "get_wikipedia_section_content": get_wikipedia_section_content,1195 "search_arxiv_papers": search_arxiv_papers,1196 "auto_behavior": auto_behavior,1197 "list_pages": list_pages,1198 "new_page": new_page,1199 "close_page": close_page,1200 "select_page": select_page,