hatimqman commited on
Commit
ad55850
·
verified ·
1 Parent(s): a3e79be

bayan-proxy: app.py

Browse files
Files changed (1) hide show
  1. app.py +399 -0
app.py ADDED
@@ -0,0 +1,399 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ # -*- coding: utf-8 -*-
2
+ """الوسيط الرسميّ لـ«بيان» (FastAPI) — المصدر الحقيقيّ لِـHF Space `hatimqman-muaalem-proxy`.
3
+
4
+ يجلس بين عميل «بيان» المنشور على Cloudflare Pages وبين ثلاث نقاط HF Inference، فيؤدّي:
5
+ • حَمْل توكن HF سرًّا (من متغيّر بيئة `HF_TOKEN`، لا من العميل) وحقنه في نداء النقطة.
6
+ • تجاوز حجب Cloudflare 1010 (العميل ينادي Space لا خادم HF مباشرةً).
7
+ • فان-آوت لثلاث نقاط علويّة:
8
+ POST /stt → نقطة الطبقة الصوتية (hatimqman/quran-stt-endpoint، handler_full.py)
9
+ POST /analyze → نقطة المُعلِّم QPS (يمرّر mode/scope/sifat للطبقة الدقيقة)
10
+ WS /stream → بثّ احتياطيّ (نموذج البثّ على GPU، finetune/stream_server.py)
11
+ • قبول جسمٍ **مسطّح** {pcm, surah, ayahs, mode?, scope?, sifat?} وتغليفه {inputs:{…}} للنقطة،
12
+ فيبقى STT_URL / MUAALEM_URL في bayan/web/app.js بلا تغيير عنوان.
13
+ • CORS لنطاق Pages، حدّ معدّل بسيط لكلّ IP، إعادة محاولة على البدء البارد (scale-to-zero)،
14
+ وتدهورٌ آمن: أيّ فشلٍ للنقطة يُرجَع كخطأ HTTP نظيف (لا استثناء) فيتجاهله العميل بصمتٍ ولا يتعطّل.
15
+
16
+ كلّ الإعدادات من متغيّرات البيئة (تُضبَط أسرارًا على HF Space) — راجِع README.md.
17
+ النشر: صورة Docker على HF Space؛ لا شيء هنا يَنشُر تلقائيًّا (كودٌ فقط).
18
+ """
19
+ from __future__ import annotations
20
+
21
+ import asyncio
22
+ import json
23
+ import logging
24
+ import os
25
+ import time
26
+ from collections import deque
27
+ from contextlib import asynccontextmanager
28
+ from typing import Any, Deque, Dict, Optional
29
+
30
+ import httpx
31
+ from fastapi import FastAPI, Request, WebSocket
32
+ from fastapi.middleware.cors import CORSMiddleware
33
+ from fastapi.responses import JSONResponse, Response
34
+
35
+ # عميل WebSocket (للبثّ الاحتياطيّ فقط). عزل استيراده حتى لا يُسقِط غيابُه مسارَي /stt و/analyze.
36
+ try:
37
+ from websockets.asyncio.client import connect as ws_connect # websockets >= 13
38
+
39
+ _WS_OK = True
40
+ except Exception: # pragma: no cover - بيئةٌ بلا websockets
41
+ _WS_OK = False
42
+
43
+ log = logging.getLogger("muaalem-proxy")
44
+ logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO"))
45
+
46
+
47
+ # ───────────────────────── الإعدادات (من البيئة) ─────────────────────────
48
+ def _env(name: str, default: str = "") -> str:
49
+ return (os.environ.get(name) or default).strip()
50
+
51
+
52
+ HF_TOKEN = _env("HF_TOKEN") # سرّ — Bearer لنقاط HF. لا يُقرأ من العميل أبدًا.
53
+
54
+ # عناوين النقاط العلويّة (URL كامل لنقطة HF Inference). فارغٌ ⇒ المسار المقابل يُرجِع 503 نظيفًا.
55
+ STT_ENDPOINT_URL = _env("STT_ENDPOINT_URL").rstrip("/")
56
+ ANALYZE_ENDPOINT_URL = _env("ANALYZE_ENDPOINT_URL").rstrip("/")
57
+ STREAM_ENDPOINT_URL = _env("STREAM_ENDPOINT_URL") # ws:// أو wss:// (اختياريّ)
58
+
59
+ # إعادة المحاولة على البدء البارد (النقطة تستيقظ من scale-to-zero فتُرجِع 502/503/504 لحظيًّا).
60
+ UPSTREAM_TIMEOUT = float(_env("UPSTREAM_TIMEOUT", "150")) # سقف الميزانية الكلّيّة للطلب (ث)
61
+ ATTEMPT_TIMEOUT = float(_env("ATTEMPT_TIMEOUT", "60")) # مهلة قراءة المحاولة الواحدة (ث)
62
+ MAX_RETRIES = int(_env("MAX_RETRIES", "6")) # حتى 6 محاولات
63
+ RETRY_BACKOFF = float(_env("RETRY_BACKOFF", "4")) # تراجع خطّيّ أساس (ث)
64
+ RETRY_BACKOFF_MAX = float(_env("RETRY_BACKOFF_MAX", "20"))
65
+ COLD_START_STATUSES = {500, 502, 503, 504} # حالاتٌ عابرة → أعِد المحاولة
66
+
67
+ # حدّ المعدّل لكلّ IP (نافذة منزلقة/دقيقة). 0 = مُعطَّل.
68
+ RATE_LIMIT_PER_MIN = int(_env("RATE_LIMIT_PER_MIN", "40"))
69
+
70
+ # CORS: قائمة صريحة و/أو تعبيرٌ نمطيّ. الافتراضيّ يسمح بنطاقات Pages/Workers وlocalhost للتطوير.
71
+ ALLOWED_ORIGINS = [o.strip() for o in _env("ALLOWED_ORIGINS").split(",") if o.strip()]
72
+ ALLOWED_ORIGIN_REGEX = _env(
73
+ "ALLOWED_ORIGIN_REGEX",
74
+ r"https://([a-z0-9-]+\.)*pages\.dev|https://([a-z0-9-]+\.)*workers\.dev|"
75
+ r"http://(localhost|127\.0\.0\.1)(:\d+)?",
76
+ )
77
+
78
+ MAX_BODY_BYTES = int(_env("MAX_BODY_BYTES", str(24 * 1024 * 1024))) # ~24MB سقف الجسم
79
+
80
+ _client: Optional[httpx.AsyncClient] = None
81
+
82
+
83
+ # ───────────────────────── حدّ المعدّل (في الذاكرة) ─────────────────────────
84
+ class RateLimiter:
85
+ """نافذة منزلقة بسيطة لكلّ IP داخل حاوية Space واحدة. تُقلّم المداخل الخاملة تلقائيًّا."""
86
+
87
+ def __init__(self, per_min: int) -> None:
88
+ self.per_min = per_min
89
+ self._hits: Dict[str, Deque[float]] = {}
90
+ self._lock = asyncio.Lock()
91
+
92
+ async def allow(self, ip: str) -> bool:
93
+ if self.per_min <= 0:
94
+ return True
95
+ now = time.monotonic()
96
+ async with self._lock:
97
+ dq = self._hits.get(ip)
98
+ if dq is None:
99
+ dq = deque()
100
+ self._hits[ip] = dq
101
+ while dq and now - dq[0] > 60.0:
102
+ dq.popleft()
103
+ if len(dq) >= self.per_min:
104
+ return False
105
+ dq.append(now)
106
+ # تقليمٌ انتهازيّ لِـIPات خاملة كي لا تنمو الذاكرة بلا حدّ.
107
+ if len(self._hits) > 4096:
108
+ stale = [k for k, v in self._hits.items() if not v or now - v[-1] > 120.0]
109
+ for k in stale:
110
+ self._hits.pop(k, None)
111
+ return True
112
+
113
+
114
+ limiter = RateLimiter(RATE_LIMIT_PER_MIN)
115
+
116
+
117
+ def client_ip(request: Request) -> str:
118
+ xff = request.headers.get("x-forwarded-for")
119
+ if xff:
120
+ return xff.split(",")[0].strip()
121
+ return request.client.host if request.client else "unknown"
122
+
123
+
124
+ # ───────────────────────── دورة الحياة (عميل HTTP مشترك) ─────────────────────────
125
+ @asynccontextmanager
126
+ async def lifespan(_app: FastAPI):
127
+ global _client
128
+ _client = httpx.AsyncClient(
129
+ timeout=httpx.Timeout(connect=15.0, read=ATTEMPT_TIMEOUT, write=60.0, pool=15.0),
130
+ limits=httpx.Limits(max_connections=100, max_keepalive_connections=20),
131
+ )
132
+ log.info(
133
+ "muaalem-proxy up | stt=%s analyze=%s stream=%s token=%s",
134
+ bool(STT_ENDPOINT_URL), bool(ANALYZE_ENDPOINT_URL), bool(STREAM_ENDPOINT_URL), bool(HF_TOKEN),
135
+ )
136
+ try:
137
+ yield
138
+ finally:
139
+ if _client is not None:
140
+ await _client.aclose()
141
+
142
+
143
+ app = FastAPI(title="muaalem-proxy", version="1.0.0", lifespan=lifespan)
144
+ app.add_middleware(
145
+ CORSMiddleware,
146
+ allow_origins=ALLOWED_ORIGINS,
147
+ allow_origin_regex=ALLOWED_ORIGIN_REGEX or None,
148
+ allow_credentials=False,
149
+ allow_methods=["GET", "POST", "OPTIONS"],
150
+ allow_headers=["*"],
151
+ max_age=86400,
152
+ )
153
+
154
+
155
+ # ───────────────────────── تمرير النداء إلى النقطة (مع إعادة محاولة) ─────────────────────────
156
+ def _upstream_headers() -> Dict[str, str]:
157
+ h = {"Content-Type": "application/json", "Accept": "application/json"}
158
+ if HF_TOKEN:
159
+ h["Authorization"] = f"Bearer {HF_TOKEN}"
160
+ return h
161
+
162
+
163
+ async def forward_json(payload: Dict[str, Any], upstream: str, label: str) -> Response:
164
+ """يمرّر `payload` (مغلّفًا بـ{inputs:{…}} من المُنادي) إلى النقطة، ويعيد ردّها كما هو.
165
+
166
+ تدهورٌ آمن: كلّ حالة فشل تُرجَع كـJSONResponse بحالة HTTP مناسبة (لا استثناء يتسرّب)،
167
+ فيرى العميل `!res.ok` ويتراجع بصمت. يُعيد المحاولة على البدء البارد (502/503/504) ضمن ميزانية.
168
+ """
169
+ if not upstream:
170
+ return JSONResponse(status_code=503, content={"error": f"{label}_upstream_not_configured"})
171
+ if _client is None: # pragma: no cover - لا يحدث بعد الإقلاع
172
+ return JSONResponse(status_code=503, content={"error": "proxy_not_ready"})
173
+
174
+ headers = _upstream_headers()
175
+ deadline = time.monotonic() + UPSTREAM_TIMEOUT
176
+ last_status: Optional[int] = None
177
+ last_detail = ""
178
+
179
+ for attempt in range(MAX_RETRIES):
180
+ remaining = deadline - time.monotonic()
181
+ if remaining <= 0:
182
+ break
183
+ read_to = max(5.0, min(ATTEMPT_TIMEOUT, remaining))
184
+ try:
185
+ resp = await _client.post(
186
+ upstream, json=payload, headers=headers,
187
+ timeout=httpx.Timeout(connect=15.0, read=read_to, write=60.0, pool=15.0),
188
+ )
189
+ except (httpx.TimeoutException, httpx.TransportError) as exc:
190
+ last_status, last_detail = 504, type(exc).__name__
191
+ else:
192
+ if resp.status_code < 400:
193
+ # نجاح — مرّر جسم النقطة كما هو (JSON عادةً؛ وإلّا الخام بنوعه).
194
+ ctype = resp.headers.get("content-type", "")
195
+ if "json" in ctype.lower():
196
+ try:
197
+ return JSONResponse(status_code=200, content=resp.json())
198
+ except Exception:
199
+ pass
200
+ return Response(content=resp.content, status_code=200,
201
+ media_type=ctype or "application/json")
202
+ last_status, last_detail = resp.status_code, resp.text[:300]
203
+ if resp.status_code not in COLD_START_STATUSES:
204
+ # خطأٌ حتميّ من النقطة (مثلًا 400/422) — لا فائدة من إعادة المحاولة.
205
+ return JSONResponse(
206
+ status_code=resp.status_code,
207
+ content={"error": f"{label}_upstream", "status": resp.status_code, "detail": last_detail},
208
+ )
209
+
210
+ # مسارٌ قابل لإعادة المحاولة: تراجَع ما لم تنفد المحاولات/الميزانية.
211
+ # ملاحظة: تراجعٌ صفريّ (backoff=0) يعني «أعِد فورًا»، لا «توقّف» — الإيقاف فقط عند نفاد الوقت.
212
+ if attempt + 1 >= MAX_RETRIES:
213
+ break
214
+ remaining = deadline - time.monotonic()
215
+ if remaining <= 0:
216
+ break
217
+ backoff = min(RETRY_BACKOFF * (attempt + 1), RETRY_BACKOFF_MAX, remaining)
218
+ log.info("%s retry %d/%d after %.0fs (last=%s)", label, attempt + 1, MAX_RETRIES, backoff, last_status)
219
+ if backoff > 0:
220
+ await asyncio.sleep(backoff)
221
+
222
+ return JSONResponse(
223
+ status_code=503,
224
+ content={"error": f"{label}_unavailable", "upstream_status": last_status, "detail": last_detail},
225
+ )
226
+
227
+
228
+ class _PayloadTooLarge(Exception):
229
+ """جسمٌ يتجاوز MAX_BODY_BYTES — يُميَّز عن خطأ فكّ JSON (وهو ValueError أيضًا)."""
230
+
231
+
232
+ async def _read_json(request: Request) -> Dict[str, Any]:
233
+ body = await request.body()
234
+ if len(body) > MAX_BODY_BYTES:
235
+ raise _PayloadTooLarge()
236
+ if not body:
237
+ return {}
238
+ data = json.loads(body) # JSONDecodeError (⊂ ValueError) عند الفساد ⇒ يعالَج كـ bad_json
239
+ if not isinstance(data, dict):
240
+ raise TypeError("expected_json_object")
241
+ return data
242
+
243
+
244
+ async def _guard(request: Request) -> Optional[JSONResponse]:
245
+ """حدّ المعدّل + قراءة/تحقّق الجسم المشتركان. يعيد استجابة خطأ أو None (مرّ)."""
246
+ if not await limiter.allow(client_ip(request)):
247
+ return JSONResponse(status_code=429, content={"error": "rate_limited"})
248
+ return None
249
+
250
+
251
+ # ───────────────────────── المسارات ─────────────────────────
252
+ @app.get("/")
253
+ @app.get("/health")
254
+ async def health() -> Dict[str, Any]:
255
+ return {
256
+ "ok": True,
257
+ "service": "muaalem-proxy",
258
+ "routes": ["/stt", "/analyze", "/stream"],
259
+ "upstream": {
260
+ "stt": bool(STT_ENDPOINT_URL),
261
+ "analyze": bool(ANALYZE_ENDPOINT_URL),
262
+ "stream": bool(STREAM_ENDPOINT_URL) and _WS_OK,
263
+ },
264
+ "token": bool(HF_TOKEN),
265
+ }
266
+
267
+
268
+ @app.post("/stt")
269
+ async def stt(request: Request) -> Response:
270
+ """الطبقة الصوتية. العميل يرسل {pcm} أو {pcm, ref}؛ نغلّفها {inputs:{pcm, ref?}} للنقطة."""
271
+ blocked = await _guard(request)
272
+ if blocked is not None:
273
+ return blocked
274
+ try:
275
+ data = await _read_json(request)
276
+ except _PayloadTooLarge:
277
+ return JSONResponse(status_code=413, content={"error": "payload_too_large"})
278
+ except Exception:
279
+ return JSONResponse(status_code=400, content={"error": "bad_json"})
280
+
281
+ pcm = data.get("pcm")
282
+ if not pcm or not isinstance(pcm, str):
283
+ return JSONResponse(status_code=400, content={"error": "missing_pcm"})
284
+
285
+ inner: Dict[str, Any] = {"pcm": pcm}
286
+ if data.get("ref"):
287
+ inner["ref"] = data["ref"]
288
+ return await forward_json({"inputs": inner}, STT_ENDPOINT_URL, "stt")
289
+
290
+
291
+ @app.post("/analyze")
292
+ async def analyze(request: Request) -> Response:
293
+ """المُعلِّم QPS. مسطّح {pcm, surah, ayahs|ayah, mode?, scope?, sifat?} → {inputs:{…}} للنقطة."""
294
+ blocked = await _guard(request)
295
+ if blocked is not None:
296
+ return blocked
297
+ try:
298
+ data = await _read_json(request)
299
+ except _PayloadTooLarge:
300
+ return JSONResponse(status_code=413, content={"error": "payload_too_large"})
301
+ except Exception:
302
+ return JSONResponse(status_code=400, content={"error": "bad_json"})
303
+
304
+ pcm = data.get("pcm")
305
+ if not pcm or not isinstance(pcm, str):
306
+ return JSONResponse(status_code=400, content={"error": "missing_pcm"})
307
+ if "surah" not in data or data.get("surah") is None:
308
+ return JSONResponse(status_code=400, content={"error": "missing_surah"})
309
+
310
+ inner: Dict[str, Any] = {"pcm": pcm, "surah": data["surah"]}
311
+ if data.get("ayahs") is not None:
312
+ inner["ayahs"] = data["ayahs"]
313
+ if data.get("ayah") is not None:
314
+ inner["ayah"] = data["ayah"]
315
+ if "ayahs" not in inner and "ayah" not in inner:
316
+ return JSONResponse(status_code=400, content={"error": "missing_ayah_or_ayahs"})
317
+
318
+ # معامِلات الطبقة الدقيقة (P1/P2) — تُمرَّر شفّافةً؛ النقطة الحاليّة تتجاهلها بأمان.
319
+ for k in ("mode", "scope", "sifat"):
320
+ if k in data and data[k] is not None:
321
+ inner[k] = data[k]
322
+
323
+ return await forward_json({"inputs": inner}, ANALYZE_ENDPOINT_URL, "analyze")
324
+
325
+
326
+ @app.websocket("/stream")
327
+ async def stream(ws: WebSocket) -> None:
328
+ """بثّ احتياطيّ: يمرّر إطارات WS ثنائيّة الاتّجاه بين العميل والنقطة العلويّة بشفافية.
329
+
330
+ البروتوكول (كما في finetune/stream_server.py): العميل يرسل إطارات int16 PCM 16k + نصّ "end"؛
331
+ النقطة تُرجِع {t, text, recent}. نمرّر البايتات والنصّ كما هي، ونحقن التوكن + معامل kind.
332
+ """
333
+ await ws.accept()
334
+ ip = ws.headers.get("x-forwarded-for", "").split(",")[0].strip() or (
335
+ ws.client.host if ws.client else "unknown"
336
+ )
337
+ if not await limiter.allow(ip):
338
+ await ws.close(code=4029, reason="rate_limited")
339
+ return
340
+ if not (_WS_OK and STREAM_ENDPOINT_URL):
341
+ await ws.close(code=1011, reason="stream_upstream_not_configured")
342
+ return
343
+
344
+ kind = ws.query_params.get("kind", "tlog")
345
+ sep = "&" if "?" in STREAM_ENDPOINT_URL else "?"
346
+ up_url = f"{STREAM_ENDPOINT_URL}{sep}kind={kind}"
347
+ headers = {"Authorization": f"Bearer {HF_TOKEN}"} if HF_TOKEN else {}
348
+
349
+ try:
350
+ async with ws_connect(
351
+ up_url, additional_headers=headers, max_size=None, open_timeout=UPSTREAM_TIMEOUT
352
+ ) as upstream:
353
+
354
+ async def client_to_upstream() -> None:
355
+ try:
356
+ while True:
357
+ msg = await ws.receive()
358
+ if msg.get("type") == "websocket.disconnect":
359
+ break
360
+ if msg.get("bytes") is not None:
361
+ await upstream.send(msg["bytes"])
362
+ elif msg.get("text") is not None:
363
+ await upstream.send(msg["text"])
364
+ except Exception:
365
+ pass
366
+ finally:
367
+ try:
368
+ await upstream.close()
369
+ except Exception:
370
+ pass
371
+
372
+ async def upstream_to_client() -> None:
373
+ try:
374
+ async for frame in upstream:
375
+ if isinstance(frame, (bytes, bytearray)):
376
+ await ws.send_bytes(bytes(frame))
377
+ else:
378
+ await ws.send_text(frame)
379
+ except Exception:
380
+ pass
381
+ finally:
382
+ try:
383
+ await ws.close()
384
+ except Exception:
385
+ pass
386
+
387
+ await asyncio.gather(client_to_upstream(), upstream_to_client())
388
+ except Exception as exc: # فشل الاتصال بالنقطة العلويّة → إغلاقٌ نظيف (تراجع صامت في العميل).
389
+ log.info("stream upstream error: %s", type(exc).__name__)
390
+ try:
391
+ await ws.close(code=1011, reason="stream_upstream_error")
392
+ except Exception:
393
+ pass
394
+
395
+
396
+ if __name__ == "__main__":
397
+ import uvicorn
398
+
399
+ uvicorn.run(app, host="0.0.0.0", port=int(_env("PORT", "7860")), log_level="info")