Spaces:
Paused
Paused
| # -*- coding: utf-8 -*- | |
| """الوسيط الرسميّ لـ«بيان» (FastAPI) — المصدر الحقيقيّ لِـHF Space `hatimqman-muaalem-proxy`. | |
| يجلس بين عميل «بيان» المنشور على Cloudflare Pages وبين ثلاث نقاط HF Inference، فيؤدّي: | |
| • حَمْل توكن HF سرًّا (من متغيّر بيئة `HF_TOKEN`، لا من العميل) وحقنه في نداء النقطة. | |
| • تجاوز حجب Cloudflare 1010 (العميل ينادي Space لا خادم HF مباشرةً). | |
| • فان-آوت لثلاث نقاط علويّة: | |
| POST /stt → نقطة الطبقة الصوتية (hatimqman/quran-stt-endpoint، handler_full.py) | |
| POST /analyze → نقطة المُعلِّم QPS (يمرّر mode/scope/sifat للطبقة الدقيقة) | |
| WS /stream → بثّ احتياطيّ (نموذج البثّ على GPU، finetune/stream_server.py) | |
| • قبول جسمٍ **مسطّح** {pcm, surah, ayahs, mode?, scope?, sifat?} وتغليفه {inputs:{…}} للنقطة، | |
| فيبقى STT_URL / MUAALEM_URL في bayan/web/app.js بلا تغيير عنوان. | |
| • CORS لنطاق Pages، حدّ معدّل بسيط لكلّ IP، إعادة محاولة على البدء البارد (scale-to-zero)، | |
| وتدهورٌ آمن: أيّ فشلٍ للنقطة يُرجَع كخطأ HTTP نظيف (لا استثناء) فيتجاهله العميل بصمتٍ ولا يتعطّل. | |
| كلّ الإعدادات من متغيّرات البيئة (تُضبَط أسرارًا على HF Space) — راجِع README.md. | |
| النشر: صورة Docker على HF Space؛ لا شيء هنا يَنشُر تلقائيًّا (كودٌ فقط). | |
| """ | |
| from __future__ import annotations | |
| import asyncio | |
| import json | |
| import logging | |
| import os | |
| import time | |
| from collections import deque | |
| from contextlib import asynccontextmanager | |
| from typing import Any, Deque, Dict, Optional | |
| import httpx | |
| from fastapi import FastAPI, Request, WebSocket | |
| from fastapi.middleware.cors import CORSMiddleware | |
| from fastapi.responses import JSONResponse, Response | |
| # عميل WebSocket (للبثّ الاحتياطيّ فقط). عزل استيراده حتى لا يُسقِط غيابُه مسارَي /stt و/analyze. | |
| try: | |
| from websockets.asyncio.client import connect as ws_connect # websockets >= 13 | |
| _WS_OK = True | |
| except Exception: # pragma: no cover - بيئةٌ بلا websockets | |
| _WS_OK = False | |
| log = logging.getLogger("muaalem-proxy") | |
| logging.basicConfig(level=os.environ.get("LOG_LEVEL", "INFO")) | |
| # ───────────────────────── الإعدادات (من البيئة) ───────────────────────── | |
| def _env(name: str, default: str = "") -> str: | |
| return (os.environ.get(name) or default).strip() | |
| HF_TOKEN = _env("HF_TOKEN") # سرّ — Bearer لنقاط HF. لا يُقرأ من العميل أبدًا. | |
| # عناوين النقاط العلويّة (URL كامل لنقطة HF Inference). فارغٌ ⇒ المسار المقابل يُرجِع 503 نظيفًا. | |
| STT_ENDPOINT_URL = _env("STT_ENDPOINT_URL").rstrip("/") | |
| ANALYZE_ENDPOINT_URL = _env("ANALYZE_ENDPOINT_URL").rstrip("/") | |
| STREAM_ENDPOINT_URL = _env("STREAM_ENDPOINT_URL") # ws:// أو wss:// (اختياريّ) | |
| # إعادة المحاولة على البدء البارد (النقطة تستيقظ من scale-to-zero فتُرجِع 502/503/504 لحظيًّا). | |
| UPSTREAM_TIMEOUT = float(_env("UPSTREAM_TIMEOUT", "150")) # سقف الميزانية الكلّيّة للطلب (ث) | |
| ATTEMPT_TIMEOUT = float(_env("ATTEMPT_TIMEOUT", "60")) # مهلة قراءة المحاولة الواحدة (ث) | |
| MAX_RETRIES = int(_env("MAX_RETRIES", "6")) # حتى 6 محاولات | |
| RETRY_BACKOFF = float(_env("RETRY_BACKOFF", "4")) # تراجع خطّيّ أساس (ث) | |
| RETRY_BACKOFF_MAX = float(_env("RETRY_BACKOFF_MAX", "20")) | |
| COLD_START_STATUSES = {500, 502, 503, 504} # حالاتٌ عابرة → أعِد المحاولة | |
| # حدّ المعدّل لكلّ IP (نافذة منزلقة/دقيقة). 0 = مُعطَّل. | |
| RATE_LIMIT_PER_MIN = int(_env("RATE_LIMIT_PER_MIN", "40")) | |
| # CORS: قائمة صريحة و/أو تعبيرٌ نمطيّ. الافتراضيّ يسمح بنطاقات Pages/Workers وlocalhost للتطوير. | |
| ALLOWED_ORIGINS = [o.strip() for o in _env("ALLOWED_ORIGINS").split(",") if o.strip()] | |
| ALLOWED_ORIGIN_REGEX = _env( | |
| "ALLOWED_ORIGIN_REGEX", | |
| r"https://([a-z0-9-]+\.)*pages\.dev|https://([a-z0-9-]+\.)*workers\.dev|" | |
| r"http://(localhost|127\.0\.0\.1)(:\d+)?", | |
| ) | |
| MAX_BODY_BYTES = int(_env("MAX_BODY_BYTES", str(24 * 1024 * 1024))) # ~24MB سقف الجسم | |
| _client: Optional[httpx.AsyncClient] = None | |
| # ───────────────────────── حدّ المعدّل (في الذاكرة) ───────────────────────── | |
| class RateLimiter: | |
| """نافذة منزلقة بسيطة لكلّ IP داخل حاوية Space واحدة. تُقلّم المداخل الخاملة تلقائيًّا.""" | |
| def __init__(self, per_min: int) -> None: | |
| self.per_min = per_min | |
| self._hits: Dict[str, Deque[float]] = {} | |
| self._lock = asyncio.Lock() | |
| async def allow(self, ip: str) -> bool: | |
| if self.per_min <= 0: | |
| return True | |
| now = time.monotonic() | |
| async with self._lock: | |
| dq = self._hits.get(ip) | |
| if dq is None: | |
| dq = deque() | |
| self._hits[ip] = dq | |
| while dq and now - dq[0] > 60.0: | |
| dq.popleft() | |
| if len(dq) >= self.per_min: | |
| return False | |
| dq.append(now) | |
| # تقليمٌ انتهازيّ لِـIPات خاملة كي لا تنمو الذاكرة بلا حدّ. | |
| if len(self._hits) > 4096: | |
| stale = [k for k, v in self._hits.items() if not v or now - v[-1] > 120.0] | |
| for k in stale: | |
| self._hits.pop(k, None) | |
| return True | |
| limiter = RateLimiter(RATE_LIMIT_PER_MIN) | |
| def client_ip(request: Request) -> str: | |
| xff = request.headers.get("x-forwarded-for") | |
| if xff: | |
| return xff.split(",")[0].strip() | |
| return request.client.host if request.client else "unknown" | |
| # ───────────────────────── دورة الحياة (عميل HTTP مشترك) ───────────────────────── | |
| async def lifespan(_app: FastAPI): | |
| global _client | |
| _client = httpx.AsyncClient( | |
| timeout=httpx.Timeout(connect=15.0, read=ATTEMPT_TIMEOUT, write=60.0, pool=15.0), | |
| limits=httpx.Limits(max_connections=100, max_keepalive_connections=20), | |
| ) | |
| log.info( | |
| "muaalem-proxy up | stt=%s analyze=%s stream=%s token=%s", | |
| bool(STT_ENDPOINT_URL), bool(ANALYZE_ENDPOINT_URL), bool(STREAM_ENDPOINT_URL), bool(HF_TOKEN), | |
| ) | |
| try: | |
| yield | |
| finally: | |
| if _client is not None: | |
| await _client.aclose() | |
| app = FastAPI(title="muaalem-proxy", version="1.0.0", lifespan=lifespan) | |
| app.add_middleware( | |
| CORSMiddleware, | |
| allow_origins=ALLOWED_ORIGINS, | |
| allow_origin_regex=ALLOWED_ORIGIN_REGEX or None, | |
| allow_credentials=False, | |
| allow_methods=["GET", "POST", "OPTIONS"], | |
| allow_headers=["*"], | |
| max_age=86400, | |
| ) | |
| # ───────────────────────── تمرير النداء إلى النقطة (مع إعادة محاولة) ───────────────────────── | |
| def _upstream_headers() -> Dict[str, str]: | |
| h = {"Content-Type": "application/json", "Accept": "application/json"} | |
| if HF_TOKEN: | |
| h["Authorization"] = f"Bearer {HF_TOKEN}" | |
| return h | |
| async def forward_json(payload: Dict[str, Any], upstream: str, label: str) -> Response: | |
| """يمرّر `payload` (مغلّفًا بـ{inputs:{…}} من المُنادي) إلى النقطة، ويعيد ردّها كما هو. | |
| تدهورٌ آمن: كلّ حالة فشل تُرجَع كـJSONResponse بحالة HTTP مناسبة (لا استثناء يتسرّب)، | |
| فيرى العميل `!res.ok` ويتراجع بصمت. يُعيد المحاولة على البدء البارد (502/503/504) ضمن ميزانية. | |
| """ | |
| if not upstream: | |
| return JSONResponse(status_code=503, content={"error": f"{label}_upstream_not_configured"}) | |
| if _client is None: # pragma: no cover - لا يحدث بعد الإقلاع | |
| return JSONResponse(status_code=503, content={"error": "proxy_not_ready"}) | |
| headers = _upstream_headers() | |
| deadline = time.monotonic() + UPSTREAM_TIMEOUT | |
| last_status: Optional[int] = None | |
| last_detail = "" | |
| for attempt in range(MAX_RETRIES): | |
| remaining = deadline - time.monotonic() | |
| if remaining <= 0: | |
| break | |
| read_to = max(5.0, min(ATTEMPT_TIMEOUT, remaining)) | |
| try: | |
| resp = await _client.post( | |
| upstream, json=payload, headers=headers, | |
| timeout=httpx.Timeout(connect=15.0, read=read_to, write=60.0, pool=15.0), | |
| ) | |
| except (httpx.TimeoutException, httpx.TransportError) as exc: | |
| last_status, last_detail = 504, type(exc).__name__ | |
| else: | |
| if resp.status_code < 400: | |
| # نجاح — مرّر جسم النقطة كما هو (JSON عادةً؛ وإلّا الخام بنوعه). | |
| ctype = resp.headers.get("content-type", "") | |
| if "json" in ctype.lower(): | |
| try: | |
| return JSONResponse(status_code=200, content=resp.json()) | |
| except Exception: | |
| pass | |
| return Response(content=resp.content, status_code=200, | |
| media_type=ctype or "application/json") | |
| last_status, last_detail = resp.status_code, resp.text[:300] | |
| if resp.status_code not in COLD_START_STATUSES: | |
| # خطأٌ حتميّ من النقطة (مثلًا 400/422) — لا فائدة من إعادة المحاولة. | |
| return JSONResponse( | |
| status_code=resp.status_code, | |
| content={"error": f"{label}_upstream", "status": resp.status_code, "detail": last_detail}, | |
| ) | |
| # مسارٌ قابل لإعادة المحاولة: تراجَع ما لم تنفد المحاولات/الميزانية. | |
| # ملاحظة: تراجعٌ صفريّ (backoff=0) يعني «أعِد فورًا»، لا «توقّف» — الإيقاف فقط عند نفاد الوقت. | |
| if attempt + 1 >= MAX_RETRIES: | |
| break | |
| remaining = deadline - time.monotonic() | |
| if remaining <= 0: | |
| break | |
| backoff = min(RETRY_BACKOFF * (attempt + 1), RETRY_BACKOFF_MAX, remaining) | |
| log.info("%s retry %d/%d after %.0fs (last=%s)", label, attempt + 1, MAX_RETRIES, backoff, last_status) | |
| if backoff > 0: | |
| await asyncio.sleep(backoff) | |
| return JSONResponse( | |
| status_code=503, | |
| content={"error": f"{label}_unavailable", "upstream_status": last_status, "detail": last_detail}, | |
| ) | |
| class _PayloadTooLarge(Exception): | |
| """جسمٌ يتجاوز MAX_BODY_BYTES — يُميَّز عن خطأ فكّ JSON (وهو ValueError أيضًا).""" | |
| async def _read_json(request: Request) -> Dict[str, Any]: | |
| body = await request.body() | |
| if len(body) > MAX_BODY_BYTES: | |
| raise _PayloadTooLarge() | |
| if not body: | |
| return {} | |
| data = json.loads(body) # JSONDecodeError (⊂ ValueError) عند الفساد ⇒ يعالَج كـ bad_json | |
| if not isinstance(data, dict): | |
| raise TypeError("expected_json_object") | |
| return data | |
| async def _guard(request: Request) -> Optional[JSONResponse]: | |
| """حدّ المعدّل + قراءة/تحقّق الجسم المشتركان. يعيد استجابة خطأ أو None (مرّ).""" | |
| if not await limiter.allow(client_ip(request)): | |
| return JSONResponse(status_code=429, content={"error": "rate_limited"}) | |
| return None | |
| # ───────────────────────── المسارات ───────────────────────── | |
| async def health() -> Dict[str, Any]: | |
| return { | |
| "ok": True, | |
| "service": "muaalem-proxy", | |
| "routes": ["/stt", "/analyze", "/stream"], | |
| "upstream": { | |
| "stt": bool(STT_ENDPOINT_URL), | |
| "analyze": bool(ANALYZE_ENDPOINT_URL), | |
| "stream": bool(STREAM_ENDPOINT_URL) and _WS_OK, | |
| }, | |
| "token": bool(HF_TOKEN), | |
| } | |
| async def stt(request: Request) -> Response: | |
| """الطبقة الصوتية. العميل يرسل {pcm} أو {pcm, ref}؛ نغلّفها {inputs:{pcm, ref?}} للنقطة.""" | |
| blocked = await _guard(request) | |
| if blocked is not None: | |
| return blocked | |
| try: | |
| data = await _read_json(request) | |
| except _PayloadTooLarge: | |
| return JSONResponse(status_code=413, content={"error": "payload_too_large"}) | |
| except Exception: | |
| return JSONResponse(status_code=400, content={"error": "bad_json"}) | |
| pcm = data.get("pcm") | |
| if not pcm or not isinstance(pcm, str): | |
| return JSONResponse(status_code=400, content={"error": "missing_pcm"}) | |
| inner: Dict[str, Any] = {"pcm": pcm} | |
| if data.get("ref"): | |
| inner["ref"] = data["ref"] | |
| return await forward_json({"inputs": inner}, STT_ENDPOINT_URL, "stt") | |
| async def analyze(request: Request) -> Response: | |
| """المُعلِّم QPS. مسطّح {pcm, surah, ayahs|ayah, mode?, scope?, sifat?} → {inputs:{…}} للنقطة.""" | |
| blocked = await _guard(request) | |
| if blocked is not None: | |
| return blocked | |
| try: | |
| data = await _read_json(request) | |
| except _PayloadTooLarge: | |
| return JSONResponse(status_code=413, content={"error": "payload_too_large"}) | |
| except Exception: | |
| return JSONResponse(status_code=400, content={"error": "bad_json"}) | |
| pcm = data.get("pcm") | |
| if not pcm or not isinstance(pcm, str): | |
| return JSONResponse(status_code=400, content={"error": "missing_pcm"}) | |
| if "surah" not in data or data.get("surah") is None: | |
| return JSONResponse(status_code=400, content={"error": "missing_surah"}) | |
| inner: Dict[str, Any] = {"pcm": pcm, "surah": data["surah"]} | |
| if data.get("ayahs") is not None: | |
| inner["ayahs"] = data["ayahs"] | |
| if data.get("ayah") is not None: | |
| inner["ayah"] = data["ayah"] | |
| if "ayahs" not in inner and "ayah" not in inner: | |
| return JSONResponse(status_code=400, content={"error": "missing_ayah_or_ayahs"}) | |
| # معامِلات الطبقة الدقيقة (P1/P2) — تُمرَّر شفّافةً؛ النقطة الحاليّة تتجاهلها بأمان. | |
| for k in ("mode", "scope", "sifat"): | |
| if k in data and data[k] is not None: | |
| inner[k] = data[k] | |
| return await forward_json({"inputs": inner}, ANALYZE_ENDPOINT_URL, "analyze") | |
| async def stream(ws: WebSocket) -> None: | |
| """بثّ احتياطيّ: يمرّر إطارات WS ثنائيّة الاتّجاه بين العميل والنقطة العلويّة بشفافية. | |
| البروتوكول (كما في finetune/stream_server.py): العميل يرسل إطارات int16 PCM 16k + نصّ "end"؛ | |
| النقطة تُرجِع {t, text, recent}. نمرّر البايتات والنصّ كما هي، ونحقن التوكن + معامل kind. | |
| """ | |
| await ws.accept() | |
| ip = ws.headers.get("x-forwarded-for", "").split(",")[0].strip() or ( | |
| ws.client.host if ws.client else "unknown" | |
| ) | |
| if not await limiter.allow(ip): | |
| await ws.close(code=4029, reason="rate_limited") | |
| return | |
| if not (_WS_OK and STREAM_ENDPOINT_URL): | |
| await ws.close(code=1011, reason="stream_upstream_not_configured") | |
| return | |
| kind = ws.query_params.get("kind", "tlog") | |
| sep = "&" if "?" in STREAM_ENDPOINT_URL else "?" | |
| up_url = f"{STREAM_ENDPOINT_URL}{sep}kind={kind}" | |
| headers = {"Authorization": f"Bearer {HF_TOKEN}"} if HF_TOKEN else {} | |
| try: | |
| async with ws_connect( | |
| up_url, additional_headers=headers, max_size=None, open_timeout=UPSTREAM_TIMEOUT | |
| ) as upstream: | |
| async def client_to_upstream() -> None: | |
| try: | |
| while True: | |
| msg = await ws.receive() | |
| if msg.get("type") == "websocket.disconnect": | |
| break | |
| if msg.get("bytes") is not None: | |
| await upstream.send(msg["bytes"]) | |
| elif msg.get("text") is not None: | |
| await upstream.send(msg["text"]) | |
| except Exception: | |
| pass | |
| finally: | |
| try: | |
| await upstream.close() | |
| except Exception: | |
| pass | |
| async def upstream_to_client() -> None: | |
| try: | |
| async for frame in upstream: | |
| if isinstance(frame, (bytes, bytearray)): | |
| await ws.send_bytes(bytes(frame)) | |
| else: | |
| await ws.send_text(frame) | |
| except Exception: | |
| pass | |
| finally: | |
| try: | |
| await ws.close() | |
| except Exception: | |
| pass | |
| await asyncio.gather(client_to_upstream(), upstream_to_client()) | |
| except Exception as exc: # فشل الاتصال بالنقطة العلويّة → إغلاقٌ نظيف (تراجع صامت في العميل). | |
| log.info("stream upstream error: %s", type(exc).__name__) | |
| try: | |
| await ws.close(code=1011, reason="stream_upstream_error") | |
| except Exception: | |
| pass | |
| if __name__ == "__main__": | |
| import uvicorn | |
| uvicorn.run(app, host="0.0.0.0", port=int(_env("PORT", "7860")), log_level="info") | |