CoolFace
Apppublic

Pamudu13/gemma-3-chat

sourceHugging Faceupdated 25d agoView on Hugging Face
0likes
server.py12002 linesDownload Raw Back to root
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,

Showing the first 1,200 of 12002 lines. Download the file for the rest.