Add HTTP resync endpoint for manual rotation triggering
Build and Push Docker Images / build-and-push (push) Successful in 1m30s
Build and Push Docker Images / build-and-push (push) Successful in 1m30s
Expose port 8080 with `/resync` endpoint to trigger on-demand WireGuard rotations. Add RESYNC_PORT and RESYNC_BIND environment variables. Implement thread-safe rotation locking to prevent concurrent runs. Update README with endpoint usage examples.
This commit is contained in:
1 parent
b20ebbfca6
commit
2dbaa0ebb0
3 files changed
+121
-7
No files matched your search
@@ -8,9 +8,12 @@ import os
|
||||
import re
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime
|
||||
import http.server
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from zoneinfo import ZoneInfo
|
||||
|
||||
from croniter import croniter
|
||||
@@ -19,6 +22,9 @@ import pia
|
||||
|
||||
EXIT_RATE_LIMITED = pia.EXIT_RATE_LIMITED
|
||||
|
||||
_ROTATION_LOCK = threading.Lock()
|
||||
_RESYNC_EVENT = threading.Event()
|
||||
|
||||
CRON_MACROS = {
|
||||
"@yearly": "0 0 1 1 *",
|
||||
"@annually": "0 0 1 1 *",
|
||||
@@ -81,6 +87,68 @@ def unhealthy_cooldown_seconds() -> float:
|
||||
return max(0.0, float(os.environ.get("UNHEALTHY_ROTATE_COOLDOWN", "60")))
|
||||
|
||||
|
||||
def resync_port() -> int | None:
|
||||
raw = os.environ.get("RESYNC_PORT", "8080").strip()
|
||||
if not raw:
|
||||
return None
|
||||
try:
|
||||
port = int(raw)
|
||||
except ValueError as exc:
|
||||
raise SystemExit(f"Invalid RESYNC_PORT '{raw}' (expected integer)") from exc
|
||||
if not (1 <= port <= 65535):
|
||||
raise SystemExit(f"Invalid RESYNC_PORT {port} (expected 1-65535)")
|
||||
return port
|
||||
|
||||
|
||||
def resync_bind_address() -> str:
|
||||
return os.environ.get("RESYNC_BIND", "0.0.0.0").strip() or "0.0.0.0"
|
||||
|
||||
|
||||
class ResyncHandler(http.server.BaseHTTPRequestHandler):
|
||||
def _respond(self, status: int, body: dict[str, Any]) -> None:
|
||||
self.send_response(status)
|
||||
self.send_header("Content-Type", "application/json")
|
||||
self.end_headers()
|
||||
self.wfile.write(json.dumps(body).encode("utf-8"))
|
||||
|
||||
def _path(self) -> str:
|
||||
return self.path.split("?", 1)[0]
|
||||
|
||||
def do_GET(self) -> None: # noqa: N802
|
||||
path = self._path()
|
||||
if path == "/resync":
|
||||
_RESYNC_EVENT.set()
|
||||
self._respond(202, {"status": "accepted", "message": "Resync requested"})
|
||||
elif path == "/":
|
||||
self._respond(
|
||||
200,
|
||||
{
|
||||
"service": "gluetun-pia-wireguard-rotator",
|
||||
"endpoints": ["/resync"],
|
||||
},
|
||||
)
|
||||
else:
|
||||
self._respond(404, {"error": "not found"})
|
||||
|
||||
def do_POST(self) -> None: # noqa: N802
|
||||
path = self._path()
|
||||
if path == "/resync":
|
||||
_RESYNC_EVENT.set()
|
||||
self._respond(202, {"status": "accepted", "message": "Resync requested"})
|
||||
else:
|
||||
self._respond(404, {"error": "not found"})
|
||||
|
||||
def log_message(self, fmt: str, *args: Any) -> None:
|
||||
log(fmt % args)
|
||||
|
||||
|
||||
def start_resync_server(port: int, bind: str) -> None:
|
||||
server = http.server.HTTPServer((bind, port), ResyncHandler)
|
||||
thread = threading.Thread(target=server.serve_forever, daemon=True, name="resync-http")
|
||||
thread.start()
|
||||
log(f"Resync endpoint listening on http://{bind}:{port}/resync")
|
||||
|
||||
|
||||
def gluetun_container_name() -> str:
|
||||
return os.environ.get("GLUETUN_CONTAINER", "m3u-filter-vpn").strip() or "m3u-filter-vpn"
|
||||
|
||||
@@ -142,6 +210,12 @@ def run_rotation(reason: str, extra_args: list[str] | None = None) -> None:
|
||||
raise SystemExit(f"Rotation failed with exit code {result.returncode}")
|
||||
|
||||
|
||||
def run_rotation_locked(reason: str, extra_args: list[str] | None = None) -> None:
|
||||
"""Run a rotation while holding the global lock to avoid concurrent runs."""
|
||||
with _ROTATION_LOCK:
|
||||
run_rotation(reason, extra_args)
|
||||
|
||||
|
||||
def wait_until_healthy_or_timeout(timeout_s: float) -> str | None:
|
||||
"""Sleep up to timeout_s, returning early if Gluetun becomes non-unhealthy."""
|
||||
deadline = time.time() + timeout_s
|
||||
@@ -161,7 +235,7 @@ def handle_unhealthy() -> None:
|
||||
"""Two-step recovery: re-auth best server, then runner-up if still unhealthy."""
|
||||
log(f"Gluetun container {gluetun_container_name()} is unhealthy; starting recovery")
|
||||
|
||||
run_rotation(
|
||||
run_rotation_locked(
|
||||
"Unhealthy recovery step 1: force rotate to best server (new keypair)",
|
||||
["--force"],
|
||||
)
|
||||
@@ -180,7 +254,7 @@ def handle_unhealthy() -> None:
|
||||
else:
|
||||
log("Still unhealthy; step 2 without exclude (no server_ip in state)")
|
||||
|
||||
run_rotation(
|
||||
run_rotation_locked(
|
||||
"Unhealthy recovery step 2: force rotate to runner-up",
|
||||
exclude_args,
|
||||
)
|
||||
@@ -193,30 +267,50 @@ def main() -> None:
|
||||
expr = resolve_cron_expr()
|
||||
tz = zone()
|
||||
interval = health_check_interval()
|
||||
port = resync_port()
|
||||
bind = resync_bind_address()
|
||||
# Fail fast on bad cron / TZ before rotating.
|
||||
nxt = next_run(expr, datetime.now(tz))
|
||||
log(
|
||||
f"Starting gluetun PIA WireGuard rotator "
|
||||
f"(TZ={tz.key}, ROTATE_CRON='{expr}', HEALTH_CHECK_INTERVAL={interval:g}s)"
|
||||
f"(TZ={tz.key}, ROTATE_CRON='{expr}', HEALTH_CHECK_INTERVAL={interval:g}s, "
|
||||
f"RESYNC_PORT={port or 'off'})"
|
||||
)
|
||||
log(f"Next scheduled rotation at {nxt.isoformat(timespec='seconds')}")
|
||||
|
||||
run_rotation("Running rotation on startup", ["--force"])
|
||||
if port:
|
||||
start_resync_server(port, bind)
|
||||
|
||||
run_rotation_locked("Running rotation on startup", ["--force"])
|
||||
nxt = next_run(expr, datetime.now(tz))
|
||||
|
||||
step = min(interval, 1.0)
|
||||
while True:
|
||||
time.sleep(interval)
|
||||
# Sleep in short chunks so a resync request wakes us early.
|
||||
deadline = time.time() + interval
|
||||
while time.time() < deadline and not _RESYNC_EVENT.is_set():
|
||||
remaining = deadline - time.time()
|
||||
time.sleep(min(step, remaining, 1.0))
|
||||
|
||||
now = datetime.now(tz)
|
||||
|
||||
if _RESYNC_EVENT.is_set():
|
||||
_RESYNC_EVENT.clear()
|
||||
run_rotation_locked("Resync requested via HTTP", ["--force"])
|
||||
nxt = next_run(expr, datetime.now(tz))
|
||||
log(f"Next scheduled rotation at {nxt.isoformat(timespec='seconds')}")
|
||||
continue
|
||||
|
||||
status = gluetun_health_status()
|
||||
if status == "unhealthy":
|
||||
handle_unhealthy()
|
||||
_RESYNC_EVENT.clear()
|
||||
nxt = next_run(expr, datetime.now(tz))
|
||||
log(f"Next scheduled rotation at {nxt.isoformat(timespec='seconds')}")
|
||||
continue
|
||||
|
||||
if now >= nxt:
|
||||
run_rotation("Running scheduled rotation", ["--skip-same-server"])
|
||||
run_rotation_locked("Running scheduled rotation", ["--skip-same-server"])
|
||||
nxt = next_run(expr, datetime.now(tz))
|
||||
log(f"Next scheduled rotation at {nxt.isoformat(timespec='seconds')}")
|
||||
|
||||
|
||||
Reference in new issue
Block a user