1+ #!/usr/bin/env python3
2+ """
3+ ONNX Client — Discover-or-Launch клиент для ONNX Singleton Server.
4+ Использует Windows Named Mutex для предотвращения гонки при запуске сервера.
5+ """
6+
7+ import os
8+ import sys
9+ import json
10+ import time
11+ import socket
12+ import subprocess
13+ import urllib .request
14+ import threading
15+ from pathlib import Path
16+ from typing import List , Optional
17+
18+
19+ class OnnxEmbedderClient :
20+ """
21+ Клиент для ONNX Singleton Server.
22+
23+ Паттерн:
24+ 1. Проверяет, запущен ли сервер на порту (health check)
25+ 2. Если нет — пытается захватить Named Mutex
26+ 3. Если мутекс захвачен — запускает subprocess сервер
27+ 4. Ждёт готовности сервера
28+ 5. Отдаёт мутекс
29+ 6. Если мутекс занят — ждёт пока другой процесс запустит сервер
30+ """
31+
32+ def __init__ (self , port : int = 9876 , model_name : str = "bge-m3" ):
33+ self .port = port
34+ self .model_name = model_name
35+ self .base_url = f"http://127.0.0.1:{ port } "
36+ self ._mutex_name = f"Global\\ MSCodeBase_OnnxServer_{ model_name } "
37+ self ._mutex_handle = None
38+ self ._server_started_by_us = False
39+
40+ def _is_server_running (self ) -> bool :
41+ """Быстрая проверка: отвечает ли сервер на /health."""
42+ try :
43+ with socket .socket (socket .AF_INET , socket .SOCK_STREAM ) as s :
44+ s .settimeout (0.5 )
45+ return s .connect_ex (("127.0.0.1" , self .port )) == 0
46+ except Exception :
47+ return False
48+
49+ def _health_check (self ) -> bool :
50+ """Полный health check через HTTP."""
51+ try :
52+ req = urllib .request .Request (f"{ self .base_url } /health" , method = "GET" )
53+ with urllib .request .urlopen (req , timeout = 2 ) as resp :
54+ data = json .loads (resp .read ().decode ())
55+ return data .get ("status" ) == "ok"
56+ except Exception :
57+ return False
58+
59+ def _acquire_launch_mutex (self ) -> bool :
60+ """
61+ Пытается захватить Named Mutex для запуска сервера.
62+ Возвращает True если мы стали владельцем и должны запускать сервер.
63+ Возвращает False если другой процесс уже запускает/запустил сервер.
64+ """
65+ if sys .platform != 'win32' :
66+ # На Unix используем файловый лок (упрощённо)
67+ return True
68+
69+ try :
70+ import ctypes
71+ kernel32 = ctypes .windll .kernel32
72+
73+ # CreateMutex с bInitialOwner=FALSE — не захватываем сразу
74+ self ._mutex_handle = kernel32 .CreateMutexW (None , False , self ._mutex_name )
75+ last_error = kernel32 .GetLastError ()
76+
77+ # ERROR_ALREADY_EXISTS = 183 — мутекс уже есть
78+ if last_error == 183 :
79+ # Мутекс существует — другой процесс уже запускает сервер
80+ kernel32 .CloseHandle (self ._mutex_handle )
81+ self ._mutex_handle = None
82+ return False
83+
84+ # Мутекс создан нами, захватываем его
85+ result = kernel32 .WaitForSingleObject (self ._mutex_handle , 10000 ) # 10 сек таймаут
86+ if result not in (0 , 128 ): # WAIT_OBJECT_0 = 0, WAIT_ABANDONED = 128
87+ kernel32 .CloseHandle (self ._mutex_handle )
88+ self ._mutex_handle = None
89+ return False
90+
91+ return True
92+
93+ except Exception as e :
94+ print (f"[ONNX Client] Mutex error: { e } " , file = sys .stderr )
95+ return True # На всякий случай пробуем запустить
96+
97+ def _release_launch_mutex (self ):
98+ """Освобождает мутекс запуска."""
99+ if self ._mutex_handle and sys .platform == 'win32' :
100+ try :
101+ import ctypes
102+ ctypes .windll .kernel32 .ReleaseMutex (self ._mutex_handle )
103+ ctypes .windll .kernel32 .CloseHandle (self ._mutex_handle )
104+ except Exception :
105+ pass
106+ self ._mutex_handle = None
107+
108+ def _launch_server (self ) -> bool :
109+ """Запускает onnx_server.py как detached процесс."""
110+ server_script = PROJECT_ROOT / "onnx_server.py"
111+ if not server_script .exists ():
112+ print (f"[ONNX Client] Server script not found: { server_script } " , file = sys .stderr )
113+ return False
114+
115+ # Флаги для Windows: DETACHED_PROCESS | CREATE_NO_WINDOW
116+ creation_flags = 0
117+ if sys .platform == 'win32' :
118+ creation_flags = 0x00000008 | 0x08000000 # DETACHED_PROCESS | CREATE_NO_WINDOW
119+
120+ env = os .environ .copy ()
121+ env ["ONNX_PORT" ] = str (self .port )
122+ env ["ONNX_IDLE_TIMEOUT" ] = "600"
123+ env ["PYTHONPATH" ] = str (PROJECT_ROOT ) + os .pathsep + env .get ("PYTHONPATH" , "" )
124+
125+ try :
126+ proc = subprocess .Popen (
127+ [sys .executable , str (server_script ), f"--port={ self .port } " ],
128+ stdout = subprocess .DEVNULL ,
129+ stderr = subprocess .DEVNULL ,
130+ stdin = subprocess .DEVNULL ,
131+ env = env ,
132+ creationflags = creation_flags ,
133+ cwd = str (PROJECT_ROOT )
134+ )
135+ self ._server_pid = proc .pid
136+ self ._server_started_by_us = True
137+ print (f"[ONNX Client] Launched server PID={ proc .pid } " , file = sys .stderr )
138+ return True
139+ except Exception as e :
140+ print (f"[ONNX Client] Failed to launch server: { e } " , file = sys .stderr )
141+ return False
142+
143+ def _wait_for_server (self , timeout : float = 30.0 ) -> bool :
144+ """Ждёт готовности сервера."""
145+ start = time .time ()
146+ while time .time () - start < timeout :
147+ if self ._health_check ():
148+ return True
149+ time .sleep (0.5 )
150+ return False
151+
152+ def ensure_server_running (self ) -> bool :
153+ """
154+ Гарантирует, что сервер запущен.
155+ Thread-safe: использует модульный lock для координации между потоками.
156+ """
157+ # Быстрая проверка без лока
158+ if self ._health_check ():
159+ return True
160+
161+ # Межпоточный lock (на случай если несколько потоков одновременно вызывают)
162+ with _client_lock :
163+ # Double-check под локом
164+ if self ._health_check ():
165+ return True
166+
167+ # Пытаемся захватить межпроцессный мутекс
168+ if not self ._acquire_launch_mutex ():
169+ # Другой процесс запускает — ждём
170+ print (f"[ONNX Client] Waiting for another process to start server..." , file = sys .stderr )
171+ return self ._wait_for_server ()
172+
173+ # Мы владелец мутекса — запускаем сервер
174+ try :
175+ if not self ._launch_server ():
176+ return False
177+ return self ._wait_for_server ()
178+ finally :
179+ self ._release_launch_mutex ()
180+
181+ def embed (self , text : str ) -> List [float ]:
182+ """Эмбеддинг одного текста."""
183+ if not self .ensure_server_running ():
184+ raise RuntimeError ("ONNX server not available" )
185+
186+ data = json .dumps ({"text" : text }).encode ("utf-8" )
187+ req = urllib .request .Request (
188+ f"{ self .base_url } /embed" ,
189+ data = data ,
190+ headers = {"Content-Type" : "application/json" },
191+ method = "POST"
192+ )
193+ with urllib .request .urlopen (req , timeout = 30 ) as resp :
194+ result = json .loads (resp .read ().decode ())
195+ if "error" in result :
196+ raise RuntimeError (result ["error" ])
197+ return result ["vector" ]
198+
199+ def embed_batch (self , texts : List [str ]) -> List [List [float ]]:
200+ """Эмбеддинг батча текстов."""
201+ if not self .ensure_server_running ():
202+ raise RuntimeError ("ONNX server not available" )
203+
204+ data = json .dumps ({"texts" : texts }).encode ("utf-8" )
205+ req = urllib .request .Request (
206+ f"{ self .base_url } /embed_batch" ,
207+ data = data ,
208+ headers = {"Content-Type" : "application/json" },
209+ method = "POST"
210+ )
211+ with urllib .request .urlopen (req , timeout = 60 ) as resp :
212+ result = json .loads (resp .read ().decode ())
213+ if "error" in result :
214+ raise RuntimeError (result ["error" ])
215+ return result ["vectors" ]
216+
217+ def shutdown (self ):
218+ """Явное завершение (для тестов)."""
219+ if self ._server_started_by_us and hasattr (self , '_server_pid' ):
220+ try :
221+ if sys .platform == 'win32' :
222+ subprocess .run (['taskkill' , '/F' , '/PID' , str (self ._server_pid )],
223+ capture_output = True )
224+ else :
225+ os .kill (self ._server_pid , 15 )
226+ except Exception :
227+ pass
228+
229+
230+ # Глобальный lock для координации внутри процесса
231+ _client_lock = threading .Lock ()
232+
233+ # Путь к корню проекта
234+ PROJECT_ROOT = Path (__file__ ).resolve ().parent
235+
236+ # ─── Singleton accessor ────────────────────────────────────
237+
238+ _client_instance : Optional [OnnxEmbedderClient ] = None
239+ _client_lock_module = threading .Lock ()
240+
241+ def get_onnx_client (port : int = 9876 , model_name : str = "bge-m3" ) -> OnnxEmbedderClient :
242+ """Возвращает singleton клиент ONNX."""
243+ global _client_instance
244+ with _client_lock_module :
245+ if _client_instance is None :
246+ _client_instance = OnnxEmbedderClient (port = port , model_name = model_name )
247+ return _client_instance
248+
249+
250+ # ─── Удобные функции ──────────────────────────────────────
251+
252+ def embed_text (text : str , port : int = 9876 ) -> List [float ]:
253+ """Быстрый вызов: эмбеддинг одного текста."""
254+ return get_onnx_client (port ).embed (text )
255+
256+ def embed_batch (texts : List [str ], port : int = 9876 ) -> List [List [float ]]:
257+ """Быстрый вызов: эмбеддинг батча."""
258+ return get_onnx_client (port ).embed_batch (texts )
259+
260+
261+ if __name__ == "__main__" :
262+ # CLI для тестов
263+ import sys
264+ client = get_onnx_client ()
265+ if len (sys .argv ) > 1 :
266+ text = " " .join (sys .argv [1 :])
267+ vec = client .embed (text )
268+ print (f"Dim: { len (vec )} , first 5: { vec [:5 ]} " )
269+ else :
270+ print ("Usage: python onnx_client.py <text>" )
0 commit comments