Files
ObsiGate/backend/webhooks.py
T

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)