Compare commits
2 Commits
be37048145
...
bbe57d5f07
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
bbe57d5f07 | ||
|
|
0be35d257b |
@@ -40,28 +40,58 @@ def _start_heartbeat(
|
|||||||
register_payload: dict,
|
register_payload: dict,
|
||||||
interval: float = 20.0,
|
interval: float = 20.0,
|
||||||
) -> threading.Thread:
|
) -> threading.Thread:
|
||||||
"""Daemon thread: sends heartbeats; re-registers automatically on 404 (tracker restart)."""
|
"""Daemon thread: sends heartbeats and re-registers automatically after tracker restarts."""
|
||||||
|
def _reregister() -> bool:
|
||||||
|
nonlocal node_id
|
||||||
|
try:
|
||||||
|
resp = _post_json(f"{tracker_url}/v1/nodes/register", register_payload)
|
||||||
|
node_id = resp.get("node_id", node_id)
|
||||||
|
return True
|
||||||
|
except Exception:
|
||||||
|
return False
|
||||||
|
|
||||||
def _loop() -> None:
|
def _loop() -> None:
|
||||||
nonlocal node_id
|
nonlocal node_id
|
||||||
hb_url = f"{tracker_url}/v1/nodes/{node_id}/heartbeat"
|
hb_url = f"{tracker_url}/v1/nodes/{node_id}/heartbeat"
|
||||||
|
outage_streak = 0 # consecutive intervals where tracker was unreachable
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
time.sleep(interval)
|
time.sleep(interval)
|
||||||
|
|
||||||
|
if outage_streak > 0:
|
||||||
|
# Tracker was down — attempt re-registration first (it may have restarted
|
||||||
|
# with a clean slate and won't know this node).
|
||||||
|
if _reregister():
|
||||||
|
hb_url = f"{tracker_url}/v1/nodes/{node_id}/heartbeat"
|
||||||
|
print(f" [node] re-registered after outage — node ID: {node_id}", flush=True)
|
||||||
|
outage_streak = 0
|
||||||
|
else:
|
||||||
|
outage_streak += 1
|
||||||
|
if outage_streak <= 3 or outage_streak % 10 == 0:
|
||||||
|
print(
|
||||||
|
f" [node] WARNING: tracker still unreachable "
|
||||||
|
f"({outage_streak * interval:.0f}s)",
|
||||||
|
flush=True,
|
||||||
|
)
|
||||||
|
continue
|
||||||
|
|
||||||
try:
|
try:
|
||||||
_post_json(hb_url, {})
|
_post_json(hb_url, {})
|
||||||
except urllib.error.HTTPError as exc:
|
except urllib.error.HTTPError as exc:
|
||||||
if exc.code == 404:
|
if exc.code == 404:
|
||||||
|
# Node was purged (e.g. long gap before restart noticed) — re-register now.
|
||||||
print(" [node] tracker lost registration — re-registering...", flush=True)
|
print(" [node] tracker lost registration — re-registering...", flush=True)
|
||||||
try:
|
if _reregister():
|
||||||
resp = _post_json(f"{tracker_url}/v1/nodes/register", register_payload)
|
|
||||||
node_id = resp.get("node_id", node_id)
|
|
||||||
hb_url = f"{tracker_url}/v1/nodes/{node_id}/heartbeat"
|
hb_url = f"{tracker_url}/v1/nodes/{node_id}/heartbeat"
|
||||||
print(f" [node] re-registered — node ID: {node_id}", flush=True)
|
print(f" [node] re-registered — node ID: {node_id}", flush=True)
|
||||||
except Exception as re_exc:
|
else:
|
||||||
print(f" [node] WARNING: re-registration failed: {re_exc}", flush=True)
|
print(" [node] WARNING: re-registration failed", flush=True)
|
||||||
|
outage_streak = 1
|
||||||
else:
|
else:
|
||||||
print(f" [node] WARNING: heartbeat failed: {exc}", flush=True)
|
print(f" [node] WARNING: heartbeat failed ({exc.code}): {exc}", flush=True)
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
print(f" [node] WARNING: heartbeat failed: {exc}", flush=True)
|
outage_streak = 1
|
||||||
|
print(f" [node] WARNING: tracker unreachable: {exc}", flush=True)
|
||||||
|
|
||||||
t = threading.Thread(target=_loop, daemon=True, name="heartbeat")
|
t = threading.Thread(target=_loop, daemon=True, name="heartbeat")
|
||||||
t.start()
|
t.start()
|
||||||
@@ -98,6 +128,19 @@ def run_startup(
|
|||||||
tracker_url = tracker_url.rstrip("/")
|
tracker_url = tracker_url.rstrip("/")
|
||||||
|
|
||||||
# 1. Hardware detection
|
# 1. Hardware detection
|
||||||
|
if advertise_host is None and host == "0.0.0.0":
|
||||||
|
# socket.getfqdn() returns an mDNS name (.local / .localdomain) that remote
|
||||||
|
# machines on a different OS or subnet often can't resolve. Instead, probe the
|
||||||
|
# outbound IP by opening a UDP socket toward the tracker — no data is sent.
|
||||||
|
try:
|
||||||
|
_tracker_host = urllib.parse.urlparse(tracker_url).hostname or "8.8.8.8"
|
||||||
|
_s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||||
|
_s.connect((_tracker_host, 80))
|
||||||
|
advertise_host = _s.getsockname()[0]
|
||||||
|
_s.close()
|
||||||
|
except Exception:
|
||||||
|
advertise_host = socket.getfqdn()
|
||||||
|
|
||||||
print("Detecting hardware...", flush=True)
|
print("Detecting hardware...", flush=True)
|
||||||
hw = detect_hardware()
|
hw = detect_hardware()
|
||||||
device: str = hw["device"]
|
device: str = hw["device"]
|
||||||
|
|||||||
Reference in New Issue
Block a user