From b602eec3ef518fcc58faf58c4e772e4dcc131669 Mon Sep 17 00:00:00 2001 From: toderian Date: Tue, 4 Aug 2026 22:07:53 +0000 Subject: [PATCH 1/2] fix: address RM-050 backend review blockers What changed: - harden excerpt redaction for multi-token and URL-safe credential shapes - bound DNS and streamed HTTP capture with existing timeout constants - prove per-node evidence isolation through final comparison aggregation Why: - close privacy, availability, and contract blockers found by independent review Checks: - focused RM-050 backend suite: 300 passed, 15 subtests passed - full RedMesh suite: 2045 passed, 4 unrelated baseline failures, 3 skipped --- .../tests/test_finalization_aggregation.py | 18 ++- .../tests/test_response_fingerprint.py | 115 ++++++++++++++ .../red_mesh/worker/response_fingerprint.py | 144 ++++++++++++++---- 3 files changed, 245 insertions(+), 32 deletions(-) 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..09fab4a4 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,14 @@ 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_bearer_tokens(self): excerpt = sanitize_excerpt("Authorization: Bearer eyJhbGciOiJIUzI1NiJ9.payload") self.assertNotIn("eyJhbGciOiJIUzI1NiJ9", excerpt) @@ -76,6 +91,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 +156,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 +214,91 @@ 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"}) + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.get", + return_value=response, + ) as get: + http, excerpt = worker._fingerprint_http("https", 443) + + worker._target_timeout.assert_called_once_with(FINGERPRINT_HTTP_TIMEOUT) + self.assertTrue(get.call_args.kwargs["stream"]) + self.assertEqual(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, + **{"Content-Length": str(RESPONSE_BODY_MAX_BYTES + 100)}, + ) + worker = _make_worker(comparison_ports=[443]) + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.get", + return_value=response, + ): + 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]) + with patch( + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.get", + return_value=response, + ): + http, excerpt = worker._fingerprint_http("https", 443) + + self.assertEqual(http["status"], 200) + self.assertEqual(excerpt, "plain text") + + 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) + + 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..95845cca 100644 --- a/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py +++ b/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py @@ -19,11 +19,16 @@ import hashlib import ipaddress +import queue import re import socket +import threading +import time import requests +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 # enough for that; storing more would commit third-party content to an @@ -37,6 +42,7 @@ CAPTURED_HEADERS = ("server", "via", "x-cache", "cf-ray", "x-powered-by", "location") TITLE_MAX_CHARS = 200 +RESPONSE_BODY_MAX_BYTES = 1024 * 1024 _TITLE_RE = re.compile(r"(.*?)", re.IGNORECASE | re.DOTALL) @@ -54,13 +60,20 @@ r"(?i)\b(authorization|api[-_]?key|apikey|access[-_]?token|token|" r"session[-_]?id|sessionid|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"(?= deadline: + complete = False + break + if not chunk: + continue + remaining = max_bytes - total + if len(chunk) >= remaining: + chunks.append(chunk[:remaining]) + total += min(len(chunk), remaining) + complete = False + break + chunks.append(chunk) + total += len(chunk) + + return b"".join(chunks), complete + + def certificate_identity(cert_der): """ Extract comparable identity fields from a DER-encoded certificate. @@ -183,7 +238,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 +262,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 @@ -232,33 +293,56 @@ def _fingerprint_http(self, scheme, port): try: user_agent = getattr(self, "scanner_user_agent", "") headers = {"User-Agent": user_agent} if user_agent else {} + timeout = self._target_timeout(FINGERPRINT_HTTP_TIMEOUT) resp = requests.get( url, - timeout=self._target_timeout(5), + timeout=timeout, verify=False, allow_redirects=True, headers=headers, + stream=True, ) except Exception as exc: self.P(f"Response fingerprint GET failed on {url}: {exc}", color='y') 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=timeout, + ) + 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": len(resp.history), + "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() From 918be218ec89313242891e7c8ac613755cdc043c Mon Sep 17 00:00:00 2001 From: toderian Date: Tue, 4 Aug 2026 22:15:50 +0000 Subject: [PATCH 2/2] fix: close RM-050 backend review round two What changed: - cover compound credential key spellings in excerpt redaction - follow redirects manually without reading intermediate bodies - enforce a wall deadline around streamed reads and preserve exact-cap hashes Why: - close residual privacy, availability, and evidence-loss findings Checks: - response fingerprint and aggregation suites: 55 passed, 9 subtests passed --- .../tests/test_response_fingerprint.py | 80 +++++++++++-- .../red_mesh/worker/response_fingerprint.py | 110 +++++++++++++----- 2 files changed, 152 insertions(+), 38 deletions(-) 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 09fab4a4..d8d4d58a 100644 --- a/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py +++ b/extensions/business/cybersec/red_mesh/tests/test_response_fingerprint.py @@ -72,6 +72,12 @@ 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) @@ -231,15 +237,18 @@ 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.get", - return_value=response, - ) as get: + "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(get.call_args.kwargs["stream"]) - self.assertEqual(get.call_args.kwargs["timeout"], 12) + 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"]) @@ -248,13 +257,15 @@ def test_http_capture_streams_and_uses_existing_timeout_constant(self): def test_oversized_body_is_capped_without_claiming_a_partial_hash(self): response = self._response( - body=b"a" * RESPONSE_BODY_MAX_BYTES, + 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.get", - return_value=response, + "extensions.business.cybersec.red_mesh.worker.response_fingerprint.requests.Session", + return_value=session, ): http, excerpt = worker._fingerprint_http("https", 443) @@ -266,15 +277,38 @@ 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.get", - return_value=response, + "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) @@ -298,6 +332,32 @@ def test_bounded_reader_stops_at_the_byte_cap(self): 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): """ diff --git a/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py b/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py index 95845cca..78fb7411 100644 --- a/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py +++ b/extensions/business/cybersec/red_mesh/worker/response_fingerprint.py @@ -26,6 +26,8 @@ import time import requests +from requests.compat import urljoin +from requests.models import DEFAULT_REDIRECT_LIMIT from ..constants import FINGERPRINT_HTTP_TIMEOUT, FINGERPRINT_TIMEOUT @@ -40,6 +42,7 @@ # 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 @@ -57,8 +60,10 @@ (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"(?P[\"'])?(?P\s*[:=]\s*)" r"(?P\"(?:\\.|[^\"\\])*\"|'(?:\\.|[^'\\])*'|[^\r\n,;&<>}]+)" @@ -172,25 +177,52 @@ def read_bounded_response_body(response, max_bytes, max_seconds): """Read a streamed response without allowing a hostile peer to grow memory forever.""" chunks = [] total = 0 - complete = True deadline = time.monotonic() + max_seconds + events = queue.Queue(maxsize=1) + stopped = threading.Event() - for chunk in response.iter_content(chunk_size=64 * 1024): - if time.monotonic() >= deadline: - complete = False - break - if not chunk: - continue - remaining = max_bytes - total - if len(chunk) >= remaining: - chunks.append(chunk[:remaining]) - total += min(len(chunk), remaining) - complete = False - break - chunks.append(chunk) - total += len(chunk) + def _read(): + try: + for chunk in response.iter_content(chunk_size=64 * 1024): + if stopped.is_set(): + return + events.put(("chunk", chunk)) + if not stopped.is_set(): + events.put(("done", None)) + except Exception: + if not stopped.is_set(): + events.put(("error", None)) - return b"".join(chunks), complete + threading.Thread(target=_read, daemon=True).start() + + while True: + remaining_seconds = deadline - time.monotonic() + if remaining_seconds <= 0: + stopped.set() + response.close() + try: + events.get_nowait() + except queue.Empty: + pass + return b"".join(chunks), False + try: + kind, payload = events.get(timeout=remaining_seconds) + except queue.Empty: + stopped.set() + response.close() + return b"".join(chunks), False + if kind == "done": + return b"".join(chunks), True + if kind == "error": + return b"".join(chunks), False + if not payload: + continue + chunks.append(payload) + total += len(payload) + if total > max_bytes: + stopped.set() + response.close() + return b"".join(chunks)[:max_bytes], False def certificate_identity(cert_der): @@ -290,27 +322,48 @@ 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 {} timeout = self._target_timeout(FINGERPRINT_HTTP_TIMEOUT) - resp = requests.get( - url, - timeout=timeout, - verify=False, - allow_redirects=True, - headers=headers, - stream=True, - ) + 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 try: body, body_complete = read_bounded_response_body( resp, max_bytes=RESPONSE_BODY_MAX_BYTES, - max_seconds=timeout, + max_seconds=max(deadline - time.monotonic(), 0), ) encoding = resp.encoding or "utf-8" try: @@ -330,7 +383,7 @@ def _fingerprint_http(self, scheme, port): http = { "status": resp.status_code, "final_url": resp.url, - "redirect_count": len(resp.history), + "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, @@ -346,3 +399,4 @@ def _fingerprint_http(self, scheme, port): return http, excerpt finally: resp.close() + session.close()