diff --git a/extensions/business/cybersec/red_mesh/tests/test_finalization_aggregation.py b/extensions/business/cybersec/red_mesh/tests/test_finalization_aggregation.py index e33a2ed7..31c029ed 100644 --- a/extensions/business/cybersec/red_mesh/tests/test_finalization_aggregation.py +++ b/extensions/business/cybersec/red_mesh/tests/test_finalization_aggregation.py @@ -444,9 +444,17 @@ def _latest_pass(self): "worker_reports": { "0xUS": {"start_port": 1, "end_port": 443, "open_ports": [80, 443], "node_ip": "1.1.1.1", "country": "US", "nr_findings": 1, - "finding_counts": {"HIGH": 1}, "finding_signatures": ["sig1"]}, + "finding_counts": {"HIGH": 1}, "finding_signatures": ["sig1"], + "response_evidence": { + "target_host": "example.test", "resolved_ips": ["1.1.1.1"], + "ports": {"443": {"reachable": True}}, + }}, "0xIN": {"start_port": 1, "end_port": 443, "open_ports": [80], - "node_ip": "2.2.2.2", "country": "in", "nr_findings": 0}, + "node_ip": "2.2.2.2", "country": "in", "nr_findings": 0, + "response_evidence": { + "target_host": "example.test", "resolved_ips": ["2.2.2.2"], + "ports": {"443": {"reachable": False}}, + }}, }, "worker_scan_metrics": { "0xUS": {"scan_metrics": { @@ -528,6 +536,12 @@ def test_node_comparison_includes_failed_and_timed_out_nodes(self): self.assertEqual(comp["0xUS"]["metrics"]["total_duration"], 19.0) self.assertEqual(comp["0xUS"]["metrics"]["traffic_windows"][0]["attempts"], 8) self.assertEqual(comp["0xUS"]["metrics"]["threads"][0]["local_worker_id"], "thread-1") + # Response evidence remains attributable per node rather than merging the + # two vantages' values into a shared aggregate. + self.assertEqual(comp["0xUS"]["response_evidence"]["resolved_ips"], ["1.1.1.1"]) + self.assertEqual(comp["0xIN"]["response_evidence"]["resolved_ips"], ["2.2.2.2"]) + self.assertTrue(comp["0xUS"]["response_evidence"]["ports"]["443"]["reachable"]) + self.assertFalse(comp["0xIN"]["response_evidence"]["ports"]["443"]["reachable"]) # Metrics-only node with all-timeout connections -> timeout. self.assertEqual(comp["0xBR"]["status"], "timeout") # Selected peer that never reported at all -> failed (China-timeout case). diff --git a/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py b/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py index 69552579..d8d4d58a 100644 --- a/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py +++ b/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py @@ -7,17 +7,24 @@ """ import unittest +import time from unittest.mock import MagicMock, patch from extensions.business.cybersec.red_mesh.worker import PentestLocalWorker from extensions.business.cybersec.red_mesh.worker.response_fingerprint import ( EXCERPT_MAX_BYTES, + RESPONSE_BODY_MAX_BYTES, certificate_identity, excerpt_allowed, normalize_content_type, + read_bounded_response_body, resolve_host, sanitize_excerpt, ) +from extensions.business.cybersec.red_mesh.constants import ( + FINGERPRINT_HTTP_TIMEOUT, + FINGERPRINT_TIMEOUT, +) from .conftest import DummyOwner @@ -57,6 +64,20 @@ def test_redacts_credential_key_values(self): self.assertNotIn("sk-abc123def456", excerpt) self.assertIn("REDACTED", excerpt) + def test_redacts_multiword_credential_values(self): + excerpt = sanitize_excerpt('{"password": "correct horse battery staple", "ok": 1}') + self.assertNotIn("correct horse battery staple", excerpt) + + def test_redacts_basic_authorization_values(self): + excerpt = sanitize_excerpt("Authorization: Basic dXNlcjpwYXNz") + self.assertNotIn("dXNlcjpwYXNz", excerpt) + + def test_redacts_common_compound_credential_keys(self): + for key in ("client_secret", "refresh_token", "session_token", "clientSecret"): + with self.subTest(key=key): + excerpt = sanitize_excerpt(f'{key}="short-sensitive"') + self.assertNotIn("short-sensitive", excerpt) + def test_redacts_bearer_tokens(self): excerpt = sanitize_excerpt("Authorization: Bearer eyJhbGciOiJIUzI1NiJ9.payload") self.assertNotIn("eyJhbGciOiJIUzI1NiJ9", excerpt) @@ -76,6 +97,11 @@ def test_redacts_long_base64_runs(self): excerpt = sanitize_excerpt(f"state {blob} end") self.assertNotIn(blob, excerpt) + def test_redacts_long_urlsafe_base64_runs(self): + blob = "0123456789_abcdefghijklmnopqrstuvwxyz-ABCDE" + excerpt = sanitize_excerpt(f"state {blob} end") + self.assertNotIn(blob, excerpt) + def test_redaction_precedes_truncation(self): # A credential straddling the byte cap must not survive as a prefix. padding = "x" * (EXCERPT_MAX_BYTES - 20) @@ -136,6 +162,16 @@ def test_addresses_are_sorted_and_deduplicated(self): self.assertEqual(addresses, ["1.2.3.4", "93.184.216.34"]) self.assertIsNone(error) + def test_resolution_timeout_is_recorded(self): + blocker = MagicMock(side_effect=lambda *_args: time.sleep(0.2)) + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.socket.getaddrinfo", + blocker, + ): + addresses, error = resolve_host("slow.example", timeout=0.01) + self.assertEqual(addresses, []) + self.assertEqual(error, "DNS resolution timed out") + class TestCertificateIdentity(unittest.TestCase): @@ -184,6 +220,145 @@ def test_evidence_present_in_status_after_capture(self): self.assertIn("response_evidence", worker.get_status()) +class TestHttpCapture(unittest.TestCase): + + @staticmethod + def _response(body=b"Acmehello", **headers): + response = MagicMock() + response.status_code = 200 + response.url = "https://example.test/" + response.history = [] + response.encoding = "utf-8" + response.headers = {"Content-Type": "text/html", **headers} + response.iter_content.return_value = [body] + return response + + def test_http_capture_streams_and_uses_existing_timeout_constant(self): + worker = _make_worker(comparison_ports=[443]) + worker._target_timeout = MagicMock(return_value=12) + response = self._response(**{"Content-Length": "37", "server": "nginx"}) + session = MagicMock() + session.get.return_value = response + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.Session", + return_value=session, + ): + http, excerpt = worker._fingerprint_http("https", 443) + + worker._target_timeout.assert_called_once_with(FINGERPRINT_HTTP_TIMEOUT) + self.assertTrue(session.get.call_args.kwargs["stream"]) + self.assertFalse(session.get.call_args.kwargs["allow_redirects"]) + self.assertLessEqual(session.get.call_args.kwargs["timeout"], 12) + self.assertEqual(http["status"], 200) + self.assertEqual(http["title"], "Acme") + self.assertIsNotNone(http["body_sha256"]) + self.assertIn("hello", excerpt) + response.close.assert_called_once() + + def test_oversized_body_is_capped_without_claiming_a_partial_hash(self): + response = self._response( + body=b"a" * (RESPONSE_BODY_MAX_BYTES + 1), + **{"Content-Length": str(RESPONSE_BODY_MAX_BYTES + 100)}, + ) + worker = _make_worker(comparison_ports=[443]) + session = MagicMock() + session.get.return_value = response + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.Session", + return_value=session, + ): + http, excerpt = worker._fingerprint_http("https", 443) + + self.assertEqual(http["body_length"], RESPONSE_BODY_MAX_BYTES + 100) + self.assertIsNone(http["body_sha256"]) + self.assertLessEqual(len(excerpt.encode("utf-8")), EXCERPT_MAX_BYTES) + + def test_unknown_response_encoding_falls_back_safely(self): + response = self._response(body=b"plain text") + response.encoding = "not-a-real-codec" + worker = _make_worker(comparison_ports=[443]) + session = MagicMock() + session.get.return_value = response + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.Session", + return_value=session, + ): + http, excerpt = worker._fingerprint_http("https", 443) + + self.assertEqual(http["status"], 200) + self.assertEqual(excerpt, "plain text") + + def test_redirect_bodies_are_never_buffered(self): + redirect = self._response(**{"Location": "/final"}) + redirect.status_code = 302 + redirect.url = "https://example.test/start" + redirect.iter_content.side_effect = AssertionError("redirect body must not be read") + final = self._response() + final.url = "https://example.test/final" + session = MagicMock() + session.get.side_effect = [redirect, final] + worker = _make_worker(comparison_ports=[443]) + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.Session", + return_value=session, + ): + http, _excerpt = worker._fingerprint_http("https", 443) + + redirect.iter_content.assert_not_called() + redirect.close.assert_called_once() + self.assertEqual(http["redirect_count"], 1) + self.assertEqual(http["final_url"], "https://example.test/final") + + def test_port_probe_uses_existing_fingerprint_timeout_constant(self): + worker = _make_worker(comparison_ports=[443]) + worker._target_timeout = MagicMock(return_value=6) + worker._tls_unverified_connect = MagicMock(return_value=(None, None, None)) + worker._fingerprint_http = MagicMock(return_value=(None, None)) + connection = MagicMock() + connection.__enter__.return_value = connection + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.socket.create_connection", + return_value=connection, + ) as create_connection: + worker._fingerprint_port(443) + + worker._target_timeout.assert_called_once_with(FINGERPRINT_TIMEOUT) + self.assertEqual(create_connection.call_args.kwargs["timeout"], 6) + + def test_bounded_reader_stops_at_the_byte_cap(self): + response = MagicMock() + response.iter_content.return_value = [b"abc", b"def"] + body, complete = read_bounded_response_body(response, max_bytes=4, max_seconds=10) + self.assertEqual(body, b"abcd") + self.assertFalse(complete) + + def test_bounded_reader_hashes_an_exact_cap_complete_body(self): + response = MagicMock() + response.iter_content.return_value = [b"a" * RESPONSE_BODY_MAX_BYTES] + body, complete = read_bounded_response_body( + response, + max_bytes=RESPONSE_BODY_MAX_BYTES, + max_seconds=10, + ) + self.assertEqual(len(body), RESPONSE_BODY_MAX_BYTES) + self.assertTrue(complete) + + def test_bounded_reader_enforces_wall_deadline(self): + response = MagicMock() + + def delayed_chunks(**_kwargs): + time.sleep(0.2) + yield b"late" + + response.iter_content.side_effect = delayed_chunks + started = time.monotonic() + body, complete = read_bounded_response_body(response, max_bytes=10, max_seconds=0.01) + elapsed = time.monotonic() - started + self.assertEqual(body, b"") + self.assertFalse(complete) + self.assertLess(elapsed, 0.1) + + class TestAggregationAttribution(unittest.TestCase): """ Response evidence describes one vantage. Registering it as a cross-worker diff --git a/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py b/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py index d89ef221..78fb7411 100644 --- a/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py +++ b/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py @@ -19,10 +19,17 @@ import hashlib import ipaddress +import queue import re import socket +import threading +import time import requests +from requests.compat import urljoin +from requests.models import DEFAULT_REDIRECT_LIMIT + +from ..constants import FINGERPRINT_HTTP_TIMEOUT, FINGERPRINT_TIMEOUT # Excerpts exist so an operator can see *that* two vantages were served # different content. A few hundred bytes of the head of the document is @@ -35,8 +42,10 @@ # Headers that identify *which* infrastructure answered. These are what # actually differ between a CDN edge in Brazil and one in China. CAPTURED_HEADERS = ("server", "via", "x-cache", "cf-ray", "x-powered-by", "location") +REDIRECT_STATUSES = frozenset({301, 302, 303, 307, 308}) TITLE_MAX_CHARS = 200 +RESPONSE_BODY_MAX_BYTES = 1024 * 1024 _TITLE_RE = re.compile(r"(.*?)", re.IGNORECASE | re.DOTALL) @@ -51,16 +60,25 @@ (re.compile(r"(?i)\bbearer\s+[A-Za-z0-9._~+/\-]+=*"), "bearer [REDACTED]"), ( re.compile( - r"(?i)\b(authorization|api[-_]?key|apikey|access[-_]?token|token|" - r"session[-_]?id|sessionid|session|secret|password|passwd|pwd)\b" + r"(?i)\b(authorization|" + r"(?:api|access|refresh|session|client|auth|id|private|secret)[-_]?" + r"(?:key|token|secret|id)|" + r"apikey|token|session|secret|password|passwd|pwd)\b" # An optional closing quote covers JSON keys such as {"api_key": "..."}. - r"[\"']?(\s*[:=]\s*)[\"']?[^\s\"',;&<>]+" + r"(?P[\"'])?(?P\s*[:=]\s*)" + r"(?P\"(?:\\.|[^\"\\])*\"|'(?:\\.|[^'\\])*'|[^\r\n,;&<>}]+)" ), - r"\1\2[REDACTED]", + r"\1\g\g[REDACTED]", ), (re.compile(r"[A-Za-z0-9._%+\-]+@[A-Za-z0-9.\-]+\.[A-Za-z]{2,}"), "[REDACTED_EMAIL]"), (re.compile(r"\b[A-Fa-f0-9]{32,}\b"), "[REDACTED_HEX]"), - (re.compile(r"\b[A-Za-z0-9+/]{32,}={0,2}"), "[REDACTED_B64]"), + ( + re.compile( + r"(? max_bytes: + stopped.set() + response.close() + return b"".join(chunks)[:max_bytes], False + + def certificate_identity(cert_der): """ Extract comparable identity fields from a DER-encoded certificate. @@ -183,7 +270,10 @@ def _capture_response_fingerprint(self): if not ports: return - resolved_ips, resolver_error = resolve_host(self.target) + resolved_ips, resolver_error = resolve_host( + self.target, + timeout=self._target_timeout(FINGERPRINT_HTTP_TIMEOUT), + ) evidence = { "target_host": self.target, "resolved_ips": resolved_ips, @@ -204,7 +294,10 @@ def _fingerprint_port(self, port): entry = {"reachable": False, "tls": None, "http": None, "excerpt": None} try: - with socket.create_connection((self.target, port), timeout=self._target_timeout(3)): + with socket.create_connection( + (self.target, port), + timeout=self._target_timeout(FINGERPRINT_TIMEOUT), + ): entry["reachable"] = True except Exception: # An unreachable port is itself a comparable result: a target that @@ -229,36 +322,81 @@ def _fingerprint_port(self, port): def _fingerprint_http(self, scheme, port): """Issue one GET and reduce the response to comparable attributes.""" url = f"{scheme}://{self.target}:{port}/" + session = requests.Session() + resp = None try: user_agent = getattr(self, "scanner_user_agent", "") headers = {"User-Agent": user_agent} if user_agent else {} - resp = requests.get( - url, - timeout=self._target_timeout(5), - verify=False, - allow_redirects=True, - headers=headers, - ) + timeout = self._target_timeout(FINGERPRINT_HTTP_TIMEOUT) + deadline = time.monotonic() + timeout + current_url = url + redirect_count = 0 + while True: + remaining = deadline - time.monotonic() + if remaining <= 0: + raise requests.Timeout("response fingerprint deadline exceeded") + resp = session.get( + current_url, + timeout=remaining, + verify=False, + allow_redirects=False, + headers=headers, + stream=True, + ) + location = resp.headers.get("Location") + if resp.status_code not in REDIRECT_STATUSES or not location: + break + if redirect_count >= DEFAULT_REDIRECT_LIMIT: + raise requests.TooManyRedirects("response fingerprint redirect limit exceeded") + current_url = urljoin(resp.url or current_url, location) + redirect_count += 1 + resp.close() + resp = None except Exception as exc: self.P(f"Response fingerprint GET failed on {url}: {exc}", color='y') + if resp is not None: + resp.close() + session.close() return None, None - content_type = normalize_content_type(resp.headers.get("Content-Type")) - title_match = _TITLE_RE.search(resp.text[:5000]) - http = { - "status": resp.status_code, - "final_url": resp.url, - "redirect_count": len(resp.history), - "title": title_match.group(1).strip()[:TITLE_MAX_CHARS] if title_match else None, - "content_type": content_type, - "body_length": len(resp.content), - "body_sha256": hashlib.sha256(resp.content).hexdigest(), - "headers": { - name: resp.headers.get(name) - for name in CAPTURED_HEADERS - if resp.headers.get(name) - }, - } - - excerpt = sanitize_excerpt(resp.text) if excerpt_allowed(content_type) else None - return http, excerpt + try: + body, body_complete = read_bounded_response_body( + resp, + max_bytes=RESPONSE_BODY_MAX_BYTES, + max_seconds=max(deadline - time.monotonic(), 0), + ) + encoding = resp.encoding or "utf-8" + try: + body_text = body.decode(encoding, errors="replace") + except LookupError: + body_text = body.decode("utf-8", errors="replace") + content_type = normalize_content_type(resp.headers.get("Content-Type")) + title_match = _TITLE_RE.search(body_text[:5000]) + declared_length = resp.headers.get("Content-Length") + try: + declared_length = int(declared_length) if declared_length is not None else None + except (TypeError, ValueError): + declared_length = None + if declared_length is not None and declared_length < 0: + declared_length = None + body_length = len(body) if body_complete else declared_length + http = { + "status": resp.status_code, + "final_url": resp.url, + "redirect_count": redirect_count, + "title": title_match.group(1).strip()[:TITLE_MAX_CHARS] if title_match else None, + "content_type": content_type, + "body_length": body_length, + "body_sha256": hashlib.sha256(body).hexdigest() if body_complete else None, + "headers": { + name: resp.headers.get(name) + for name in CAPTURED_HEADERS + if resp.headers.get(name) + }, + } + + excerpt = sanitize_excerpt(body_text) if excerpt_allowed(content_type) else None + return http, excerpt + finally: + resp.close() + session.close()