From 4e1829a25c24a39b4c88408f0cac0fa1614a5946 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=83=AD=E5=90=89=E6=B5=A9?= <1625567290@qq.com> Date: Fri, 14 Aug 2026 19:20:07 +0800 Subject: [PATCH] fix(eval): classify CONNECT hosts and kill raw TCP tunnels MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Refuse blocklisted CONNECT targets before the tunnel opens. Raw TCP fallbacks — including SSH over 443 — are audited as raw_tunnel, stripped of payload, and closed. HTTPS interception for unrelated hosts is unchanged. Fixes #2977 Generated-by: Grok --- packages/eval/README.md | 4 +- packages/eval/harbor/egress_filter.py | 96 +++++++++++++++++++++- packages/eval/harbor/test_egress_filter.py | 83 +++++++++++++++++++ 3 files changed, 180 insertions(+), 3 deletions(-) diff --git a/packages/eval/README.md b/packages/eval/README.md index 039673e051..015eac9fdb 100644 --- a/packages/eval/README.md +++ b/packages/eval/README.md @@ -86,8 +86,8 @@ cohort. `MAKA_EVAL_EGRESS_NAMESPACE_TEST=1 python3 harbor/test_cell_egress_names the overlay and the checked-in policy and asserts that contract in a real cell namespace; it needs a Docker daemon and outbound network, and skips otherwise. This URL policy is a blocklist for known benchmark and public-solution contamination surfaces, not a complete defense against a deliberately -invented lookup channel. It classifies what it can read: a `CONNECT` tunnel carrying something other -than TLS or HTTP reaches no rule and no audit record, which is tracked in issue #2977. Collected Maka runtime files +invented lookup channel. It classifies HTTP(S) requests and `CONNECT` hosts against the blocklist, and +kills tunnels that fall back to raw TCP. Collected Maka runtime files and egress audit logs are represented in attempt artifacts with byte counts and SHA-256 digests. The local image tag remains a machine deployment identity rather than a registry digest; digest pinning is tracked in issue #2953. diff --git a/packages/eval/harbor/egress_filter.py b/packages/eval/harbor/egress_filter.py index ff6f601b81..685c5246d9 100644 --- a/packages/eval/harbor/egress_filter.py +++ b/packages/eval/harbor/egress_filter.py @@ -101,12 +101,50 @@ def public_trajectory_repository(host: str, path_query: str) -> bool: except ImportError: http = None +try: + from mitmproxy.proxy import commands as proxy_commands +except ImportError: + proxy_commands = None + def request(flow: object) -> None: + apply_http_policy(flow, flow.request.pretty_url) + + +def http_connect(flow: object) -> None: + try: + raw_url = connect_target_url(flow) + except Exception: + raw_url = "" + apply_http_policy(flow, raw_url) + + +def tcp_start(flow: object) -> None: + record_raw_tunnel(flow) + kill_flow(flow) + + +def tcp_message(flow: object) -> None: + messages = getattr(flow, "messages", None) + if messages: + messages[-1].content = b"" + kill_flow(flow) + + +def next_layer(nextlayer: object) -> None: + current = getattr(nextlayer, "layer", None) + if current is None or type(current).__name__ != "TCPLayer": + return + context = getattr(nextlayer, "context", None) + record_raw_tunnel(context) + nextlayer.layer = CloseRawLayer(context) + + +def apply_http_policy(flow: object, raw_url: str) -> None: if http is None: raise RuntimeError("mitmproxy is required to run the Eval egress filter") try: - matched = contamination_rule(flow.request.pretty_url) + matched = contamination_rule(raw_url) if not matched: return rule_id, host, normalized_path = matched @@ -127,6 +165,62 @@ def request(flow: object) -> None: pass +def connect_target_url(flow: object) -> str: + request = flow.request + host = (getattr(request, "pretty_host", None) or getattr(request, "host", "") or "").strip() + if not host: + raise ValueError("empty CONNECT host") + if ":" in host and not host.startswith("["): + host = f"[{host}]" + port = getattr(request, "port", None) + if port in (None, 443): + return f"https://{host}/" + if port == 80: + return f"http://{host}/" + return f"https://{host}:{port}/" + + +def tcp_peer(flow: object) -> tuple[str, str]: + server = getattr(flow, "server_conn", None) or getattr(flow, "server", None) + address = getattr(server, "address", None) if server is not None else None + if isinstance(address, (tuple, list)) and address: + host = str(address[0])[:255] + port = address[1] if len(address) > 1 else "" + return host, f":{port}" if port != "" else "" + return "", "" + + +def record_raw_tunnel(flow: object) -> None: + host, path = tcp_peer(flow) + try: + append_audit("raw_tunnel", host, path) + except Exception: + pass + + +def kill_flow(flow: object) -> None: + kill = getattr(flow, "kill", None) + if callable(kill) and getattr(flow, "killable", True): + try: + kill() + except Exception: + pass + + +class CloseRawLayer: + def __init__(self, context: object) -> None: + self.context = context + + def handle_event(self, event: object): + if proxy_commands is None: + return + yield + for name in ("client", "server"): + connection = getattr(self.context, name, None) + if connection is not None: + yield proxy_commands.CloseConnection(connection) + + def blocked_response(rule_id: str): return http.Response.make( 451, diff --git a/packages/eval/harbor/test_egress_filter.py b/packages/eval/harbor/test_egress_filter.py index 6efe000191..ca0e5a27c8 100644 --- a/packages/eval/harbor/test_egress_filter.py +++ b/packages/eval/harbor/test_egress_filter.py @@ -78,6 +78,89 @@ def make(status, body, headers): self.assertIn("host", record) self.assertIn("normalizedPath", record) + def test_http_connect_refuses_blocklisted_hosts_before_the_tunnel_opens(self) -> None: + class Response: + @staticmethod + def make(status, body, headers): + return {"status": status, "body": body, "headers": headers} + + with tempfile.TemporaryDirectory() as directory: + MODULE.http = SimpleNamespace(Response=Response) + MODULE.AUDIT_PATH = Path(directory) / "hits.jsonl" + blocked = type( + "Flow", + (), + { + "request": type( + "Request", + (), + {"pretty_host": "tbench.ai", "host": "tbench.ai", "port": 443}, + )() + }, + )() + MODULE.http_connect(blocked) + self.assertEqual(blocked.response["status"], 451) + self.assertEqual(blocked.response["headers"]["X-Maka-Eval-Egress-Rule"], "tbench_domain") + record = json.loads(MODULE.AUDIT_PATH.read_text().splitlines()[0]) + self.assertEqual(record["ruleId"], "tbench_domain") + self.assertEqual(record["host"], "tbench.ai") + + for host in ("example.com", "github.com", "ssh.github.com"): + allowed = type( + "Flow", + (), + { + "request": type( + "Request", + (), + {"pretty_host": host, "host": host, "port": 443}, + )(), + "response": None, + }, + )() + MODULE.http_connect(allowed) + self.assertIsNone(allowed.response, host) + + def test_tcp_start_kills_raw_tunnels_and_records_them(self) -> None: + with tempfile.TemporaryDirectory() as directory: + MODULE.AUDIT_PATH = Path(directory) / "hits.jsonl" + killed: list[str] = [] + flow = type( + "Flow", + (), + {"server_conn": SimpleNamespace(address=("ssh.github.com", 443))}, + )() + flow.kill = lambda: killed.append("killed") + MODULE.tcp_start(flow) + self.assertEqual(killed, ["killed"]) + record = json.loads(MODULE.AUDIT_PATH.read_text().splitlines()[0]) + self.assertEqual(record["ruleId"], "raw_tunnel") + self.assertEqual(record["host"], "ssh.github.com") + self.assertEqual(record["normalizedPath"], ":443") + + def test_tcp_message_drops_raw_payloads(self) -> None: + message = SimpleNamespace(content=b"SSH-2.0-test\r\n") + killed: list[str] = [] + flow = SimpleNamespace(messages=[message], killable=True) + flow.kill = lambda: killed.append("killed") + MODULE.tcp_message(flow) + self.assertEqual(message.content, b"") + self.assertEqual(killed, ["killed"]) + + def test_next_layer_replaces_raw_tcp_with_a_closer(self) -> None: + class TCPLayer: + pass + + with tempfile.TemporaryDirectory() as directory: + MODULE.AUDIT_PATH = Path(directory) / "hits.jsonl" + context = SimpleNamespace(server=SimpleNamespace(address=("ssh.github.com", 443))) + nextlayer = SimpleNamespace(layer=TCPLayer(), context=context) + MODULE.next_layer(nextlayer) + self.assertIsInstance(nextlayer.layer, MODULE.CloseRawLayer) + record = json.loads(MODULE.AUDIT_PATH.read_text().splitlines()[0]) + self.assertEqual(record["ruleId"], "raw_tunnel") + self.assertEqual(record["host"], "ssh.github.com") + if __name__ == "__main__": unittest.main()