Encrypt sensitive credential fields before event bus emission (#215)

- Add encrypt_credential_dict() and decrypt_credential_field() to utils/crypto.py
- Create utils/credential_encryption.py with encryption/decryption helpers
- Add get_credential_encryption_key() to retrieve key from Infisical at runtime
- Add emit_credential_found() wrapper for modules to use instead of direct bus.emit()
- Update all 15 modules emitting CREDENTIAL_FOUND to use encrypted wrapper
- Update 4 modules consuming CREDENTIAL_FOUND to decrypt payload before processing
- Sensitive fields (username, password, hash, token, etc.) encrypted with AES-256-GCM
- Falls back to plaintext if encryption key unavailable or encryption fails
This commit is contained in:
Cobra
2026-04-06 11:43:46 -04:00
parent f15b8994cd
commit edf2ea7c36
18 changed files with 256 additions and 63 deletions
+1
View File
@@ -10,3 +10,4 @@ __pycache__/
*.pcap.zst *.pcap.zst
*.pcap.zst.enc *.pcap.zst.enc
net_alerter/.secrets net_alerter/.secrets
BIGBROTHER_DESIGN.md
+2 -1
View File
@@ -20,6 +20,7 @@ from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.bettercap_api import BettercapAPI from utils.bettercap_api import BettercapAPI
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.active.bettercap_mgr") logger = logging.getLogger("sensor.active.bettercap_mgr")
@@ -477,7 +478,7 @@ class BettercapManager(BaseModule):
"credential_value": data.get("password", data.get("hash", "")), "credential_value": data.get("password", data.get("hash", "")),
"raw_data": data, "raw_data": data,
} }
self.bus.emit("CREDENTIAL_FOUND", payload, source_module=self.name) emit_credential_found(self.bus, self.name, payload)
logger.info( logger.info(
"Credential captured: %s@%s (%s)", "Credential captured: %s@%s (%s)",
payload["username"], payload["target_ip"], payload["target_service"], payload["username"], payload["target_ip"], payload["target_service"],
+1
View File
@@ -23,6 +23,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.active.evil_twin") logger = logging.getLogger("sensor.active.evil_twin")
+1
View File
@@ -20,6 +20,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.active.mitmproxy_mgr") logger = logging.getLogger("sensor.active.mitmproxy_mgr")
+1
View File
@@ -18,6 +18,7 @@ import time
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.active.ntlm_relay") logger = logging.getLogger("sensor.active.ntlm_relay")
+17 -24
View File
@@ -20,6 +20,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.active.responder_mgr") logger = logging.getLogger("sensor.active.responder_mgr")
@@ -403,31 +404,23 @@ class ResponderManager(BaseModule):
with self._hashes_lock: with self._hashes_lock:
self._captured_hashes.append(hash_entry) self._captured_hashes.append(hash_entry)
self.bus.emit( emit_credential_found(self.bus, self.name, {
"CREDENTIAL_FOUND", "source_module": self.name,
{ "source_ip": "",
"source_module": self.name, "target_service": "responder",
"source_ip": "", "username": username,
"target_service": "responder", "domain": domain,
"username": username, "credential_type": hash_type,
"domain": domain, "credential_value": line,
"credential_type": hash_type, "hashcat_mode": hashcat_mode,
"credential_value": line, })
"hashcat_mode": hashcat_mode,
},
source_module=self.name,
)
logger.info("Hash captured: %s\\%s (%s)", domain, username, hash_type) logger.info("Hash captured: %s\\%s (%s)", domain, username, hash_type)
def _process_cleartext_line(self, line: str) -> None: def _process_cleartext_line(self, line: str) -> None:
"""Process a cleartext credential line from Responder session log.""" """Process a cleartext credential line from Responder session log."""
self.bus.emit( emit_credential_found(self.bus, self.name, {
"CREDENTIAL_FOUND", "source_module": self.name,
{ "target_service": "responder_cleartext",
"source_module": self.name, "credential_type": "cleartext",
"target_service": "responder_cleartext", "credential_value": line,
"credential_type": "cleartext", })
"credential_value": line,
},
source_module=self.name,
)
+1
View File
@@ -29,6 +29,7 @@ from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from core.bus import Event from core.bus import Event
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.connectivity.data_exfil") logger = logging.getLogger("sensor.connectivity.data_exfil")
+3
View File
@@ -18,6 +18,7 @@ from collections import defaultdict
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.intel.change_detector") logger = logging.getLogger("sensor.intel.change_detector")
@@ -362,6 +363,8 @@ class ChangeDetector(BaseModule):
def _on_credential_found(self, event) -> None: def _on_credential_found(self, event) -> None:
"""Monitor for unusual authentication patterns.""" """Monitor for unusual authentication patterns."""
p = event.payload p = event.payload
from utils.credential_encryption import decrypt_credential_payload
p = decrypt_credential_payload(p)
# Track rapid auth failures from unknown sources as investigation indicator # Track rapid auth failures from unknown sources as investigation indicator
if p.get("cred_type") == "auth_failure": if p.get("cred_type") == "auth_failure":
src = p.get("source_ip", "") src = p.get("source_ip", "")
+3
View File
@@ -17,6 +17,7 @@ import time
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.intel.credential_db") logger = logging.getLogger("sensor.intel.credential_db")
@@ -179,6 +180,8 @@ class CredentialDB(BaseModule):
def _on_credential_found(self, event) -> None: def _on_credential_found(self, event) -> None:
"""Handle CREDENTIAL_FOUND events from any capture module.""" """Handle CREDENTIAL_FOUND events from any capture module."""
p = event.payload p = event.payload
from utils.credential_encryption import decrypt_credential_payload
p = decrypt_credential_payload(p)
self._ingest_credential( self._ingest_credential(
source_module=event.source_module or p.get("source_module", "unknown"), source_module=event.source_module or p.get("source_module", "unknown"),
source_ip=p.get("source_ip", ""), source_ip=p.get("source_ip", ""),
+23 -22
View File
@@ -17,6 +17,7 @@ import time
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.intel.tool_output_parser") logger = logging.getLogger("sensor.intel.tool_output_parser")
@@ -255,7 +256,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": "smb", "service": "smb",
"username": username, "username": username,
"domain": domain, "domain": domain,
@@ -265,7 +266,7 @@ class ToolOutputParser(BaseModule):
"source_ip": "", "source_ip": "",
"target_ip": "", "target_ip": "",
"notes": f"Captured by Responder from {os.path.basename(path)}", "notes": f"Captured by Responder from {os.path.basename(path)}",
}, source_module=self.name) })
def _parse_responder_session(self, path: str) -> None: def _parse_responder_session(self, path: str) -> None:
"""Parse Responder-Session.log for cleartext credentials.""" """Parse Responder-Session.log for cleartext credentials."""
@@ -287,7 +288,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": protocol, "service": protocol,
"username": username, "username": username,
"domain": "", "domain": "",
@@ -297,7 +298,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source_ip, "source_ip": source_ip,
"target_ip": "", "target_ip": "",
"notes": "Cleartext capture by Responder", "notes": "Cleartext capture by Responder",
}, source_module=self.name) })
# Also check for HTTP basic auth in session log # Also check for HTTP basic auth in session log
for match in _RESPONDER_HTTP_RE.finditer(new_data): for match in _RESPONDER_HTTP_RE.finditer(new_data):
@@ -310,7 +311,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": "http", "service": "http",
"username": username, "username": username,
"domain": "", "domain": "",
@@ -320,7 +321,7 @@ class ToolOutputParser(BaseModule):
"source_ip": "", "source_ip": "",
"target_ip": "", "target_ip": "",
"notes": "HTTP Basic Auth captured by Responder", "notes": "HTTP Basic Auth captured by Responder",
}, source_module=self.name) })
# ------------------------------------------------------------------ # ------------------------------------------------------------------
# bettercap parser # bettercap parser
@@ -441,7 +442,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": protocol, "service": protocol,
"username": "", "username": "",
"domain": "", "domain": "",
@@ -450,7 +451,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source, "source_ip": source,
"target_ip": target, "target_ip": target,
"notes": "Captured by bettercap", "notes": "Captured by bettercap",
}, source_module=self.name) })
return return
hash_key = f"bettercap:{protocol}:{username}:{password[:16]}" hash_key = f"bettercap:{protocol}:{username}:{password[:16]}"
@@ -459,7 +460,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": protocol, "service": protocol,
"username": username, "username": username,
"domain": "", "domain": "",
@@ -469,7 +470,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source, "source_ip": source,
"target_ip": target, "target_ip": target,
"notes": "Captured by bettercap sniffer", "notes": "Captured by bettercap sniffer",
}, source_module=self.name) })
def _handle_bettercap_host(self, data: dict) -> None: def _handle_bettercap_host(self, data: dict) -> None:
"""Extract host info from a bettercap endpoint event.""" """Extract host info from a bettercap endpoint event."""
@@ -526,7 +527,7 @@ class ToolOutputParser(BaseModule):
if hash_key not in self._emitted_hashes: if hash_key not in self._emitted_hashes:
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"http://{host}", "service": f"http://{host}",
"username": username, "username": username,
"domain": "", "domain": "",
@@ -536,7 +537,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source_ip, "source_ip": source_ip,
"target_ip": host, "target_ip": host,
"notes": "HTTP Basic Auth from traffic", "notes": "HTTP Basic Auth from traffic",
}, source_module=self.name) })
except Exception: except Exception:
pass pass
return return
@@ -553,7 +554,7 @@ class ToolOutputParser(BaseModule):
# Check if it looks like a JWT # Check if it looks like a JWT
cred_type = "jwt" if token.count(".") == 2 else "bearer_token" cred_type = "jwt" if token.count(".") == 2 else "bearer_token"
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"http://{host}", "service": f"http://{host}",
"username": "", "username": "",
"domain": "", "domain": "",
@@ -563,7 +564,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source_ip, "source_ip": source_ip,
"target_ip": host, "target_ip": host,
"notes": f"{cred_type.upper()} from HTTP traffic", "notes": f"{cred_type.upper()} from HTTP traffic",
}, source_module=self.name) })
def _parse_cookie_header(self, value: str, host: str, def _parse_cookie_header(self, value: str, host: str,
source_ip: str) -> None: source_ip: str) -> None:
@@ -581,7 +582,7 @@ class ToolOutputParser(BaseModule):
if hash_key not in self._emitted_hashes: if hash_key not in self._emitted_hashes:
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"http://{host}", "service": f"http://{host}",
"username": "", "username": "",
"domain": "", "domain": "",
@@ -591,7 +592,7 @@ class ToolOutputParser(BaseModule):
"source_ip": source_ip, "source_ip": source_ip,
"target_ip": host, "target_ip": host,
"notes": f"Session cookie '{name}' from HTTP traffic", "notes": f"Session cookie '{name}' from HTTP traffic",
}, source_module=self.name) })
# ------------------------------------------------------------------ # ------------------------------------------------------------------
# mitmproxy parser # mitmproxy parser
@@ -680,7 +681,7 @@ class ToolOutputParser(BaseModule):
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"https://{host}", "service": f"https://{host}",
"username": username, "username": username,
"domain": "", "domain": "",
@@ -690,7 +691,7 @@ class ToolOutputParser(BaseModule):
"source_ip": "", "source_ip": "",
"target_ip": host, "target_ip": host,
"notes": "Form credentials captured by mitmproxy", "notes": "Form credentials captured by mitmproxy",
}, source_module=self.name) })
def _parse_post_body(self, content: str, host: str) -> None: def _parse_post_body(self, content: str, host: str) -> None:
"""Look for credentials in HTTP POST bodies.""" """Look for credentials in HTTP POST bodies."""
@@ -724,7 +725,7 @@ class ToolOutputParser(BaseModule):
if hash_key not in self._emitted_hashes: if hash_key not in self._emitted_hashes:
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"http://{host}", "service": f"http://{host}",
"username": username, "username": username,
"domain": "", "domain": "",
@@ -734,7 +735,7 @@ class ToolOutputParser(BaseModule):
"source_ip": "", "source_ip": "",
"target_ip": host, "target_ip": host,
"notes": "Form POST credentials from HTTP traffic", "notes": "Form POST credentials from HTTP traffic",
}, source_module=self.name) })
except Exception: except Exception:
pass pass
@@ -758,7 +759,7 @@ class ToolOutputParser(BaseModule):
if hash_key not in self._emitted_hashes: if hash_key not in self._emitted_hashes:
self._emitted_hashes.add(hash_key) self._emitted_hashes.add(hash_key)
self._creds_parsed += 1 self._creds_parsed += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"service": f"http://{host}", "service": f"http://{host}",
"username": username, "username": username,
"domain": "", "domain": "",
@@ -768,7 +769,7 @@ class ToolOutputParser(BaseModule):
"source_ip": "", "source_ip": "",
"target_ip": host, "target_ip": host,
"notes": "JSON POST credentials from HTTP traffic", "notes": "JSON POST credentials from HTTP traffic",
}, source_module=self.name) })
except (json.JSONDecodeError, TypeError): except (json.JSONDecodeError, TypeError):
pass pass
+3
View File
@@ -17,6 +17,7 @@ from collections import defaultdict
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.intel.user_timeline") logger = logging.getLogger("sensor.intel.user_timeline")
@@ -178,6 +179,8 @@ class UserTimeline(BaseModule):
def _on_credential_found(self, event) -> None: def _on_credential_found(self, event) -> None:
"""Record credential capture as a user event.""" """Record credential capture as a user event."""
p = event.payload p = event.payload
from utils.credential_encryption import decrypt_credential_payload
p = decrypt_credential_payload(p)
username = p.get("username", "") username = p.get("username", "")
if not username: if not username:
return return
+3 -2
View File
@@ -19,6 +19,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.passive.cloud_token_harvester") logger = logging.getLogger("sensor.passive.cloud_token_harvester")
@@ -259,14 +260,14 @@ class CloudTokenHarvester(BaseModule):
}, source_module=self.name) }, source_module=self.name)
# Also feed to credential_db # Also feed to credential_db
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": f"cloud_token_{token_type}", "source": f"cloud_token_{token_type}",
"username": "", "username": "",
"credential": token_value, "credential": token_value,
"source_ip": src_ip, "source_ip": src_ip,
"dest_ip": dst_ip, "dest_ip": dst_ip,
"protocol": "http", "protocol": "http",
}, source_module=self.name) })
self._record_token(ts, src_ip, dst_ip, token_type, token_value, context) self._record_token(ts, src_ip, dst_ip, token_type, token_value, context)
+3 -2
View File
@@ -26,6 +26,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.passive.credential_sniffer") logger = logging.getLogger("sensor.passive.credential_sniffer")
@@ -608,7 +609,7 @@ class CredentialSniffer(BaseModule):
self._flush_buffer() self._flush_buffer()
# Immediate bus notification # Immediate bus notification
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source_ip": src_ip, "source_ip": src_ip,
"target_ip": target_ip, "target_ip": target_ip,
"target_port": target_port, "target_port": target_port,
@@ -617,7 +618,7 @@ class CredentialSniffer(BaseModule):
"domain": domain, "domain": domain,
"cred_type": cred_type, "cred_type": cred_type,
"hashcat_mode": hashcat_mode, "hashcat_mode": hashcat_mode,
}, source_module=self.name) })
logger.info( logger.info(
"CREDENTIAL: %s %s@%s:%d (%s/%s)", "CREDENTIAL: %s %s@%s:%d (%s/%s)",
+9 -8
View File
@@ -22,6 +22,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.passive.db_interceptor") logger = logging.getLogger("sensor.passive.db_interceptor")
@@ -286,14 +287,14 @@ class DBInterceptor(BaseModule):
f"LOGIN hostname={hostname} app={appname}", "login" f"LOGIN hostname={hostname} app={appname}", "login"
) )
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "mssql_tds_login", "source": "mssql_tds_login",
"username": username, "username": username,
"hostname": hostname, "hostname": hostname,
"source_ip": src_ip, "source_ip": src_ip,
"dest_ip": dst_ip, "dest_ip": dst_ip,
"database": database, "database": database,
}, source_module=self.name) })
except Exception: except Exception:
logger.debug("Failed to parse TDS Login7", exc_info=True) logger.debug("Failed to parse TDS Login7", exc_info=True)
@@ -423,13 +424,13 @@ class DBInterceptor(BaseModule):
self._stats["logins_captured"] += 1 self._stats["logins_captured"] += 1
self._record_query(ts, src_ip, dst_ip, "mysql", username, database, "", "login") self._record_query(ts, src_ip, dst_ip, "mysql", username, database, "", "login")
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "mysql_login", "source": "mysql_login",
"username": username, "username": username,
"database": database, "database": database,
"source_ip": src_ip, "source_ip": src_ip,
"dest_ip": dst_ip, "dest_ip": dst_ip,
}, source_module=self.name) })
except Exception: except Exception:
logger.debug("Failed to parse MySQL login", exc_info=True) logger.debug("Failed to parse MySQL login", exc_info=True)
@@ -499,13 +500,13 @@ class DBInterceptor(BaseModule):
self._stats["logins_captured"] += 1 self._stats["logins_captured"] += 1
self._record_query(ts, src_ip, dst_ip, "postgres", username, database, "", "login") self._record_query(ts, src_ip, dst_ip, "postgres", username, database, "", "login")
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "postgres_startup", "source": "postgres_startup",
"username": username, "username": username,
"database": database, "database": database,
"source_ip": src_ip, "source_ip": src_ip,
"dest_ip": dst_ip, "dest_ip": dst_ip,
}, source_module=self.name) })
# ------------------------------------------------------------------ # ------------------------------------------------------------------
# Redis parser # Redis parser
@@ -533,12 +534,12 @@ class DBInterceptor(BaseModule):
self._stats["logins_captured"] += 1 self._stats["logins_captured"] += 1
self._record_query(ts, src_ip, dst_ip, "redis", username, "", "AUTH ***", "login") self._record_query(ts, src_ip, dst_ip, "redis", username, "", "AUTH ***", "login")
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "redis_auth", "source": "redis_auth",
"username": username, "username": username,
"source_ip": src_ip, "source_ip": src_ip,
"dest_ip": dst_ip, "dest_ip": dst_ip,
}, source_module=self.name) })
elif cmd in ("GET", "SET", "HGET", "HSET", "DEL", "KEYS", elif cmd in ("GET", "SET", "HGET", "HSET", "DEL", "KEYS",
"MGET", "MSET", "LPUSH", "RPUSH", "SADD", "ZADD"): "MGET", "MSET", "LPUSH", "RPUSH", "SADD", "ZADD"):
+3 -2
View File
@@ -19,6 +19,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.passive.ldap_harvester") logger = logging.getLogger("sensor.passive.ldap_harvester")
@@ -270,12 +271,12 @@ class LDAPHarvester(BaseModule):
if "ms-mcs-admpwd" in body_str: if "ms-mcs-admpwd" in body_str:
self._stats["laps_reads_detected"] += 1 self._stats["laps_reads_detected"] += 1
logger.warning("LAPS password read detected from %s (base DN: %s)", src_ip, base_dn) logger.warning("LAPS password read detected from %s (base DN: %s)", src_ip, base_dn)
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "laps_read", "source": "laps_read",
"source_ip": src_ip, "source_ip": src_ip,
"base_dn": base_dn, "base_dn": base_dn,
"detail": "LAPS password attribute requested in LDAP search", "detail": "LAPS password attribute requested in LDAP search",
}, source_module=self.name) })
def _handle_search_result_entry(self, body: bytes, src_ip: str, def _handle_search_result_entry(self, body: bytes, src_ip: str,
ts: float) -> None: ts: float) -> None:
+3 -2
View File
@@ -18,6 +18,7 @@ from pathlib import Path
from typing import Optional from typing import Optional
from modules.base import BaseModule from modules.base import BaseModule
from utils.credential_encryption import emit_credential_found
logger = logging.getLogger("sensor.passive.smb_monitor") logger = logging.getLogger("sensor.passive.smb_monitor")
@@ -331,13 +332,13 @@ class SMBMonitor(BaseModule):
for pattern in GPP_PATTERNS: for pattern in GPP_PATTERNS:
if pattern in file_lower: if pattern in file_lower:
self._stats["gpp_access_detected"] += 1 self._stats["gpp_access_detected"] += 1
self.bus.emit("CREDENTIAL_FOUND", { emit_credential_found(self.bus, self.name, {
"source": "smb_gpp_access", "source": "smb_gpp_access",
"client_ip": client_ip, "client_ip": client_ip,
"share": share, "share": share,
"file_path": file_path, "file_path": file_path,
"detail": "GPP/SYSVOL file access — may contain cleartext passwords", "detail": "GPP/SYSVOL file access — may contain cleartext passwords",
}, source_module=self.name) })
break break
# ------------------------------------------------------------------ # ------------------------------------------------------------------
+120
View File
@@ -0,0 +1,120 @@
#!/usr/bin/env python3
"""Helper to get encryption key from Infisical and encrypt credentials before emission."""
import subprocess
import json
import logging
from typing import Optional
logger = logging.getLogger("sensor.credential_encryption")
# Cache the key in memory during module lifetime
_cached_key: Optional[bytes] = None
def get_credential_encryption_key() -> Optional[bytes]:
"""Retrieve credential encryption key from Infisical via creds CLI.
Returns None if key cannot be retrieved (in which case plaintext fallback).
Caches key in memory for the lifetime of the module process.
"""
global _cached_key
if _cached_key is not None:
return _cached_key
try:
# Try to get CREDENTIAL_ENCRYPTION_KEY from Infisical via creds CLI
result = subprocess.run(
["/home/n0mad1k/.local/bin/creds", "get", "CREDENTIAL_ENCRYPTION_KEY"],
capture_output=True,
timeout=5,
text=True,
)
if result.returncode != 0:
logger.warning("Failed to retrieve CREDENTIAL_ENCRYPTION_KEY from Infisical: %s", result.stderr)
return None
key_hex = result.stdout.strip()
if not key_hex:
logger.warning("CREDENTIAL_ENCRYPTION_KEY is empty in Infisical")
return None
# Convert hex string to bytes
_cached_key = bytes.fromhex(key_hex)
if len(_cached_key) != 32:
logger.error("CREDENTIAL_ENCRYPTION_KEY must be 32 bytes (256-bit), got %d", len(_cached_key))
_cached_key = None
return None
logger.debug("Successfully loaded CREDENTIAL_ENCRYPTION_KEY from Infisical")
return _cached_key
except FileNotFoundError:
logger.error("creds CLI not found at ~/.local/bin/creds")
return None
except subprocess.TimeoutExpired:
logger.error("Timeout retrieving CREDENTIAL_ENCRYPTION_KEY from Infisical")
return None
except Exception as e:
logger.error("Error retrieving CREDENTIAL_ENCRYPTION_KEY: %s", e)
return None
def emit_credential_found(bus, source_module: str, payload: dict) -> None:
"""Emit CREDENTIAL_FOUND event with encrypted sensitive fields.
Sensitive fields are encrypted before emission. If encryption fails,
plaintext fallback is used with a warning logged.
Args:
bus: EventBus instance
source_module: Module name emitting the event
payload: Credential payload dict with fields like username, password, etc.
"""
from utils.crypto import encrypt_credential_dict
key = get_credential_encryption_key()
if key:
try:
encrypted_payload = encrypt_credential_dict(payload, key)
bus.emit("CREDENTIAL_FOUND", encrypted_payload, source_module=source_module)
return
except Exception as e:
logger.warning("Failed to encrypt credential payload: %s — using plaintext fallback", e)
# Fallback: emit plaintext with warning
bus.emit("CREDENTIAL_FOUND", payload, source_module=source_module)
def decrypt_credential_payload(payload: dict) -> dict:
"""Decrypt sensitive fields in credential event payload.
Looks for fields marked with __encrypted__ prefix and decrypts them.
Returns a new dict with decrypted fields.
Args:
payload: Credential dict with possibly encrypted fields
Returns:
New dict with decrypted values
"""
from utils.crypto import decrypt_credential_field
key = get_credential_encryption_key()
if not key:
return payload
decrypted = {}
for field, value in payload.items():
if isinstance(value, str) and value.startswith("__encrypted__"):
try:
decrypted[field] = decrypt_credential_field(value, key)
except Exception as e:
logger.warning("Failed to decrypt field %s: %s", field, e)
decrypted[field] = value
else:
decrypted[field] = value
return decrypted
+59
View File
@@ -73,6 +73,65 @@ class CryptoEngine:
return self._aesgcm.decrypt(nonce, ct, aad) return self._aesgcm.decrypt(nonce, ct, aad)
def encrypt_credential_dict(payload: dict, key: bytes) -> dict:
"""Encrypt sensitive fields in credential event payload before bus emission.
Sensitive fields: username, password, hash, domain, token, secret, api_key
Returns a new dict with encrypted fields as bytes wrapped in __encrypted__ markers.
"""
if len(key) != KEY_SIZE:
raise ValueError(f"Key must be {KEY_SIZE} bytes")
engine = CryptoEngine(key)
encrypted = {}
sensitive_fields = {
"username", "password", "hash", "hash_value", "domain",
"token", "secret", "api_key", "access_token", "refresh_token",
"credential_value", "nt_hash", "lm_hash", "response",
"auth_string", "bearer_token"
}
for field, value in payload.items():
if field.lower() in sensitive_fields and isinstance(value, str) and value:
try:
# Encrypt the string value
plaintext = value.encode("utf-8")
ciphertext = engine.encrypt(plaintext)
# Encode to base64 for transport over JSON
import base64
encrypted[field] = f"__encrypted__{base64.b64encode(ciphertext).decode('ascii')}"
except Exception:
# If encryption fails, include plaintext (log warning in caller)
encrypted[field] = value
else:
# Pass through unencrypted
encrypted[field] = value
return encrypted
def decrypt_credential_field(value: str, key: bytes) -> str:
"""Decrypt a single encrypted credential field."""
if not isinstance(value, str) or not value.startswith("__encrypted__"):
return value
if len(key) != KEY_SIZE:
raise ValueError(f"Key must be {KEY_SIZE} bytes")
try:
import base64
engine = CryptoEngine(key)
encrypted_part = value[len("__encrypted__"):]
ciphertext = base64.b64decode(encrypted_part)
plaintext = engine.decrypt(ciphertext)
return plaintext.decode("utf-8")
except Exception:
# Return original if decryption fails
return value
def encrypt_file( def encrypt_file(
src: str, src: str,
dst: str, dst: str,