File size: 3,410 Bytes
a0a2f38
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
"""SA3 Sampler beta telemetry relay.

The plugin POSTs its daily performance log here without any credential. This service holds
the Hugging Face token (Space secret HF_TOKEN) and writes the file into the private dataset
repo. Nothing else is exposed: one route, one file per (installation id, day), size-capped,
rate-limited per client address.
"""
import os
import re
import threading
import time

from fastapi import FastAPI, Request
from fastapi.responses import JSONResponse
from huggingface_hub import HfApi

REPO = os.environ.get("TELEMETRY_REPO", "VortexSamples/ReverseBass-Beta-Telemetry")
TOKEN = os.environ.get("HF_TOKEN", "")
MAX_BYTES = 1_000_000            # a day of events is a few KB; anything bigger is not ours
PER_HOUR = 120                   # uploads per client address per hour
ID_RE = re.compile(r"^[0-9a-zA-Z-]{6,64}$")
DAY_RE = re.compile(r"^\d{8}$")

app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None)
api = HfApi(token=TOKEN or None)
commit_lock = threading.Lock()
recent: dict[str, list[float]] = {}
recent_lock = threading.Lock()
stats = {"accepted": 0, "rejected": 0, "started": int(time.time())}


def client_ip(request: Request) -> str:
    fwd = request.headers.get("x-forwarded-for", "")
    return (fwd.split(",")[0].strip() if fwd else (request.client.host if request.client else "?"))


def rate_limited(ip: str) -> bool:
    now = time.time()
    with recent_lock:
        hits = [t for t in recent.get(ip, []) if now - t < 3600]
        if len(hits) >= PER_HOUR:
            recent[ip] = hits
            return True
        hits.append(now)
        recent[ip] = hits
        if len(recent) > 10000:      # keep the table bounded
            for k in [k for k, v in recent.items() if not v or now - v[-1] > 3600]:
                recent.pop(k, None)
    return False


def reject(status: int, why: str):
    stats["rejected"] += 1
    return JSONResponse({"ok": False, "why": why}, status_code=status)


@app.get("/")
def root():
    return {"ok": True, "service": "sa3-telemetry-relay", "repo": REPO, **stats}


@app.post("/logs/{install_id}/{day}")
async def put_log(install_id: str, day: str, request: Request):
    if not TOKEN:
        return reject(503, "relay has no token configured")
    if not ID_RE.match(install_id) or not DAY_RE.match(day):
        return reject(400, "bad id or day")
    if rate_limited(client_ip(request)):
        return reject(429, "too many uploads")
    body = await request.body()
    if not body or len(body) > MAX_BYTES:
        return reject(413 if body else 400, "bad size")
    try:
        text = body.decode("utf-8")
    except UnicodeDecodeError:
        return reject(400, "not utf-8")
    if not text.lstrip().startswith("{"):
        return reject(400, "not a log")

    path = f"logs/{install_id}/{day}.jsonl"
    try:
        with commit_lock:
            api.upload_file(path_or_fileobj=body, path_in_repo=path, repo_id=REPO, repo_type="dataset",
                            commit_message=f"telemetry {install_id[:8]} {day}")
    except Exception as e:   # noqa: BLE001
        msg = str(e)
        if "No files have been modified" in msg or "no changes" in msg.lower():
            stats["accepted"] += 1
            return {"ok": True, "bytes": len(body), "unchanged": True}
        return reject(502, "upload failed: " + msg[:200])
    stats["accepted"] += 1
    return {"ok": True, "bytes": len(body)}