diff --git a/api/bridge.py b/api/bridge.py index a3b190b..ef8b3f6 100644 --- a/api/bridge.py +++ b/api/bridge.py @@ -26,6 +26,7 @@ try: from ..services.execution_budgets import BudgetExceededError from ..services.idempotency_store import IdempotencyStore from ..services.rate_limit import build_rate_limit_response, check_rate_limit + from ..services.redaction import stable_redaction_tag from ..services.sidecar.auth import is_bridge_enabled, require_bridge_auth from ..services.sidecar.bridge_contract import ( BRIDGE_ENDPOINTS, @@ -44,6 +45,7 @@ except ImportError: from services.execution_budgets import BudgetExceededError from services.idempotency_store import IdempotencyStore from services.rate_limit import build_rate_limit_response, check_rate_limit + from services.redaction import stable_redaction_tag from services.sidecar.auth import is_bridge_enabled, require_bridge_auth from services.sidecar.bridge_contract import ( BRIDGE_ENDPOINTS, @@ -82,6 +84,10 @@ MAX_FILES_COUNT = 10 _startup_time = time.time() +def _bridge_sensitive_tag(value: Optional[str], *, label: str) -> str: + return stable_redaction_tag(value, label=label) + + class BridgeHandlers: """Handlers for bridge API endpoints.""" @@ -107,7 +113,7 @@ class BridgeHandlers: def _bridge_token(self, device_id: Optional[str], scope: Optional[str] = None): scopes = {scope} if scope else set() return SimpleNamespace( - token_id=f"bridge:{device_id or 'unknown'}", + token_id=f"bridge:{_bridge_sensitive_tag(device_id, label='device')}", role="bridge", scopes=scopes, ) @@ -716,7 +722,10 @@ class BridgeHandlers: store_key = f"wr:{idempotency_key}" is_dup, _ = self._idempotency_store.check_and_record(store_key, ttl=86400) if is_dup: - logger.info(f"Duplicate worker result suppressed: {idempotency_key}") + logger.info( + "Duplicate worker result suppressed for %s", + _bridge_sensitive_tag(idempotency_key, label="idem"), + ) return web.json_response( { "ok": True, @@ -745,7 +754,8 @@ class BridgeHandlers: self._worker_results[job_id] = { "status": data.get("status", "completed"), "outputs": data.get("outputs", {}), - "worker_id": device_id, + # IMPORTANT: keep worker identity redacted in cached bridge state. + "worker_id": _bridge_sensitive_tag(device_id, label="device"), "timestamp": time.time(), } @@ -761,7 +771,11 @@ class BridgeHandlers: scope=BridgeScope.JOB_SUBMIT.value, details={"status": data.get("status", "completed")}, ) - logger.info(f"F46: Worker result accepted for job={job_id} from={device_id}") + logger.info( + "F46: Worker result accepted for job=%s from=%s", + job_id, + _bridge_sensitive_tag(device_id, label="device"), + ) return web.json_response(response_data, status=201) @endpoint_metadata( diff --git a/services/audit.py b/services/audit.py index 6f226a9..736189a 100644 --- a/services/audit.py +++ b/services/audit.py @@ -12,6 +12,8 @@ import time import uuid from typing import Any, Dict, Iterable, Optional, Tuple +from .redaction import redact_json, stable_redaction_tag + logger = logging.getLogger("ComfyUI-OpenClaw.services.audit") _TRUTHY = {"1", "true", "yes", "on"} @@ -136,6 +138,18 @@ def _resolve_trace_id(request: Any, details: Dict[str, Any]) -> str: return uuid.uuid4().hex +def _sanitize_audit_details(details: Optional[Dict[str, Any]]) -> Any: + safe_details = _json_safe(details or {}) + if isinstance(safe_details, dict): + sanitized = dict(safe_details) + actor_ip = sanitized.pop("actor_ip", None) + if actor_ip is not None: + # IMPORTANT: keep network provenance correlatable without storing raw client IPs. + sanitized["actor_ip_tag"] = stable_redaction_tag(actor_ip, label="ip") + return redact_json(sanitized) + return redact_json(safe_details) + + def _chain_hash(prev_hash: str, entry: Dict[str, Any]) -> str: payload = json.dumps( entry, sort_keys=True, separators=(",", ":"), ensure_ascii=True @@ -237,7 +251,7 @@ def _emit_modern( request: Optional[Any] = None, source: str = "openclaw", ) -> Dict[str, Any]: - details_dict = _json_safe(details or {}) + details_dict = _sanitize_audit_details(details or {}) token = token_info or _resolve_request_token_info(request) token_id = "anonymous" role = "unknown" @@ -283,7 +297,7 @@ def _emit_legacy( error: Optional[str] = None, metadata: Optional[Dict[str, Any]] = None, ) -> Dict[str, Any]: - details = {"actor_ip": actor_ip} + details = {"actor_ip_tag": stable_redaction_tag(actor_ip, label="ip")} if provider: details["provider"] = provider if error: @@ -341,7 +355,9 @@ def audit_config_write(actor_ip: str, ok: bool, error: Optional[str] = None) -> outcome="allow" if ok else "error", status_code=200 if ok else 400, details=( - {"actor_ip": actor_ip, "error": error} if error else {"actor_ip": actor_ip} + {"actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), "error": error} + if error + else {"actor_ip_tag": stable_redaction_tag(actor_ip, label="ip")} ), ) @@ -355,9 +371,16 @@ def audit_secret_write( outcome="allow" if ok else "error", status_code=200 if ok else 500, details=( - {"actor_ip": actor_ip, "provider": provider, "error": error} + { + "actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), + "provider": provider, + "error": error, + } if error - else {"actor_ip": actor_ip, "provider": provider} + else { + "actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), + "provider": provider, + } ), ) @@ -371,9 +394,16 @@ def audit_secret_delete( outcome="allow" if ok else "error", status_code=200 if ok else 404, details=( - {"actor_ip": actor_ip, "provider": provider, "error": error} + { + "actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), + "provider": provider, + "error": error, + } if error - else {"actor_ip": actor_ip, "provider": provider} + else { + "actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), + "provider": provider, + } ), ) @@ -385,6 +415,8 @@ def audit_llm_test(actor_ip: str, ok: bool, error: Optional[str] = None) -> None outcome="allow" if ok else "error", status_code=200 if ok else 500, details=( - {"actor_ip": actor_ip, "error": error} if error else {"actor_ip": actor_ip} + {"actor_ip_tag": stable_redaction_tag(actor_ip, label="ip"), "error": error} + if error + else {"actor_ip_tag": stable_redaction_tag(actor_ip, label="ip")} ), ) diff --git a/services/redaction.py b/services/redaction.py index 3079964..2000ed5 100644 --- a/services/redaction.py +++ b/services/redaction.py @@ -7,6 +7,7 @@ Prevents sensitive data leakage in observability outputs. from __future__ import annotations +import hashlib import logging import re from typing import Any, Dict, List, Optional, Set, Tuple @@ -94,6 +95,22 @@ SENSITIVE_KEYS: Set[str] = { } +def stable_redaction_tag(value: Any, *, label: str = "value") -> str: + """ + Build a deterministic, non-cleartext correlation tag for sensitive identifiers. + + The output is intentionally one-way and short so operators can correlate repeated + values across logs without exposing the original identifier. + """ + if value is None: + return f"{label}:none" + text = str(value).strip() + if not text: + return f"{label}:empty" + digest = hashlib.sha256(text.encode("utf-8")).hexdigest()[:12] + return f"{label}:{digest}" + + def redact_text( text: str, patterns: Optional[List[Tuple[re.Pattern, str]]] = None ) -> str: diff --git a/services/request_ip.py b/services/request_ip.py index de58da5..d354409 100644 --- a/services/request_ip.py +++ b/services/request_ip.py @@ -41,7 +41,8 @@ def get_trusted_proxies() -> List[ipaddress.IPv4Network]: # Try as network CIDR networks.append(ipaddress.ip_network(part, strict=False)) except ValueError: - logger.warning(f"Invalid trusted proxy entry: {part}") + # IMPORTANT: never echo malformed proxy config verbatim into logs. + logger.warning("Invalid trusted proxy entry ignored.") return networks diff --git a/services/safe_io.py b/services/safe_io.py index 78b7491..7d14b1c 100644 --- a/services/safe_io.py +++ b/services/safe_io.py @@ -679,7 +679,8 @@ def safe_request_json( ): request.add_header(key, value) else: - logger.debug(f"Skipping disallowed header: {key}") + # IMPORTANT: do not log caller-supplied header names verbatim here. + logger.debug("Skipping disallowed outbound header.") # Build Pinned Opener opener = _build_pinned_opener(pinned_ips) @@ -790,7 +791,7 @@ def safe_request_text_stream( ): request.add_header(key, value) else: - logger.debug(f"Skipping disallowed header: {key}") + logger.debug("Skipping disallowed outbound header.") opener = _build_pinned_opener(pinned_ips) diff --git a/services/sidecar/auth.py b/services/sidecar/auth.py index 24842df..1d1b2f8 100644 --- a/services/sidecar/auth.py +++ b/services/sidecar/auth.py @@ -17,6 +17,11 @@ except ModuleNotFoundError: # pragma: no cover (optional for unit tests) from .bridge_contract import BridgeScope +try: + from ..redaction import stable_redaction_tag +except ImportError: # pragma: no cover + from services.redaction import stable_redaction_tag # type: ignore + logger = logging.getLogger("ComfyUI-OpenClaw.sidecar.auth") # Environment configuration @@ -43,6 +48,14 @@ LEGACY_HEADER_SCOPES = "X-Moltbot-Scopes" HEADER_CLIENT_CERT_HASH = "X-Client-Cert-Hash" # SHA256 fingerprint from proxy +def _device_tag(device_id: Optional[str]) -> str: + return stable_redaction_tag(device_id, label="device") + + +def _cert_tag(cert_hash: Optional[str]) -> str: + return stable_redaction_tag(cert_hash, label="cert") + + def _env_get(primary: str, legacy: str, default: str = "") -> str: """Get env var value (prefers new names, falls back to legacy). Respects empty-string overrides.""" if primary in os.environ: @@ -123,12 +136,16 @@ def validate_mtls_binding(request: web.Request, device_id: str) -> Tuple[bool, s if not expected_hash: # Strict mode: mTLS enabled implies explicit device binding - return False, f"Device not bound to a certificate (Device ID: {device_id})" + return False, "Device not bound to a certificate" # Constant-time comparison not strictly required for public fingerprints but good practice if not hmac.compare_digest(cert_hash, expected_hash): + # IMPORTANT: keep bound-device diagnostics redacted; raw fingerprints are sensitive. logger.warning( - f"mTLS violation: Device {device_id} presented {cert_hash}, expected {expected_hash}" + "mTLS violation for %s presented=%s expected=%s", + _device_tag(device_id), + _cert_tag(cert_hash), + _cert_tag(expected_hash), ) return False, "Certificate fingerprint mismatch" @@ -204,8 +221,8 @@ def validate_device_token( lifecycle_result = get_token_store().validate_token( device_token, required_scope=required_scope_value ) - except Exception as e: - logger.warning(f"S58 lifecycle validation unavailable, falling back: {e}") + except Exception: + logger.warning("S58 lifecycle validation unavailable, falling back.") lifecycle_result = None if lifecycle_result and lifecycle_result.ok and lifecycle_result.token: @@ -240,7 +257,7 @@ def validate_device_token( # Constant-time token comparison if not hmac.compare_digest(device_token, expected_token): - logger.warning(f"Invalid device token from: {device_id[:8]}...") + logger.warning("Invalid device token from %s", _device_tag(device_id)) return False, "Invalid device token", None # Scope validation (legacy header contract) @@ -250,7 +267,8 @@ def validate_device_token( ) or request.headers.get(LEGACY_HEADER_SCOPES, "") if not scopes_header: logger.warning( - f"Device {device_id[:8]} missing required scopes header." + "Bridge device %s missing required scopes header.", + _device_tag(device_id), ) return False, "Missing X-OpenClaw-Scopes header", None @@ -259,20 +277,23 @@ def validate_device_token( ) if required_scope not in granted_scopes: logger.warning( - f"Device {device_id[:8]} missing scope {required_scope}. Granted: {granted_scopes}" + "Bridge device %s missing scope %s (granted_count=%d).", + _device_tag(device_id), + required_scope, + len(granted_scopes), ) return False, f"Missing required scope: {required_scope}", None # Check allowlist if configured allowed_ids = get_allowed_device_ids() if allowed_ids is not None and device_id not in allowed_ids: - logger.warning(f"Device ID not in allowlist: {device_id[:8]}...") + logger.warning("Device ID not in allowlist: %s", _device_tag(device_id)) return False, "Device not authorized", None # R104: mTLS Binding Check is_mtls_valid, mtls_error = validate_mtls_binding(request, device_id) if not is_mtls_valid: - logger.warning(f"mTLS validation failed for {device_id}: {mtls_error}") + logger.warning("mTLS validation failed for %s: %s", _device_tag(device_id), mtls_error) return False, mtls_error, None return True, "", device_id diff --git a/tests/security/test_s78_redaction.py b/tests/security/test_s78_redaction.py new file mode 100644 index 0000000..ca804a8 --- /dev/null +++ b/tests/security/test_s78_redaction.py @@ -0,0 +1,272 @@ +import asyncio +import importlib +import json +import os +import unittest +from unittest.mock import AsyncMock, MagicMock, patch + +from services.audit import emit_audit_event +from services.idempotency_store import IdempotencyStore +from services.request_ip import get_trusted_proxies +from services.redaction import stable_redaction_tag +from services.safe_io import safe_request_json, safe_request_text_stream + + +class TestS78BridgeAuthRedaction(unittest.TestCase): + def setUp(self): + self.env_keys = [ + "OPENCLAW_BRIDGE_ENABLED", + "OPENCLAW_BRIDGE_DEVICE_TOKEN", + "OPENCLAW_BRIDGE_MTLS_ENABLED", + "OPENCLAW_BRIDGE_DEVICE_CERT_MAP", + ] + for key in self.env_keys: + os.environ.pop(key, None) + + def tearDown(self): + for key in self.env_keys: + os.environ.pop(key, None) + + def _reload_auth_module(self): + import services.sidecar.auth as auth_module + + return importlib.reload(auth_module) + + def test_invalid_token_log_redacts_device_id(self): + os.environ["OPENCLAW_BRIDGE_ENABLED"] = "1" + os.environ["OPENCLAW_BRIDGE_DEVICE_TOKEN"] = "expected-token" + auth_module = self._reload_auth_module() + + req = MagicMock() + req.headers = { + "X-OpenClaw-Device-Id": "worker-1", + "X-OpenClaw-Device-Token": "wrong-token", + } + + with self.assertLogs("ComfyUI-OpenClaw.sidecar.auth", level="WARNING") as logs: + is_valid, error, _ = auth_module.validate_device_token(req) + + self.assertFalse(is_valid) + self.assertIn("invalid", error.lower()) + output = "\n".join(logs.output) + self.assertNotIn("worker-1", output) + self.assertIn("device:", output) + + def test_mtls_mismatch_log_redacts_cert_fingerprints(self): + os.environ["OPENCLAW_BRIDGE_ENABLED"] = "1" + os.environ["OPENCLAW_BRIDGE_DEVICE_TOKEN"] = "expected-token" + os.environ["OPENCLAW_BRIDGE_MTLS_ENABLED"] = "1" + os.environ["OPENCLAW_BRIDGE_DEVICE_CERT_MAP"] = "worker-1:sha256_expected" + auth_module = self._reload_auth_module() + + req = MagicMock() + req.headers = { + "X-OpenClaw-Device-Id": "worker-1", + "X-OpenClaw-Device-Token": "expected-token", + "X-Client-Cert-Hash": "sha256_actual", + } + + with self.assertLogs("ComfyUI-OpenClaw.sidecar.auth", level="WARNING") as logs: + is_valid, error, _ = auth_module.validate_device_token(req) + + self.assertFalse(is_valid) + self.assertIn("fingerprint mismatch", error.lower()) + output = "\n".join(logs.output) + self.assertNotIn("worker-1", output) + self.assertNotIn("sha256_actual", output) + self.assertNotIn("sha256_expected", output) + self.assertIn("cert:", output) + + +class TestS78BridgeWorkerRedaction(unittest.TestCase): + def setUp(self): + os.environ["OPENCLAW_BRIDGE_ENABLED"] = "1" + os.environ["OPENCLAW_BRIDGE_DEVICE_TOKEN"] = "test-token-secret" + IdempotencyStore().clear() + import services.sidecar.auth as auth_module + + importlib.reload(auth_module) + + def tearDown(self): + for key in [ + "OPENCLAW_BRIDGE_ENABLED", + "OPENCLAW_BRIDGE_DEVICE_TOKEN", + "MOLTBOT_BRIDGE_ENABLED", + "MOLTBOT_BRIDGE_DEVICE_TOKEN", + ]: + os.environ.pop(key, None) + IdempotencyStore().clear() + + def _make_auth_request(self, *, job_id: str, idempotency_key: str): + req = MagicMock() + req.method = "POST" + req.path = f"/bridge/worker/result/{job_id}" + req.headers = { + "X-OpenClaw-Device-Id": "worker-1", + "X-OpenClaw-Device-Token": "test-token-secret", + "X-OpenClaw-Scopes": "job:submit,job:status", + "X-Idempotency-Key": idempotency_key, + } + req.query = {} + req.match_info = {"job_id": job_id} + req.json = AsyncMock(return_value={"status": "completed", "outputs": {}}) + return req + + def test_worker_result_cache_and_logs_redact_sensitive_ids(self): + from api.bridge import BridgeHandlers + + handlers = BridgeHandlers() + req = self._make_auth_request(job_id="job-1", idempotency_key="idem-001") + + with self.assertLogs("ComfyUI-OpenClaw.api.bridge", level="INFO") as logs: + resp = asyncio.run(handlers.worker_result_handler(req)) + + self.assertEqual(resp.status, 201) + self.assertEqual( + handlers._worker_results["job-1"]["worker_id"], + stable_redaction_tag("worker-1", label="device"), + ) + output = "\n".join(logs.output) + self.assertNotIn("worker-1", output) + self.assertIn("device:", output) + + def test_duplicate_result_log_redacts_idempotency_key(self): + from api.bridge import BridgeHandlers + + handlers = BridgeHandlers() + req1 = self._make_auth_request(job_id="job-2", idempotency_key="idem-dup") + req2 = self._make_auth_request(job_id="job-2", idempotency_key="idem-dup") + + asyncio.run(handlers.worker_result_handler(req1)) + + with self.assertLogs("ComfyUI-OpenClaw.api.bridge", level="INFO") as logs: + resp = asyncio.run(handlers.worker_result_handler(req2)) + + self.assertEqual(resp.status, 200) + output = "\n".join(logs.output) + self.assertNotIn("idem-dup", output) + self.assertIn("idem:", output) + + +class TestS78AuditRedaction(unittest.TestCase): + def setUp(self): + self.test_log = "test_s78_audit.log" + self.path_patcher = patch("services.audit.AUDIT_LOG_PATH", self.test_log) + self.hash_patcher = patch("services.audit._LAST_HASH", None) + self.path_patcher.start() + self.hash_patcher.start() + if os.path.exists(self.test_log): + os.remove(self.test_log) + + def tearDown(self): + self.path_patcher.stop() + self.hash_patcher.stop() + if os.path.exists(self.test_log): + os.remove(self.test_log) + + def _read_entries(self): + with open(self.test_log, "r", encoding="utf-8") as handle: + return [json.loads(line) for line in handle if line.strip()] + + def test_audit_storage_redacts_actor_ip(self): + emit_audit_event( + action="config.update", + target="config.json", + outcome="allow", + status_code=200, + details={"actor_ip": "1.2.3.4", "error": "Authorization: Bearer sk-secret"}, + ) + + entries = self._read_entries() + details = entries[0]["details"] + self.assertNotIn("actor_ip", details) + self.assertEqual(details["actor_ip_tag"], "ip:6694f83c9f47") + self.assertIn("***REDACTED***", details["error"]) + self.assertNotIn("1.2.3.4", json.dumps(entries[0])) + + +class TestS78DefensiveLogRedaction(unittest.TestCase): + def test_invalid_trusted_proxy_log_omits_raw_value(self): + with patch.dict( + os.environ, + {"OPENCLAW_TRUSTED_PROXIES": "127.0.0.1,not-a-network"}, + clear=False, + ): + with self.assertLogs( + "ComfyUI-OpenClaw.services.request_ip", level="WARNING" + ) as logs: + get_trusted_proxies() + + output = "\n".join(logs.output) + self.assertIn("Invalid trusted proxy entry ignored.", output) + self.assertNotIn("not-a-network", output) + + @patch("services.safe_io._build_pinned_opener") + @patch("services.safe_io.validate_outbound_url") + def test_safe_request_json_log_omits_disallowed_header_name( + self, mock_validate, mock_build + ): + mock_validate.return_value = ("https", "example.com", 443, ["93.184.216.34"]) + + mock_response = MagicMock() + mock_response.getcode.return_value = 200 + mock_response.read.return_value = b'{"ok": true}' + + mock_opener = MagicMock() + mock_opener.open.return_value.__enter__.return_value = mock_response + mock_build.return_value = mock_opener + + with self.assertLogs("ComfyUI-OpenClaw.services.safe_io", level="DEBUG") as logs: + out = safe_request_json( + method="POST", + url="https://example.com/test", + json_body={"x": 1}, + headers={"Secret-Header": "blocked"}, + allow_hosts={"example.com"}, + ) + + self.assertTrue(out["ok"]) + output = "\n".join(logs.output) + self.assertIn("Skipping disallowed outbound header.", output) + self.assertNotIn("Secret-Header", output) + + @patch("services.safe_io._build_pinned_opener") + @patch("services.safe_io.validate_outbound_url") + def test_safe_request_stream_log_omits_disallowed_header_name( + self, mock_validate, mock_build + ): + mock_validate.return_value = ("https", "example.com", 443, ["93.184.216.34"]) + + class _FakeStreamResponse: + def __init__(self): + self.headers = {} + self._lines = [b"data: ok\n", b""] + + def getcode(self): + return 200 + + def readline(self, _max_bytes): + return self._lines.pop(0) + + def close(self): + return None + + mock_opener = MagicMock() + mock_opener.open.return_value = _FakeStreamResponse() + mock_build.return_value = mock_opener + + with self.assertLogs("ComfyUI-OpenClaw.services.safe_io", level="DEBUG") as logs: + lines = list( + safe_request_text_stream( + method="POST", + url="https://example.com/stream", + json_body={"x": 1}, + headers={"Secret-Header": "blocked"}, + allow_hosts={"example.com"}, + ) + ) + + self.assertEqual(lines, ["data: ok\n"]) + output = "\n".join(logs.output) + self.assertIn("Skipping disallowed outbound header.", output) + self.assertNotIn("Secret-Header", output)