Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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": {
Expand Down Expand Up @@ -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).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand Down Expand Up @@ -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):

Expand Down Expand Up @@ -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"<html><title>Acme</title>hello</html>", **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
Expand Down
Loading