Repository navigation
Expand file tree
/
Copy pathproxy.py
More file actions
254 lines (212 loc) · 11.7 KB
/
Copy pathproxy.py
File metadata and controls
254 lines (212 loc) · 11.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
from flask import request, Response, jsonify, stream_with_context
import requests
from werkzeug.exceptions import ClientDisconnected
from threading import Lock
from collections import defaultdict
from protocols import PROTOCOLS
class ProxyServer:
def __init__(self, zookeeper):
self.zookeeper = zookeeper
self.app = zookeeper.app
self.active_connections = defaultdict(int) # tracks connections by full URL
self.connections_lock = Lock()
self.setup_routes()
def setup_routes(self):
# Health and Identity
self.app.route('/health', methods=['GET'])(self.health_check)
self.app.route('/.well-known/serviceinfo', methods=['GET'])(self.service_info)
# OpenAI API routes
self.app.route('/v1/models', methods=['GET'])(self.get_models)
# OpenAI API routes
self.app.route('/v1/completions', methods=['POST'])(self.handle_completions)
self.app.route('/v1/chat/completions', methods=['POST'])(self.handle_chat_completions)
# A1111 API routes
self.app.route('/sdapi/v1/sd-models', methods=['GET'])(self.get_sd_models)
self.app.route('/sdapi/v1/txt2img', methods=['POST'])(self.handle_txt2img)
self.app.route('/sdapi/v1/img2img', methods=['POST'])(self.handle_img2img)
def get_models(self):
unique_models = {}
# Get all available models and filter for those with text capabilities
available_models = self.zookeeper.get_available_models()
for model in available_models:
protocol = model['listener']['protocol']
if (protocol in PROTOCOLS and
(PROTOCOLS[protocol]['completions'] or PROTOCOLS[protocol]['chat_completions'])):
model_name = model['model_name']
if model_name not in unique_models:
unique_models[model_name] = {
"id": model_name,
"owned_by": model['source'],
"running": model.get('status', {}).get('running', False),
"ready": model.get('status', {}).get('ready', False),
"metrics": []
}
else:
# Update running/ready status (true if any instance is running/ready)
existing = unique_models[model_name]
existing["running"] = existing["running"] or model.get('status', {}).get('running', False)
existing["ready"] = existing["ready"] or model.get('status', {}).get('ready', False)
# Collect metrics
if 'metrics' in model and model['metrics']:
unique_models[model_name]["metrics"].append(model['metrics'])
# Aggregate metrics
for model_data in unique_models.values():
if model_data["metrics"]:
metrics_list = model_data["metrics"]
model_data["num_running"] = sum(m.num_running for m in metrics_list)
model_data["num_waiting"] = sum(m.num_waiting for m in metrics_list)
model_data["kv_utilization"] = sum(m.kv_utilization for m in metrics_list) / len(metrics_list)
del model_data["metrics"] # Remove temporary metrics list
return jsonify({"data": list(unique_models.values())})
def get_sd_models(self):
"""Return list of available SD models in A1111 format"""
image_models = []
# Get all available models and filter for those with image capabilities
available_models = self.zookeeper.get_available_models()
for model in available_models:
protocol = model['listener']['protocol']
if protocol in PROTOCOLS and (PROTOCOLS[protocol]['txt2img'] or PROTOCOLS[protocol]['img2img']):
image_models.append({
"title": model['model_name'],
"model_name": model['model_name'],
"hash": "0000000000", # Placeholder
"sha256": "0" * 64, # Placeholder
"filename": model.get('path', ""),
"config": None
})
return jsonify(image_models)
def handle_completions(self):
"""Handle /v1/completions endpoint"""
return self.handle_request('completions')
def handle_chat_completions(self):
"""Handle /v1/chat/completions endpoint"""
return self.handle_request('chat_completions')
def handle_txt2img(self):
"""Handle /sdapi/v1/txt2img endpoint"""
return self.handle_request('txt2img', required_keys=['prompt'])
def handle_img2img(self):
"""Handle /sdapi/v1/img2img endpoint"""
return self.handle_request('img2img', required_keys=['prompt'])
def handle_request(self, endpoint_type, required_keys=None):
"""Unified request handler for all endpoints.
Args:
endpoint_type: The type of endpoint to handle (completions, chat_completions, etc)
required_keys: List of required keys in the request payload
"""
if not request.is_json:
return jsonify({"error": "Request must be JSON"}), 400
data = request.get_json()
if required_keys:
missing_keys = [key for key in required_keys if key not in data]
if missing_keys:
return jsonify({"error": f"Missing required keys: {', '.join(missing_keys)}"}), 400
return self._handle_request(endpoint_type, data)
def health_check(self):
# Check local running models
if len(self.zookeeper.get_available_models(local_models=True, remote_models=False)) > 0:
return '', PROTOCOLS['openai']['health_status']
def service_info(self):
"""Return service information according to the serviceinfo spec."""
return jsonify({
"version": "0.2",
"software": {
"name": "ModelZoo",
"version": "0.0.1",
"repository": "https://github.com/the-crypt-keeper/modelzoo",
"homepage": "https://github.com/the-crypt-keeper/modelzoo"
},
"api": {
"openai": {
"name": "OpenAI API",
"rel_url": "/v1",
"documentation": "https://openai.com/documentation/api",
"version": "0.0.1"
}
}
})
def _handle_request(self, endpoint_type, data):
try:
requested_model = data.get('model')
if not requested_model: return jsonify({"error": "Model not specified in the request"}), 400
# Get all available instances of the requested model
model_instances = []
# Get all available model instances, build potential endpoint map
available_models = self.zookeeper.get_available_models()
for model in available_models:
# Match on display name or model name
if model['model_name'] == requested_model:
protocol = model['listener']['protocol']
endpoint = PROTOCOLS.get(protocol, {})[endpoint_type]
if not endpoint: # Protocol doesn't support this endpoint
continue
host = model['listener']['host']
port = model['listener']['port']
model_instances.append({
'model_name': model['model_name'],
'model_id': model['model_id'],
'protocol': protocol,
'url': f"http://{host}:{port}{endpoint}"
})
if not model_instances:
return jsonify({"error": f"Model {requested_model} not found or not running"}), 404
# Select instance with least connections
with self.connections_lock:
selected = min(model_instances,
key=lambda x: self.active_connections[x['url']])
self.active_connections[selected['url']] += 1
target_url = selected['url']
protocol = selected['protocol']
headers = {k: v for k, v in request.headers.items() if k.lower() != 'host'}
try:
# Apply sampler mapping if this is an image endpoint
if endpoint_type in ['txt2img', 'img2img'] and 'sampler_name' in data:
sampler_map = PROTOCOLS[protocol].get('image_sampler_map', {})
if data['sampler_name'] in sampler_map:
data['sampler_name'] = sampler_map[data['sampler_name']]
# Replace model name with model ID in request
data = data.copy()
data['model'] = selected['model_id']
# Apply request adapter if specified in protocol
request_adapter = PROTOCOLS[protocol].get(f'{endpoint_type}_request_adapter')
if request_adapter is not None:
data = request_adapter(data, target_url)
# Determine if we should stream based on request
should_stream = data.get('stream', False)
resp = requests.post(target_url, json=data, headers=headers, stream=should_stream)
content_type = resp.headers.get('Content-Type', 'application/json')
if should_stream:
@stream_with_context
def generate():
try:
for chunk in resp.iter_content(chunk_size=4096):
if chunk:
yield chunk
except GeneratorExit:
print("Client disconnected. Stopping stream.")
finally:
resp.close()
return Response(generate(),
direct_passthrough=True,
status=resp.status_code,
content_type=content_type)
else:
# For non-streaming requests, get full response
response_data = resp.json() if content_type == 'application/json' else resp.content
# Apply response adapter if specified in protocol
if content_type == 'application/json':
response_adapter = PROTOCOLS[protocol].get(f'{endpoint_type}_response_adapter')
if response_adapter is not None:
response_data = response_adapter(response_data, target_url)
return Response(response_data if isinstance(response_data, bytes) else jsonify(response_data).data,
status=resp.status_code,
content_type=content_type)
finally:
# Always decrement connection count
with self.connections_lock:
self.active_connections[target_url] -= 1
except requests.RequestException as e:
print(f"Network error occurred: {str(e)}")
return jsonify({"error": f"Network error: {str(e)}"}), 500
except Exception as e:
print(f"Unexpected error occurred: {str(e)}")
return jsonify({"error": f"Unexpected error: {str(e)}"}), 500