File size: 18,894 Bytes
ad55850
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
# -*- 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 مشترك) ─────────────────────────
@asynccontextmanager
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


# ───────────────────────── المسارات ─────────────────────────
@app.get("/")
@app.get("/health")
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),
    }


@app.post("/stt")
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")


@app.post("/analyze")
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")


@app.websocket("/stream")
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")