Spaces:
Running
Running
Update app.py
Browse files
app.py
CHANGED
|
@@ -3,6 +3,135 @@ import os
|
|
| 3 |
from flask import Flask, request, jsonify
|
| 4 |
|
| 5 |
app = Flask(__name__)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 6 |
|
| 7 |
# ู
ุณุงุฑ ุงุฎุชุจุงุฑ ุณุฑูุน
|
| 8 |
@app.get("/")
|
|
|
|
| 3 |
from flask import Flask, request, jsonify
|
| 4 |
|
| 5 |
app = Flask(__name__)
|
| 6 |
+
# โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
|
| 7 |
+
# ู
ุญุฑู ุงุชุตุงู ุฎุงุฑุฌู: Direct ุฃููุงูุ ุซู
Relay (WebSocket) ูุฎุทุฉ B
|
| 8 |
+
# ุฃูุตูู ุจุนุฏ ุงูุงุณุชูุฑุงุฏุงุช ููุจู ุจุฏุก ุฎุฏุงู
ู (RPC/Flask)
|
| 9 |
+
# โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
|
| 10 |
+
import os, time, threading, socket, random, requests, importlib
|
| 11 |
+
MODE = None # "direct" ุฃู "relay"
|
| 12 |
+
CURRENT_SERVER = None
|
| 13 |
+
PORT = None # ู
ููุฐู ุงูู
ุญูู ูุฎุฏู
ุฉ /run ุฃู RPC
|
| 14 |
+
CONNECTED = threading.Event()
|
| 15 |
+
RELAY_SIO = None
|
| 16 |
+
|
| 17 |
+
# โ ุฌูุจ ุงูู
ุฑุดุญูู ู
ู peer_discovery ูู ูู ุฏูุฑุฉ (ุญุชู ูู ุจูููุฏูู
ุฏููุงู
ูููุงู)
|
| 18 |
+
def _load_candidates():
|
| 19 |
+
import peer_discovery as pd
|
| 20 |
+
importlib.reload(pd)
|
| 21 |
+
servers = list(getattr(pd, "CENTRAL_REGISTRY_SERVERS", []))
|
| 22 |
+
# ู
ููุฐู ุงูุฐู ุณุชุดุบู ุนููู ุงูู RPC ู
ุญูููุง (ูููุณ ู
ููุฐ ุงูุณูุฑูุฑ ุงูุนุงู
)
|
| 23 |
+
# ููุถู ุฃู ูููู ู
ุง ูููุฏู peer_discoveryุ ูุฅู ูู
ููุฌุฏ ุงุณุชุฎุฏู
7860/8000/5000 ุงุญุชูุงุท
|
| 24 |
+
ports = list(getattr(pd, "rport", [])) or ["7860", "8000", "5000"]
|
| 25 |
+
get_ip = getattr(pd, "get_local_ip", lambda: "127.0.0.1")
|
| 26 |
+
random.shuffle(servers); random.shuffle(ports)
|
| 27 |
+
return servers, [int(p) for p in ports], get_ip
|
| 28 |
+
|
| 29 |
+
def _is_port_free(p, host="0.0.0.0"):
|
| 30 |
+
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
|
| 31 |
+
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
| 32 |
+
try:
|
| 33 |
+
s.bind((host, int(p)))
|
| 34 |
+
return True
|
| 35 |
+
except OSError:
|
| 36 |
+
return False
|
| 37 |
+
|
| 38 |
+
def _try_direct_register(server, port, get_ip):
|
| 39 |
+
"""
|
| 40 |
+
ูุญุงูู ุงูุชุณุฌูู ุงูู
ุจุงุดุฑ ุนูุฏ ุงูุณูุฑูุฑ: POST /register
|
| 41 |
+
ุฅุฐุง ูุฌุญุ ูุณุชุฎุฏู
ูุฐุง ุงูู
ููุฐ ูุชุดุบูู ุฎุงุฏู
ูุง ุงูู
ุญูู ูุงุญููุง.
|
| 42 |
+
"""
|
| 43 |
+
info = {
|
| 44 |
+
"node_id": os.getenv("NODE_ID", socket.gethostname()),
|
| 45 |
+
"ip": get_ip(),
|
| 46 |
+
"port": int(port),
|
| 47 |
+
}
|
| 48 |
+
r = requests.post(f"{server}/register", json=info, timeout=6)
|
| 49 |
+
r.raise_for_status() # ูุฑู
ู ุงุณุชุซูุงุก ูู ูุดู
|
| 50 |
+
return True
|
| 51 |
+
|
| 52 |
+
def _start_relay_client(server):
|
| 53 |
+
"""
|
| 54 |
+
ูุถุน Relay: ููุงุฉ Socket.IO ุฏุงุฆู
ุฉ ู
ุน ุงูุณูุฑูุฑ (ุนูู 443 ุบุงูุจูุง).
|
| 55 |
+
ูุง ุชุญุชุงุฌ ุฃู ู
ููุฐ inbound. ุฅุฐุง ุงูุณูุฑูุฑ ูุฏุนู
relayุ ููุนุชุจุฑ ุงุชุตุงู ุฎุงุฑุฌู ูุงุฌุญ.
|
| 56 |
+
"""
|
| 57 |
+
global RELAY_SIO
|
| 58 |
+
import socketio # python-socketio client
|
| 59 |
+
sio = socketio.Client(reconnection=True, reconnection_attempts=0)
|
| 60 |
+
|
| 61 |
+
node_id = os.getenv("NODE_ID", socket.gethostname())
|
| 62 |
+
|
| 63 |
+
@sio.event
|
| 64 |
+
def connect():
|
| 65 |
+
print(f"๐ Relay connected to {server}")
|
| 66 |
+
sio.emit("register_node", {"node_id": node_id})
|
| 67 |
+
# ุงุนุชุจุฑ ุงูุงุชุตุงู ูุงุฌุญ ุฎุงุฑุฌููุง
|
| 68 |
+
# ููุจูู PORT ูุงุณุชุฎุฏุงู
LAN/RPC ู
ุญูู (ูู ุงุญุชุฌุชู). ุชุณุชุทูุน ุถุจุทู ูุงุญููุง.
|
| 69 |
+
|
| 70 |
+
@sio.event
|
| 71 |
+
def disconnect():
|
| 72 |
+
print("โ ๏ธ Relay disconnected; will auto-reconnectโฆ")
|
| 73 |
+
|
| 74 |
+
# ุฃู
ุซูุฉ ุฃุญุฏุงุซ (ุนุฏูู ุนูู ู
ุฒุงุฌู):
|
| 75 |
+
@sio.on("ping_node")
|
| 76 |
+
def _on_ping(msg):
|
| 77 |
+
sio.emit("pong_node", {"node_id": node_id, "t": time.time()})
|
| 78 |
+
|
| 79 |
+
try:
|
| 80 |
+
# Socket.IO URL โ ุนุงุฏุฉ ููุณ ุงูุฏูู
ูู ุนูู /socket.io
|
| 81 |
+
sio.connect(server, transports=["websocket", "polling"], wait_timeout=10)
|
| 82 |
+
RELAY_SIO = sio
|
| 83 |
+
return True
|
| 84 |
+
except Exception as e:
|
| 85 |
+
print(f"โ Relay connect failed: {e}")
|
| 86 |
+
return False
|
| 87 |
+
|
| 88 |
+
def _connect_until_success():
|
| 89 |
+
global MODE, CURRENT_SERVER, PORT
|
| 90 |
+
backoff = 1
|
| 91 |
+
while True:
|
| 92 |
+
servers, ports, get_ip = _load_candidates()
|
| 93 |
+
|
| 94 |
+
# 1) ุฌุฑูุจ DIRECT ุนูู ูู ุงูุณูุฑูุฑุงุช ู
ุน ู
ุฌู
ูุนุฉ ู
ูุงูุฐู
|
| 95 |
+
for s in servers:
|
| 96 |
+
for p in ports:
|
| 97 |
+
# ุชุฃูุฏ ุฃู ุงูู
ููุฐ ู
ุชุงุญ ู
ุญูููุง ูุชุดุบูู ุฎุงุฏู
ู ูุงุญููุง
|
| 98 |
+
if not _is_port_free(p):
|
| 99 |
+
continue
|
| 100 |
+
try:
|
| 101 |
+
_try_direct_register(s, p, get_ip)
|
| 102 |
+
MODE, CURRENT_SERVER, PORT = "direct", s, int(p)
|
| 103 |
+
print(f"โ
DIRECT connected: {s} (node port {PORT})")
|
| 104 |
+
CONNECTED.set()
|
| 105 |
+
return
|
| 106 |
+
except Exception as e:
|
| 107 |
+
print(f"โ ๏ธ DIRECT failed on {s} / {p}: {e}")
|
| 108 |
+
|
| 109 |
+
# 2) ุฅุฐุง ูุดูุช ุงูู
ุจุงุดุฑุฉ ุนูู ุงููู โ ุฌุฑูุจ RELAY ููู ุณูุฑูุฑ ุจุงูุชุฑุชูุจ
|
| 110 |
+
for s in servers:
|
| 111 |
+
ok = _start_relay_client(s)
|
| 112 |
+
if ok:
|
| 113 |
+
MODE, CURRENT_SERVER = "relay", s
|
| 114 |
+
# PORT ุงุฎุชูุงุฑู ููุง (ูู LAN ููุท). ูู ุนูุฏู RPC ู
ุญููู ุดุบููู ุนูู ุฃู ู
ููุฐ ู
ุชุงุญ.
|
| 115 |
+
if ports:
|
| 116 |
+
# ุงุฎุชุฑ ุฃูู ู
ููุฐ ู
ุชุงุญ ู
ุญูููุง ููุชุดุบูู ุงูุฏุงุฎูู
|
| 117 |
+
for p in ports:
|
| 118 |
+
if _is_port_free(p):
|
| 119 |
+
PORT = int(p)
|
| 120 |
+
break
|
| 121 |
+
print(f"โ
RELAY connected via {s} (no inbound port needed)")
|
| 122 |
+
CONNECTED.set()
|
| 123 |
+
return
|
| 124 |
+
|
| 125 |
+
# ูุง ูุฌุงุญ ุจุนุฏ ุงูุชุฌุฑุจุชูู โ ุงูุชุธุฑ ูุฒูุฏ ุงูู
ููุฉ ูุฃุนุฏ ุงูู
ุญุงููุฉ
|
| 126 |
+
time.sleep(backoff)
|
| 127 |
+
backoff = min(backoff * 2, 30)
|
| 128 |
+
|
| 129 |
+
def start_connect_loop():
|
| 130 |
+
threading.Thread(target=_connect_until_success, daemon=True).start()
|
| 131 |
+
|
| 132 |
+
# ุดุบูู ุญููุฉ ุงูุงุชุตุงู ู
ุจูุฑูุง:
|
| 133 |
+
start_connect_loop()
|
| 134 |
+
# โโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโโ
|
| 135 |
|
| 136 |
# ู
ุณุงุฑ ุงุฎุชุจุงุฑ ุณุฑูุน
|
| 137 |
@app.get("/")
|