#!/usr/bin/env python3 # -*- coding: utf-8 -*- """ ZERO COST PROJECT v3.0.0 — RELAY SERVER (Docker + FastAPI) ============================================================ CPU relay yang handle SEMUA management: - Token pool (85 akun HF untuk ZeroGPU auth) - User token economy (quota harian, ban, admin) - Safety filter (hard-block + soft-score) - Spy log dengan TTL - Broadcast system - Admin panel endpoints - Routing ke Backend 1 (normal) / Backend 2 (uncensored) - Blur processing + passkey unlock - Legacy frontend guard middleware Backend ZeroGPU 1 & 2 = PURE PEKERJA (image generation only). """ import os import io import re import json import base64 import time import random import threading from datetime import datetime, timedelta, timezone from typing import Optional, Dict, List, Any from fastapi import FastAPI, Request, HTTPException from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import JSONResponse import httpx from PIL import Image, ImageFilter from pydantic import BaseModel # ╔══════════════════════════════════════════════════════════════╗ # ║ [0] CONFIG — HARDCODED ║ # ╚══════════════════════════════════════════════════════════════╝ VERSION = "3.0.0-relay-docker" # Backend ZeroGPU workers (pure image generators) BACKEND1_URL = "https://bl4ckspaces-zimageturbo-backend.hf.space" # normal BACKEND2_URL = "https://bl4ckspaces-mps.hf.space" # uncensored MPS # Shared secret antara relay ↔ worker (untuk verify request internal) WORKER_SECRET = "zc_worker_k7x2m9q4_bl4ck" # Passkeys UNCENSORED_PASSKEY = "Rusdi6967" ADMIN_PASSKEY = "Ambalabu69" # Generation params (dikirim ke worker) LOCKED_STEPS = 8 LOCKED_CFG_SCALE = 0.0 LOCKED_MAX_SEQUENCE_LENGTH = 256 MIN_BATCH = 1 MAX_BATCH = 2 MAX_PIXELS = 1024 * 1024 MIN_PIXELS = 512 * 512 MIN_SIDE = 512 MAX_SIDE = 2048 RESOLUTION_STEP = 8 MAX_ASPECT = 3.0 # Token economy TOKEN_DAILY_QUOTA = 15 COST_NORMAL = 1 COST_UNCENSORED = 3 BLUR_RADIUS = 30 # Spy log SPY_MAX_ENTRIES = 40 SPY_FULL_TTL_SECONDS = 900 # 15 menit # Worker timeouts (second) WORKER_CONNECT_TIMEOUT = 20 WORKER_READ_TIMEOUT = 600 # 10 menit (generate butuh waktu) WORKER_POOL_LIMIT = 50 # max koneksi ke worker paralel # ╔══════════════════════════════════════════════════════════════╗ # ║ [0a] HF TOKEN POOL (85 akun) ║ # ╚══════════════════════════════════════════════════════════════╝ _RAW_PAYLOAD = ( "PiRCDDtPcPFMLWkTkVaZmzoleHOunXnLIA" "-BHvZXGICstaktSwycmwNmzHGrTNmKxnlRZ" "-ZdgawyTPzXIpwhnRYIteUKSMsWnEDtGKtM" "-nMiFYAFsINxAJWPwiCQlaunmdgmrcxKoaT" "-PccpUIbTckCiafwErDLkRlsvqhgtfZaBHL" "-faGyXBPfBkaHXDMUSJtxEggonhhZbomFIz" "-SndsPaRWsevDXCgZcSjTUlBYUJqOkSfFmn" "-CqobFdUpeVCeuhUaiuXwvdczBUmoUHXRGa" "-JKCQYUhhHPPkpucegqkNSyureLdXpmeXRF" "-tBYfslUwHNiNMufzwAYIlrDVovEWmOQulC" "-LKLdrdUxyUyKODSUthmqHXqDMfHrQueera" "-ivSBboJYQVcifWkCNcOTOnxUQrZOtOglnU" "-jiSbBMUmAniRpJOmVIlczuqpRjwSeuizLk" "-VcXaKQLEawBWZbNrBOSLTjrVtTuSvobhLL" "-ZrlTPvhDmYqZGGFuIqDDCrQRcWRhYcuyOI" "-FCambosUqUQJrThbIveHglnvjoNpOGWBsW" "-kUyoiWTbZlNfSrdTNaVINuwlNTQseFCfZB" "-WGarKlgPBzpJeKxpqirFgnKKAtOFBFomSe" "-IZwzmRBCALYfvYtmtvTWsIQYvHuRGUiGyr" "-NtijfwwAPQRknELkhIWjMQQUUqzgwhIjeu" "-obVKYRMqECBoLsBWOKyfVWtHlugAhhuaIH" "-EsDAvVqRZCbigQrpDFNinlVeijagnAjETW" "-yuMifxRJoXWKPRGgYFrXHXTGdoKBuCZCUU" "-YthKrdEtrmyDbBteZcGzNeoDqGAxzeEinv" "-JgNjfcunLsOBcZIaOYcFqgcZIZWjbnocJn" "-cINBgwvihyKiTpxwDTXjnHTnlQivLCluGJ" "-jnciPeeWUwQbHNITBRtOgPjnqWkgAqZDhq" "-uwTpdbUKmiKUZkpWgsOJyaywRBSPeSHcLY" "-ddTsgQvyVXUSRsYrcAioNbGnTsOVugcSpK" "-BoDGLtNQdSSiIvuveGBiCFJKzEEqQmZTbz" "-aKBQewTCodMpeijniAvkcsqoSEhMrbhiKE" "-hOrZAwkLlabOCjjZVStFoufImunhcjlEhz" "-NkfolaYSzHbBLkqkalzvbLJWHxDQfUEwsr" "-OoFVTidQQHCpKkDBnTsTvPeOogJkVfEWAZ" "-plFpieuUJmDJBhyPIdjBmDKKTVNcNpulXJ" "-AyHpfXRVxcrmiOFUYMLiNAaVtyLwwkdvWJ" "-pHNLBEtYhQwfUQJFAxtLcklIaIXyQQkIAt" "-CpXxqjXPCvnLEyfdlsbgdqXdotkMkzvmBU" "-PerSDYxXPlPviQNmACCXyzyIfNBnJgpXst" "-PKxgUunqjxFjeMUrhgPuHAfUYhrGsFPlBh" "-QaVcpabAfUrqncLNBwgXCkpksoGWWxgQxK" "-PdxIUKbxDsVyMeDtHozMnbQrQXwFxZjKdy" "-aMejPqEuMzjwYCKuUVQWOrcvfoWKtKbbgv" "-CVBXSUlqCHUOPgjbLtmkxXQieFECclgZkT" "-pMRYEHGgfGjuxgmKVFobnNhPpxmEoWgsZb" "-HGSvrRpieiWTbSqRVtdKVgeMHhxVQmakAE" "-yrUmRlsTkRjREynarJfHTAOGjcbCQAJDQd" "-PtPUYmOAhSihYOJpQClGZUsWXBBVzQLOrU" "-PLNxMJmXrCgzMNZtIJahbKQfmGFeNceLTE" "-zLAfTPijfcCqWHOLBuTjEeNXIcWVxsPbbq" "-yjgUeEDychEqlzSMSsRZnUFVvOtpAdWtAx" "-fKcKRTqWTLhgIqiCquYvYyqksBBwvXTlSm" "-uUdhNJtMIsIWxZZQmjOAoRwkbuhVbDRkUn" "-MunRYwrfqwMdLZEbZOGHpPPYnDmuDimopu" "-UWAwKSXCNzWvGmsdOhnmqeUIyTJfiNXIZt" "-iheoERRZTkhdAdnDHJatYgOiLapIgVTFJw" "-LLnfrrHopubyFJllMntnFliVrXFFXaAwoj" "-PGjVbqVHxzglCMkvteIIbdVtPJlCPxefLA" "-hkrmPSIlSUTnqsmQLsGECXCwxSmQLnJmEJ" "-gqxPdLdhmngjpzBznOrzIQqtOrsUGNYbBJ" "-QwsXWJiLNOCRMYuZGNsYxYTLpXygHmQjvW" "-yQBGhxEAztptTIQXgOYOoVqcAEfiCLdsFd" "-pFGVSiujNvgifJoBOHLLPogiOgOZyiQlqu" "-HwAhtTKtaGUOEjDcjfXVNOMFHKHyZlsQpI" "-wbnWdXOgZrkfktwYDlbiqCuhQRZLCsBvak" "-LDuJYfWVljMjsFFcbKahhsglPTjsvTJfqy" "-DAmhXglyPbLZAlbiekBRHnjWZrzdtJSbAL" "-SByeDprwukcociMPOhJKEtAodUOWSiETeq" "-oqfZPYkBcrgJiuWOPLINQImroLKjQDmbRZ" "-aJLwnooVCfTUssMXtAWIYfgsKhZPthlkYK" "-UqUQinbBqBVPuewmsxNImHMQCTIUWOpCdA" "-EiEjqpvtdAtxpUpnxIgRrKPBkUJuWaQJnA" "-cvXjppPnFhwdCbBFukSwNLVOkDOQEPUSqd" "-PoSgaPVVsdesVCucFerlntYWiJpgreWdGJ" "-UELvMUNsxlMemLmYiBMLDXIUHxTiWqIYqd" "-FSVAPeAzGEfLJIYRbSNxQxfTvseIkWVCZw" "-StCLoInFnlseUbPzqIFpFCCUCBsQjxzuIc" "-oHKZuJHonaPBaiGvYPMSENeXOcYIzcJTJy" "-cjmPeQZixhFdRdImgagBonMbwQhuTitAqM" "-dPztnFlmVxzpdbFbCWWfjWAZGTcahucGIi" "-NMEbWuaiULTGpoxulWSHwubeGeopaxiPIH" "-sGJhTjdDVTBsPPZRlXPYHlmBbPpGkymfVH" "-HwkuSfVGIzfUrnGXOYonaLgBqKMksniSmG" "-vnJYppejvDHxpxDIHgUoNBshieASIRvtXD" "-JybhNvxSGmhEMJTxnLktzLdempTyAEXYuu" ) POOL_85 = ["hf_" + seg.strip() for seg in _RAW_PAYLOAD.split("-")] assert len(POOL_85) == 85, f"Token pool should have 85, got {len(POOL_85)}" # ╔══════════════════════════════════════════════════════════════╗ # ║ [1] UTILITY FUNCTIONS ║ # ╚══════════════════════════════════════════════════════════════╝ def _utcnow(): return datetime.now(timezone.utc).replace(tzinfo=None) def _utc_str(): return _utcnow().strftime("%Y-%m-%d %H:%M:%S UTC") _log_seen = {} _log_throttle_lock = threading.Lock() def _log_throttled(key: str, msg: str, interval: int = 300) -> bool: now = time.time() with _log_throttle_lock: last = _log_seen.get(key, 0) if now - last < interval: return False _log_seen[key] = now print(msg, flush=True) return True def _check_quota_error(error_msg: str) -> bool: msg = str(error_msg).lower() return any(kw in msg for kw in [ "exceeded your zerogpu quota", "quota", "0s left", "authenticate with a hugging face token" ]) # ╔══════════════════════════════════════════════════════════════╗ # ║ [2] TOKEN POOL MANAGER ║ # ╚══════════════════════════════════════════════════════════════╝ class TokenPoolManager: def __init__(self, tokens: List[str]): self.tokens = tokens self._lock = threading.Lock() self._current_index = 0 self._status = { i: {"state": "active", "rest_until": None, "exhausted_at": None, "use_count": 0, "error_count": 0} for i in range(len(tokens)) } def _reactivate_rested(self): now = time.time() for i, s in self._status.items(): if s["state"] == "resting" and s["rest_until"] and now >= s["rest_until"]: s["state"] = "active" s["rest_until"] = None s["exhausted_at"] = None print(f"[TokenPool] Token #{i} reactivated", flush=True) def get_next_token(self): with self._lock: self._reactivate_rested() for i in range(len(self.tokens)): idx = (self._current_index + i) % len(self.tokens) if self._status[idx]["state"] == "active": self._current_index = (idx + 1) % len(self.tokens) self._status[idx]["use_count"] += 1 return idx, self.tokens[idx] return None, None def mark_exhausted(self, index: int, error_msg: str = ""): with self._lock: rest = self._parse_rest_time(error_msg) self._status[index].update({ "state": "resting", "exhausted_at": time.time(), "rest_until": time.time() + rest, "error_count": self._status[index]["error_count"] + 1, }) print(f"[TokenPool] Token #{index} exhausted → resting {self._format_duration(rest)}", flush=True) def is_token_active(self, index: int) -> bool: with self._lock: self._reactivate_rested() return self._status[index]["state"] == "active" def get_token_index(self, token: str): try: return self.tokens.index(token) except ValueError: return None def _parse_rest_time(self, error_msg: str) -> int: msg = str(error_msg).lower() mult = {"minute": 60, "hour": 3600, "day": 86400} m = re.search(r"in\s+(\d+)\s+(minute|hour|day)s?", msg) if m: return int(m.group(1)) * mult.get(m.group(2), 3600) m = re.search(r"(\d+)\s+(minute|hour|day)s?\s+left", msg) if m: return int(m.group(1)) * mult.get(m.group(2), 3600) now = _utcnow() tomorrow = (now + timedelta(days=1)).replace(hour=0, minute=0, second=0, microsecond=0) return int((tomorrow - now).total_seconds()) @staticmethod def _format_duration(seconds: float) -> str: if seconds <= 0: return "0s" seconds = int(seconds) h, m, s = seconds // 3600, (seconds % 3600) // 60, seconds % 60 if h > 0: return f"{h}h {m}m" if m > 0: return f"{m}m {s}s" return f"{s}s" def get_pool_status(self) -> Dict: with self._lock: self._reactivate_rested() now = time.time() active = sum(1 for s in self._status.values() if s["state"] == "active") resting = sum(1 for s in self._status.values() if s["state"] == "resting") resting_details = [] for i, s in self._status.items(): if s["state"] == "resting" and s["rest_until"]: rem = max(0, s["rest_until"] - now) resting_details.append({ "index": i, "remaining_seconds": round(rem), "remaining_human": self._format_duration(rem), }) return { "total_tokens": len(self.tokens), "active": active, "resting": resting, "total_uses": sum(s["use_count"] for s in self._status.values()), "total_errors": sum(s["error_count"] for s in self._status.values()), "resting_details": resting_details, "all_exhausted": active == 0, } pool_manager = TokenPoolManager(POOL_85) # ╔══════════════════════════════════════════════════════════════╗ # ║ [3] USER TOKEN MANAGER ║ # ╚══════════════════════════════════════════════════════════════╝ STORAGE_DIR = os.environ.get("RELAY_STORAGE", "/data") USER_TOKENS_FILE = os.path.join(STORAGE_DIR, "user_tokens.json") BROADCAST_FILE = os.path.join(STORAGE_DIR, "broadcast.json") os.makedirs(STORAGE_DIR, exist_ok=True) class UserTokenManager: def __init__(self, filepath: str): self.filepath = filepath self._lock = threading.Lock() self._data = {} self._load() def _load(self): try: if os.path.exists(self.filepath): with open(self.filepath, "r") as f: self._data = json.load(f) print(f"[UserTokens] Loaded {len(self._data)} users", flush=True) except Exception as e: print(f"[UserTokens] Load failed: {e}", flush=True) self._data = {} def _save(self): try: tmp = self.filepath + ".tmp" with open(tmp, "w") as f: json.dump(self._data, f, indent=2) os.replace(tmp, self.filepath) except Exception as e: _log_throttled("ut_save_fail", f"[UserTokens] Save failed: {e}", 300) def _should_reset(self, user: Dict) -> bool: lr = user.get("last_reset", "") if not lr: return True try: dt = datetime.strptime(lr, "%Y-%m-%d %H:%M:%S UTC") return (_utcnow() - dt).total_seconds() > 86400 except Exception: return True def get_user(self, ip: str) -> Dict: with self._lock: if ip not in self._data: self._data[ip] = { "tokens": TOKEN_DAILY_QUOTA, "is_admin": False, "banned_until": None, "ban_reason": None, "ban_by": None, "last_reset": _utc_str(), "total_used": 0, "uncensored_used": 0, "first_seen": _utc_str(), "gens": 0, } self._save() user = self._data[ip] if not user.get("is_admin") and self._should_reset(user): user["tokens"] = TOKEN_DAILY_QUOTA user["last_reset"] = _utc_str() self._save() return dict(user) def is_banned(self, ip: str): user = self.get_user(ip) if not user.get("banned_until"): return False, None try: until = datetime.strptime(user["banned_until"], "%Y-%m-%d %H:%M:%S UTC") if _utcnow() < until: return True, {"until": user["banned_until"], "reason": user.get("ban_reason", "No reason"), "by": user.get("ban_by", "system")} with self._lock: self._data[ip]["banned_until"] = None self._data[ip]["ban_reason"] = None self._save() return False, None except Exception: return False, None def consume(self, ip: str, cost: int): with self._lock: user = self._data.get(ip) if not user: return False, 0 if user.get("is_admin"): return True, 999999 if user.get("tokens", 0) < cost: return False, user.get("tokens", 0) user["tokens"] -= cost user["total_used"] = user.get("total_used", 0) + cost user["gens"] = user.get("gens", 0) + 1 self._save() return True, user["tokens"] def mark_uncensored(self, ip: str): with self._lock: if ip in self._data: self._data[ip]["uncensored_used"] = self._data[ip].get("uncensored_used", 0) + 1 self._save() def grant_tokens(self, ip: str, amount: int, by: str = "system"): with self._lock: self.get_user(ip) self._data[ip]["tokens"] = self._data[ip].get("tokens", 0) + int(amount) self._data[ip]["last_grant"] = {"amount": amount, "by": by, "at": _utc_str()} self._save() return self._data[ip]["tokens"] def ban_ip(self, ip: str, hours: float, reason: str = "", by: str = "admin"): with self._lock: self.get_user(ip) until = _utcnow() + timedelta(hours=float(hours)) self._data[ip]["banned_until"] = until.strftime("%Y-%m-%d %H:%M:%S UTC") self._data[ip]["ban_reason"] = str(reason) self._data[ip]["ban_by"] = str(by) self._save() return self._data[ip]["banned_until"] def unban_ip(self, ip: str): with self._lock: if ip in self._data: self._data[ip]["banned_until"] = None self._data[ip]["ban_reason"] = None self._save() return True return False def promote_admin(self, ip: str, by: str = "system"): with self._lock: self.get_user(ip) self._data[ip]["is_admin"] = True self._data[ip]["promoted_by"] = str(by) self._data[ip]["promoted_at"] = _utc_str() self._save() def demote_admin(self, ip: str): with self._lock: if ip in self._data: self._data[ip]["is_admin"] = False self._data[ip].pop("promoted_by", None) self._data[ip].pop("promoted_at", None) self._save() return True return False def list_users(self, limit: int = 300): with self._lock: items = [{"ip": ip, **data} for ip, data in self._data.items()] items.sort(key=lambda x: x.get("gens", 0), reverse=True) return items[:limit] def stats(self): with self._lock: return { "total_users": len(self._data), "admins": sum(1 for u in self._data.values() if u.get("is_admin")), "banned": sum(1 for u in self._data.values() if u.get("banned_until")), "total_tokens_circulating": sum(u.get("tokens", 0) for u in self._data.values()), "total_tokens_consumed": sum(u.get("total_used", 0) for u in self._data.values()), "total_generations": sum(u.get("gens", 0) for u in self._data.values()), } user_tokens = UserTokenManager(USER_TOKENS_FILE) # ╔══════════════════════════════════════════════════════════════╗ # ║ [4] BROADCAST ║ # ╚══════════════════════════════════════════════════════════════╝ _broadcast_lock = threading.Lock() _broadcast = {"message": None, "type": "info", "active": False, "set_at": None, "set_by": None} def _load_broadcast(): global _broadcast try: if os.path.exists(BROADCAST_FILE): with open(BROADCAST_FILE, "r") as f: _broadcast = json.load(f) except Exception: pass def _save_broadcast(): try: with open(BROADCAST_FILE, "w") as f: json.dump(_broadcast, f, indent=2) except Exception: pass _load_broadcast() def broadcast_set(message, btype: str = "info", by: str = "admin"): global _broadcast with _broadcast_lock: _broadcast = { "message": str(message) if message else None, "type": btype, "active": bool(message), "set_at": _utc_str(), "set_by": str(by), } _save_broadcast() return dict(_broadcast) def broadcast_clear(): return broadcast_set(None) def broadcast_get(): with _broadcast_lock: return dict(_broadcast) # ╔══════════════════════════════════════════════════════════════╗ # ║ [5] SAFETY FILTER ║ # ╚══════════════════════════════════════════════════════════════╝ class SmartSafetyFilter: HARD_BLOCK_PHRASES = [ "child porn", "child pornography", "underage sex", "minor sex", "kid sex", "children sex", "preteen sex", "loli hentai", "loli nude", "loli sex", "loli porn", "shota nude", "shota sex", "shota hentai", "schoolgirl nude", "schoolgirl sex", "school girl sex", "underage nude", "underage naked", "underage nsfw", "cub porn", "cub sex", "infant sex", "baby sex", "baby nude", "bestiality", "zoophilia", "animal sex", "rape ", "raped ", "raping ", "forced sex", "non-con sex", "nonconsensual", "sexual assault", "necrophilia", "gore porn", "snuff", "deepfake nude", "deepfake sex", "deepfake porn", "revenge porn", ] HARD_BLOCK_WORDS = {"loli", "shota", "shotacon", "lolicon", "pedophile", "pedophilia", "pedo", "zoophile", "incest"} SOFT_TERMS = { "nude": 8, "naked": 8, "nudity": 8, "topless": 7, "bottomless": 7, "bare breasts": 7, "exposed breasts": 7, "genitalia": 12, "penis": 10, "vagina": 10, "vulva": 10, "sex ": 10, "fucking": 10, "intercourse": 10, "orgasm": 7, "ejaculation": 10, "cumshot": 10, "creampie": 10, "masturbat": 8, "porn": 10, "porno": 10, "pornography": 10, "xxx": 8, "nsfw": 4, "hentai": 6, "erotic": 4, "lingerie": 4, "underwear": 3, "panties": 4, "thong": 3, "see-through": 5, "wet t-shirt": 6, "upskirt": 9, "downblouse": 7, "spread legs": 6, "ahegao": 7, "tentacle": 5, "striptease": 6, "stripper": 5, "camgirl": 5, "onlyfans": 3, "bondage": 6, "bdsm": 5, "fetish": 4, "gore": 8, "murder": 5, "r18": 6, "r-18": 6, "dildo": 7, "sex toy": 7, } SAFE_CONTEXT = { "bikini": -4, "swimsuit": -4, "swimwear": -4, "bathing suit": -4, "beachwear": -5, "beach": -5, "pool": -4, "resort": -4, "tropical": -3, "summer": -3, "vacation": -4, "fashion": -5, "runway": -5, "editorial": -5, "magazine": -4, "photoshoot": -3, "crop top": -3, "shorts": -2, "mini skirt": -2, "gym": -4, "workout": -4, "fitness": -4, "athletic": -4, "yoga": -4, "swimming": -4, "fine art": -5, "artistic": -3, "sculpture": -4, "anatomy": -5, "medical": -5, "educational": -5, "portrait": -3, "headshot": -5, } THRESHOLD = 15 @classmethod def check(cls, prompt: str, negative_prompt: str = "", hard_only: bool = False): combined = ((prompt or "") + " " + (negative_prompt or "")).lower() for phrase in cls.HARD_BLOCK_PHRASES: if phrase in combined: return (False, f"Prohibited content detected: '{phrase.strip()}'. Not allowed in any mode.", 100, [phrase.strip()]) for word in cls.HARD_BLOCK_WORDS: if re.search(r"\b" + re.escape(word) + r"\b", combined): return (False, f"Prohibited content detected: '{word}'. Not allowed in any mode.", 100, [word]) if hard_only: return True, "OK (hard-blocks only)", 0, [] score = 0 hits = [] for term, weight in cls.SOFT_TERMS.items(): if " " in term: if term in combined: score += weight hits.append((term, weight)) elif re.search(r"\b" + re.escape(term) + r"\b", combined): score += weight hits.append((term, weight)) for term, weight in cls.SAFE_CONTEXT.items(): if " " in term: if term in combined: score += weight elif re.search(r"\b" + re.escape(term) + r"\b", combined): score += weight score = max(0, score) if score >= cls.THRESHOLD: top = sorted(hits, key=lambda x: -x[1])[:5] top_str = ", ".join(f"'{t}'(+{w})" for t, w in top) return (False, f"Content scored {score}/{cls.THRESHOLD}. Flagged: {top_str}. Try milder prompt or Uncensored mode.", score, [t for t, _ in top]) return True, "OK", score, [] safety_filter = SmartSafetyFilter() # ╔══════════════════════════════════════════════════════════════╗ # ║ [6] STATS + SPY LOG ║ # ╚══════════════════════════════════════════════════════════════╝ _generation_stats = { "total_generations": 0, "total_time_seconds": 0.0, "last_generation_time": None, "last_resolution": None, "last_seed": None, "started_at": time.time(), "errors": 0, "safety_blocks": 0, "normal_gens": 0, "uncensored_gens": 0, "tokens_consumed": 0, "relay_failures": 0, "legacy_blocked": 0, } def _update_stats(width, height, elapsed, seed, mode="normal", tokens=1): _generation_stats["total_generations"] += 1 _generation_stats["total_time_seconds"] += elapsed _generation_stats["last_generation_time"] = round(elapsed, 2) _generation_stats["last_resolution"] = f"{width}x{height}" _generation_stats["last_seed"] = seed _generation_stats["tokens_consumed"] += tokens if mode == "uncensored": _generation_stats["uncensored_gens"] += 1 else: _generation_stats["normal_gens"] += 1 _spy_lock = threading.Lock() _spy_log = [] _spy_counter = 0 def _blur_data_url(data_url: str, radius: int = BLUR_RADIUS): try: _, b64 = data_url.split(",", 1) raw = base64.b64decode(b64) img = Image.open(io.BytesIO(raw)) img = img.filter(ImageFilter.GaussianBlur(radius=radius)) if img.mode != "RGB": img = img.convert("RGB") buf = io.BytesIO() img.save(buf, format="JPEG", quality=90) return "data:image/jpeg;base64," + base64.b64encode(buf.getvalue()).decode("utf-8") except Exception as e: print(f"[Relay] blur failed: {e}", flush=True) return None def _make_thumbnail(data_url: str, max_size: int = 224): try: _, b64 = data_url.split(",", 1) raw = base64.b64decode(b64) img = Image.open(io.BytesIO(raw)) img.thumbnail((max_size, max_size), Image.LANCZOS) if img.mode != "RGB": img = img.convert("RGB") buf = io.BytesIO() img.save(buf, format="JPEG", quality=60) return "data:image/jpeg;base64," + base64.b64encode(buf.getvalue()).decode("utf-8") except Exception: return None def _add_spy_entry(entry: Dict) -> int: global _spy_counter entry["time_ts"] = time.time() with _spy_lock: _spy_counter += 1 entry["id"] = _spy_counter _spy_log.append(entry) if len(_spy_log) > SPY_MAX_ENTRIES: del _spy_log[:len(_spy_log) - SPY_MAX_ENTRIES] now = time.time() for e in _spy_log: exp = e.get("_full_expires") if exp and now > exp: for im in e.get("images", []): im["full"] = None e["_full_expires"] = None return entry["id"] def _log_to_spy(client_ip, client_ua, token_idx, prompt, neg, width, height, batch, result_dict, stage="generate", mode="normal", worker_key="backend1"): try: images_spy = [] if result_dict.get("success"): for img in result_dict.get("images", []): item = { "seed": img.get("seed"), "thumb": _make_thumbnail(img.get("data", "")), "full": None, "is_blurred": bool(img.get("is_blurred", False)), } if mode == "uncensored": item["full"] = img.get("data") images_spy.append(item) md = result_dict.get("metadata", {}) if isinstance(result_dict, dict) else {} entry = { "time": _utc_str(), "stage": stage, "mode": mode, "worker": worker_key, "ip": client_ip, "user_agent": client_ua, "prompt": str(prompt)[:1000], "negative": str(neg)[:500] if neg else "", "width": width, "height": height, "batch": batch, "seed": md.get("seed"), "duration_s": md.get("duration"), "hf_token_index": token_idx, "success": bool(result_dict.get("success")), "error": None if result_dict.get("success") else str(result_dict.get("error", ""))[:300], "images": images_spy, "safety_score": result_dict.get("_safety_score"), "tokens_consumed": md.get("tokens_consumed"), "user_is_admin": result_dict.get("_user_is_admin"), "_full_expires": (time.time() + SPY_FULL_TTL_SECONDS) if mode == "uncensored" else None, } return _add_spy_entry(entry) except Exception as e: _log_throttled("spy_err", f"[Relay] spy log error: {e}", 300) return None def _build_feed_payload(since: int = 0): with _spy_lock: snapshot = list(_spy_log) if since > 0: picked = [e for e in snapshot if e.get("id", 0) > since] picked.reverse() else: picked = list(reversed(snapshot)) feed = [] for e in picked: ec = dict(e) ec.pop("_full_expires", None) ec["images"] = [{"seed": im.get("seed"), "thumb": im.get("thumb"), "is_blurred": im.get("is_blurred", False)} for im in e.get("images", [])] feed.append(ec) return feed # ╔══════════════════════════════════════════════════════════════╗ # ║ [7] HTTPX CLIENT (ASYNC, CONNECTION POOLING) ║ # ╚══════════════════════════════════════════════════════════════╝ _http_client: Optional[httpx.AsyncClient] = None async def get_http_client() -> httpx.AsyncClient: global _http_client if _http_client is None or _http_client.is_closed: _http_client = httpx.AsyncClient( timeout=httpx.Timeout( connect=WORKER_CONNECT_TIMEOUT, read=WORKER_READ_TIMEOUT, write=30.0, pool=10.0, ), limits=httpx.Limits( max_connections=WORKER_POOL_LIMIT, max_keepalive_connections=20, keepalive_expiry=30, ), follow_redirects=True, ) return _http_client _worker_status = { "backend1": {"reachable": None, "last_check": 0, "last_error": None, "ok": 0, "fail": 0}, "backend2": {"reachable": None, "last_check": 0, "last_error": None, "ok": 0, "fail": 0}, } _ws_lock = asyncio.Lock() if False else threading.Lock() # placeholder, pakai sync lock async def _check_worker(backend_key: str, force: bool = False) -> Dict: url = BACKEND1_URL if backend_key == "backend1" else BACKEND2_URL now = time.time() with _ws_lock: if not force and now - _worker_status[backend_key]["last_check"] < 120: return dict(_worker_status[backend_key]) ok, err = False, None try: client = await get_http_client() r = await client.get(url + "/", timeout=8.0) ok = (r.status_code == 200) if not ok: err = f"HTTP {r.status_code}" except Exception as e: err = str(e)[:150] with _ws_lock: _worker_status[backend_key].update({ "reachable": ok, "last_check": time.time(), "last_error": err }) return dict(_worker_status[backend_key]) async def _relay_to_worker(backend_key: str, prompt: str, neg: str, width: int, height: int, seed: int, batch: int, hf_token: str) -> Dict: """ Relay request ke backend ZeroGPU. Backend expose: POST /gradio_api/call/worker_generate Args: [WORKER_SECRET, prompt, neg, width, height, seed, batch] """ url = BACKEND1_URL if backend_key == "backend1" else BACKEND2_URL data = [WORKER_SECRET, str(prompt), str(neg or ""), int(width), int(height), int(seed), int(batch)] headers = { "Content-Type": "application/json", "Authorization": f"Bearer {hf_token}", } post_url = f"{url}/gradio_api/call/worker_generate" client = await get_http_client() # POST — get event_id last_err = None event_id = None for attempt in range(3): try: r = await client.post(post_url, json={"data": data}, headers=headers) if r.status_code in (502, 503): last_err = f"{backend_key} cold-starting (HTTP {r.status_code})" print(f"[Worker] {last_err} — retry {attempt+1}/3 in 10s", flush=True) await asyncio.sleep(10) continue r.raise_for_status() body = r.json() event_id = body.get("event_id") if not event_id: raise RuntimeError(f"No event_id from {backend_key}: {r.text[:200]}") break except httpx.TimeoutException: last_err = f"{backend_key} POST timeout" print(f"[Worker] {last_err} — retry {attempt+1}/3", flush=True) await asyncio.sleep(3) continue except httpx.HTTPStatusError as e: raise RuntimeError(f"{backend_key} POST failed: HTTP {e.response.status_code}") except Exception as e: raise RuntimeError(f"{backend_key} POST error: {type(e).__name__}: {e}") if not event_id: raise RuntimeError(last_err or f"{backend_key} unreachable after 3 retries") # SSE stream get_url = f"{post_url}/{event_id}" current_event = None try: async with client.stream("GET", get_url, headers=headers) as resp: resp.raise_for_status() async for line in resp.aiter_lines(): if not line: current_event = None continue if line.startswith("event:"): current_event = line[6:].strip() continue if not line.startswith("data:"): continue payload = line[5:].strip() if not payload: continue try: float(payload) continue # heartbeat except ValueError: pass try: parsed = json.loads(payload) except json.JSONDecodeError: continue if current_event == "complete" or isinstance(parsed, list): if parsed and isinstance(parsed[0], str): return json.loads(parsed[0]) raise RuntimeError(f"{backend_key} unexpected output format") if isinstance(parsed, dict) and parsed.get("msg"): msg = parsed["msg"] if msg == "process_completed": if parsed.get("success") is False: raise RuntimeError(f"{backend_key}: {parsed.get('error') or 'process failed'}") out = parsed.get("output") or {} arr = out.get("data") or parsed.get("data") or [] if arr and isinstance(arr[0], str): return json.loads(arr[0]) raise RuntimeError(f"{backend_key} empty output") if msg == "process_failed": raise RuntimeError(f"{backend_key}: {parsed.get('error') or 'process_failed'}") if msg == "queue_full": raise RuntimeError(f"{backend_key} queue full") continue except httpx.TimeoutException: raise RuntimeError(f"{backend_key} SSE timeout") except httpx.HTTPStatusError as e: raise RuntimeError(f"{backend_key} SSE HTTP {e.response.status_code}") raise RuntimeError(f"{backend_key} stream ended without output") # ╔══════════════════════════════════════════════════════════════╗ # ║ [8] FASTAPI APP ║ # ╚══════════════════════════════════════════════════════════════╝ import asyncio app = FastAPI( title="ZeroCost Relay", version=VERSION, docs_url=None, # disable Swagger UI redoc_url=None, # disable ReDoc ) # CORS app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], expose_headers=["*"], ) # ── Helper: extract client info ── def _extract_client_info(request: Request): ip = "unknown" ua = "unknown" try: xff = request.headers.get("x-forwarded-for", "") if xff: ip = xff.split(",")[0].strip() else: ip = request.headers.get("x-real-ip", "") or (request.client.host if request.client else "unknown") ua = request.headers.get("user-agent", "unknown") if len(ua) > 300: ua = ua[:300] + "..." except Exception: pass return ip or "unknown", ua or "unknown" # ── Helper: extract token from Authorization header ── def _extract_token(request: Request): auth = request.headers.get("authorization", "") if auth.lower().startswith("bearer "): return auth[7:].strip() return auth.strip() if auth else None # ── Helper: admin auth ── def _admin_auth(key: str) -> bool: return str(key).strip() == ADMIN_PASSKEY # ╔══════════════════════════════════════════════════════════════╗ # ║ [9] MIDDLEWARE — LEGACY FRONTEND GUARD ║ # ╚══════════════════════════════════════════════════════════════╝ @app.middleware("http") async def legacy_frontend_guard(request: Request, call_next): """Block request dari frontend v1.5.1 di pintu masuk.""" path = request.url.path method = request.method # Block /get_tokens (plural) — endpoint v1.5.1 if method == "POST" and path == "/get_tokens": _generation_stats["legacy_blocked"] += 1 _log_throttled( "guard_legacy_get_tokens", f"[Guard] 🚫 Blocked legacy /get_tokens from {request.client.host if request.client else '?'}", 60, ) return JSONResponse( status_code=410, content={ "error": "OUTDATED_FRONTEND", "message": "Your frontend is outdated (v1.5.1). Hard-refresh: Ctrl+Shift+R", "required_version": "3.0.0", "blocked_by": "legacy_guard", } ) # Block /generate dengan 6 args (hf_token di arg ke-6) if method == "POST" and path == "/generate": try: body_bytes = await request.body() if body_bytes: body = json.loads(body_bytes) # Gradio-style: { "data": [...] } if isinstance(body, dict) and "data" in body: payload = body["data"] if isinstance(payload, list) and len(payload) == 6: sixth = payload[5] if isinstance(sixth, str) and sixth.startswith("hf_"): _generation_stats["legacy_blocked"] += 1 _log_throttled( "guard_legacy_generate", f"[Guard] 🚫 Blocked legacy generate (6 args) from {request.client.host if request.client else '?'}", 30, ) return JSONResponse( status_code=410, content={ "error": "OUTDATED_FRONTEND", "message": "Frontend v1.5.1 detected. Hard-refresh to load v3.0.0", "required_version": "3.0.0", "blocked_by": "legacy_guard", } ) except (json.JSONDecodeError, Exception): pass # Restore body untuk downstream handler async def receive_wrapper(): return {"type": "http.request", "body": body_bytes, "more_body": False} request._receive = receive_wrapper return await call_next(request) # ╔══════════════════════════════════════════════════════════════╗ # ║ [10] PUBLIC ENDPOINTS ║ # ╚══════════════════════════════════════════════════════════════╝ @app.get("/") async def root(): return {"service": "ZeroCost Relay", "version": VERSION, "status": "online"} @app.get("/info") async def info(): """Compatibility endpoint (frontend pakai ini untuk detect API).""" return { "version": VERSION, "named_endpoints": { "/generate": {}, "/get_token": {}, "/get_tokens": {}, "/pool_status": {}, "/user_status": {}, "/get_config": {}, "/health": {}, "/get_broadcast": {}, "/unlock_image": {}, "/diagnose": {}, "/get_stats": {}, "/reset_stats": {}, "/admin_feed": {}, "/admin_feed_since": {}, "/admin_users": {}, "/admin_ban": {}, "/admin_unban": {}, "/admin_grant": {}, "/admin_promote": {}, "/admin_demote": {}, "/admin_broadcast": {}, "/admin_image": {}, "/admin_log": {}, } } @app.post("/get_token") async def get_token(): idx, token = pool_manager.get_next_token() if idx is None: status = pool_manager.get_pool_status() soonest = None for d in status.get("resting_details", []): if soonest is None or d["remaining_seconds"] < soonest: soonest = d["remaining_seconds"] return { "success": False, "error": f"All HF tokens resting. Next in ~{TokenPoolManager._format_duration(soonest or 0)}.", "pool_status": status, } return { "success": True, "token": token, "token_index": idx, "pool_active": pool_manager.get_pool_status()["active"], "pool_total": len(POOL_85), } @app.post("/get_tokens") async def get_tokens_legacy(): """Stub untuk frontend v1.5.1 — return POOL_85 array.""" _log_throttled("legacy_get_tokens", "[Relay] Legacy 'get_tokens' called (served pool)", 600) return POOL_85 @app.post("/pool_status") async def pool_status(): return pool_manager.get_pool_status() @app.post("/user_status") async def user_status(request: Request): client_ip, _ = _extract_client_info(request) user = user_tokens.get_user(client_ip) banned, ban_info = user_tokens.is_banned(client_ip) return { "ip": client_ip, "tokens": user.get("tokens", 0), "is_admin": user.get("is_admin", False), "total_used": user.get("total_used", 0), "uncensored_used": user.get("uncensored_used", 0), "gens": user.get("gens", 0), "last_reset": user.get("last_reset"), "banned": banned, "ban_info": ban_info, "costs": {"normal": COST_NORMAL, "uncensored": COST_UNCENSORED}, "daily_quota": TOKEN_DAILY_QUOTA, } @app.post("/get_config") async def get_config(): b1 = await _check_worker("backend1") b2 = await _check_worker("backend2") return { "version": VERSION, "architecture": "Docker relay + ZeroGPU workers", "modes": { "normal": {"label": "Normal (HF Official)", "cost": COST_NORMAL, "nsfw_filter": "full", "blur": False, "worker": "backend1"}, "uncensored": {"label": "Uncensored (CivitAI)", "cost": COST_UNCENSORED, "nsfw_filter": "hard-blocks only", "blur": True, "passkey_required": True, "worker": "backend2"}, }, "token_economy": {"daily_quota": TOKEN_DAILY_QUOTA, "cost_normal": COST_NORMAL, "cost_uncensored": COST_UNCENSORED, "reset": "daily (24h)"}, "resolution_policy": {"max_pixels": MAX_PIXELS, "min_pixels": MIN_PIXELS, "min_side": MIN_SIDE, "max_side": MAX_SIDE, "max_aspect": MAX_ASPECT, "step": RESOLUTION_STEP}, "batch_policy": {"min": MIN_BATCH, "max": MAX_BATCH}, "workers": {"backend1": b1, "backend2": b2}, } @app.post("/health") async def health(): b1 = await _check_worker("backend1") b2 = await _check_worker("backend2") with _spy_lock: spy_count = len(_spy_log) return { "status": "healthy", "version": VERSION, "architecture": "Docker relay + ZeroGPU workers", "backend1": b1, "backend2": b2, "hf_pool": pool_manager.get_pool_status(), "token_economy": user_tokens.stats(), "generation_stats": { "total": _generation_stats["total_generations"], "normal": _generation_stats["normal_gens"], "uncensored": _generation_stats["uncensored_gens"], "safety_blocks": _generation_stats["safety_blocks"], "tokens_consumed": _generation_stats["tokens_consumed"], "relay_failures": _generation_stats["relay_failures"], "legacy_blocked": _generation_stats["legacy_blocked"], }, "spy_entries": spy_count, } @app.post("/get_broadcast") async def get_broadcast(): b = broadcast_get() if not b.get("active"): return {"active": False} return b @app.post("/diagnose") async def diagnose(request: Request): headers_dict = {} try: for key, value in request.headers.items(): k = key.lower() if k in ("authorization", "cookie") and len(value) > 20: headers_dict[key] = value[:20] + "..." else: headers_dict[key] = value except Exception as e: headers_dict["_error"] = str(e) client_ip, _ = _extract_client_info(request) b1 = await _check_worker("backend1") b2 = await _check_worker("backend2") return { "all_headers": headers_dict, "client_ip": client_ip, "user": user_tokens.get_user(client_ip), "pool_status": pool_manager.get_pool_status(), "broadcast": broadcast_get(), "workers": {"backend1": b1, "backend2": b2}, } @app.post("/get_stats") async def get_stats(): stats = dict(_generation_stats) stats["uptime_seconds"] = round(time.time() - _generation_stats["started_at"], 1) stats["avg_time_seconds"] = (round(stats["total_time_seconds"] / stats["total_generations"], 2) if stats["total_generations"] > 0 else 0) return stats @app.post("/reset_stats") async def reset_stats(): global _generation_stats _generation_stats = { "total_generations": 0, "total_time_seconds": 0.0, "last_generation_time": None, "last_resolution": None, "last_seed": None, "started_at": time.time(), "errors": 0, "safety_blocks": 0, "normal_gens": 0, "uncensored_gens": 0, "tokens_consumed": 0, "relay_failures": 0, "legacy_blocked": 0, } return {"success": True, "message": "Stats reset"} @app.post("/unlock_image") async def unlock_image(request: Request): body = await request.json() passkey = body.get("passkey", "") entry_id = body.get("entry_id") if str(passkey).strip() != UNCENSORED_PASSKEY: return {"success": False, "error": "Invalid passkey"} try: entry_id = int(entry_id) except (ValueError, TypeError): return {"success": False, "error": "Bad entry id"} with _spy_lock: for e in _spy_log: if e["id"] == entry_id: imgs = [{"seed": im.get("seed"), "full": im.get("full")} for im in e.get("images", []) if im.get("full")] if not imgs: return {"success": False, "error": "Image expired (15-min TTL). Please regenerate."} return {"success": True, "entry_id": entry_id, "images": imgs} return {"success": False, "error": "Image expired from buffer. Please regenerate."} # ╔══════════════════════════════════════════════════════════════╗ # ║ [11] VALIDATION HELPERS ║ # ╚══════════════════════════════════════════════════════════════╝ def _validate_resolution(width: int, height: int): if width % RESOLUTION_STEP != 0 or height % RESOLUTION_STEP != 0: return f"Dimensions must be multiples of {RESOLUTION_STEP}" if not (MIN_SIDE <= width <= MAX_SIDE) or not (MIN_SIDE <= height <= MAX_SIDE): return f"Each side must be between {MIN_SIDE} and {MAX_SIDE}" pixels = width * height if pixels > MAX_PIXELS: return f"REJECTED: {pixels:,} px exceeds 1024×1024 spec (max {MAX_PIXELS:,})" if pixels < MIN_PIXELS: return f"REJECTED: {pixels:,} px below minimum ({MIN_PIXELS:,})" ratio = max(width, height) / max(1, min(width, height)) if ratio > MAX_ASPECT: return f"REJECTED: aspect ratio {ratio:.2f}:1 too extreme (max {MAX_ASPECT}:1)" return None # ╔══════════════════════════════════════════════════════════════╗ # ║ [12] /generate — MAIN ROUTER ║ # ╚══════════════════════════════════════════════════════════════╝ @app.post("/generate") async def generate(request: Request): body = await request.json() # Accept both array and object format if isinstance(body, dict) and "data" in body: payload = body["data"] elif isinstance(body, list): payload = body else: payload = [ body.get("prompt", ""), body.get("negative_prompt", ""), body.get("width", 1024), body.get("height", 1024), body.get("seed", -1), body.get("batch", 1), body.get("mode", "normal"), ] # Backward compat: detect legacy v1.5.1 (6 args dengan hf_token di posisi ke-6) legacy_mode = False if len(payload) == 6: sixth = payload[5] if isinstance(sixth, str) and sixth.startswith("hf_"): legacy_mode = True if legacy_mode: prompt, neg, width, height, seed, _ = payload batch, mode = 1, "normal" _log_throttled("legacy_generate", "[Relay] Legacy frontend (v1.5.1) — auto-adapting to normal mode", 60) _generation_stats["legacy_blocked"] += 1 else: prompt = payload[0] if len(payload) > 0 else "" neg = payload[1] if len(payload) > 1 else "" width = payload[2] if len(payload) > 2 else 1024 height = payload[3] if len(payload) > 3 else 1024 seed = payload[4] if len(payload) > 4 else -1 batch = payload[5] if len(payload) > 5 else 1 mode = payload[6] if len(payload) > 6 else "normal" client_ip, client_ua = _extract_client_info(request) mode = str(mode or "normal").lower().strip() if mode not in ("normal", "uncensored"): mode = "normal" # ── Ban check ── banned, ban_info = user_tokens.is_banned(client_ip) if banned: rej = {"success": False, "error": f"Banned until {ban_info['until']}. Reason: {ban_info['reason']}", "retryable": False, "banned": True, "ban_info": ban_info} _log_to_spy(client_ip, client_ua, -1, prompt, neg, 0, 0, 0, rej, stage="banned", mode=mode) return rej # ── Batch validation ── try: batch_val = int(batch) except (ValueError, TypeError): batch_val = 1 if batch_val < MIN_BATCH or batch_val > MAX_BATCH: batch_val = 1 # ── Resolution ── try: width = int(width) height = int(height) seed = int(seed) if seed is not None else -1 except (ValueError, TypeError) as e: return {"success": False, "error": f"Bad param: {e}", "retryable": False} res_err = _validate_resolution(width, height) if res_err: rej = {"success": False, "error": res_err, "retryable": False, "resolution_rejected": True} _log_to_spy(client_ip, client_ua, -1, prompt, neg, width, height, batch_val, rej, stage="resolution_rejected", mode=mode) return rej # ── User tokens ── user = user_tokens.get_user(client_ip) is_admin = bool(user.get("is_admin")) cost_each = COST_UNCENSORED if mode == "uncensored" else COST_NORMAL total_cost = cost_each * batch_val if not is_admin and user.get("tokens", 0) < total_cost: rej = {"success": False, "error": f"Insufficient tokens. Need {total_cost}, have {user.get('tokens', 0)}. Resets daily.", "retryable": False, "insufficient_tokens": True, "tokens_have": user.get("tokens", 0), "tokens_needed": total_cost} _log_to_spy(client_ip, client_ua, -1, prompt, neg, width, height, batch_val, rej, stage="insufficient_tokens", mode=mode) return rej # ── Safety filter ── is_safe, safety_reason, safety_score, safety_hits = safety_filter.check( prompt, neg, hard_only=(mode == "uncensored")) if not is_safe: _generation_stats["safety_blocks"] += 1 print(f"[Relay] 🛡️ SAFETY BLOCKED ({mode}) from {client_ip}", flush=True) rej = {"success": False, "error": safety_reason, "retryable": False, "safety_blocked": True, "safety_score": safety_score, "safety_hits": safety_hits, "_safety_score": safety_score} _log_to_spy(client_ip, client_ua, -1, prompt, neg, width, height, batch_val, rej, stage="safety_filtered", mode=mode) return rej # ── Consume user tokens ── remaining = 999999 if not is_admin: ok, remaining = user_tokens.consume(client_ip, total_cost) if not ok: return {"success": False, "error": f"Insufficient tokens (race). Need {total_cost}.", "retryable": True, "insufficient_tokens": True} if mode == "uncensored": user_tokens.mark_uncensored(client_ip) def _refund_and_respond(err_msg, stage, extra=None): if not is_admin: try: user_tokens.grant_tokens(client_ip, total_cost, by="refund") except Exception: pass fail = {"success": False, "error": err_msg, "retryable": True, "tokens_refunded": total_cost} if extra: fail.update(extra) _log_to_spy(client_ip, client_ua, -1, prompt, neg, width, height, batch_val, fail, stage=stage, mode=mode) return fail # ── Pick worker ── worker_key = "backend2" if mode == "uncensored" else "backend1" worker_label = "MPS (uncensored)" if mode == "uncensored" else "ZeroGPU (normal)" # ── Get HF token for ZeroGPU auth ── hf_idx, hf_token = pool_manager.get_next_token() if hf_idx is None: status = pool_manager.get_pool_status() return _refund_and_respond( f"All HF tokens resting. Pool: {status['active']}/{status['total_tokens']}.", "tokens_exhausted", {"pool_status": status}) if seed < 0: seed = random.randint(0, 2**32 - 1) # ── Relay to worker ── print(f"[Relay] 🚀 {worker_label} | {width}x{height} b={batch_val} | {client_ip} | token #{hf_idx}", flush=True) t0 = time.time() try: result = await _relay_to_worker(worker_key, prompt, neg, width, height, seed, batch_val, hf_token) except Exception as e: emsg = str(e) _generation_stats["relay_failures"] += 1 with _ws_lock: _worker_status[worker_key]["fail"] += 1 if _check_quota_error(emsg): print(f"[Relay] ⚠️ ZeroGPU quota exhausted for HF token #{hf_idx}", flush=True) pool_manager.mark_exhausted(hf_idx, emsg) if not is_admin: try: user_tokens.grant_tokens(client_ip, total_cost, by="refund") except Exception: pass fail = {"success": False, "error": f"GPU quota exhausted (token #{hf_idx}). Auto-rotating.", "retryable": True, "quota_exhausted": True, "try_different_token": True, "pool_status": pool_manager.get_pool_status(), "tokens_refunded": total_cost} _log_to_spy(client_ip, client_ua, hf_idx, prompt, neg, width, height, batch_val, fail, stage="quota_exhausted", mode=mode, worker_key=worker_key) return fail print(f"[Relay] ❌ Worker failed: {emsg[:180]}", flush=True) return _refund_and_respond(f"{worker_label} worker error: {emsg[:300]}", "error") if not isinstance(result, dict) or not result.get("success"): err = (result or {}).get("error", f"{worker_key} returned failure") if _check_quota_error(err): pool_manager.mark_exhausted(hf_idx, err) fail = {"success": False, "error": f"GPU quota exhausted: {err}", "retryable": True, "quota_exhausted": True, "try_different_token": True, "pool_status": pool_manager.get_pool_status(), "tokens_refunded": total_cost} _log_to_spy(client_ip, client_ua, hf_idx, prompt, neg, width, height, batch_val, fail, stage="quota_exhausted", mode=mode, worker_key=worker_key) return fail return _refund_and_respond(str(err), "error") with _ws_lock: _worker_status[worker_key]["ok"] += 1 elapsed = time.time() - t0 # ── Enrich metadata ── md = result.setdefault("metadata", {}) md["mode"] = mode md["worker"] = worker_key md["model"] = ("ZimageTurbo-CivitAI-3025713" if mode == "uncensored" else "ZimageTurbo-HF-Official") md["tokens_consumed"] = total_cost md["duration"] = round(elapsed, 2) result["_safety_score"] = safety_score result["_user_is_admin"] = is_admin _update_stats(width, height, elapsed, md.get("seed", seed), mode=mode, tokens=0 if is_admin else total_cost) # ── Spy log ── entry_id = _log_to_spy(client_ip, client_ua, hf_idx, prompt, neg, width, height, batch_val, result, stage="generate", mode=mode, worker_key=worker_key) # ── Blur if uncensored & not admin ── if mode == "uncensored" and not is_admin: for img in result.get("images", []): blurred = _blur_data_url(img.get("data", "")) if blurred is None: return _refund_and_respond("Blur failed. Tokens refunded.", "error") img["data"] = blurred img["is_blurred"] = True md["is_blurred"] = True md["blur_passkey_required"] = True md["unlock_id"] = entry_id else: md["is_blurred"] = False md["unlock_id"] = None result["tokens_remaining"] = remaining print(f"[Relay] ✅ OK | {elapsed:.2f}s | entry_id={entry_id} | blurred={mode=='uncensored' and not is_admin}", flush=True) return result # ╔══════════════════════════════════════════════════════════════╗ # ║ [13] ADMIN ENDPOINTS ║ # ╚══════════════════════════════════════════════════════════════╝ _admin_action_log = [] _admin_log_lock = threading.Lock() def _log_admin_action(action: str, details: Dict): with _admin_log_lock: _admin_action_log.append({"time": _utc_str(), "action": action, "details": details}) if len(_admin_action_log) > 200: del _admin_action_log[:len(_admin_action_log) - 200] print(f"[Admin] {action}: {details}", flush=True) @app.post("/admin_feed") async def admin_feed(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} return { "success": True, "entries": _build_feed_payload(0), "incremental": False, "latest_id": (_spy_counter if _spy_log else 0), "stats": _generation_stats, "user_stats": user_tokens.stats(), "pool": pool_manager.get_pool_status(), "workers": {"backend1": await _check_worker("backend1"), "backend2": await _check_worker("backend2")}, "broadcast": broadcast_get(), "server_time": _utc_str(), } @app.post("/admin_feed_since") async def admin_feed_since(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} try: since = int(body.get("since_id", 0)) except (ValueError, TypeError): since = 0 return { "success": True, "entries": _build_feed_payload(since), "incremental": since > 0, "latest_id": (_spy_counter if _spy_log else 0), "stats": _generation_stats, "user_stats": user_tokens.stats(), "pool": pool_manager.get_pool_status(), "workers": {"backend1": await _check_worker("backend1"), "backend2": await _check_worker("backend2")}, "broadcast": broadcast_get(), "server_time": _utc_str(), } @app.post("/admin_users") async def admin_users(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} return {"success": True, "users": user_tokens.list_users(limit=300)} @app.post("/admin_ban") async def admin_ban(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} ip = body.get("ip", "") try: hours = float(body.get("hours", 0)) except (ValueError, TypeError): return {"success": False, "error": "hours must be a number"} if hours <= 0: return {"success": False, "error": "hours must be > 0"} reason = body.get("reason", "") until = user_tokens.ban_ip(str(ip), hours, str(reason or "No reason"), by="admin_panel") _log_admin_action("ban", {"ip": ip, "hours": hours, "reason": reason, "until": until}) return {"success": True, "ip": ip, "banned_until": until, "reason": reason} @app.post("/admin_unban") async def admin_unban(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} ip = body.get("ip", "") ok = user_tokens.unban_ip(str(ip)) _log_admin_action("unban", {"ip": ip}) return {"success": ok, "ip": ip} @app.post("/admin_grant") async def admin_grant(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} ip = body.get("ip", "") try: amount = int(body.get("amount", 0)) except (ValueError, TypeError): return {"success": False, "error": "amount must be integer"} bal = user_tokens.grant_tokens(str(ip), amount, by="admin_panel") _log_admin_action("grant", {"ip": ip, "amount": amount, "new_balance": bal}) return {"success": True, "ip": ip, "amount": amount, "new_balance": bal} @app.post("/admin_promote") async def admin_promote(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} ip = body.get("ip", "") user_tokens.promote_admin(str(ip), by="admin_panel") _log_admin_action("promote", {"ip": ip}) return {"success": True, "ip": ip, "is_admin": True} @app.post("/admin_demote") async def admin_demote(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} ip = body.get("ip", "") user_tokens.demote_admin(str(ip)) _log_admin_action("demote", {"ip": ip}) return {"success": True, "ip": ip, "is_admin": False} @app.post("/admin_broadcast") async def admin_broadcast(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} message = body.get("message", "") btype = body.get("type", "info") if not message or not str(message).strip(): result = broadcast_clear() _log_admin_action("broadcast_clear", {}) else: result = broadcast_set(str(message), str(btype or "info"), by="admin_panel") _log_admin_action("broadcast_set", {"message": str(message)[:200], "type": btype}) return {"success": True, "broadcast": result} @app.post("/admin_image") async def admin_image(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} try: entry_id = int(body.get("entry_id", 0)) except (ValueError, TypeError): return {"success": False, "error": "Bad entry id"} with _spy_lock: for e in _spy_log: if e["id"] == entry_id: imgs = [{"seed": im.get("seed"), "full": im.get("full")} for im in e.get("images", [])] if all(not im.get("full") for im in imgs): return {"success": False, "error": f"Full image expired ({SPY_FULL_TTL_SECONDS//60}-min TTL)."} return {"success": True, "entry_id": entry_id, "images": imgs, "entry": {k: v for k, v in e.items() if k not in ("images", "_full_expires")}} return {"success": False, "error": "Entry not found (rotated out)"} @app.post("/admin_log") async def admin_log(request: Request): body = await request.json() key = body.get("admin_key", "") if isinstance(body, dict) else "" if not _admin_auth(key): return {"success": False, "error": "Invalid admin key"} with _admin_log_lock: return {"success": True, "log": list(reversed(_admin_action_log))} # ╔══════════════════════════════════════════════════════════════╗ # ║ [14] STARTUP ║ # ╚══════════════════════════════════════════════════════════════╝ @app.on_event("startup") async def startup(): print("\n" + "=" * 60, flush=True) print(f" ZERO COST v3.0.0 RELAY (Docker + FastAPI) - STARTING", flush=True) print("=" * 60, flush=True) print(f" Storage: {STORAGE_DIR}", flush=True) print(f" Backend 1: {BACKEND1_URL}", flush=True) print(f" Backend 2: {BACKEND2_URL}", flush=True) print(f" Worker Secret: set ({len(WORKER_SECRET)} chars)", flush=True) print(f" HF Pool: {len(POOL_85)} tokens", flush=True) print(f" Token economy: {TOKEN_DAILY_QUOTA}/user/day | normal={COST_NORMAL} | uncensored={COST_UNCENSORED}", flush=True) print(f" Safety: hard-block ALL modes; soft-score normal only", flush=True) print(f" Spy log: {SPY_MAX_ENTRIES} entries, TTL {SPY_FULL_TTL_SECONDS//60} min", flush=True) print(f" Worker timeout: connect={WORKER_CONNECT_TIMEOUT}s, read={WORKER_READ_TIMEOUT}s", flush=True) print("-" * 60, flush=True) print(" Endpoints:", flush=True) print(" POST /generate, /get_token, /pool_status, /user_status", flush=True) print(" POST /get_config, /health, /get_broadcast, /unlock_image", flush=True) print(" POST /diagnose, /get_stats, /reset_stats", flush=True) print(" POST /admin_feed, /admin_users, /admin_ban, etc.", flush=True) print("-" * 60, flush=True) print(" Initial worker health check:", flush=True) b1 = await _check_worker("backend1", force=True) b2 = await _check_worker("backend2", force=True) print(f" Backend 1: {'✅ reachable' if b1['reachable'] else '❌ ' + str(b1['last_error'])}", flush=True) print(f" Backend 2: {'✅ reachable' if b2['reachable'] else '❌ ' + str(b2['last_error'])}", flush=True) print("=" * 60 + "\n", flush=True) @app.on_event("shutdown") async def shutdown(): global _http_client if _http_client and not _http_client.is_closed: await _http_client.aclose() # ╔══════════════════════════════════════════════════════════════╗ # ║ [15] LAUNCH (for local testing) ║ # ╚══════════════════════════════════════════════════════════════╝ if __name__ == "__main__": import uvicorn uvicorn.run( "app:app", host="0.0.0.0", port=7860, workers=4, limit_concurrency=500, timeout_keep_alive=65, )