295 lines
9.7 KiB
Python
295 lines
9.7 KiB
Python
"""
|
|
Webhook management and dispatch for ObsiGate.
|
|
|
|
Webhooks are HTTP POST callbacks triggered on file/directory events.
|
|
Configuration is persisted in data/webhooks.json.
|
|
|
|
Security (BUG-026):
|
|
* target URLs are validated against SSRF (scheme + resolved IP must be public);
|
|
* HTTP is refused unless ``OBSIGATE_WEBHOOK_ALLOW_HTTP=true``;
|
|
* private/loopback/link-local targets are refused unless
|
|
``OBSIGATE_WEBHOOK_ALLOW_PRIVATE=true``;
|
|
* redirects are never followed;
|
|
* signing secrets are **not** stored in the public config file — they live in
|
|
``data/webhook_secrets.json`` (0600) or in an environment variable named
|
|
``OBSIGATE_WEBHOOK_SECRET_<ID>``.
|
|
|
|
Events: file_created, file_deleted, file_modified, file_renamed,
|
|
directory_created, directory_deleted, directory_renamed
|
|
"""
|
|
|
|
import asyncio
|
|
import hashlib
|
|
import hmac
|
|
import ipaddress
|
|
import json
|
|
import logging
|
|
import os
|
|
import socket
|
|
import threading
|
|
import uuid
|
|
from datetime import datetime, timezone
|
|
from pathlib import Path
|
|
from urllib.parse import urlparse
|
|
|
|
import aiohttp
|
|
|
|
logger = logging.getLogger("obsigate.webhooks")
|
|
|
|
WEBHOOKS_FILE = Path("data/webhooks.json")
|
|
WEBHOOK_SECRETS_FILE = Path("data/webhook_secrets.json")
|
|
|
|
VALID_EVENTS = {
|
|
"file_created", "file_deleted", "file_modified", "file_renamed",
|
|
"directory_created", "directory_deleted", "directory_renamed",
|
|
}
|
|
|
|
|
|
def _allow_http() -> bool:
|
|
return os.environ.get("OBSIGATE_WEBHOOK_ALLOW_HTTP", "false").lower() == "true"
|
|
|
|
|
|
def _allow_private() -> bool:
|
|
return os.environ.get("OBSIGATE_WEBHOOK_ALLOW_PRIVATE", "false").lower() == "true"
|
|
|
|
|
|
def _is_public_ip(ip_str: str) -> bool:
|
|
"""True when *ip_str* is a globally routable unicast address."""
|
|
try:
|
|
ip = ipaddress.ip_address(ip_str)
|
|
except ValueError:
|
|
return False
|
|
return not (
|
|
ip.is_private or ip.is_loopback or ip.is_link_local or ip.is_reserved
|
|
or ip.is_multicast or ip.is_unspecified
|
|
)
|
|
|
|
|
|
def validate_webhook_url(url: str) -> str:
|
|
"""Validate the URL syntax/scheme and literal-IP safety at config time.
|
|
|
|
Raises:
|
|
ValueError: When the URL is malformed or points at an obviously
|
|
forbidden scheme/host.
|
|
"""
|
|
if not url or not isinstance(url, str):
|
|
raise ValueError("URL requise")
|
|
parsed = urlparse(url)
|
|
if parsed.scheme not in ("http", "https"):
|
|
raise ValueError("L'URL doit utiliser http ou https")
|
|
if parsed.scheme == "http" and not _allow_http():
|
|
raise ValueError("HTTPS requis (définir OBSIGATE_WEBHOOK_ALLOW_HTTP=true pour autoriser http)")
|
|
host = parsed.hostname
|
|
if not host:
|
|
raise ValueError("Hôte manquant dans l'URL")
|
|
# Reject literal private/loopback IPs immediately (no DNS needed).
|
|
try:
|
|
ip = ipaddress.ip_address(host)
|
|
except ValueError:
|
|
return url # hostname — resolved and checked at dispatch time
|
|
if not _allow_private() and not _is_public_ip(str(ip)):
|
|
raise ValueError("Adresse privée/interne refusée")
|
|
return url
|
|
|
|
|
|
def is_safe_target(url: str) -> bool:
|
|
"""Full SSRF check performed right before dispatch (resolves the host).
|
|
|
|
Returns False when the URL is malformed, the scheme is forbidden, or any
|
|
resolved address is private/loopback/reserved.
|
|
"""
|
|
try:
|
|
validate_webhook_url(url)
|
|
except ValueError:
|
|
return False
|
|
if _allow_private():
|
|
return True
|
|
parsed = urlparse(url)
|
|
host = parsed.hostname or ""
|
|
port = parsed.port or (443 if parsed.scheme == "https" else 80)
|
|
try:
|
|
infos = socket.getaddrinfo(host, port, proto=socket.IPPROTO_TCP)
|
|
except socket.gaierror:
|
|
logger.warning(f"Webhook target host could not be resolved: {host}")
|
|
return False
|
|
for info in infos:
|
|
addr = str(info[4][0])
|
|
if not _is_public_ip(addr):
|
|
logger.warning(f"Webhook target resolves to a non-public address ({addr}); blocked")
|
|
return False
|
|
return True
|
|
|
|
|
|
def _read() -> list:
|
|
if not WEBHOOKS_FILE.exists():
|
|
return []
|
|
try:
|
|
return json.loads(WEBHOOKS_FILE.read_text(encoding="utf-8"))
|
|
except (json.JSONDecodeError, OSError):
|
|
return []
|
|
|
|
|
|
def _write(webhooks: list):
|
|
WEBHOOKS_FILE.parent.mkdir(parents=True, exist_ok=True)
|
|
tmp = WEBHOOKS_FILE.with_suffix(".tmp")
|
|
tmp.write_text(json.dumps(webhooks, indent=2, default=str), encoding="utf-8")
|
|
tmp.replace(WEBHOOKS_FILE)
|
|
|
|
|
|
def _read_secrets() -> dict:
|
|
if not WEBHOOK_SECRETS_FILE.exists():
|
|
return {}
|
|
try:
|
|
return json.loads(WEBHOOK_SECRETS_FILE.read_text(encoding="utf-8"))
|
|
except (json.JSONDecodeError, OSError):
|
|
return {}
|
|
|
|
|
|
# ROADMAP #85 T10a — verrou autour des read-modify-write des deux stores
|
|
# (webhooks + secrets) : perte de mises à jour en cas de mutations
|
|
# concurrentes.
|
|
_lock = threading.RLock()
|
|
|
|
|
|
def _write_secrets(secrets: dict):
|
|
WEBHOOK_SECRETS_FILE.parent.mkdir(parents=True, exist_ok=True)
|
|
tmp = WEBHOOK_SECRETS_FILE.with_suffix(".tmp")
|
|
tmp.write_text(json.dumps(secrets, indent=2), encoding="utf-8")
|
|
tmp.replace(WEBHOOK_SECRETS_FILE)
|
|
try:
|
|
WEBHOOK_SECRETS_FILE.chmod(0o600)
|
|
except OSError:
|
|
pass # Windows doesn't support Unix permissions
|
|
|
|
|
|
def _store_secret(wh_id: str, secret: str | None) -> None:
|
|
with _lock:
|
|
secrets = _read_secrets()
|
|
if secret:
|
|
secrets[wh_id] = secret
|
|
else:
|
|
secrets.pop(wh_id, None)
|
|
_write_secrets(secrets)
|
|
|
|
|
|
def _get_secret(wh: dict) -> str | None:
|
|
"""Resolve a webhook secret from env, dedicated store, or legacy record."""
|
|
env_key = "OBSIGATE_WEBHOOK_SECRET_" + wh["id"].replace("-", "_").upper()
|
|
env_val = os.environ.get(env_key)
|
|
if env_val:
|
|
return env_val
|
|
stored = _read_secrets().get(wh["id"])
|
|
if stored:
|
|
return stored
|
|
return wh.get("secret") # legacy inline secret
|
|
|
|
|
|
def _public_view(wh: dict) -> dict:
|
|
"""Return a webhook record safe to expose through the API."""
|
|
clean = {k: v for k, v in wh.items() if k != "secret"}
|
|
clean["has_secret"] = bool(_get_secret(wh))
|
|
return clean
|
|
|
|
|
|
def get_webhooks() -> list:
|
|
return [_public_view(wh) for wh in _read()]
|
|
|
|
|
|
def create_webhook(name: str, url: str, events: list[str], secret: str | None = None) -> dict:
|
|
validate_webhook_url(url)
|
|
with _lock:
|
|
webhooks = _read()
|
|
wh_id = str(uuid.uuid4())
|
|
wh = {
|
|
"id": wh_id,
|
|
"name": name,
|
|
"url": url,
|
|
"events": [e for e in events if e in VALID_EVENTS],
|
|
"enabled": True,
|
|
"created_at": datetime.now(timezone.utc).isoformat(),
|
|
"last_fired_at": None,
|
|
}
|
|
webhooks.append(wh)
|
|
_write(webhooks)
|
|
if secret:
|
|
_store_secret(wh_id, secret)
|
|
logger.info(f"Created webhook '{name}' → {url}")
|
|
return _public_view(wh)
|
|
|
|
|
|
def update_webhook(wh_id: str, updates: dict) -> dict | None:
|
|
with _lock:
|
|
webhooks = _read()
|
|
for wh in webhooks:
|
|
if wh["id"] == wh_id:
|
|
if updates.get("url"):
|
|
validate_webhook_url(updates["url"])
|
|
if "secret" in updates:
|
|
_store_secret(wh_id, updates["secret"])
|
|
safe_updates = {
|
|
k: v for k, v in updates.items()
|
|
if k not in ("id", "secret")
|
|
}
|
|
wh.update(safe_updates)
|
|
_write(webhooks)
|
|
return _public_view(wh)
|
|
return None
|
|
|
|
|
|
def delete_webhook(wh_id: str) -> bool:
|
|
with _lock:
|
|
webhooks = _read()
|
|
new_list = [wh for wh in webhooks if wh["id"] != wh_id]
|
|
if len(new_list) == len(webhooks):
|
|
return False
|
|
_write(new_list)
|
|
secrets = _read_secrets()
|
|
if secrets.pop(wh_id, None) is not None:
|
|
_write_secrets(secrets)
|
|
return True
|
|
|
|
|
|
async def dispatch_webhooks(event_type: str, data: dict):
|
|
"""Fire all enabled webhooks subscribed to event_type."""
|
|
webhooks = _read()
|
|
targets = [wh for wh in webhooks if wh.get("enabled", True) and event_type in wh.get("events", [])]
|
|
if not targets:
|
|
return
|
|
|
|
payload = {
|
|
"event": event_type,
|
|
"timestamp": datetime.now(timezone.utc).isoformat(),
|
|
"data": data,
|
|
}
|
|
body = json.dumps(payload, default=str)
|
|
|
|
async def _post(wh):
|
|
try:
|
|
# BUG-026: re-check the target right before connecting (DNS rebinding).
|
|
if not is_safe_target(wh["url"]):
|
|
logger.warning(f"Webhook '{wh['name']}' blocked by SSRF policy")
|
|
return
|
|
headers = {"Content-Type": "application/json", "X-ObsiGate-Event": event_type}
|
|
secret = _get_secret(wh)
|
|
if secret:
|
|
sig = hmac.new(secret.encode(), body.encode(), hashlib.sha256).hexdigest()
|
|
headers["X-ObsiGate-Signature"] = f"sha256={sig}"
|
|
|
|
timeout = aiohttp.ClientTimeout(total=5)
|
|
async with aiohttp.ClientSession(timeout=timeout) as session, session.post(
|
|
wh["url"], data=body, headers=headers, allow_redirects=False
|
|
) as resp:
|
|
if resp.status < 400:
|
|
logger.debug(f"Webhook '{wh['name']}' OK ({resp.status})")
|
|
else:
|
|
logger.warning(f"Webhook '{wh['name']}' failed ({resp.status})")
|
|
update_webhook(wh["id"], {"last_fired_at": datetime.now(timezone.utc).isoformat()})
|
|
except Exception as e:
|
|
logger.warning(f"Webhook '{wh['name']}' error: {e}")
|
|
|
|
tasks = [asyncio.create_task(_post(wh)) for wh in targets]
|
|
# Don't await — fire and forget (webhooks should not block the main thread)
|
|
# But we do register them so they run in background
|
|
for task in tasks:
|
|
task.add_done_callback(lambda t: t.exception() if not t.cancelled() else None)
|