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"
[\"'])?(?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()