Distributed Computation
Each authorized PC computes locally,
broadcasts its result, discovers other
nodes automatically, and merges their
latest results.
1. Wallet Authorization
ETPY Requirement
Loading...
Token Contract
Loading...
Wallet authorization required.
2. Computation
Computation locked until wallet authorization succeeds.
4. Peers
No peers discovered.
5. Merged Results
No results received.
6. Wallet Proof
No wallet authorization yet.
"""
# ============================================================
# HTTP HANDLER
# ============================================================
class RequestHandler(
BaseHTTPRequestHandler
):
def log_message(
self,
format_string,
*args
):
return
def send_json(
self,
payload,
status_code=200
):
raw = json_bytes(
payload
)
self.send_response(
status_code
)
self.send_header(
"Content-Type",
"application/json; charset=utf-8"
)
self.send_header(
"Content-Length",
str(
len(raw)
)
)
self.send_header(
"Cache-Control",
"no-store"
)
self.end_headers()
self.wfile.write(
raw
)
def read_json_body(self):
try:
length = int(
self.headers.get(
"Content-Length",
"0"
)
)
except Exception:
return None
if length <= 0:
return None
if length > 1024 * 1024:
return None
try:
raw = self.rfile.read(
length
)
return json.loads(
raw.decode(
"utf-8"
)
)
except Exception:
return None
def do_GET(self):
path = urlparse(
self.path
).path
if path == "/status":
self.send_json(
status()
)
return
if path == "/":
config = {
"node_id":
NODE_ID,
"token_address":
ETPY_TOKEN_ADDRESS,
"required_etpy":
REQUIRED_ETPY,
"required_chain_id":
REQUIRED_CHAIN_ID
}
page = PAGE.replace(
"__ETHERS_CDN__",
ETHERS_CDN
)
page = page.replace(
"__CONFIG__",
json.dumps(
config,
separators=(",", ":")
)
)
raw = page.encode(
"utf-8"
)
self.send_response(
200
)
self.send_header(
"Content-Type",
"text/html; charset=utf-8"
)
self.send_header(
"Content-Length",
str(
len(raw)
)
)
self.send_header(
"Cache-Control",
"no-store"
)
self.end_headers()
self.wfile.write(
raw
)
return
if path == "/health":
self.send_json({
"ok": True,
"node_id": NODE_ID,
"token_gate":
token_gate_passed,
"computation":
token_gate_passed
})
return
self.send_json(
{
"error": "Not found"
},
404
)
def do_POST(self):
path = urlparse(
self.path
).path
# ----------------------------------------------------
# Wallet authorization
# ----------------------------------------------------
if path == "/authorize":
payload = (
self.read_json_body()
)
ok, message = (
authorize_wallet(
payload
)
)
if not ok:
self.send_json(
{
"ok": False,
"error": message
},
403
)
return
self.send_json(
{
"ok": True,
"authorized": True,
"node_id": NODE_ID
}
)
print(
"[AUTH]"
" Wallet authorized:"
f" {authorized_wallet}"
)
return
# ----------------------------------------------------
# Wallet revocation
# ----------------------------------------------------
if path == "/revoke":
revoke_wallet()
self.send_json(
{
"ok": True,
"authorized": False
}
)
print(
"[AUTH] Wallet authorization revoked."
)
return
# ----------------------------------------------------
# Peer result
# ----------------------------------------------------
if path == "/peer/result":
payload = (
self.read_json_body()
)
if not isinstance(
payload,
dict
):
self.send_json(
{
"ok": False,
"error":
"Invalid JSON"
},
400
)
return
node_id = str(
payload.get(
"node_id",
""
)
)
register_peer(
node_id,
self.client_address[0],
PEER_PORT,
payload.get(
"round",
0
)
)
accepted = merge_result(
payload
)
self.send_json(
{
"ok": True,
"merged": accepted,
"node_id": NODE_ID
}
)
return
self.send_json(
{
"error":
"Not found"
},
404
)
# ============================================================
# SERVER
# ============================================================
def start_server():
server = None
try:
server = ThreadingHTTPServer(
(
"0.0.0.0",
PEER_PORT
),
RequestHandler
)
server.daemon_threads = True
print(
"[SERVER]"
f" http://127.0.0.1:{PEER_PORT}"
)
print(
"[LAN]"
f" http://{local_ip}:{PEER_PORT}"
)
while running:
server.handle_request()
except Exception as error:
print(
"[SERVER ERROR]",
error
)
finally:
if server:
try:
server.server_close()
except Exception:
pass
# ============================================================
# AUTOMATIC BROWSER
# ============================================================
def open_browser():
time.sleep(
1.2
)
try:
webbrowser.open(
"http://127.0.0.1:"
+ str(PEER_PORT)
+ "/"
)
except Exception:
pass
# ============================================================
# CLEAN SHUTDOWN
# ============================================================
def shutdown():
global running
running = False
print()
print(
"[SYSTEM] Shutting down."
)
# ============================================================
# MAIN
# ============================================================
def main():
global running
print()
print(
"=" * 72
)
print(
" ETPY DISTRIBUTED LOCAL COMPUTATION"
)
print(
"=" * 72
)
print()
print(
"Node ID:"
)
print(
" ",
NODE_ID
)
print()
print(
"Local IP:"
)
print(
" ",
local_ip
)
print()
print(
"Peer server:"
)
print(
" ",
f"{local_ip}:{PEER_PORT}"
)
print()
print(
"Discovery:"
)
print(
" ",
f"UDP {DISCOVERY_PORT}"
)
print()
print(
"ETPY contract:"
)
print(
" ",
ETPY_TOKEN_ADDRESS
)
print()
print(
"Required ETPY:"
)
print(
" ",
f"{REQUIRED_ETPY:,}"
)
print()
if (
ETPY_TOKEN_ADDRESS
==
"0xc8ff7af0b48116922d16a839b16a4d91c57bc3b0"
):
print(
"[WARNING]"
)
print(
"Set ETPY_TOKEN_ADDRESS to the official"
)
print(
"ETPY ERC-20 contract address."
)
print()
print(
"Wallet gate:"
)
print(
" LOCKED until wallet balance and"
)
print(
" off-chain wallet signature are verified."
)
print()
print(
"Computation:"
)
print(
" LOCKED until authorization."
)
print()
print(
"=" * 72
)
print()
# --------------------------------------------------------
# HTTP / peer server
# --------------------------------------------------------
threading.Thread(
target=start_server,
daemon=True
).start()
# --------------------------------------------------------
# UDP listener
# --------------------------------------------------------
threading.Thread(
target=discovery_listener,
daemon=True
).start()
# --------------------------------------------------------
# UDP broadcaster
# --------------------------------------------------------
threading.Thread(
target=discovery_broadcaster,
daemon=True
).start()
# --------------------------------------------------------
# Computation
# --------------------------------------------------------
threading.Thread(
target=computation_worker,
daemon=True
).start()
# --------------------------------------------------------
# Result broadcasting
# --------------------------------------------------------
threading.Thread(
target=result_worker,
daemon=True
).start()
# --------------------------------------------------------
# Browser
# --------------------------------------------------------
if OPEN_BROWSER:
threading.Thread(
target=open_browser,
daemon=True
).start()
# --------------------------------------------------------
# Main monitoring
# --------------------------------------------------------
try:
while running:
time.sleep(
5
)
with state_lock:
gate = (
token_gate_passed
)
peers_count = len(
peers
)
nodes_count = len(
distributed_results
)
round_number = (
computation_round
)
average = (
network_average()
)
wallet = (
authorized_wallet
)
print()
print(
"[STATUS]"
f" gate={'OPEN' if gate else 'LOCKED'}"
f" round={round_number}"
f" peers={peers_count}"
f" nodes={nodes_count}"
f" average={average:.10f}"
)
if wallet:
print(
"[WALLET]"
f" {wallet}"
)
except KeyboardInterrupt:
shutdown()
except Exception as error:
print(
"[SYSTEM ERROR]",
error
)
shutdown()
# ============================================================
# ENTRY POINT
# ============================================================
if __name__ == "__main__":
main()