Merge pull request 'feat: authenticated DNS RRSet upsert' (#57) from feat/dns-rrset-upsert-20260907 into main
This commit is contained in:
commit
f24809a22a
1 changed files with 188 additions and 2 deletions
190
app.py
190
app.py
|
|
@ -12,14 +12,17 @@ from contextlib import asynccontextmanager
|
||||||
from contextvars import ContextVar
|
from contextvars import ContextVar
|
||||||
|
|
||||||
log = logging.getLogger("butler")
|
log = logging.getLogger("butler")
|
||||||
VERSION = "2.4.1"
|
VERSION = "2.4.6"
|
||||||
|
|
||||||
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
|
API_DIR = os.environ.get("API_KEY_DIR", "/data/api")
|
||||||
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
|
VAULT_CACHE_DIR = os.environ.get("VAULT_CACHE_DIR", "/data/vault-cache")
|
||||||
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
|
BUTLER_TOKEN = os.environ.get("BUTLER_TOKEN", "")
|
||||||
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
|
CONFIG_PATH = os.environ.get("BUTLER_CONFIG", "/data/butler.yaml")
|
||||||
UI_PATH = os.environ.get("BUTLER_UI_PATH", os.path.join(os.path.dirname(__file__), "ui.html"))
|
UI_PATH = os.environ.get("BUTLER_UI_PATH", os.path.join(os.path.dirname(__file__), "ui.html"))
|
||||||
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get("MEDIA_HANDOFF_ALLOWED_NETWORKS", "10.2.1.119/32")
|
MEDIA_HANDOFF_ALLOWED_NETWORKS = os.environ.get(
|
||||||
|
"MEDIA_HANDOFF_ALLOWED_NETWORKS",
|
||||||
|
"10.2.1.119/32,10.5.85.12/32",
|
||||||
|
)
|
||||||
|
|
||||||
# --- Config loading ---
|
# --- Config loading ---
|
||||||
|
|
||||||
|
|
@ -1571,6 +1574,13 @@ class ProxyRouteRequest(BaseModel):
|
||||||
dns_token: str | None = None
|
dns_token: str | None = None
|
||||||
|
|
||||||
|
|
||||||
|
class DnsRrsetUpsertRequest(BaseModel):
|
||||||
|
record_type: Literal["A", "AAAA", "CNAME"] = "A"
|
||||||
|
value: str
|
||||||
|
ttl: int = Field(default=300, ge=60, le=86400)
|
||||||
|
comment: str = Field(default="Managed by Homelab Butler", max_length=200)
|
||||||
|
|
||||||
|
|
||||||
def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]:
|
def _remote_python(script: str, timeout: int = 30) -> tuple[int, str, str]:
|
||||||
encoded = base64.b64encode(script.encode()).decode()
|
encoded = base64.b64encode(script.encode()).decode()
|
||||||
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
|
return _ssh(VPS_SSH, f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"', timeout=timeout)
|
||||||
|
|
@ -3237,6 +3247,53 @@ async def _hetzner_zone_and_rrsets(zone_name: str):
|
||||||
return zone, rr_response.json().get("rrsets", []), headers
|
return zone, rr_response.json().get("rrsets", []), headers
|
||||||
|
|
||||||
|
|
||||||
|
@app.put("/dns/rrset/{zone_name}/{record_name}")
|
||||||
|
async def dns_rrset_upsert(zone_name: str, record_name: str, req: DnsRrsetUpsertRequest, _=Depends(_verify)):
|
||||||
|
zone_name = zone_name.strip().lower().rstrip(".")
|
||||||
|
hostname = _validate_managed_hostname(f"{record_name}.{zone_name}")
|
||||||
|
record_name = hostname[: -(len(zone_name) + 1)]
|
||||||
|
record_type = req.record_type.upper()
|
||||||
|
value = req.value.strip().rstrip("." if record_type == "CNAME" else "")
|
||||||
|
try:
|
||||||
|
if record_type == "A" and ipaddress.ip_address(value).version != 4:
|
||||||
|
raise ValueError
|
||||||
|
if record_type == "AAAA" and ipaddress.ip_address(value).version != 6:
|
||||||
|
raise ValueError
|
||||||
|
if record_type == "CNAME":
|
||||||
|
_validate_managed_hostname(value)
|
||||||
|
except ValueError as exc:
|
||||||
|
raise HTTPException(400, f"Invalid {record_type} record value") from exc
|
||||||
|
zone, rrsets, headers = await _hetzner_zone_and_rrsets(zone_name)
|
||||||
|
existing = next((item for item in rrsets if item.get("name") == record_name and item.get("type") == record_type), None)
|
||||||
|
payload = {
|
||||||
|
"name": record_name,
|
||||||
|
"type": record_type,
|
||||||
|
"ttl": req.ttl,
|
||||||
|
"records": [{"value": value, "comment": req.comment}],
|
||||||
|
}
|
||||||
|
async with httpx.AsyncClient(timeout=30) as client:
|
||||||
|
if existing:
|
||||||
|
response = await client.put(
|
||||||
|
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets/{record_name}/{record_type}",
|
||||||
|
headers=headers,
|
||||||
|
json=payload,
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
response = await client.post(
|
||||||
|
f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets",
|
||||||
|
headers=headers,
|
||||||
|
json=payload,
|
||||||
|
)
|
||||||
|
if response.status_code not in {200, 201}:
|
||||||
|
raise HTTPException(502, f"Hetzner RRSet upsert failed: HTTP {response.status_code}")
|
||||||
|
verify = await client.get(f"{HETZNER_DNS_API}/zones/{zone['id']}/rrsets", headers=headers)
|
||||||
|
current = [item for item in verify.json().get("rrsets", []) if item.get("name") == record_name and item.get("type") == record_type] if verify.status_code == 200 else []
|
||||||
|
if not current or not any(record.get("value") == value for record in current[0].get("records", [])):
|
||||||
|
raise HTTPException(502, "RRSet read-back does not contain requested value")
|
||||||
|
_audit(f"/dns/rrset/{zone_name}/{record_name}", "PUT", 200, f"type={record_type} value={value}")
|
||||||
|
return {"zone": zone_name, "zone_id": zone.get("id"), "name": record_name, "created": existing is None, "rrset": current[0]}
|
||||||
|
|
||||||
|
|
||||||
@app.get("/dns/rrset/{zone_name}/{record_name}")
|
@app.get("/dns/rrset/{zone_name}/{record_name}")
|
||||||
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
|
async def dns_rrset_get(zone_name: str, record_name: str, _=Depends(_verify)):
|
||||||
zone_name = zone_name.strip().lower().rstrip(".")
|
zone_name = zone_name.strip().lower().rstrip(".")
|
||||||
|
|
@ -3902,6 +3959,133 @@ async def tts_health(_=Depends(_verify)):
|
||||||
return results
|
return results
|
||||||
|
|
||||||
|
|
||||||
|
def _sab_history_command() -> str:
|
||||||
|
container_script = r'''import json, re, subprocess, urllib.parse, urllib.request
|
||||||
|
text = open('/config/sabnzbd.ini', encoding='utf-8', errors='replace').read()
|
||||||
|
match = re.search(r'^api_key\s*=\s*(\S+)', text, re.M)
|
||||||
|
port_match = re.search(r'^port\s*=\s*(\d+)', text, re.M)
|
||||||
|
api_key = match.group(1) if match else ''
|
||||||
|
port = port_match.group(1) if port_match else '7777'
|
||||||
|
params = urllib.parse.urlencode({'mode': 'history', 'limit': 100, 'output': 'json', 'apikey': api_key})
|
||||||
|
with urllib.request.urlopen('http://127.0.0.1:' + port + '/api?' + params, timeout=20) as response:
|
||||||
|
data = json.load(response)
|
||||||
|
slots = data.get('history', {}).get('slots', [])
|
||||||
|
allowed = ('nzo_id', 'name', 'category', 'status', 'script', 'script_line', 'fail_message', 'completed', 'storage', 'path')
|
||||||
|
print(json.dumps([{key: item.get(key) for key in allowed} for item in slots]))
|
||||||
|
'''
|
||||||
|
container_encoded = base64.b64encode(container_script.encode()).decode()
|
||||||
|
host_script = f'''import subprocess, sys
|
||||||
|
command = ["sudo", "-n", "docker", "exec", "sabnzbd", "python3", "-c", "import base64;exec(base64.b64decode('{container_encoded}'))"]
|
||||||
|
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
|
||||||
|
if proc.returncode != 0:
|
||||||
|
command = command[2:]
|
||||||
|
proc = subprocess.run(command, capture_output=True, text=True, timeout=30)
|
||||||
|
if proc.returncode != 0:
|
||||||
|
print((proc.stderr or proc.stdout)[-500:], file=sys.stderr)
|
||||||
|
raise SystemExit(proc.returncode)
|
||||||
|
print(proc.stdout)
|
||||||
|
'''
|
||||||
|
encoded = base64.b64encode(host_script.encode()).decode()
|
||||||
|
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/media/handoff/sab-history")
|
||||||
|
async def media_handoff_sab_history(_=Depends(_verify)):
|
||||||
|
"""Return sanitized SAB history without exposing the SAB API key."""
|
||||||
|
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
|
||||||
|
if not inventory:
|
||||||
|
raise HTTPException(404, "Host sabnzbd not found")
|
||||||
|
target = f'{inventory["user"]}@{inventory["ip"]}'
|
||||||
|
rc, out, err = await asyncio.to_thread(_ssh, target, _sab_history_command(), 40)
|
||||||
|
if rc != 0:
|
||||||
|
raise HTTPException(502, (err or out).strip()[-500:] or "SAB history failed")
|
||||||
|
try:
|
||||||
|
return {"history": json.loads(out)}
|
||||||
|
except json.JSONDecodeError as exc:
|
||||||
|
raise HTTPException(502, "invalid SAB history response") from exc
|
||||||
|
|
||||||
|
|
||||||
|
def _media_handoff_diagnostics_command() -> str:
|
||||||
|
script = r'''import hashlib, json, os, subprocess, urllib.error, urllib.request
|
||||||
|
|
||||||
|
|
||||||
|
def run(args):
|
||||||
|
proc = subprocess.run(args, capture_output=True, text=True, timeout=20)
|
||||||
|
return proc.returncode, proc.stdout.strip(), proc.stderr.strip()
|
||||||
|
|
||||||
|
|
||||||
|
def docker_exec(command):
|
||||||
|
for prefix in (["sudo", "-n", "docker"], ["docker"]):
|
||||||
|
rc, out, err = run(prefix + ["exec", "sabnzbd", "sh", "-lc", command])
|
||||||
|
if rc == 0 or "not found" not in (err + out).lower():
|
||||||
|
return rc, out, err
|
||||||
|
return rc, out, err
|
||||||
|
|
||||||
|
inspect_rc, inspect_out, _ = run(["sudo", "-n", "docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
|
||||||
|
if inspect_rc != 0:
|
||||||
|
inspect_rc, inspect_out, _ = run(["docker", "inspect", "-f", "{{.State.Running}}", "sabnzbd"])
|
||||||
|
py_rc, py_out, _ = docker_exec("command -v python3")
|
||||||
|
curl_rc, curl_out, _ = docker_exec("command -v curl")
|
||||||
|
stat_rc, stat_out, _ = docker_exec("test -f /usenet/scripts/movetdarr.sh && stat -c '%a %s' /usenet/scripts/movetdarr.sh && sha256sum /usenet/scripts/movetdarr.sh")
|
||||||
|
log_rc, log_out, _ = docker_exec("tail -n 200 /usenet/scripts/postprocess.log 2>/dev/null || true")
|
||||||
|
source_rc, source_out, _ = docker_exec("find /usenet/complete -mindepth 2 -maxdepth 2 -type d ! -name '_UNPACK_*' -print 2>/dev/null | sort | tail -100")
|
||||||
|
target_rc, target_out, _ = docker_exec("find /tdarr/complete -mindepth 2 -maxdepth 2 -type d -print 2>/dev/null | sort | tail -100")
|
||||||
|
probe_code = 0
|
||||||
|
probe_body = ""
|
||||||
|
probe = urllib.request.Request(
|
||||||
|
"http://10.5.85.2:8888/media/handoff",
|
||||||
|
method="POST",
|
||||||
|
headers={"Content-Type": "application/json"},
|
||||||
|
data=b'{"action":"status","jobId":"diagnostic-probe"}',
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
with urllib.request.urlopen(probe, timeout=10) as response:
|
||||||
|
probe_code = response.status
|
||||||
|
probe_body = response.read(500).decode(errors="replace")
|
||||||
|
except urllib.error.HTTPError as exc:
|
||||||
|
probe_code = exc.code
|
||||||
|
probe_body = exc.read(500).decode(errors="replace")
|
||||||
|
except Exception as exc:
|
||||||
|
probe_body = type(exc).__name__
|
||||||
|
script_info = {"exists": stat_rc == 0, "executable": False, "sha256": None, "mode": None, "size": None}
|
||||||
|
if stat_rc == 0:
|
||||||
|
lines = stat_out.splitlines()
|
||||||
|
if lines:
|
||||||
|
parts = lines[0].split()
|
||||||
|
if len(parts) >= 2:
|
||||||
|
script_info.update(mode=parts[0], size=int(parts[1]), executable=any(ch in parts[0][-3:] for ch in "1357"))
|
||||||
|
if len(lines) > 1:
|
||||||
|
script_info["sha256"] = lines[1].split()[0]
|
||||||
|
print(json.dumps({
|
||||||
|
"container_running": inspect_rc == 0 and inspect_out == "true",
|
||||||
|
"tools": {"python3": py_rc == 0 and bool(py_out), "curl": curl_rc == 0 and bool(curl_out)},
|
||||||
|
"script": script_info,
|
||||||
|
"caller_probe": {"http_status": probe_code, "body": probe_body},
|
||||||
|
"source_directories": source_out.splitlines() if source_rc == 0 else [],
|
||||||
|
"target_directories": target_out.splitlines() if target_rc == 0 else [],
|
||||||
|
"recent_log": log_out.splitlines()[-200:],
|
||||||
|
}))
|
||||||
|
'''
|
||||||
|
encoded = base64.b64encode(script.encode()).decode()
|
||||||
|
return f'python3 -c "import base64;exec(base64.b64decode(\'{encoded}\'))"'
|
||||||
|
|
||||||
|
|
||||||
|
@app.get("/media/handoff/diagnostics")
|
||||||
|
async def media_handoff_diagnostics(_=Depends(_verify)):
|
||||||
|
"""Read-only diagnostics for the fixed SABnzbd handoff script and caller path."""
|
||||||
|
inventory = await asyncio.to_thread(_find_inventory_host, "sabnzbd")
|
||||||
|
if not inventory:
|
||||||
|
raise HTTPException(404, "Host sabnzbd not found")
|
||||||
|
target = f'{inventory["user"]}@{inventory["ip"]}'
|
||||||
|
rc, out, err = await asyncio.to_thread(_ssh, target, _media_handoff_diagnostics_command(), 30)
|
||||||
|
if rc != 0:
|
||||||
|
raise HTTPException(502, (err or out).strip()[-500:] or "media handoff diagnostics failed")
|
||||||
|
try:
|
||||||
|
return json.loads(out)
|
||||||
|
except json.JSONDecodeError as exc:
|
||||||
|
raise HTTPException(502, "invalid media handoff diagnostic response") from exc
|
||||||
|
|
||||||
|
|
||||||
class MediaHandoffPayload(BaseModel):
|
class MediaHandoffPayload(BaseModel):
|
||||||
action: Literal["start", "moved", "status", "fail"]
|
action: Literal["start", "moved", "status", "fail"]
|
||||||
category: str | None = Field(None, max_length=32)
|
category: str | None = Field(None, max_length=32)
|
||||||
|
|
@ -3933,6 +4117,8 @@ def _media_handoff_caller_allowed(request: Request) -> bool:
|
||||||
async def media_handoff(payload: MediaHandoffPayload, request: Request):
|
async def media_handoff(payload: MediaHandoffPayload, request: Request):
|
||||||
"""Narrow SABnzbd-to-n8n bridge; no generic unauthenticated proxy access."""
|
"""Narrow SABnzbd-to-n8n bridge; no generic unauthenticated proxy access."""
|
||||||
if not _media_handoff_caller_allowed(request):
|
if not _media_handoff_caller_allowed(request):
|
||||||
|
caller = request.client.host if request.client else "unknown"
|
||||||
|
_audit("/media/handoff", "POST", 403, f"caller={caller} action={payload.action}")
|
||||||
raise HTTPException(403, "Media handoff caller is not allowed")
|
raise HTTPException(403, "Media handoff caller is not allowed")
|
||||||
|
|
||||||
allowed_categories = {"serien4k", "serien", "video4k", "video"}
|
allowed_categories = {"serien4k", "serien", "video4k", "video"}
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue